系列导航:

上一篇讲的是消费者如何应对重复投递;这一篇再往前后各退一步,看整条链路真正的可靠性边界。因为在生产环境里,最容易犯的错不是“不知道怎么开 publisher confirms”,而是把“Broker 收到了消息”误当成“业务已经可靠完成”。

本文继续使用同一个共享示例:订单服务创建订单后发布 order.created,库存服务消费该事件创建库存预留。

{
  "eventId": "01J7X8A6YQ4Q9M6Q4QG7H2J8F1",
  "eventType": "order.created",
  "occurredAt": "2026-07-15T10:00:00Z",
  "orderId": "ORD-20260715-001",
  "customerId": "C10086",
  "amount": 199.00,
  "currency": "CNY"
}

本篇只回答 6 件事:四层可靠性边界是什么、为什么 publisher confirms 不够、双写窗口在哪里、Outbox 解决什么、Inbox/dedup 解决什么,以及 at-least-once 与 effectively-once 的真正区别。

四层可靠性边界,不要把不同问题混成一个“消息可靠”

order.created 从订单服务送到库存服务,至少会跨过四层边界:

边界代表动作典型失败主要治理
应用事务边界订单库写入成功订单已提交,但消息还没发Outbox
Publisher → Broker消息被实际路由到的队列接受发布结果未知、消息不可路由mandatory return + publisher confirms + 稳定 eventId 重发
Broker → Consumer消息送到库存服务消费者崩溃、重新投递manual ack + at-least-once
消费事务边界库存库写入成功DB 已提交,但 ack 还没发Inbox / dedup

这四层之所以要拆开,是因为每一层的“确认成功”只对本层负责:

  • 数据库事务提交成功,不代表 RabbitMQ 已经接管消息。
  • publisher confirm 成功,不代表消费者已经执行业务。
  • 消费者收到了消息,不代表库存写入已经提交。
  • 业务库提交成功,也不代表 RabbitMQ 已经收到最后的 ack

所以“端到端可靠”从来不是一个开关,而是每一层都承认自己只能保证到哪里

publisher confirms 很重要,但它解决不了双写

publisher confirms 的价值,是让生产者知道 RabbitMQ 是否已经接管某次发布。对于已经成功路由的消息,它能回答的问题是:

  • RabbitMQ 是否接受了这次发布。
  • 对发送为 persistent 且路由到 durable 队列的消息,Broker 是否已经按相应队列类型的持久化语义承担责任。
  • 发布结果未知时,生产者是否应该重发。

positive publisher confirm 不能单独证明消息已经进入目标队列。AMQP 0-9-1 中,不可路由消息也可能收到 positive confirm;如果发布时没有设置 mandatory: true,它会被直接丢弃。因此 Outbox Publisher 只有同时满足下面两个条件,才能把记录标记成 Published

  1. 没有收到与本次发布对应的 basic.return
  2. 收到与本次发布对应的 positive publisher confirm。

RabbitMQ 保证 mandatory 消息的 basic.return 先于对应的 confirm 发出,但客户端仍须用消息的 MessageId/eventId 关联 return,用 Channel 发布序号关联 confirm,并最终映射回同一条 Outbox 记录。具体语义见 RabbitMQ Publisher ConfirmsUnroutable Message Handling

还要注意,mandatory 只证明至少一个队列接受了消息,不能证明某个指定队列或某条预期绑定一定存在。示例只有一个目标队列时这两个判断可以重合;存在多条绑定时,还必须通过部署期拓扑声明、权限校验和运行时监控保证预期目标没有缺席。

但它回答不了另一个更难的问题:订单库提交和消息发布之间,谁来保证原子性。

如果你的代码是这样:

  1. 向数据库插入订单。
  2. 提交数据库事务。
  3. 调用 RabbitMQ 发布 order.created
  4. 等待 confirm。

那只要进程在第 2 步和第 3 步之间崩溃,就会出现订单已经存在,但消息根本没发出去的情况。publisher confirms 再可靠,也无法回到过去补上这次没发生的发布。

反过来如果改成“先发消息,再提交订单事务”,窗口也没有消失,只是反过来了:消息可能已经被 RabbitMQ 接管,但订单事务随后回滚,下游就看到了一个本不该存在的 order.created

这就是双写窗口。

双写窗口的本质:同一件业务要写两个系统

只要一个业务动作既要写业务数据库,又要写消息系统,而且这两个写入不在同一个本地事务里,双写窗口就一定存在。

围绕 order.created,最典型的风险有两种:

窗口 A:先提交订单,再发消息
结果:订单存在,但事件缺失

窗口 B:先发消息,再提交订单
结果:事件已发出,但订单不存在

很多团队第一次做 RabbitMQ 时,会希望通过“重试 publish”“等 confirm”“多加几次日志”把这个问题糊住。但这些措施只能降低某些失败的不可见性,不能消灭双写本身。

