Outbox:把数据库事务可靠地接到消息队列

Outbox:把数据库事务可靠地接到消息队列

在电商订单、支付、库存、积分等系统中,经常需要在修改数据库之后通知另一个服务:

  • 创建订单后通知库存服务
  • 支付成功后通知履约服务
  • 扣减库存后通知营销服务
  • 修改用户状态后刷新搜索索引

最直观的代码通常是:

1
2
3
4
5
@Transactional
public void createOrder(OrderCommand command) {
orderRepository.insert(createOrder(command));
messageBroker.publish(new OrderCreated(...));
}

但数据库事务只管理数据库,不能自动管理外部的消息队列。只要数据库提交和消息发送之间存在两个独立的操作,就存在双写不一致的故障窗口。

Outbox 的核心做法是:

在同一个数据库事务中保存业务数据和待发送的事件;事务提交后,再由独立的投递程序把事件发送到 MQ;MQ 确认成功后,才把事件标记为已发送。

它通常实现的是:

1
业务事务可靠提交 + 消息至少一次投递 + 消费端幂等

数据库和 MQ 为什么会双写不一致

假设代码按下面的顺序执行:

1
2
3
保存订单
提交数据库事务
发送 OrderCreated 消息

可能出现以下情况:

  • 数据库提交成功,应用在发送消息前宕机:订单已经存在,但库存服务永远收不到消息
  • MQ 发送失败:订单存在,但下游状态没有推进
  • MQ 实际已经接收消息,生产者因为网络超时认为发送失败:重试后出现重复消息
  • MQ 发送成功,数据库事务随后回滚:下游收到一个不存在的订单事件
  • 数据库提交结果未知:客户端或应用超时,重试可能创建重复订单

即使把发送消息放进带有事务注解的方法中,普通本地事务也不会自动覆盖 MQ。Publisher Confirm 也只能告诉生产者 Broker 是否接收了消息,不能保证数据库事务和消息发送是同一个原子操作。

Outbox 的核心原理

Outbox 会在业务数据库中增加一张事件表。业务表和事件表必须参与同一个本地数据库事务:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
BEGIN;

INSERT INTO orders(order_id, status, total_amount)
VALUES ('O1001', 'PENDING_PAY', 199.00);

INSERT INTO outbox_events(
event_id,
event_type,
aggregate_id,
payload,
status,
attempts,
created_at
)
VALUES (
'E9001',
'OrderCreated',
'O1001',
'{...}',
'PENDING',
0,
CURRENT_TIMESTAMP
);

COMMIT;

提交后,订单和 Outbox 事件只有两种结果:

  • 两者一起提交成功
  • 两者一起回滚

之后由后台投递程序读取 PENDING 事件:

1
2
3
4
5
6
7
8
9
读取 PENDING 事件

发送到 MQ

等待 Publisher Confirm

确认成功

标记为 SENT

所以 Outbox 的顺序是:

1
先保存,后发送;发送确认后再标记。

不能先发送 MQ,再补写 Outbox。否则如果 MQ 已经发送成功,而数据库保存失败,就会产生一个没有业务事实支撑的孤立事件。

Outbox 事件的状态

常见状态如下:

1
2
3
4
5
6
7
8
PENDING
↓ 领取
SENDING
↓ Broker 确认
SENT

SENDING --租约超时--> PENDING
PENDING --确定性错误--> FAILED 或 DEAD

状态的含义需要提前约定:

  • PENDING:事件已经和业务数据一起提交,但还没有完成可靠投递
  • SENDING:某个投递程序暂时领取了事件
  • SENT:Broker 已经按配置确认接收,不代表消费者已经处理
  • FAILED:本次投递失败,可以继续重试或等待人工处理
  • DEAD:超过重试上限,进入死信或异常处理流程

SENT 不应该被理解成“订单已经扣库存”或“消息已经被业务消费”。如果需要知道消费结果,应由消费者更新自己的业务状态,或者发布新的成功/失败事件。

