【发布时间】:2017-07-11 12:35:45
【问题描述】:
我在使用 Apache Spark 实施一个工作流方面需要帮助。我的任务在下一个:
- 我有几个 CSV 文件作为源数据。注意:这些文件可能有不同的布局
- 我有元数据,其中包含我需要如何解析每个文件的信息(这不是问题)
- 主要目标:结果是带有几个附加列的源文件。我必须更新每个源文件而不加入一个输出范围。例如:源 10 个文件 -> 10 个结果文件,每个结果文件只有对应源文件的数据。
据我所知,Spark 可以通过掩码打开许多文件:
var source = sc.textFile("/source/data*.gz");
但在这种情况下,我无法识别文件的哪一行。如果我得到源文件列表并尝试通过以下场景进行处理:
JavaSparkContext sc = new JavaSparkContext(...);
List<String> files = new ArrayList() //list of source files full name's
for(String f : files)
{
JavaRDD<String> data = sc.textFile(f);
//process this file with Spark
outRdd.coalesce(1, true).saveAsTextFile(f + "_out");
}
但在这种情况下,我将以顺序模式处理所有文件。
接下来是我的问题:如何以并行模式处理多个文件?例如:一个文件 - 一个执行者?
非常感谢您的帮助!
【问题讨论】:
标签: apache-spark parallel-processing