Dataflow流式作业扩容时未确认消息的处理机制及重复疑问
Dataflow流式作业扩缩容时的消息重复处理与确认规则
重复处理是否为预期行为?
- 这是预期行为。当Dataflow触发自动扩容时,原Worker持有的未确认Pub/Sub消息,会因为消费组重新平衡、Worker租约过期等原因,被Pub/Sub重新分发给新加入的Worker,进而出现同一条消息被新旧Worker同时处理的情况。
- 从设计逻辑来看,Pub/Sub本身是至少一次投递的语义,Dataflow流式作业也天生具备容忍重复处理的能力,毕竟分布式系统里节点故障、扩缩容这类场景都是常态,重复投递是保证消息不丢失的必要机制。
消息确认及下游处理规则
- 消息确认逻辑:不管新旧Worker谁先完成处理,只要其中一个Worker成功把结果写入BQ,并向Pub/Sub发送确认(ACK),Pub/Sub就会把这条消息标记为已处理,后续其他Worker的ACK会被直接忽略。如果先处理的Worker发送ACK失败,后处理的Worker完成后发送的ACK依然会生效。
- 下游BQ的处理:如果BQ表没做去重逻辑,大概率会出现重复数据。解决方式要么在Dataflow的DoFn里加去重逻辑(比如基于消息ID做判断),要么在BQ侧设置主键约束,或者查询时做去重处理。
- 处理顺序的影响:即使新旧Worker的处理完成顺序打乱,Pub/Sub只认第一个成功的ACK。但BQ的写入顺序不一定和消息原始顺序一致,除非你的作业启用了恰好一次语义(需要满足一系列条件,比如使用可重入的DoFn、开启BQ流式写入的恰好一次模式等),否则要接受可能的乱序和重复情况。
内容的提问来源于stack exchange,提问作者Pav3k
相关产品推荐
相关产品推荐