业务数据和 Outbox 的保存

保存成功

业务数据和 Outbox 事件在一个事务中提交成功:

1
2
orders = 成功
outbox_events = PENDING

此时即使 MQ 完全不可用,事件也仍然保存在数据库中,投递程序恢复后可以继续发送。

业务校验失败

例如商品不存在、金额不合法、用户没有权限:

1
2
业务数据不保存
Outbox 事件也不保存

这种情况不应该写一条需要重试的 Outbox 事件。业务拒绝和基础设施故障要区分开。

事务回滚

数据库死锁、连接断开、唯一键冲突或业务主动抛出异常时,业务数据和 Outbox 事件都应回滚。

如果只保存了业务数据,之后再用 afterCommit 回调发送消息,应用可能在回调执行前宕机。afterCommit 只能降低一部分延迟,不能代替可靠持久化。

数据库提交结果未知

应用提交事务后,响应在网络中丢失,应用不知道数据库究竟成功还是失败。客户端重试时不能简单地重新插入一条订单。

应使用请求幂等键或业务唯一键:

  • 同一个请求键只允许创建一笔订单
  • 重试时先查询已有订单和状态
  • 不能仅依赖订单号由客户端随机生成后直接插入

Outbox 保存失败

如果业务数据和 Outbox 在同一个数据库事务中,Outbox 写失败应该让整个事务回滚。不能接受“订单已经创建,但 Outbox 写失败后继续返回成功”。

如果业务库和 Outbox 表不在同一个事务边界内,Outbox 就不能解决这两个存储之间的原子性问题。此时需要使用 CDC、事务消息、事件溯源或重新划分事实来源。

投递程序如何领取事件

投递程序可以定时轮询,也可以通过数据库 CDC 发现新事件。最简单的轮询逻辑是:

1
2
3
4
5
查询 status = PENDING 且 available_at <= 当前时间的事件
领取事件
提交领取状态
发送 MQ
根据结果更新状态

多个投递实例同时运行时,不能让它们无限制地读取同一批事件。常见做法有:

  • 数据库行锁和 SKIP LOCKED
  • 条件更新:只有 PENDING 才能改成 SENDING
  • 租约:记录 lease_until 和 worker_id
  • 分片扫描:按事件 ID、聚合 ID 或时间范围分配任务

网络发送期间不建议一直持有数据库行锁。更常见的方式是先用短事务领取事件并提交,再在事务外发送;发送完后用带条件的更新标记结果。

例如:

1
2
3
4
5
6
UPDATE outbox_events
SET status = 'SENDING',
worker_id = 'worker-a',
lease_until = CURRENT_TIMESTAMP + INTERVAL '1' MINUTE
WHERE event_id = 'E9001'
AND status = 'PENDING';

更新成功的实例才有资格继续发送。

没有可发送事件

这是正常状态,不应该当成错误。投递程序可以等待下一次轮询或由 CDC 触发。

领取成功后投递程序宕机

事件可能停留在 SENDING。租约到期后,其他投递程序应把它重新放回 PENDING:

1
SENDING + lease_until < now → PENDING

不能依赖人工修改数据库,否则少量异常就会变成永久积压。

多个投递程序重复领取

即使使用了锁或条件更新,租约过短、网络抖动、进程暂停等情况仍可能造成同一事件被多个实例处理。因此领取机制只能减少重复发送,不能从根本上消除重复,event_id 幂等仍然需要存在。

消息发送阶段

发送成功并收到 Confirm

流程如下:

1
2
3
4
事件仍为 SENDING
发送 MQ
Broker 返回成功确认
事件更新为 SENT

此时可以认为 MQ 已经按照自己的持久化和副本配置接收了消息,但不能认为消费者已经完成业务处理。

发送立即失败

例如 Topic 不存在、交换机不存在、权限错误、消息格式不合法:

  • 临时性错误:网络抖动、Broker 过载、连接断开,可以重试
  • 确定性错误:权限错误、序列化错误、配置错误,不应该无限重试

