系列导航:

前一篇解决的是“消息怎样可靠到达 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 的视角看,这次投递没有被确认,所以重新投递完全合理;从业务的视角看,这却已经是第二次收到同一事件。如果消费者设计成“默认每次到达都执行一次”,那重复消费就会直接变成重复扣减、重复建单、重复发券。

因此消费可靠性的第一原则不是“阻止重复”,而是:

  1. 接受重复投递一定会发生。
  2. 让同一业务事件重复到达时,最终结果仍然只生效一次。
  3. 把真正不能自动恢复的消息尽快隔离出去,而不是无限空转。

幂等键先回答一个问题:什么算同一件事

幂等实现经常失败,不是因为代码不会写,而是因为团队没有先定义“重复”的判定标准。对 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 时,一个安全顺序通常是:

  1. 开启数据库事务。
  2. 尝试插入 (inventory-reservation, eventId)
  3. 如果唯一键冲突,说明这条消息已经处理过,直接安全 ack
  4. 如果插入成功,在同一个事务里写入库存预留记录。
  5. 事务提交成功后,再向 RabbitMQ 发送 ack
重复投递下的幂等事务边界同一个事件先写入去重记录并提交库存预留。确认前崩溃后,重投递会因唯一键冲突而安全确认,不重复写库存。ONE EVENT ID · ONE BUSINESS RESULT正常事务:业务副作用与去重记录一起提交Brokerorder.createdeventId=evt…稳定业务事件 IDINSERT inbox_messages唯一约束库存预留同一数据库事务COMMITack故障窗口:COMMIT 后、ack 前退出ack 前崩溃redeliveryeventId=evt…第二次投递唯一键冲突已处理过ack不重复写库存安全退出同一 eventId 的第二次到达不会再次越过库存写入边界。移动端幂等事务边界先完成正常事务,再说明确认前崩溃,最后展示重投递因唯一键冲突而安全确认。ONE EVENT ID · ONE RESULT1 · 正常事务Broker → eventId=evt…INSERT inbox_messages库存预留 → COMMIT提交后才可以 ack2 · ack 前崩溃COMMIT 已完成,Broker 未收到 ack所以同一事件会 redelivery3 · 重投递安全退出redelivery · eventId=evt…唯一键冲突ack · 不重复写库存effectively-once 业务结果
重复投递如何收敛为一次业务结果RabbitMQ 可以 at-least-once 投递;事务边界和唯一约束把业务观察结果收敛为 effectively-once。
同一个 eventId 在数据库事务中先写入去重记录,再写库存预留并提交;如果 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”,那“订单不存在”更可能是数据缺失、顺序错误或脏消息,继续原地重试只会放大问题。

因此,有限重试的核心不是“做几个重试队列”,而是:

  1. 只有确认可能恢复的失败才进入重试链路。
  2. 每次重试都要保留 eventId 和当前尝试次数。
  3. 达到上限后必须退出主消费路径。

如果瞬时失败来自 HTTP 调用,还要给请求携带稳定的幂等键;否则消费者本身虽然有 Inbox,下游仍可能在超时重试时重复执行。

一个够用的分级退避可以是:

有限重试与死信隔离瞬时失败进入一档退避队列,等待后经 DLX 回主队列再消费;再次失败才进入下一档。第 4 次仍失败进入隔离路径,毒消息绕过常规重试。RETRY BUDGET · RECOVER OR ISOLATE主队列inventory.order-created失败分类瞬时失败毒消息 / 不可恢复第 1 次失败retry.30s30s 后回 main短时恢复窗口第 2 次失败retry.5m5m 后回 main继续保留 eventId第 3 次失败retry.30m30m 后回 main最后一次退避每档:TTL / DLX 回主队列再消费,仍失败才进入下一档第 4 次仍失败inventory.dlq调查与受控重放毒消息 / 不可恢复 → 不进入常规重试只把可能恢复的故障放入退避梯;无法恢复的问题应尽快从主消费路径隔离。移动端有限重试与死信隔离每档等待结束后都回主队列重新消费,再次失败才进入下一档;毒消息绕过常规重试进入隔离路径。FINITE RETRY LADDER主队列 · inventory.order-created开始处理失败分类瞬时失败 → 有限退避毒消息 / 不可恢复第 1 次失败retry.30s等待 30s · DLX 回 main第 2 次失败retry.5m等待 5m · DLX 回 main第 3 次失败retry.30m等待 30m · DLX 回 main每档等待后回主队列再消费再次失败才进入下一档第 4 次仍失败inventory.dlq隔离、调查、修复后重放不进入常规重试
有限重试与死信隔离横向箭头压缩表示“等待后回主队列再消费,仍失败才进入下一档”;重试预算用尽后,消息退出主消费路径。
短时可恢复的失败进入 30 秒、5 分钟和 30 分钟退避队列;每档等待结束后都由 DLX 路由回主队列重新消费,而不是在重试队列之间直接转发。第 4 次仍失败进入死信队列。

