【问题标题】:How to do parallel pipeline?如何做并行流水线?
【发布时间】:2016-04-03 21:50:35
【问题描述】:

我在 Spark v.1.6.0 中构建了一个 scala 应用程序,它实际上结合了各种功能。我有用于扫描数据帧以查找某些条目的代码,我有对数据帧执行某些计算的代码,我有用于创建输出的代码,等等。

此时组件是“静态”组合的,即在我的代码中,我从组件 X 调用代码进行计算,我获取结果数据并调用组件 Y 的方法,该方法采用数据作为输入。

我想让这个更灵活,让用户简单地指定一个管道(可能是一个具有并行执行的管道)。我会假设工作流程相当小且简单,如下图所示:

但是,我不知道如何最好地解决这个问题。

  • 我可以自己构建整个管道逻辑,这可能会导致相当多的工作,也可能会导致一些错误......
  • 我已经看到 Apache Spark 在 ML 包中带有一个 Pipeline 类,但是,如果我理解正确,它不支持并行执行(在示例中,两个 ParquetReader 可以同时读取和处理数据)
  • 显然有Luigi project 可以做到这一点(但是,它在页面上说 Luigi 用于长时间运行的工作流程,而我只需要短期运行的工作流程;Luigi 可能是矫枉过正?)?李>

对于在 Spark 中构建工作/数据流,您有什么建议?

【问题讨论】:

    标签: scala apache-spark dataflow apache-spark-ml


    【解决方案1】:

    我建议使用 Spark 的 MLlib 管道功能,您所描述的听起来很适合这种情况。它的一个好处是它允许 Spark 以一种可能比你更智能的方式为你优化流程。

    您提到它不能并行读取两个 Parquet 文件,但它可以以分布式方式读取每个单独的文件。因此,与其让 N/2 个节点分别处理每个文件,不如让 N 个节点串行处理它们,我希望这会给您一个类似的运行时间,特别是如果到 y-c 的映射是一对一的。基本上,您不必担心 Spark 未充分利用您的资源(如果您的数据正确分区)。

    但实际上情况可能会更好,因为 Spark 在优化流程方面比您更聪明。要记住的重要一点是,Spark 可能不会完全按照您定义它们的方式和单独的步骤执行操作:当您告诉它计算 y-c 时,它实际上并没有立即执行此操作。它是懒惰的(以一种好的方式!)并等到你建立了整个流程并要求它提供答案,此时它会分析流程,应用优化(例如,一种可能性是它可以弄清楚它没有'不必读取和处理其中一个或两个 Parquet 文件的一大块,尤其是 partition discovery),然后才执行最终计划。

    【讨论】:

    • 确实是优化代码的好方法。然而,这里没有免费的午餐!我没有发现 spark 的管道是超级并行优化的。对于以逻辑回归作为估计器的简单交叉验证管道,4 个 EC2 从属设备中 3 个以上的 CPU 利用率低于 50%。不好!不过,一切都缓存在 RAM 中。我实际上正在寻找如何优化它。
    猜你喜欢
    • 2011-08-10
    • 1970-01-01
    • 1970-01-01
    • 2012-06-04
    • 1970-01-01
    • 1970-01-01
    • 2013-11-07
    • 2023-03-07
    • 2022-10-13
    相关资源
    最近更新 更多