确定性错误应记录详细错误,进入 FAILED 或 DEAD,并触发告警。

发送超时

发送超时不等于 Broker 没有收到消息。可能是:

1
2
3
消息没有到达 Broker
消息已经到达 Broker,但 Confirm 返回途中丢失
Broker 已确认,但生产者连接断开

因此超时后不能直接认为发送失败。最安全的处理是使用同一个 event_id 重试,并让消费者幂等。

Broker 已接收,投递程序还没有标记

这是最典型的重复发送窗口:

1
2
3
4
发送成功
进程在更新 SENT 前宕机
事件仍然是 SENDING 或 PENDING
恢复后再次发送

结果是 MQ 中出现两条相同 event_id 的消息。这个结果是可接受的,因为 Outbox 的设计目标是避免丢失,重复则交给消费者幂等处理。

MQ 不可用或持续过载

事件应继续保留在 PENDING 或进入带退避时间的状态:

1
2
3
attempts = attempts + 1
available_at = now + backoff(attempts)
last_error = ...

重试应使用指数退避和随机抖动,避免所有投递程序同时重试形成重试风暴。还需要限制单次扫描数量,防止积压把业务数据库读压力打满。

发送后的标记阶段

Confirm 成功,标记 SENT 成功

这是理想路径:

1
PENDING → SENDING → SENT

可以记录 sent_at、broker_message_id 和最后一次确认信息,方便排查。

Confirm 成功,标记 SENT 失败

例如数据库暂时不可用,状态更新事务失败。下一轮会再次发送同一事件。

这会产生重复,但不能为了避免重复而在没有 Confirm 时提前标记 SENT,否则更容易造成消息丢失。

标记 SENT 成功,但消费者还没有处理

这是正常的。SENT 只描述生产端和 Broker 之间的关系,不描述消费者状态。

消费者如果宕机,消息应留在 MQ 中等待重新投递;如果消费者已经 ACK 但业务状态异常,则需要由消费者自己的幂等记录和补偿机制处理。

旧投递程序覆盖新状态

如果使用租约,旧实例可能在租约过期后恢复,并继续把事件标记为 SENT 或 FAILED,覆盖新实例的结果。

更新状态时应携带 worker_id 或 lease token:

1
2
3
4
5
6
7
UPDATE outbox_events
SET status = 'SENT',
sent_at = CURRENT_TIMESTAMP
WHERE event_id = 'E9001'
AND status = 'SENDING'
AND worker_id = 'worker-a'
AND lease_token = 'token-a';

这样旧实例的过期操作不会覆盖新实例。

是否应该立即删除 SENT 记录

通常不应该立即删除。Outbox 记录可以用于:

  • 审计某个业务事件是否生成
  • 排查消息发送延迟
  • 重放消息
  • 对账和补偿
  • 分析某段时间的故障

可以在保留足够长的时间后归档或删除。删除前要确认消费者不再需要回放,也要避免清理任务和投递程序互相干扰。

消费阶段

消费者的正确顺序通常是:

1
2
3
4
5
6
7
8
9
10
11
12
13
收到消息

校验消息和版本

检查幂等记录

开启本地业务事务

写入消费记录并修改业务数据

提交事务

ACK

消费者的 Inbox 表可以记录:

1
2
3
4
5
6
event_id
consumer_name
aggregate_id
status
processed_at
last_error

event_id 和 consumer_name 通常建立唯一约束。

消费前宕机

消息还没有执行业务逻辑,消费者宕机后重新收到消息。没有产生业务副作用,直接重试即可。

业务事务执行失败

例如数据库死锁、连接超时,可以重新投递。消费者不要在失败时 ACK,否则消息会被认为已经处理完成。

业务事务提交成功,但 ACK 失败

这是消费端最常见的重复窗口:

1
2
3
业务数据已经提交
消费者还没来得及 ACK 就宕机
MQ 再次投递同一消息

