RabbitMQ 消费可靠性:幂等、重试与死信队列(十)
系列导航:
- 上一篇:RabbitMQ 发布可靠性:mandatory、publisher confirms 与不可路由消息(九)
- 当前篇:RabbitMQ 消费可靠性:幂等、重试与死信队列(十)(本文)
- 下一篇:RabbitMQ 端到端可靠性:Outbox、Inbox 与 effectively-once(十一)
前一篇解决的是“消息怎样可靠到达 RabbitMQ”;这一篇只解决下一段链路:消息进入消费者之后,怎样既不轻易丢,也不因为重复投递把业务做重。 讨论消费可靠性时,先接受一个前提:RabbitMQ 的消费语义天然更接近 at-least-once,只要你把 ack 放在业务处理成功之后,重复消费就不是异常,而是系统在故障窗口里的正常表现。
继续沿用第七至九篇的 order.created:
{
"eventId": "01J7X8A6YQ4Q9M6Q4QG7H2J8F1",
"eventType": "order.created",
"occurredAt": "2026-07-15T10:00:00Z",
"orderId": "ORD-20260715-001",
"customerId": "C10086",
"amount": 199.00,
"currency": "CNY"
}
假设库存服务消费 order.created,为订单创建一条库存预留记录。接下来所有讨论都围绕这件事展开。
重复消费不是例外,而是默认事实
最典型的重复窗口,是消费者已经把业务写入数据库,但还没来得及 ack,进程就退出了。从 RabbitMQ 的视角看,这次投递没有被确认,所以重新投递完全合理;从业务的视角看,这却已经是第二次收到同一事件。如果消费者设计成“默认每次到达都执行一次”,那重复消费就会直接变成重复扣减、重复建单、重复发券。
因此消费可靠性的第一原则不是“阻止重复”,而是:
- 接受重复投递一定会发生。
- 让同一业务事件重复到达时,最终结果仍然只生效一次。
- 把真正不能自动恢复的消息尽快隔离出去,而不是无限空转。
幂等键先回答一个问题:什么算同一件事
幂等实现经常失败,不是因为代码不会写,而是因为团队没有先定义“重复”的判定标准。对 order.created 来说,最稳妥的判断通常不是“同一个时间到达的消息”,也不是“内容看起来差不多的消息”,而是同一个业务事件 ID。
这就是为什么示例消息里要有稳定的 eventId:
orderId表示业务对象是谁。eventType表示发生了什么。eventId表示这一件事件实例本身。
如果上游重发的是同一条事件,eventId 必须保持不变;只有这样,消费者才能可靠地识别“这是同一件事又来了一次”,而不是“又发生了一次新的订单创建”。
一个够用的去重表最小可以长这样:
CREATE TABLE inbox_messages (
consumer_name varchar(128) NOT NULL,
event_id varchar(64) NOT NULL,
processed_at timestamptz NOT NULL,
PRIMARY KEY (consumer_name, event_id)
);
这里刻意把主键设计成 (consumer_name, event_id),因为同一条 order.created 允许被“库存服务”和“积分服务”各处理一次,但不应该被“库存服务”处理两次。
幂等不是一张表,而是一条事务边界
真正关键的不是“有没有 dedup 表”,而是去重记录和业务结果必须通过同一个 DbContext、同一个数据库连接,在同一个数据库事务里提交。库存服务处理 order.created 时,一个安全顺序通常是:
- 开启数据库事务。
- 尝试插入
(inventory-reservation, eventId)。 - 如果唯一键冲突,说明这条消息已经处理过,直接安全
ack。 - 如果插入成功,在同一个事务里写入库存预留记录。
- 事务提交成功后,再向 RabbitMQ 发送
ack。
下面以 PostgreSQL 和 EF Core 为例。ON CONFLICT DO NOTHING 会把重复事件变成“影响 0 行”,不会像未处理的唯一键异常那样让当前 PostgreSQL 事务进入失败状态:
const string consumerName = "inventory-reservation";
await using var tx = await db.Database.BeginTransactionAsync(cancellationToken);
var inserted = await db.Database.ExecuteSqlInterpolatedAsync($"""
INSERT INTO inbox_messages (consumer_name, event_id, processed_at)
VALUES ({consumerName}, {message.EventId}, now())
ON CONFLICT (consumer_name, event_id) DO NOTHING
""", cancellationToken);
if (inserted == 0)
{
// 先结束当前数据库事务,再确认 RabbitMQ 投递。
await tx.RollbackAsync(cancellationToken);
await channel.BasicAckAsync(deliveryTag, multiple: false, cancellationToken);
return;
}
await inventoryReservations.AddAsync(new InventoryReservation
{
OrderId = message.OrderId,
EventId = message.EventId,
CreatedAt = DateTimeOffset.UtcNow
}, cancellationToken);
await db.SaveChangesAsync(cancellationToken);
await tx.CommitAsync(cancellationToken);
await channel.BasicAckAsync(deliveryTag, multiple: false, cancellationToken);
这条顺序背后的判断非常重要:
- 不能先
ack再提交数据库,否则崩溃时消息已经从队列删掉,业务结果却没落库。 - 不能先单独写去重表,再单独写业务表,否则两次写入之间仍然有窗口。
- 不能让去重辅助方法偷偷创建另一个
DbContext或数据库连接,否则文字上写着“同一事务”,运行时却没有原子性。 - 不要捕获唯一键异常后直接在原事务里继续执行;PostgreSQL 事务可能已经失败,应使用
ON CONFLICT DO NOTHING,或显式回滚到保存点。 - 也不能把
redelivered标记当成幂等依据,因为它只能说明 RabbitMQ 认为这条消息曾经投递过,不能替代你的业务去重主键。
这条事务边界只能覆盖同一个数据库里的写入。如果处理过程中还会调用 HTTP、发短信或写另一个数据库,eventId 还必须作为下游幂等键,或者通过新的 Outbox 把外部副作用拆到事务之后;否则“下游已经成功,但响应丢失”仍可能在重试时造成重复。
有限重试,只处理“等一会儿可能会好”的错误
重试不是对所有异常一视同仁地 requeue: true。更实用的做法,是先把失败分成三类:
| 类别 | 例子 | 是否值得重试 |
|---|---|---|
| 瞬时失败 | 数据库连接瞬断、下游 HTTP 503、短时锁等待 | 是,且要有上限和退避 |
| 预期内的业务拒绝 | 订单状态不允许预留、库存规则校验不通过 | 否,记录拒绝结果后 ack |
| 毒消息或不变量异常 | JSON 损坏、必填字段缺失、理论上必须存在的订单不存在 | 否,隔离并告警 |
对 order.created 而言,一个很常见的误判是:库存服务查不到订单,就把它当成“再等等也许就好了”。但如果系统语义明确要求“订单创建成功后才会发 order.created”,那“订单不存在”更可能是数据缺失、顺序错误或脏消息,继续原地重试只会放大问题。
因此,有限重试的核心不是“做几个重试队列”,而是:
- 只有确认可能恢复的失败才进入重试链路。
- 每次重试都要保留
eventId和当前尝试次数。 - 达到上限后必须退出主消费路径。
如果瞬时失败来自 HTTP 调用,还要给请求携带稳定的幂等键;否则消费者本身虽然有 Inbox,下游仍可能在超时重试时重复执行。
一个够用的分级退避可以是:
图里为了便于阅读,把三次“回主队列、重新消费、再次失败”的循环压成了一条横向路径。真实拓扑不是 retry.30s → retry.5m → retry.30m 直接串联,而是:
main → consumer → 失败 → retry.30s
│ TTL 到期,由 DLX 路由回 main
↓
main → consumer → 再次失败 → retry.5m
│ TTL 到期,由 DLX 路由回 main
↓
main → consumer → 再次失败 → retry.30m
│ TTL 到期,由 DLX 路由回 main
↓
main → consumer → 第 4 次仍失败 → 隔离
应用可以用 retry-attempt 头选择下一档队列,但必须一直复用同一个 eventId。如果依赖 RabbitMQ 的死信机制,也可以读取 Broker 添加的 x-death;要注意它按“队列 + 原因”压缩计数,单个 x-death.count 不能直接当作整条链路的总尝试次数。
应用把原消息转入重试队列时,应先等待新消息的 publisher confirm,确认成功后再 ack 原消息。若 confirm 已成功而原消息的 ack 丢失,两份消息仍可能同时存在,这正是前文幂等事务需要处理的重复窗口。
这条链路表达的是明确态度:重试预算是有限资源,不是无限兜底。 比如把最大处理次数定成 4,主队列第一次失败后等待 30s 再消费,第二次失败等待 5m,第三次失败等待 30m,第 4 次仍失败才退出主链路。
先分清 DLX 和 DLQ,再谈死信安全
RabbitMQ 不会把消息“直接塞进某个死信队列”。队列满足死信条件后,会把消息重新发布到 Dead Letter Exchange(DLX),再由交换机按绑定把它路由到 DLQ。常见触发条件包括:消费者 nack/reject 且 requeue: false、消息 TTL 到期、队列超过长度限制,以及 Quorum Queue 超过投递次数限制。完整触发条件见 RabbitMQ Dead Letter Exchanges。
生产环境优先通过 policy 设置源队列的 dead-letter-exchange 和 dead-letter-routing-key,避免把参数硬编码进队列声明;policy 可以在线调整,硬编码参数通常需要删除并重建队列。
这里有一个容易被忽略的安全边界:RabbitMQ 默认的内部死信重新发布不启用 publisher confirms。目标交换机不存在、路由不到队列或目标不可用时,源队列中的消息可能已经被移除,因此“配置了 DLQ”并不自动等于可靠隔离。对于不能接受这类丢失窗口的链路,可以选择:
- 由应用执行带 publisher confirm 的重新发布,确认成功后再
ack原消息。 - 源队列使用 Quorum Queue,并显式设置
dead-letter-strategy=at-least-once、overflow=reject-publish和 DLX。该模式不是默认值,约束与资源开销见 Quorum Queue 的 at-least-once dead-lettering。
死信队列最常见的另一种误用,是把它当成“另一个主队列”,进去之后自动绕回去。这样做通常会把问题藏起来,而不是解决问题。
更稳妥的理解是:
- 主队列负责正常处理。
- 重试队列负责有限退避。
- 死信队列负责承接已经退出主链路的消息。
Broker 在死信时会添加 x-death,其中包含来源队列、原因、次数、交换机和 routing key 等信息;但下面这些应用诊断字段不会凭空产生。inventory.dlq 最好还能查到:
eventIdeventTypeorderId- 原始 routing key
- 失败分类
- 最后一次错误摘要
- 已尝试次数
如果需要补充失败分类和错误摘要,不能先 ack 原消息再“尽力发布”一份增强后的副本。要么确认增强副本已发布后再 ack,要么把诊断上下文写入以 eventId 为键的外部存储。这样当值班同学打开管理台或日志系统时,看到的不只是“有一条消息死了”,而是能直接判断是否需要补数据、是否应该人工重放。
毒消息、业务拒绝和不变量异常,不是同一种问题
这两个概念在生产里经常被混在一起,但它们应该被明确分开:
| 类型 | 特征 | 处理方式 |
|---|---|---|
| 毒消息 | JSON 解析失败、字段缺失、类型错误、协议不兼容 | 不进入常规重试,可靠隔离并告警 |
| 预期内的业务拒绝 | 订单状态不允许预留、库存规则校验不通过 | 提交拒绝结果,必要时通过 Outbox 发布 inventory.reservation-rejected,然后 ack |
| 不变量异常 | 按契约必然存在的订单缺失、跨服务数据矛盾 | 可靠隔离并告警,由人工调查 |
两者都不适合无限重试,但原因不同:
- 毒消息的问题在消息本身,重试多少次都不会突然变成合法 JSON。
- 预期内的业务拒绝是领域流程的正常分支,不应默认污染 DLQ;它应该成为可查询的业务结果,必要时再发出拒绝事件。
- 不变量异常意味着系统事实与契约冲突,才更适合进入运维隔离路径。
把这两类问题混成一个“消费异常”统一处理,常见结果就是:日志暴涨、消息来回打转、健康消息被拖慢,而真正需要人工介入的问题迟迟没有暴露出来。
一个足够实用的消费决策表
围绕 order.created,消费者可以按下面这张表做判断:
| 结果 | 例子 | 动作 |
|---|---|---|
| 处理成功 | 库存预留成功 | ack |
| 已处理过 | eventId 命中唯一约束 | ack |
| 瞬时失败 | DB 超时、依赖 503 | confirmed republish 到下一档重试队列,再 ack 原消息 |
| 预期业务拒绝 | 订单状态不允许预留 | 记录或发布拒绝结果,再 ack |
| 毒消息 / 不变量异常 | 反序列化失败、契约事实矛盾 | 通过可靠路径隔离并告警 |
如果你的消费者目前还是“只要报错就 nack(requeue: true)”,那最先要改的通常不是参数,而是这张决策表。
本篇解决了什么:
- 说明了为什么消费端必须接受重复投递,并围绕稳定
eventId建立幂等键与“事务提交后再 ack”的最小消费事务边界。 - 区分了有限重试、DLX/DLQ、毒消息、业务拒绝和不变量异常各自的职责,也补上了死信重新发布与 publisher confirm 的安全边界。