消息队列实战(5):消息确认、重试与幂等消费
2026/9/15 0:25:53 网站建设 项目流程

上一篇选型时反复强调"至少一次投递"是几乎所有 MQ 的默认语义——它保证不丢消息,但绝不保证不重复。本篇把"确认、重试、幂等"三件事讲成一条完整的链路:消费者怎么确认、失败后怎么重试、以及最关键的一步——如何用幂等把重复投递收敛成"一次业务结果"。这是所有 MQ 落地都必须先解决的通用问题,与具体产品无关。

一、确认与重试:为什么退避要加抖动

消费端有两种确认:确认成功(ack)和确认失败(nack)。nack 后 Broker 会按策略重投,常见策略是固定间隔或指数退避。直接固定间隔重试的隐患是"重试风暴":如果一条消息因为下游数据库抖动而失败,成千上万条同类消息会在同一时刻一起重试,把已经过载的下游再次压垮。正确的做法是指数退避加随机抖动,让重试时间错开。下面实现退避计算,并观察一条毒丸消息如何走完重试→死信的全程。

importrandomdefretry_delay(attempt,base=1.0,cap=60.0,jitter=0.3):expo=min(cap,base*(2**attempt))returnround(expo+random.uniform(0,expo*jitter),2)MAX_RETRIES=3defprocess(message):# 模拟: 毒丸消息永远失败, 其余成功returnmessage!="poison"defrun(message):attempt=0whileattempt<=MAX_RETRIES:ifprocess(message):print(f"{message}: 消费成功, 发送 ack")returnifattempt==MAX_RETRIES:print(f"{message}: 重试{MAX_RETRIES}次仍失败, 移入 DLQ")returndelay=retry_delay(attempt)print(f"{message}: 第{attempt+1}次失败,{delay}s 后重试")attempt+=1random.seed(1)formin("normal-1","poison","normal-2"):run(m)

运行输出:

normal-1: 消费成功, 发送 ack poison: 第 1 次失败, 1.04s 后重试 poison: 第 2 次失败, 2.51s 后重试 poison: 第 3 次失败, 4.92s 后重试 poison: 重试 3 次仍失败, 移入 DLQ normal-2: 消费成功, 发送 ack

退避间隔从 1.04s 增长到 4.92s,且因为加了抖动,间隔不是精确的 1、2、4 秒。抖动幅度一般取退避时间的 10%~30%,这个区间既能错峰,又不至于让延迟失控。重试次数耗尽后,消息被移入死信队列(DLQ),而不是无限重试阻塞后面的消息——这是保护主队列吞吐的关键,下一篇会专门展开死信队列。

二、幂等消费:把"至少一次"收敛成"恰好一次的业务结果"

即使有 ack 和重试,消费者在"处理成功但 ack 前崩溃"这类窗口里,仍会收到重复消息。业务层必须假设消息可能重复,并用幂等键去重。核心做法是:给每条消息赋予一个全局唯一的业务键(如订单号、支付流水号),消费者在处理前先查"这个键是否已处理过",处理过就跳过。下面用内存集合模拟去重,生产上应换成数据库唯一约束或 Redis 去重表。

classIdempotentConsumer:def__init__(self):self.seen=set()self.balance={}defhandle(self,biz_key,amount):ifbiz_keyinself.seen:print(f"{biz_key}: 重复消息, 跳过")returnself.seen.add(biz_key)self.balance["acct"]=self.balance.get("acct",0)+amountprint(f"{biz_key}: 首次处理, 入账{amount}")c=IdempotentConsumer()forbiz_key,amountin[("pay-1001",10),("pay-1001",10),("pay-1002",20)]:c.handle(biz_key,amount)print("最终余额:",c.balance["acct"])

运行输出:

pay-1001: 首次处理, 入账 10 pay-1001: 重复消息, 跳过 pay-1002: 首次处理, 入账 20 最终余额: 30

pay-1001第二次出现被直接跳过,余额没有被重复加两次。这里有个容易被忽略的细节:去重集合本身也要持久化。如果消费者进程重启,内存集合清空,重复消息又会被当成首次处理。生产上的可靠做法是让幂等键落到数据库里,用"插入唯一键"的原子性来判重——插入成功即首次处理,插入失败(唯一约束冲突)即重复,业务处理与幂等键写入放在同一个本地事务里,这样即使崩溃也不会漏判。

把 ack、重试、幂等串起来,就是一条完整的不丢不重的链路:Broker 靠 ack 确认消费者确实收到,靠重试补足偶发失败,消费者靠幂等键兜住重试带来的重复。注意这里达成的是"业务结果恰好一次",而不是严格意义上的"消息恰好一次投递"——真正端到端的 exactly-once 需要消息系统和下游存储配合事务,成本极高,第八篇讲事务消息时会再谈这个边界。

下一篇进入消息的高级能力:顺序消息、延迟消息与死信队列,看看在"不丢不重"之上,业务还如何利用 MQ 解决有序、定时和兜底三类问题。

三、幂等键落库与 exactly-once 的边界

内存集合去重只适合演示,生产上幂等键必须落到持久化存储,且和业务处理保持原子性。最可靠的做法是数据库唯一约束:把业务键作为唯一索引字段,处理业务时同时插入一条"处理记录",如果插入因唯一冲突失败,就说明这条消息处理过,直接跳过。业务写操作和这条去重记录放在同一个本地事务里,崩溃也不会出现"业务做了但没记录"或"记录了但业务没做"的裂缝。

选什么做幂等键有讲究。用消息自身的业务 ID(订单号、支付流水号)最理想,因为它天然唯一且语义清晰;如果没有业务 ID,可以用"消息体哈希 + 投递序号"合成一个键,但要小心消息体里含时间戳会导致哈希每次不同,反而去重失效。去重记录还要设置过期时间:如果无限保留,去重表会无限膨胀;但过期太早,一个延迟很久的重试消息会被误判为首次处理。生产上通常按"消息最大重试窗口"的几倍来设置过期,比如最大重试 24 小时,去重记录保留 72 小时。

最后明确 exactly-once 的边界。前面用"幂等"达成的,是"业务结果恰好一次",它假设消息本身可能重复,靠消费端兜住。而严格意义上的端到端 exactly-once(消息从生产到消费全程恰好一次投递)需要消息系统、生产端、消费端三方配合事务,比如 Kafka 的事务 + 幂等生产者 + 事务性消费,成本极高且仍有适用限制。绝大多数业务根本不需要这么强,用"至少一次投递 + 消费端幂等"就足以保证正确性,这是投入产出比最高的方案,也是本系列贯穿始终的默认假设。

ack 的时机也要注意:手动 ack 应该在业务处理成功之后、且最好和幂等键写入同一个事务内,这样"处理成功"和"确认消费"才是一致的,否则可能出现业务成功了却没 ack、消息被重复投递,或者 ack 了业务却失败、消息被永久跳过。

参考来源

  • RabbitMQ:消费者确认与重试
  • Kafka:传递语义
  • AWS:消息幂等处理

👍 觉得有用就点个赞 + 收藏,方便回头查阅;有疑问直接在评论区留言,我看到都会回。

🚀 本文属于《消息队列实战》系列,持续更新,关注不迷路。

📌 文章里的代码都能直接跑。想要可直接 clone 的完整工程 + 配套部署脚本 / 踩坑清单?评论一声或发邮件到cj2664@qq.com,我免费发你。
如果你正好在做类似系统、或有工程化难题想找人做,也欢迎邮件聊一句——我按实际情况评估,能落地的就接单或出方案。评论和邮件都能直接找到我,不用跳别的平台。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询