消费者必须通过 event_id 或业务唯一键识别重复消息,然后直接返回成功并 ACK,不能再次扣库存、发积分或创建订单。

先写 Inbox,后写业务数据

这种顺序有风险:

1
2
3
4
Inbox 标记已处理
业务数据写入失败
消息再次到来时被 Inbox 拦截
业务数据永远不会被处理

如果 Inbox 和业务数据在同一个数据库中,应让两者进入同一个事务:

1
2
3
4
5
BEGIN
写入 Inbox
修改业务数据
COMMIT
ACK

如果两者不在同一个事务边界内,就需要状态机、补偿任务或另一个 Outbox,不能只把 Inbox 当成简单的“已读标记”。

消费者调用外部服务成功,但本地事务失败

例如消费者已经调用支付、短信或发货服务成功,随后本地数据库提交失败。外部调用无法随着本地事务自动回滚。

解决方式包括:

  • 外部接口使用业务幂等键
  • 记录调用状态并可查询
  • 用本地 Outbox 推进下一步
  • 使用 Saga 补偿,例如退款、释放库存
  • 使用定时对账发现并修复异常状态

不要假设“数据库事务回滚后外部服务也会回滚”。

确定性业务失败

例如订单已关闭、库存不足、金额不一致,这类错误重试通常不会改变结果。可以:

  • 推进业务状态到失败
  • 发布失败事件
  • 进入人工处理或补偿流程
  • ACK 原消息,避免无限重试

暂时性技术失败

例如数据库短暂不可用、依赖服务超时,应该进行有限次重试。重试次数、退避间隔和最大延迟都要有上限。

毒消息和死信

如果一条消息无论如何都无法成功,持续重试会阻塞同一队列或分区。超过阈值后应进入 DLQ,并保留:

  • 原始消息
  • event_id
  • 失败原因
  • 重试次数
  • 最后一次异常堆栈
  • 首次失败时间和最后失败时间

DLQ 不是丢弃区,而是需要告警、查询、修复和安全重放的异常队列。

Outbox 表应该保存什么

典型字段如下:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
event_id        全局唯一事件 ID
event_type 事件类型
aggregate_type 聚合类型,例如 Order
aggregate_id 聚合 ID,例如 order_id
aggregate_version 业务聚合版本
payload 事件快照
headers trace_id、租户等元数据
status PENDING、SENDING、SENT、FAILED
attempts 已尝试次数
available_at 下次可投递时间
worker_id 当前领取者
lease_until 领取租约截止时间
lease_token 防止旧实例覆盖状态
last_error 最近一次错误
created_at 创建时间
sent_at 发送成功时间

payload 最好是不可变的事件快照。消费者不应该只拿 order_id,再去读取一个已经变化的订单状态来猜测事件发生时的内容,否则重放时可能得到不同结果。

事件格式还应包含 schema version。以后增加字段时尽量向后兼容,避免旧消费者因为新字段直接失败。

顺序、重复和幂等

Outbox 本身不自动保证顺序。

如果同一个订单依次产生:

1
2
3
OrderCreated
PaymentSucceeded
OrderClosed

多个投递程序、不同事务提交时间和 MQ 分区策略都可能导致事件乱序。

需要顺序时,通常采用:

  • 以 order_id、user_id 或 sku_id 作为消息 Key
  • 相同 Key 路由到同一队列或分区
  • 同一个 Key 同时只允许一个消费任务
  • 每个事件带 aggregate_version
  • 消费者拒绝旧版本,或等待缺失版本
  • 严格有序时,前一条失败要阻塞同一 Key 的后续事件

如果业务只需要最终状态,不需要每一个中间事件,那么发送状态快照或使用版本号覆盖旧事件,可能比强行保证所有事件有序更简单。

幂等也要区分范围:

  • event_id 幂等:同一条逻辑消息重复投递不会重复处理
  • 业务键幂等:同一订单、支付单或库存操作不会重复执行
  • 外部调用幂等:重试外部 API 不会重复扣款或发货

