【问题标题】:Apache Beam bundling issueApache Beam 捆绑问题
【发布时间】: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 操作?

【问题讨论】:

    标签: apache-spark apache-beam


    【解决方案1】:

    Spark 决定如何拆分输入,这里决定一次性遍历整个文件,因为它太小了。

    我在distcp application 中做过类似的事情;这使用 Spark 的 ParallelCollectionRDD 类来明确告诉 spark 将列表一一拆分。

    该类应该足以让您执行类似的操作 - 您可能必须在本地将初始文本文件读取到列表中,然后将列表传递给 ParallelCollectionRDD 构造函数

    【讨论】:

    • 感谢您的帮助,但您知道如何使用 Beam 做同样的事情吗?
    • 不,但如果您将它与 spark 一起使用,它可能只是排队。尝试在 apache 上询问 beam 用户邮件列表
    【解决方案2】:

    回复有点晚,但我查看了 Beam 在 2.16.0 版本中的功能。

    在第一个 TextIO.read() 之后,您将获得 2 个任务——我怀疑您最初的 346k 行文件列表被分成两个分区。此行为由 TextIO 中的 desiredBundleSize 控制,它被硬编码为 64MB。

    在 Spark 中,您的操作 ReadFromS3Mapper 将“融合”到到达的记录中,并且您将始终停留在两个分区中。

    如果你想保持相同的代码,你可以在两个转换之间强制重新分区:

    PCollection<String> allContents = p.apply("ReadLines", TextIO.read().from(options.getInputFile()))
            .apply("Repartition", Reshuffle.viaRandomKey())
            .apply(ParDo.of(new ReadFromS3Mapper()));
    

    作为替代方案,TextIO 和FileIO 实用程序中有很多有趣的模式可用。有一个与您的 almost exactly 匹配的示例(隐式包括 reshuffle)。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-08-04
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-04-29
      相关资源
      最近更新 更多