pub-sub模式下如何追踪消息处理完成状态实现分批发布控制
Pub/Sub模式下批量消息的状态追踪与批次触发方案
单条消息处理状态追踪实现
核心是给消息打唯一批次标识+持久化存状态,别靠内存记,服务一重启状态全乱:
- 发布端发消息时,给同批次100条消息统一生成全局唯一的
batch_id,再给每条消息分配0-99的不重复message_seq,两个值都塞到消息中间件的消息头里,不要只放在消息体中,消费端不用解析全量消息体就能拿到这两个标识,性能更高也更稳。 - 持久化存储单条消息的处理状态,服务量级小就直接在业务库建张轻量表,核心字段就
batch_id、message_seq、process_status(枚举值:待处理/处理成功/处理失败),加个更新时间、错误信息字段方便排查;量级大就用Redis的Hash结构存,key设为msg_process_status:{batch_id},field用message_seq,value存状态值,读写性能更高。 - 消费端处理逻辑加幂等校验:拿到消息先查对应
batch_id+message_seq的状态,如果已经是处理成功,直接返回ACK丢弃消息,避免重复消费搞乱状态;如果没有对应记录,先插入一条待处理状态的记录,再执行业务逻辑。 - 业务逻辑执行无报错,就把对应记录状态更新为处理成功,再提交消息ACK;如果执行报错,先把状态标记为处理失败、记录错误信息,按配置的重试策略走重试,不要无限阻塞消费队列。
批次完成判定与发布端通知
核心是用原子计数判断批次完成度,靠同体系的Pub/Sub事件传完成信号,比硬调接口稳:
- 消费端每次把单条消息状态更新为处理成功时,同步对该批次的成功计数做原子自增:用数据库就加行锁后执行count统计,用Redis直接调用
INCR batch_success_count:{batch_id},拿到自增后的返回值判断是否等于100。 - 当自增返回值刚好等于100时,说明当前批次所有消息都处理成功,此时往专门的
batch_complete_event主题发一条事件消息,内容携带batch_id、完成时间、成功消息总数即可,不要在消费线程里直接HTTP调用发布端接口,避免发布端临时不可用导致通知丢失。 - 发布端发完一批100条消息后,把当前批次状态标记为「处理中」,暂停下一批次的发布逻辑,同时订阅
batch_complete_event主题。收到完成事件后,先校验事件里的batch_id是不是当前正在处理的批次、成功数是不是100,校验通过后把当前批次标记为「已完成」,再触发下一批100条消息的发布流程即可。
几个落地必踩的坑提前避:
- 全链路必须做幂等:不管是消息中间件重复投递业务消息,还是完成事件重复投递,都要靠之前持久化的状态做校验,禁止重复执行业务逻辑、重复触发下一批次发布。
- 不要做线程阻塞等待:发布端发完批次后不要挂起工作线程等结果,靠事件驱动触发后续流程就行,不然服务并发能力会直接垮掉。
- 失败消息要加兜底策略:如果某条消息重试多次还是失败,要及时触发告警人工介入,设置合理的批次超时时间,不要因为单条消息异常导致整个批次永远卡住、后续消息全发不出去。
内容的提问来源于stack exchange,提问作者ARJAV
相关产品推荐
相关产品推荐