Outbox 解决了什么问题

Outbox 主要解决以下问题:

  • 数据库事务提交和消息发布之间的双写不一致
  • 应用在发送消息前宕机造成的事件丢失
  • MQ 暂时不可用时的可靠积压
  • 发送失败后的可控重试
  • 业务事件的审计、补偿和重放
  • 生产端消息发送与业务数据库之间的清晰事实边界

它把“必须发送的事件”先变成数据库中的事实,再把“发布到 MQ”变成一个可以重试的后台过程。

Outbox 不解决什么问题

Outbox 不是万能的分布式事务方案:

  • 不保证端到端 exactly-once
  • 不保证消息不重复
  • 不自动保证消费者顺序
  • 不自动保证多个服务最终状态一致
  • 不会回滚已经成功的外部 API 调用
  • 不会替代 MQ 的持久化、副本和确认配置
  • 不会自动处理库存、支付和订单之间的业务补偿
  • 不会消除数据库和投递程序本身的运维成本

它提供的是可靠的事件出口,消费者仍然需要 Inbox、幂等、重试、死信和对账。

Publisher Confirm 和 Outbox 的关系

Publisher Confirm 只确认生产者到 Broker 的一段链路:

1
Producer → Broker

Outbox 解决的是:

1
业务数据库事务 → 待发布事件

两者应组合使用:

1
2
3
4
5
6
7
业务数据 + Outbox:同一个数据库事务

投递程序发送 MQ

Publisher Confirm

Outbox 标记 SENT

Publisher Confirm 不能代替 Outbox,因为它无法解决:

  • 数据库提交成功但还没有发送消息
  • MQ 发送成功但数据库事务回滚
  • 应用在数据库提交和发送消息之间宕机

如果 MQ 本身就是系统的事实来源,数据库只是消费者生成的查询副本,才可以考虑不使用 Outbox,而直接依赖 MQ 的持久化、确认、重放和消费者幂等。

Outbox 和其他方案

轮询式 Outbox

投递程序定期扫描数据库,是最容易理解和落地的方式。

优点:

  • 实现简单
  • 依赖普通数据库和 MQ
  • 失败记录容易查询
  • 可以人工重放

缺点:

  • 轮询会产生额外数据库负载
  • 需要处理批量领取和并发投递
  • 高吞吐时数据库扫描可能成为瓶颈
  • 投递延迟受轮询周期影响

中小规模订单系统通常从这种方式开始。

Outbox 加 CDC

业务事务仍然先写 Outbox,但不再由应用轮询,而是由数据库日志或 Change Stream 捕获新记录,再投递到 MQ。

优点:

  • 减少轮询负载
  • 事件发现更及时
  • 更适合高吞吐系统
  • 发布器可以独立扩展

缺点:

  • 引入 CDC 组件和运维复杂度
  • 需要处理位点、断点恢复和重复投递
  • 需要维护数据库变更到业务事件的映射
  • CDC 只能发现数据库变化,不能自动设计出正确的业务事件

当 Outbox 记录量很大、轮询已经明显影响数据库时,Outbox 加 CDC 往往比继续堆轮询线程更好。

MQ 事务消息

Kafka 事务、RocketMQ 事务消息等可以把消息发布和 MQ 内部的事务协调起来。

它们适合:

  • MQ 本身是核心事实来源
  • 业务状态也在消息系统能够协调的事务范围内
  • 希望原子写入多个 Topic 或同时提交消费位点

但外部数据库通常不自动参加这个事务。数据库仍然需要 Outbox、CDC 或其他一致性方案。事务消息不是把任意数据库和任意 MQ 变成 XA 事务。

XA 或两阶段提交

XA 可以尝试让数据库和消息中间件参加一个全局事务。

优点:

  • 理论上的原子提交和回滚
  • 调用方不需要自己处理部分状态

