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