You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.25 16:20:28