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