【发布时间】:2019-03-20 08:03:20
【问题描述】:
我的问题如下,我想汇总一些存储在 S3 上的数据。作为管道的初始输入,我使用了一个文本文件,其中包含应聚合的所有 S3 文件的路径。
PCollection<String> readInputPipeline = p.apply("ReadLines", TextIO.read().from(options.getInputFile()));
readInputPipeline = readInputPipeline.apply(ParDo.of(new ReadFromS3Mapper()));
输入文件有 346k 行。当我将此代码部署到从 S3 读取的 Spark 集群时,它看起来只发生在 2 个 Spark 任务中,即使有许多内核可用。我有什么办法可以增加这个操作的并行度吗?
我在 Amazon 上的 EMR 上运行此程序,其中有一台主机 (m3.xlarge) 和一台核心机器 (R3.4xlarge),具有以下选项:
"spark-submit"
"--driver-java-options='-Dspark.yarn.app.container.log.dir=/mnt/var/log/hadoop'",
"--master", "yarn",
"--executor-cores","16",
"--executor-memory","6g"
PS:也许解决方案是我不应该在这种情况下进行这种昂贵的 IO 操作?
【问题讨论】: