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开发中必须规避的典型错误,属于可变对象复用/可变状态泄漏类问题。
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