缺点:

  • 持有资源时间长
  • 协调器和参与者故障时容易阻塞
  • 性能和可用性下降
  • 中间件支持和运维复杂
  • 网络分区时很难同时满足一致性和可用性

互联网订单系统通常更倾向于本地事务加 Outbox,再用幂等和补偿完成最终一致,而不是把所有操作放进长时间的 XA 事务。

MQ 作为第一事实来源

也可以先把订单命令或业务事件可靠写入 MQ,再由消费者写数据库。

这种模式适合事件流、日志采集和事件溯源系统,但它会改变系统的读取方式:

  • 数据库变成派生状态
  • 需要能够从消息重放构建数据库
  • 用户查询可能暂时读不到最新状态
  • 事件 schema、版本迁移和回放都变得重要

它不是 Outbox 的简单替代,而是另一种系统架构。

不使用 MQ 的同步处理

如果操作必须立即返回确定结果,且业务量并不需要削峰,直接在一个本地数据库事务中完成可能更简单:

1
2
3
4
5
校验
扣库存
创建订单
提交事务
返回结果

没有必要为了“看起来分布式”而给一个本来可以同步完成的操作增加 MQ、Outbox、重试和补偿。

电商订单中的典型用法

一个常见的订单状态流转是:

1
2
3
4
5
6
7
8
9
10
11
12
PENDING
↓ 库存锁定成功
WAIT_PAY
↓ 支付回调验签成功
PAID

FULFILLING

COMPLETED

WAIT_PAY --超时--> CLOSED
CLOSED --补偿--> 释放库存

订单服务可以在本地事务中保存订单和 OrderCreated Outbox 事件:

1
2
订单 = PENDING
Outbox = OrderCreated / PENDING

库存服务消费 OrderCreated:

  • 以 order_id 做幂等
  • 成功则锁定库存
  • 失败则发布库存不足事件或关闭订单
  • 只有库存锁定成功后,才允许用户进入支付,或者明确告诉用户订单仍在排队确认

支付回调成功后,支付服务或订单服务在同一个本地事务中:

1
2
更新订单为 PAID
保存 OrderPaid Outbox 事件

履约、积分、优惠券等服务再分别消费 OrderPaid。每个服务都有自己的本地事务和幂等边界。

库存、订单、支付跨服务时,不能指望一个 Outbox 把它们变成一个全局事务。需要 Saga 状态机、超时关闭、库存释放、退款和定时对账。

生产环境的注意事项

  • 业务数据和 Outbox 必须在同一个可提交事务中
  • 不要在业务事务提交后只依赖内存回调发送消息
  • event_id 必须稳定且唯一,消费者必须幂等
  • 发送超时按“结果未知”处理,允许重复,不要贸然丢弃
  • Publisher Confirm 成功后才能标记 SENT
  • 发送 MQ 时不要长时间持有数据库行锁
  • 使用租约避免 SENDING 永久卡住
  • 用条件更新防止旧投递程序覆盖新状态
  • 瞬时错误和确定性错误要分开处理
  • 重试应有上限、退避和随机抖动
  • 死信队列必须可查询、可告警、可安全重放
  • ACK 必须在业务事务提交之后
  • Inbox 记录和业务更新尽量放在同一个消费者事务中
  • 外部 API 必须使用幂等键,并准备查询和补偿
  • 事件 payload 要版本化,不要依赖随时变化的数据库记录重建历史
  • 对同一聚合需要顺序时,使用业务 Key 和版本号
  • 监控 Outbox 积压数量、最老事件年龄、发送延迟、失败次数、SENDING 超时、消费者重复率、消费延迟和 DLQ 数量
  • SENT 记录不要立即删除,保留足够时间用于审计、重放和对账
  • 对保存后、发送前、Confirm 后、标记前、消费提交后、ACK 前等故障窗口做故障注入测试

Outbox 最重要的理解可以压缩成一句话:

1
2
把业务事件先作为数据库事实保存,再把发布消息变成可重试的异步动作;
接受至少一次投递,用幂等和补偿处理重复与跨服务失败。