因此,生产级可靠性必须先承认一件事:双写问题首先是应用架构问题,不是 Broker 参数问题。

Outbox 的职责:保证事件一定进入“可发布集合”

Transactional Outbox 的核心思想很简单:不要在订单事务提交之后临时起意去发消息,而是在同一个订单事务里,把待发布事件也写进本地数据库。

例如:

CREATE TABLE outbox_messages (
  event_id      varchar(64) PRIMARY KEY,
  aggregate_id  varchar(64) NOT NULL,
  sequence      bigint NOT NULL,
  event_type    varchar(128) NOT NULL,
  payload       jsonb NOT NULL,
  occurred_at   timestamptz NOT NULL,
  status        varchar(16) NOT NULL DEFAULT 'Pending',
  attempt_count integer NOT NULL DEFAULT 0,
  next_attempt_at timestamptz NOT NULL DEFAULT now(),
  last_error    text NULL,
  locked_by     varchar(128) NULL,
  locked_until  timestamptz NULL,
  published_at  timestamptz NULL,
  UNIQUE (aggregate_id, sequence),
  CHECK (status IN ('Pending', 'Processing', 'Published', 'Failed')),
  CHECK ((status = 'Published') = (published_at IS NOT NULL))
);

CREATE INDEX ix_outbox_dispatch
  ON outbox_messages (occurred_at, event_id)
  WHERE status IN ('Pending', 'Processing');

这里让 status 成为调度状态的唯一事实源,published_at 只记录成功时间,并用约束防止两者互相矛盾。attempt_countnext_attempt_atlast_error 分别服务于有限重试、退避调度与故障诊断;部分索引避免发布完成后的历史记录继续拖慢待发布扫描。PostgreSQL 中裸 timestamp 表示不带时区的时间,跨时区事件时间和 Lease 更适合统一使用 timestamptz

订单服务创建订单时做两件事,而且必须在同一个事务里提交:

  1. 写入 orders
  2. 写入一行 outbox_messages,内容就是那条 order.created

这样做之后,系统至少能保证:

  • 订单事务没提交,事件也不会凭空存在。
  • 订单事务一旦提交,事件就已经进入“待发布集合”。

单个后台发布器可以扫描 published_at IS NULL 的记录,把它们发布到 RabbitMQ,并在收到 confirm 后标记 published_at;多个实例共同扫描时,则需要在发布前先认领记录,后文会展开。

这里要特别注意一个常见误解:Outbox 不是为了避免重复发布,而是为了避免“订单成功了但事件永远没机会发”。

如果发布器在收到 confirm 之后、更新 published_at 之前崩溃,恢复后它仍然可能再次读取同一条 Outbox 记录并重发。所以 Outbox 天然接受 at-least-once 发布,前提是它重发时必须继续使用同一个 eventId

多实例 Outbox Publisher:先认领,再发布

单个后台 Publisher 可以顺序扫描未发布记录;但当多个 OrderService 实例都扫描同一张 Outbox 表时,不能让它们都读到同一行再各自发布。PostgreSQL 可以在一个短事务里用 FOR UPDATE SKIP LOCKED 批量认领消息:

BEGIN;

WITH candidates AS (
  SELECT event_id
  FROM outbox_messages
  WHERE (status = 'Pending' AND next_attempt_at <= CURRENT_TIMESTAMP)
     OR (status = 'Processing' AND locked_until < CURRENT_TIMESTAMP)
  ORDER BY occurred_at, event_id
  FOR UPDATE SKIP LOCKED
  LIMIT 100
)
UPDATE outbox_messages AS outbox
SET status = 'Processing',
    locked_by = $1,
    locked_until = CURRENT_TIMESTAMP + INTERVAL '60 seconds'
FROM candidates
WHERE outbox.event_id = candidates.event_id
RETURNING outbox.event_id, outbox.event_type, outbox.payload;

COMMIT;

SKIP LOCKED 负责“同一时刻谁抢到哪条记录”;Processinglocked_until 则让认领在事务提交后继续有效。locked_by 应是每次认领唯一的 claim ID,例如 实例名:GUID,而不是固定的应用名。这样 Lease 过期并被另一实例重新认领后,旧 Worker 就不能覆盖新 Worker 的状态。PostgreSQL 官方也把 SKIP LOCKED 定位为适合多个消费者访问 queue-like table 的机制,而不是一般查询的一致性视图,参见 SELECT ... SKIP LOCKED

示例中的 100 条和 60 seconds 只是占位值,不是生产默认值。Lease 必须覆盖从认领事务提交,到发布结果判定并完成状态更新的最长时间。 批次较大或 confirm 变慢时,应减小批次、限制发布并发,或者使用带 claim ID 条件的心跳续租。即便实现续租,进程长暂停和网络分区仍可能让 Lease 过期并产生并发重发,所以消费端幂等依旧不可省略。

