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

Apache Beam测试通过但生产失败:多次产出同一实例问题排查

问题根因

你的代码核心问题是复用了可变的PubsubMessage实例。
PubsubMessage是可变对象,你在SplitUpBatches.process的循环中直接修改同一个传入的msg实例的data字段,再直接yield该实例,最终yield出去的3个元素实际上是指向同一块内存地址的引用。循环结束后这个实例的data值是最后一次修改的third,只要下游是等所有元素都收集完成后再统一处理,就会出现所有消息的data都是third的问题,也就是你观测到的同一条消息出现在多个位置的异常。

第一个测试用例的问题及调整方案

test_incorrectly_passing属于巧合通过的无效测试:
TestPipeline默认的本地执行模式下,ParDo输出的元素会被立即传递给下游算子处理,不会批量缓存,刚好每次修改msg.data之后值就被下游Map算子取走,所以测出来的结果是对的。但线上Dataflow Runner为了提升性能会做元素批处理、缓存等优化,不会实时处理每一个yield的元素,问题就会暴露。

调整测试的方法也很简单:在断言前先缓存所有输出的PubsubMessage实例,再统一提取data字段,就能复现问题:

def test_adjusted(self):
    msg = PubsubMessage(self.body, {})
    with TestPipeline() as p:
        output = (
            p
            | beam.Create([msg])
            | "split up batches" >> beam.ParDo(SplitUpBatches())
            # 先收集所有消息实例再提取data,避免顺序巧合
            | beam.combiners.ToList()
            | "map to data" >> beam.Map(lambda msgs: [m.data for m in msgs])
        )
        assert_that(output, equal_to([['first', 'second', 'third']]))

调整后的测试会直接失败,和test_correctly_failing的结果一致。

是否属于Apache Beam典型错误

这是Apache Beam开发中必须规避的典型错误,属于可变对象复用/可变状态泄漏类问题。
Beam的执行模型不承诺元素的处理时机和缓存策略,不同Runner的优化逻辑也不一样,只要你在DoFn中修改传入的元素对象、或者复用可变对象yield,就必然会在某些场景下出现非预期结果,不管是Dataflow、Spark还是Flink Runner都可能踩坑。

修复方案

不要修改传入的PubsubMessage实例,每次yield新创建的实例即可:

class SplitUpBatches(beam.DoFn):
    def process(self, msg):
        bodies = msg.data.split('\n')
        for body in bodies:
            # 新建实例,不修改原对象
            yield PubsubMessage(
                data=body.strip(),
                attributes=msg.attributes,
                message_id=msg.message_id,
                publish_time=msg.publish_time
            )

修复后两个测试都会正常通过,线上问题也会解决。

内容的提问来源于stack exchange,提问作者Frans

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.05 23:45:03