图里为了便于阅读,把三次“回主队列、重新消费、再次失败”的循环压成了一条横向路径。真实拓扑不是 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/rejectrequeue: false、消息 TTL 到期、队列超过长度限制,以及 Quorum Queue 超过投递次数限制。完整触发条件见 RabbitMQ Dead Letter Exchanges

生产环境优先通过 policy 设置源队列的 dead-letter-exchangedead-letter-routing-key,避免把参数硬编码进队列声明;policy 可以在线调整,硬编码参数通常需要删除并重建队列。

这里有一个容易被忽略的安全边界:RabbitMQ 默认的内部死信重新发布不启用 publisher confirms。目标交换机不存在、路由不到队列或目标不可用时,源队列中的消息可能已经被移除,因此“配置了 DLQ”并不自动等于可靠隔离。对于不能接受这类丢失窗口的链路,可以选择:

  • 由应用执行带 publisher confirm 的重新发布,确认成功后再 ack 原消息。
  • 源队列使用 Quorum Queue,并显式设置 dead-letter-strategy=at-least-onceoverflow=reject-publish 和 DLX。该模式不是默认值,约束与资源开销见 Quorum Queue 的 at-least-once dead-lettering

死信队列最常见的另一种误用,是把它当成“另一个主队列”,进去之后自动绕回去。这样做通常会把问题藏起来,而不是解决问题。

更稳妥的理解是:

  • 主队列负责正常处理。
  • 重试队列负责有限退避。
  • 死信队列负责承接已经退出主链路的消息。

Broker 在死信时会添加 x-death,其中包含来源队列、原因、次数、交换机和 routing key 等信息;但下面这些应用诊断字段不会凭空产生。inventory.dlq 最好还能查到:

  • eventId
  • eventType
  • orderId
  • 原始 routing key
  • 失败分类
  • 最后一次错误摘要
  • 已尝试次数

如果需要补充失败分类和错误摘要,不能先 ack 原消息再“尽力发布”一份增强后的副本。要么确认增强副本已发布后再 ack,要么把诊断上下文写入以 eventId 为键的外部存储。这样当值班同学打开管理台或日志系统时,看到的不只是“有一条消息死了”,而是能直接判断是否需要补数据、是否应该人工重放。

毒消息、业务拒绝和不变量异常,不是同一种问题

这两个概念在生产里经常被混在一起,但它们应该被明确分开:

类型特征处理方式
毒消息JSON 解析失败、字段缺失、类型错误、协议不兼容不进入常规重试,可靠隔离并告警
预期内的业务拒绝订单状态不允许预留、库存规则校验不通过提交拒绝结果,必要时通过 Outbox 发布 inventory.reservation-rejected,然后 ack
不变量异常按契约必然存在的订单缺失、跨服务数据矛盾可靠隔离并告警,由人工调查

两者都不适合无限重试,但原因不同:

  • 毒消息的问题在消息本身,重试多少次都不会突然变成合法 JSON。
  • 预期内的业务拒绝是领域流程的正常分支,不应默认污染 DLQ;它应该成为可查询的业务结果,必要时再发出拒绝事件。
  • 不变量异常意味着系统事实与契约冲突,才更适合进入运维隔离路径。

把这两类问题混成一个“消费异常”统一处理,常见结果就是:日志暴涨、消息来回打转、健康消息被拖慢,而真正需要人工介入的问题迟迟没有暴露出来。

一个足够实用的消费决策表

围绕 order.created,消费者可以按下面这张表做判断:

结果例子动作
处理成功库存预留成功ack
已处理过eventId 命中唯一约束ack
瞬时失败DB 超时、依赖 503confirmed republish 到下一档重试队列,再 ack 原消息
预期业务拒绝订单状态不允许预留记录或发布拒绝结果,再 ack
毒消息 / 不变量异常反序列化失败、契约事实矛盾通过可靠路径隔离并告警

如果你的消费者目前还是“只要报错就 nack(requeue: true)”,那最先要改的通常不是参数,而是这张决策表。

本篇解决了什么:

  • 说明了为什么消费端必须接受重复投递,并围绕稳定 eventId 建立幂等键与“事务提交后再 ack”的最小消费事务边界。
  • 区分了有限重试、DLX/DLQ、毒消息、业务拒绝和不变量异常各自的职责,也补上了死信重新发布与 publisher confirm 的安全边界。

下一篇:RabbitMQ 端到端可靠性:Outbox、Inbox 与 effectively-once(十一)