认领事务一结束,就在事务外发布 RabbitMQ。不要把网络 I/O 放在持有数据库行锁的事务里。只有发布设置了 mandatory: true、没有收到 return,并且收到 positive publisher confirm 后,才能用 claim ID 条件完成这条记录:

UPDATE outbox_messages
SET status = 'Published',
    published_at = CURRENT_TIMESTAMP,
    locked_by = NULL,
    locked_until = NULL
WHERE event_id = $1
  AND status = 'Processing'
  AND locked_by = $2;

受影响行数为 1 才代表当前 Worker 完成了它自己的认领;如果为 0,通常意味着 Lease 已过期并被其他实例接手。对 basic.return、negative confirm 或明确的连接失败,应使用同样的 claim 条件增加 attempt_count、记录 last_error,并把 next_attempt_at 推迟到下一次退避时间;达到预算后转成 Failed 并告警。对“Broker 可能已经收到、但 confirm 未返回”的未知结果,不能标记为 Published,恢复后仍要用相同 eventId 重发并接受重复。

这套机制避免了 Lease 有效期内多个实例同时认领同一条 Outbox 记录,却不能消除 confirm 结果未知、Lease 到期和状态更新中断带来的重复窗口。最典型的一种是:

Broker 已接管消息

Publisher 收到 confirm

更新 Outbox 为 Published 前进程崩溃

Lease 到期后另一实例重新发布

因此 Lease 解决的是分工与故障恢复,不是 exactly-once。重发必须复用同一个 eventId,消费端仍要依赖 Inbox、唯一约束或业务幂等把最终结果收敛为一次。发布 Channel 的并发所有权见.NET 发布端实战(七),publisher confirm 的证明范围见发布可靠性(九)

多实例认领提高吞吐,但不保证全局发布顺序

ORDER BY occurred_at, event_id 只让每次候选集的选择稳定;一旦多个 Worker 在事务外并行发布,网络延迟、confirm 顺序和 Lease 重领都可能改变消息真正进入队列的先后。因此这套 Outbox Relay 不保证全局发布顺序,也不能仅凭 SQL 中的 ORDER BY 推导出业务事件顺序。

如果同一订单可能连续产生 order.createdorder.paidorder.cancelled,就要先定义是否要求“单订单有序”。常见做法包括:

  • 在 Outbox 中保存 aggregate_id 和单调递增的 sequence,消费者拒绝或暂存越过当前版本的事件。
  • aggregate_id 分区,让同一聚合键只由一个串行发布通道处理。
  • 不依赖传输顺序,把消费者设计成基于状态版本或业务不变量收敛。

这些方案解决的是业务顺序,不会把传输语义变成 exactly-once。

Inbox / dedup 的职责:把重复投递收敛成一次业务结果

消息一旦进入 RabbitMQ,就要接受这样一个现实:无论因为发布重发、消费重投递,还是网络抖动,库存服务都有可能多次收到同一个 eventId

因此消费端也必须有自己的边界保护。上一篇已经给出了 Inbox 表、ON CONFLICT DO NOTHING 和同一数据库事务中的完整实现,参见消费可靠性(十)。这里不再重复代码,只保留它在端到端链路中的职责:先用 (consumer_name, eventId) 抢占处理权,在同一事务里提交去重记录与库存预留,事务提交后才 ack

这一步的意义在于:只要 Inbox 唯一键仍然存在,RabbitMQ 可以重复投递,但 inventory_reservations 中这次库存预留只会成功落地一次。若处理过程还包含 HTTP、短信或另一个数据库,稳定 eventId 还要继续作为下游幂等键,或通过新的 Outbox 拆分副作用。

Inbox 的保留期也是可靠性协议的一部分:它至少要覆盖 Broker 最大消息寿命、Outbox 重试周期、人工重放窗口和备份恢复可能带回的旧消息。若清理去重记录后仍允许重放同一个 eventId,effectively-once 的边界也会随之失效。

你也可以把它理解成分工:

  • Outbox 负责“不要漏发业务事件”。
  • Inbox / dedup 负责“重复到达也不要重复生效”。

两边都围绕同一个 eventId 工作,才会真正闭环。

at-least-once 和 effectively-once,不差在字面,差在边界承诺

这两个概念最容易被讲成一句口号,但工程上它们是两种完全不同的承诺:

语义真正含义能否允许重复到达
at-least-once在约定的故障范围与重试机制下,消息会被处理一次或多次
effectively-once在明确的业务结果和去重保留期内,重复到达收敛为同一结果也能

