【问题标题】:How do implement a variable similar to Spark's Accumulators in Apache Beam如何在 Apache Beam 中实现类似于 Spark 的累加器的变量
【发布时间】:2021-09-10 15:53:53
【问题描述】:

我目前在 Spark 中使用 Apache Beam 2.29.0。我的管道使用来自 Kafka 的数据,为此我有一个自定义的 KafkaConsumer,Beam 通过调用 ConsumerFactoryFn 来创建该数据。在运行期间,我需要在自定义 Kafka 消费者之间共享一段持久数据。这在 Spark 中非常简单,我将创建一个 Accumulator 变量,所有执行程序以及驱动程序都可以访问该变量。 由于 Beam 旨在在多个平台上运行 Spark、Flink、Google Dataflow,因此它不提供此功能。有谁知道实现这个的方法吗?

【问题讨论】:

    标签: apache-beam


    【解决方案1】:

    我相信边输入应该可以工作。您可以阅读有关侧面输入here 的信息。侧输入是您的 DoFn 每次处理输入 PCollection 中的元素时都可以访问的附加输入。

    Here是一个如何使用的例子。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2019-10-05
      • 1970-01-01
      • 1970-01-01
      • 2013-11-26
      • 2017-07-05
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多