【发布时间】: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