所以 effectively-once 不是 RabbitMQ 提供的独立投递模式,而是应用针对一个明确结果边界构造出来的效果。对本文示例,这个边界只是同一数据库事务中的 Inbox 记录与 inventory_reservations

  1. 生产端允许 Outbox 重发。
  2. Broker 和消费链路允许重复投递。
  3. 消费端用 Inbox / dedup / 唯一约束把最终结果收敛成一次。

如果没有最后这一步,系统最多只能叫 at-least-once。即使有 Inbox,超出同一数据库事务的外部副作用,或者超出 Inbox 保留期的历史重放,也必须单独设计幂等。RabbitMQ 官方同样建议消费者按可能发生 redelivery 来设计,参见 Reliability Guide

一个共享示例,串起来看整条链路

order.created 放回完整流程里,可以得到这样一条闭环:

Outbox 到 Inbox 的端到端可靠性闭环同一个 eventId 从订单本地事务进入 Outbox,经短事务认领和事务外发布;没有 mandatory return 且收到 positive confirm 后才完成 Outbox,最后由库存本地事务写入 Inbox 和库存预留,并在提交后确认消息。ONE EVENT ID · FOUR RELIABILITY BOUNDARIESOrderServicelocal transactionorders + outbox_messages同一事务写入eventId = 01J7X8A6YQ4Q9M6Q4QG7H2J8F1COMMIT订单提交 = 进入可发布集合Outbox publisher先认领,再在事务外发布eventId = 01J7X8A6YQ4Q9M6Q4QG7H2J8F1SKIP LOCKED + Lease短事务认领COMMIT事务外发布mandatory / return没有 returnpositive confirm → Published可能重发:confirm 结果未知 / Lease 到期或状态更新中断;eventId 不变RabbitMQBroker 边界eventId = 01J7X8A6YQ4Q9M6Q4QG7H2J8F1orders.events→ inventory.order-events同一 eventId 传输InventoryServicelocal transactioneventId = 01J7X8A6YQ4Q9M6Q4QG7H2J8F1inbox_messages +inventory_reservationsCOMMITack重复到达由 Inbox / 唯一约束收敛at-least-once 传输effectively-once 业务结果移动端 Outbox 到 Inbox 的可靠性闭环四个边界从上到下排列:订单本地事务、Outbox 发布器、RabbitMQ 与库存本地事务;同一 eventId 贯穿整条链路。ONE EVENT ID · END-TO-END1 · OrderService local transactionorders + outbox_messageseventId = 01J7X8A6YQ4Q9M6Q4QG7H2J8F1COMMIT2 · Outbox publisherSKIP LOCKED + Lease · 短事务认领eventId = 01J7X8A6YQ4Q9M6Q4QG7H2J8F1;COMMIT 后,事务外发布mandatory / return没有 return + positive confirm → Published可能重发:confirm 结果未知 / Lease 到期或状态更新中断;eventId 不变3 · RabbitMQorders.events → inventory.order-eventseventId = 01J7X8A6YQ4Q9M6Q4QG7H2J8F14 · Inventory local transactioneventId = 01J7X8A6YQ4Q9M6Q4QG7H2J8F1inbox_messages + inventory_reservationsCOMMIT → ackInbox 与唯一约束收敛重复at-least-once 传输effectively-once 业务结果
Outbox 到 Inbox 的端到端可靠性闭环Outbox 保留待发布事件,return 与 confirm 共同证明发布结果;结果未知、Lease 到期或状态更新中断都可能重发,重复由 Inbox 和业务唯一约束收敛。
订单服务在同一事务中写入 orders 和 outbox_messages;发布器短事务认领,在事务外用 mandatory / return 与 positive publisher confirm 共同判断结果;库存服务提交 Inbox 与库存预留后才 ack。全过程复用同一个 eventId。

这条链路没有承诺“物理上只出现一条消息”,它描述的是一组有条件的最终交付与结果收敛约束:

  • 订单事务提交后,事件进入本地可发布集合。
  • 在 Outbox Publisher 持续运行、记录未被提前清理、目标拓扑可路由且 RabbitMQ 最终恢复可用的前提下,事件最终进入目标队列。
  • 在 Inbox 唯一键仍处于重放窗口内的前提下,重复到达不会让库存预留重复生效。

这正是生产环境里更值得追求的结果。

本篇解决了什么:

  • 拆清了 order.created 从订单库到 RabbitMQ 再到库存库之间的四层可靠性边界,也说明了 publisher confirms 无法消除双写窗口。
  • 补齐了多实例 Outbox Publisher 的分工:短事务认领、事务外发布、Lease 续租与条件完成共同控制并发,同时明确它不提供全局发布顺序。
  • 明确了 Outbox 负责把事件保留到可发布集合,mandatory return 与 publisher confirm 共同判定发布结果,Inbox / dedup 则在指定结果和保留期内收敛重复。

下一篇:RabbitMQ 生产环境实践:Classic、Quorum、Streams、监控与高可用(十二)