Apache Beam 2.39.0+ DirectRunner AssertionError问题咨询
Apache Beam 2.39.0+版本
AssertionError: A total of 2 watermark-pending bundles did not execute问题分析与解决 问题背景
使用Python 3.9 + Apache Beam 2.38.0时,以下最小可复现示例可正常运行;但升级至2.39.0或2.44.0版本后,触发错误:AssertionError: A total of 2 watermark-pending bundles did not execute
开启DEBUG日志后可见关键信息:
- 两条bundle的
Stage input watermark: Timestamp(-9223372036854.775000)(即timestamp.MIN_TIMESTAMP) Bundle schedule watermark: Timestamp(9223372036854.775000)(即timestamp.MAX_TIMESTAMP)- 同时出现
Unable to add bundle for stage消息
复现代码:
import logging import apache_beam as beam def setup_logging(): log_format = '[%(asctime)-15s] [%(name)s] [%(levelname)s]: %(message)s' logging.basicConfig(format=log_format, level=logging.INFO) logging.info("Pipeline Started") class CreateKvPCollectWithSideInputDoFn(beam.DoFn): def __init__(self): super().__init__() def process(self, element, side_input): print(f"side_input_type: {type(side_input)}") yield "b", "2" class CreateKvPCollectDoFn(beam.DoFn): def __init__(self): super().__init__() def process(self, element): yield "a", "1" def main(): setup_logging() pipeline = beam.Pipeline() pcollect_input = ( pipeline | "Input/Create" >> beam.Create(["input"]) ) kvpcollect_1 = ( pcollect_input | "PCollection_1" >> beam.ParDo(CreateKvPCollectDoFn()) ) beamdict_1 = beam.pvalue.AsDict(kvpcollect_1) kvpcollect_2 = ( pcollect_input | "PCollection_2" >> beam.ParDo( CreateKvPCollectWithSideInputDoFn(), side_input=beamdict_1 ) ) kvpcollect_3 = ( (kvpcollect_1, kvpcollect_2) | "Flatten" >> beam.Flatten() ) beamdict_3 = beam.pvalue.AsDict(kvpcollect_3) ( pcollect_input | "UseBeamDict_3" >> beam.ParDo(CreateKvPCollectWithSideInputDoFn(), side_input=beamdict_3) | "PrintResult" >> beam.Map(print) ) result = pipeline.run() result.wait_until_finish() if __name__ == '__main__': main()
原因分析
该问题源于Apache Beam 2.39.0版本对侧输入(Side Input)的水印处理逻辑进行了严格升级:
- 2.38.0及更早版本中,对于
AsDict等创建的侧输入,系统允许其依赖的PCollection形成循环依赖(示例中kvpcollect_1→beamdict_1→kvpcollect_2→kvpcollect_3→beamdict_3→最终ParDo又依赖beamdict_3,而kvpcollect_3包含kvpcollect_1),此时水印处理会宽松跳过循环检查。 - 2.39.0+版本中,系统强化了水印传播的正确性校验,当检测到循环依赖导致水印无法推进(出现MIN/MAX_TIMESTAMP极端值)时,会直接触发断言错误,阻止无意义的bundle调度。
规避方案
方案1:打破侧输入的循环依赖
重构Pipeline,避免侧输入依赖的PCollection反向依赖自身输出。示例中可调整kvpcollect_3的来源,移除循环链路:
# 修改kvpcollect_3的定义,仅使用kvpcollect_2,打破循环 kvpcollect_3 = ( kvpcollect_2 | "Flatten" >> beam.Flatten() )
方案2:显式指定全局窗口约束侧输入
对于无需时间窗口的场景,给侧输入关联的PCollection显式指定全局窗口,强制水印正常推进:
from apache_beam.transforms.window import GlobalWindows, WindowInto # 给kvpcollect_1和kvpcollect_2添加全局窗口 kvpcollect_1 = ( pcollect_input | "PCollection_1" >> beam.ParDo(CreateKvPCollectDoFn()) | "Window1" >> WindowInto(GlobalWindows()) ) kvpcollect_2 = ( pcollect_input | "PCollection_2" >> beam.ParDo( CreateKvPCollectWithSideInputDoFn(), side_input=beamdict_1 ) | "Window2" >> WindowInto(GlobalWindows()) )
方案3:临时回退到Beam 2.38.0版本
如果暂时无法重构Pipeline,可回退到旧版本保持原有逻辑运行:
pip install apache-beam==2.38.0
内容的提问来源于stack exchange,提问作者Francis
相关产品推荐
相关产品推荐

