【问题标题】:Apache Beam/Dataflow: pipeline doesn't use most recent version of side input (streaming pipeline with global window and frequently updated side input)Apache Beam/Dataflow:管道不使用最新版本的侧输入(具有全局窗口和经常更新的侧输入的流式管道)
【发布时间】:2022-10-14 04:38:19
【问题描述】:

我将 Apache Beam (SDK 2.40.0) 与 GCP Dataflow 运行器和流式管道一起使用。我需要使用配置来处理我可以随时更改的数据。因此,我每 2 分钟加载一次(可接受的延迟),作为这样的侧面输入:

configs = (
        p
        | PeriodicImpulse(fire_interval=120, apply_windowing=False)
        | "Global Window" >> beam.WindowInto(
                window.GlobalWindows(),
                trigger=trigger.Repeatedly(trigger.AfterProcessingTime(5)),
                accumulation_mode=trigger.AccumulationMode.DISCARDING
            )
        | 'Get Side Input' >> beam.ParDo(GetConfigsFn())
)

通过附加的打印语句,我验证了配置每 2 分钟成功加载一次并输出到 PCollection。

我在另一个处理 PubSub 消息的步骤中使用配置(我省略了所有不相关的步骤,消息也在全局窗口中):

msgs_with_config = (
        pubsub_messages
        | 'Merge data and configs' >> beam.ParDo(AddConfigFromSideInputFn(), config_dict=beam.pvalue.AsDict(configs))
)

我面临的问题是,合并数据和配置步骤是使用旧版本的配置而不是最新版本。在使用更新版本的配置之前,需要任意时间(从几分钟、20 分钟到几个小时)。我的怀疑是,侧面输入缓存在某处,并且不会为每个处理的消息加载。

这是对这种行为的有效解释吗?它是预期的行为吗?还有其他可能的原因吗?

如何避免这种行为,以便始终使用最新的侧面输入版本?

【问题讨论】:

    标签: python google-cloud-dataflow apache-beam


    【解决方案1】:

    是的,辅助输入缓存在 Dataflow 工作程序中并用于捆绑包。如果您确实需要更快的重新加载,我建议重构管道以执行两个窗口 PCollections 的连接,而不是使用一个 PCollection 作为侧面输入。例如,使用CoGroupByKey 转换。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2018-05-11
      • 1970-01-01
      • 1970-01-01
      • 2021-10-19
      • 2021-11-25
      • 2019-05-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多