【问题标题】:Beam + Flink: No parallelism when using SDFBoundedSourceReaderBeam + Flink:使用 SDFBoundedSourceReader 时没有并行性
【发布时间】:2021-09-20 01:48:06
【问题描述】:

背景:我使用 TFX 管道和 Flink 作为 Beam 的运行器(使用 flink-on-k8s-operator 的 flink 会话集群)。 Flink 集群有 2 个任务管理器,每个 16 核,并行度设置为 32。TFX 组件调用beam.io.ReadFromTFRecord 来加载数据,传入一个 glob 文件模式。我有一个跨 160 个文件的 TFRecords 数据集。当我尝试运行该组件时,所有 160 个文件的处理最终都在 Flink 中的一个子任务中结束,即并行度实际上是 1。见下图:

我尝试了各种 Beam/Flink 选项和不同版本的 Beam/Flink,但行为保持不变。

此外,该行为会影响使用 apache_beam.io.iobase.SDFBoundedSourceReader 的任何内容,例如apache_beam.io.parquetio.ReadFromParquet 也有同样的问题。我的配置中是否有一些模糊的设置,或者这是 Flink 运行器的错误?我还在互联网上进行了广泛搜索,除了使用 beam.Reshuffle 的建议之外,找不到任何提及此问题的内容。

【问题讨论】:

    标签: apache-flink apache-beam tfx


    【解决方案1】:

    从运行图上看,算子节点是通过forward连接的,这样上游的一个任务的数据就只能被下游的一个任务处理。 该方案可以改变不同运营商的连接方式

    public DataStream<T> rebalance() {
        return setConnectionType(new RebalancePartitioner<T>());
    }
    

    【讨论】:

      【解决方案2】:

      似乎我缺少的模糊设置是 Beam 管道选项中的--experiments=pre_optimize=all。这将导致以下代码正在运行,并在 Splittable DoFn 扩展中包含 RESHUFFLEhttps://github.com/apache/beam/blob/v2.32.0/sdks/python/apache_beam/runners/portability/fn_api_runner/translations.py#L1433

      对于那些将来阅读本文的人来说,这适用于 Beam 2.32.0 和 Flink 1.13.2 - 这无疑会在某个时候发生变化,因此这个答案可能不再相关。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2020-10-26
        • 1970-01-01
        • 2020-09-21
        相关资源
        最近更新 更多