可靠消息与Outbox:双写、投递、消费和重放全过程
可靠消息通常以至少一次语义运行。消费者不能假设同一事件只来一次,完整幂等模型见:MQ消费幂等全过程。
一、学完本页要真正会什么
可靠消息不是“调用send()没报错”。一条业务事件要跨越数据库本地事务、应用进程、Broker存储、消费确认和消费者数据库,每个边界都可能在“动作成功、确认丢失”时产生未知结果和重复。
学完本页,你应该能够:
- 列出数据库和MQ双写的所有错误顺序与崩溃窗口。
- 解释Outbox关闭了哪个窗口、为什么仍然至少一次。
- 设计Outbox状态、租约、owner token、重试和归档字段。
- 说明多个投递器怎样安全抢占,不重复并发处理同一行。
- 处理Broker已收但发送端状态回写失败。
- 比较轮询Outbox、CDC Outbox和MQ事务消息。
- 设计消费者去重记录与业务写入同一本地事务。
- 处理外部HTTP、短信等无法与消费日志同事务的副作用。
- 解释顺序、事件版本、死信和安全重放。
- 按eventId、businessKey、Outbox状态、Broker位置和消费日志排查。
二、先理解At-most-once、At-least-once与Exactly-once
| 语义 | 含义 | 代价/风险 |
|---|---|---|
| At-most-once | 最多处理一次,不重试或先确认 | 失败时可能丢失 |
| At-least-once | 失败或确认未知时重试 | 可能重复,消费者必须幂等 |
| Exactly-once effect | 业务效果恰好一次 | 依赖稳定业务键、原子去重和副作用边界,不是一个Broker开关 |
Kafka生产者幂等、事务或RocketMQ事务消息能增强特定边界,但不会自动覆盖消费者的MySQL、HTTP、短信和第三方支付。端到端“恰好一次效果”最终仍要由业务唯一键和状态机证明。
三、数据库和MQ为什么不能直接双写
3.1 先写数据库再发消息
flowchart TD
A["数据库事务提交订单"] --> B["应用准备发送消息"]
B --> C["进程宕机或网络失败"]
C --> D["订单存在,但事件永久没有发送"]3.2 先发消息再写数据库
flowchart TD
A["消息先进入Broker并可消费"] --> B["消费者开始处理"]
A --> C["生产者数据库事务失败或回滚"]
C --> D["下游处理了一个不存在的业务事实"]3.3 在本地事务里直接send
把mqClient.send()写进@Transactional也不能让普通MQ加入数据库事务:
send成功 → 数据库后来回滚:幽灵消息
数据库准备提交 → send失败:整个长事务回滚且锁持有更久
send响应超时:不知道Broker是否已经收到四、Outbox关闭双写窗口的原理
业务数据和事件意图写到同一个数据库、同一个本地事务:
flowchart TD
A["开始本地事务"] --> B["写订单等业务事实"]
B --> C["写Outbox事件PENDING"]
C --> D{"本地事务结果"}
D -- "提交" --> E["业务事实和待发事件同时存在"]
D -- "回滚" --> F["业务事实和事件同时不存在"]
E --> G["独立投递器随后可靠发送"]Outbox保证:只要订单事务提交,数据库里就有一条可恢复的待发事件。它没有保证投递器只发送一次,因为“Broker接收”和“Outbox标SENT”仍属于两个系统。
五、Outbox表设计
create table t_outbox_event (
event_id varchar(64) not null,
aggregate_type varchar(64) not null,
aggregate_id varchar(128) not null,
aggregate_version bigint not null,
event_type varchar(128) not null,
event_version int not null,
destination varchar(128) not null,
message_key varchar(128) not null,
payload text not null,
headers text null,
status varchar(32) not null,
owner_token varchar(64) null,
lease_until datetime null,
retry_count int not null,
next_retry_at datetime null,
last_error_code varchar(64) null,
last_error_message varchar(512) null,
created_at datetime not null,
sent_at datetime null,
updated_at datetime not null,
primary key (event_id),
key idx_outbox_dispatch (status, next_retry_at, created_at),
key idx_outbox_aggregate (aggregate_type, aggregate_id, aggregate_version)
);字段意义:
| 字段 | 为什么需要 |
|---|---|
event_id | 生产、投递和消费端稳定去重键 |
aggregate_id/version | 同一业务对象排序和旧事件保护 |
event_type/version | 事件契约演进 |
destination/message_key | 路由与分区键 |
status | PENDING、SENDING、SENT、RETRY、DEAD |
owner_token/lease_until | 多投递器抢占与宕机接管 |
retry_count/next_retry_at | 有界退避和可索引扫描 |
last_error | 错误分类和排查 |
payload中不要直接放密码、Token、密钥和不必要的完整敏感信息。事件是长期复制的数据资产,脱敏和最小化应在产生时完成。
六、生产者本地事务Demo
@Transactional(rollbackFor = Exception.class)
public void createOrder(CreateOrderCommand command) {
Order order = Order.create(command.getOrderNo(), command.getUserId());
orderRepository.insert(order);
OutboxEvent event = OutboxEvent.pending(
command.getEventId(),
"Order",
command.getOrderNo(),
order.getVersion(),
"OrderCreated",
1,
JsonUtils.toJson(OrderCreatedEvent.from(order)));
outboxRepository.insert(event);
}关键约束:
eventId在业务操作开始时生成并跨重试稳定。- 订单号和eventId有唯一约束。
- 订单写入和Outbox插入使用同一数据源、同一本地事务。
- 不在提交前把事件变成可消费消息。
七、多投递器怎样安全抢占
多个实例同时扫描PENDING时,不能都读取同一批后再发送。常见方案:
7.1 数据库锁定领取
在数据库支持且版本语义明确时,可以在短事务中使用锁定读和SKIP LOCKED领取不同批次,然后快速更新为SENDING并提交。不同数据库、MySQL版本和隔离级别支持不同,必须验证。
7.2 条件更新领取
update t_outbox_event
set status = 'SENDING',
owner_token = :token,
lease_until = :leaseUntil,
updated_at = now()
where event_id = :eventId
and status in ('PENDING', 'RETRY')
and next_retry_at <= now();影响行数为1才获得发送权。领取事务必须短,不能持数据库行锁等待MQ网络发送。
flowchart TD
A["扫描可发送eventId"] --> B["用唯一ownerToken条件抢占"]
B --> C{"更新是否为1行"}
C -- "是" --> D["提交领取事务后发送MQ"]
C -- "否" --> E["其他投递器已领取,跳过"]八、SENDING投递器宕机怎样接管
投递器领取后可能在发送前或发送后宕机。使用lease_until:
- 正常投递器只更新自己ownerToken对应的记录。
- SENDING租约未到期,不允许其他实例并发接管。
- 租约过期后,新实例先判断Broker/业务是否可查询,再以新ownerToken条件接管。
- 旧实例恢复后,更新状态时ownerToken不匹配,不能覆盖新处理者。
租约只能防止大部分并发发送,不能消除“旧投递器长时间暂停后恢复继续send”的风险。由于消息发送本身允许至少一次,消费者幂等仍是最终防线。
九、Broker已收但Outbox未标SENT为什么必然重复
flowchart TD
A["投递器发送eventId=E1"] --> B["Broker持久化E1"]
B --> C["Broker响应发送成功"]
C --> D["投递器在更新Outbox前宕机"]
D --> E["Outbox仍为SENDING/RETRY"]
E --> F["恢复任务再次发送E1"]
F --> G["消费者必须按E1复用第一次结果"]即使Producer SDK自身有幂等能力,也必须理解其作用域、会话、序列和Broker版本。业务Outbox不能依赖一个不明确的传输层假设来取消消费者幂等。
十、投递状态机和错误分类
PENDING → SENDING → SENT
↓
RETRY → SENDING
↓
DEAD发送结果分类:
| 类型 | 处理 |
|---|---|
| 明确成功 | ownerToken条件更新SENT,记录Broker元数据 |
| 瞬时失败 | RETRY,指数退避和抖动 |
| 响应未知 | RETRY或查询Broker,允许重复发送 |
| 永久配置错误 | DEAD并告警,例如Topic不存在/权限拒绝 |
| Payload不合法 | DEAD,修复数据或转换器后受控重放 |
不能把所有异常无限立即重试。权限错误、Schema错误不会靠第10000次重试自动恢复。
十一、CDC Outbox怎样工作
轮询器不断查询Outbox表;CDC则读取数据库提交日志:
flowchart TD
A["订单和Outbox同一事务提交"] --> B["数据库产生Redo/Binlog/WAL变更"]
B --> C["CDC Connector按提交顺序读取"]
C --> D["只筛选Outbox表事件"]
D --> E["转换为Broker消息"]
E --> F["保存Connector位点"]优点:减少业务表轮询和延迟,按数据库提交日志捕获。仍需处理:
- Connector位点持久化和恢复。
- 位点提交前后崩溃导致重复。
- DDL/Schema演进。
- 大事务和大payload。
- 事件路由和删除/墓碑语义。
- Outbox归档与CDC读取位置关系。
- 下游幂等。
直接CDC业务表会暴露存储Schema并难表达领域事件;CDC Outbox由业务显式写事件,语义通常更稳定。
十二、RocketMQ事务消息关闭哪个窗口
抽象流程:
flowchart TD
A["Producer发送半消息"] --> B["Broker持久化但暂不可消费"]
B --> C["Producer执行本地数据库事务"]
C --> D{"本地事务事实"}
D -- "提交" --> E["Producer通知Broker Commit"]
D -- "回滚" --> F["通知Broker Rollback"]
D -- "通信未知" --> G["Broker回查本地事务"]
G --> H["回查订单/事务流水,不读内存变量"]
H --> E
H --> F事务消息把“本地事务结果与消息是否可见”的协调交给Broker回查机制。它不保证:
- 消费者只收到一次。
- 消费者数据库和消息确认原子。
- 消费者调用外部HTTP只执行一次。
- 所有MQ产品都有相同语义。
本地事务需要保存可回查状态;若回查永远返回UNKNOWN,Broker最终会按产品策略处理,不能把它当永久等待队列。具体次数、间隔和状态以RocketMQ版本配置为准。
十三、消费者至少一次处理全过程
flowchart TD
A["Consumer收到eventId"] --> B["开始本地数据库事务"]
B --> C["原子插入消费日志或占用唯一键"]
C --> D{"是否首次消费"}
D -- "是" --> E["校验事件版本和业务状态"]
E --> F["更新业务表"]
F --> G["消费日志标SUCCESS"]
G --> H["提交本地事务"]
H --> I["再向Broker ACK/提交Offset"]
D -- "否且已SUCCESS" --> J["复用结果并ACK"]业务更新和消费日志在同一数据库时,应放在同一本地事务。若业务提交后ACK丢失,消息会重投;消费日志让第二次跳过副作用并重新ACK。
十四、消费日志表
create table t_message_consume (
consumer_group varchar(128) not null,
event_id varchar(64) not null,
event_type varchar(128) not null,
request_hash varchar(128) not null,
status varchar(32) not null,
result_ref varchar(128) null,
retry_count int not null,
first_received_at datetime not null,
updated_at datetime not null,
primary key (consumer_group, event_id)
);为什么主键包含consumer_group:同一事件需要被库存、积分、通知等不同逻辑各处理一次;不同消费者不能共用一条全局“已消费”记录互相跳过。
为什么保存request_hash:同一个eventId却携带不同payload是协议冲突,不能复用旧成功结果。
十五、消费者先查再处理为什么不安全
错误流程:
exists(eventId) == false
→ 执行业务
→ insert consume_log两个并发投递都可能查到false并重复执行。应先通过唯一键原子占位,或让业务表本身以eventId/businessKey唯一;唯一冲突后读取已有状态。完整原理见分布式幂等。
十六、消费者调用HTTP、短信等外部副作用怎么办
消费日志和本地业务表可以同事务,但外部HTTP无法加入普通本地事务:
外部短信发送成功
→ 应用在记录SUCCESS前宕机
→ 消息重投
→ 短信可能再次发送可选策略:
- 下游支持幂等键:使用eventId调用并查询结果。
- 再写本地Command Outbox,由专门投递器调用外部系统。
- 对天然不可幂等通知,业务上接受至少一次并去重模板/频率。
- 高风险资金接口使用稳定业务流水和事实查询。
不能在调用外部接口前就把消费日志标SUCCESS,否则调用失败后Broker不再重投,副作用永久丢失。
十七、事件顺序和版本保护
同一订单可能产生:
OrderCreated version=1
OrderPaid version=2
OrderCanceled version=3消息网络和重试可能让旧事件晚到。消费者应:
- 使用aggregateId作为同一分区/队列路由键,尽量保持局部顺序。
- 在派生表保存last_applied_version。
- 只接受期望的下一版本,或按业务允许跳跃但拒绝旧版本。
- 缺版本时进入等待/重查事实源,而不是无条件覆盖。
update t_order_search_projection
set status = :status,
last_applied_version = :version
where order_no = :orderNo
and last_applied_version < :version;仅靠Broker顺序不够:多生产者、失败重试、分区扩容和重放都可能影响观察顺序。
十八、事件契约与Schema演进
事件是已经发生的事实,发布后可能被长期存储和重放。演进原则:
- 事件包含
eventType和eventVersion。 - 新增字段提供默认语义。
- 不随意改变已有字段含义或类型。
- 消费者支持新旧版本并在完成迁移后再淘汰旧版。
- Payload过大时存引用要考虑源数据变化和保留期。
- 敏感字段变更要考虑历史消息和归档中的删除/脱敏要求。
十九、重试、死信和停车场队列
错误分类:
| 类型 | 示例 | 处理 |
|---|---|---|
| 瞬时技术错误 | 数据库短暂不可用、网络抖动 | 指数退避重试 |
| 业务等待条件 | 依赖事件尚未到达 | 延迟重试或查事实源 |
| 永久数据错误 | 缺字段、无法反序列化 | 死信并修复转换/数据 |
| 永久业务拒绝 | 状态不允许 | 记录异常,不无限重试 |
| 结果未知 | 外部调用响应丢失 | 使用原幂等键查询事实 |
死信不是垃圾桶。必须有:
- 告警和负责人。
- 原Topic、分区/队列、Offset/消息标识。
- eventId、businessKey、错误分类。
- Payload安全存储和保留期。
- 修复后的受控重放工具。
- 重放审计和速率限制。
二十、安全重放流程
flowchart TD
A["选择死信/时间范围/业务键"] --> B["在隔离环境验证修复后的消费者"]
B --> C["确认消费者幂等和版本保护"]
C --> D["小批灰度重放并限速"]
D --> E["监控成功、重复、下游负载和新死信"]
E --> F["扩大批次或停止回滚"]
F --> G["记录操作者、原因、范围和结果"]不能把死信直接批量重新发回原Topic且不限制速度,否则可能再次打垮仍未恢复的下游。重放使用原eventId,不能生成新ID绕过幂等。
二十一、Outbox与消费日志怎样归档
永久保留所有记录会导致表膨胀和索引变慢,过早删除又失去去重和审计依据。
保留期至少考虑:
- Broker最大重试/保留/可回放窗口。
- 客户端和人工可能重放的最长时间。
- 业务对账和审计要求。
- 法规与敏感数据删除要求。
只归档已确认终态且超过安全窗口的记录。进行中的SENDING、RETRY、PROCESSING不能按创建时间粗暴删除。大表可按时间分区、冷热归档,但索引和扫描策略要随之调整。
二十二、Outbox、CDC与事务消息对比
| 方案 | 事实记录位置 | 投递机制 | 优点 | 主要代价 |
|---|---|---|---|---|
| 轮询Outbox | 业务数据库 | 定时扫描 | 简单、可见、易补偿 | 轮询延迟、表压力、抢占 |
| CDC Outbox | 业务数据库+日志 | Connector读取提交日志 | 低轮询压力、提交顺序 | CDC基础设施、位点和Schema治理 |
| MQ事务消息 | Broker半消息+本地事实 | 提交/回滚和回查 | 发送链路自然 | 依赖产品语义、回查和本地状态 |
afterCommit回调 | 仅应用内存 | 提交后立即send | 简单 | 提交后宕机仍会丢,不可靠 |
afterCommit适合非关键提示或作为低延迟尝试,但不能替代持久Outbox。可以在提交后立即尝试唤醒投递器,可靠性仍由数据库Outbox兜底。
二十三、商业场景:支付成功传播
flowchart TD
A["支付回调验签、金额和状态校验"] --> B["本地事务条件更新支付/订单"]
B --> C["同事务写PaymentSucceeded Outbox"]
C --> D["投递器或CDC发布事件"]
D --> E["积分消费者幂等入账"]
D --> F["通知消费者幂等/限频发送"]
D --> G["ES消费者按订单版本更新"]
E --> H["对账验证所有派生结果"]
F --> H
G --> H支付回调主事务不应同步等待积分、短信、ES。每个消费者使用自己的consumerGroup + eventId去重,失败互不阻塞,并可独立重试和对账。
二十四、生产排查Runbook
24.1 Outbox PENDING持续堆积
- 看待发数量、最老年龄和新增速率。
- 检查扫描索引是否命中、查询是否锁等待。
- 检查投递器实例、线程池、连接池和租约。
- 看Broker发送TPS、延迟、权限和Topic路由。
- 按错误码区分瞬时失败和永久配置错误。
- 扩容前确认数据库扫描和Broker容量,避免投递器争抢更严重。
24.2 大量SENDING租约过期
检查投递器GC/宕机、发送超时、lease是否小于正常P99、ownerToken条件更新以及旧实例是否恢复后继续发送。允许接管后重复,但消费者必须幂等。
24.3 消息已在Broker但Outbox仍RETRY
这是发送确认未知窗口。按同一eventId再次发送或使用产品可用的查询能力,不要手工直接标SENT;先确认Broker和业务事实。消费者去重会吸收重复。
24.4 消费者重复修改业务
检查是否先查后写、消费日志和业务是否同事务、eventId是否稳定、不同payload是否复用ID、ACK是否在本地事务提交前、外部副作用是否支持幂等。
24.5 某个聚合事件乱序
检查messageKey/分区、生产者数量、aggregateVersion、重试队列是否改变路由、重放方式和消费者版本条件更新。根据数据库事实源补缺失版本或重建投影,不要无条件按到达顺序覆盖。
24.6 死信重放后再次堆积
确认根因是否真正修复、重放速率是否超过下游容量、旧payload是否能被新代码解析、依赖数据是否已准备。停止全量重放,回到小批灰度。
二十五、监控指标
| 指标 | 价值 |
|---|---|
| outbox_pending_count/oldest_age | 生产传播延迟 |
| outbox_claim_latency | 数据库扫描和争抢 |
| send_success/retry/dead | 发送健康度 |
| send_unknown_total | Broker确认不确定窗口 |
| lease_expired_total | 投递器宕机或超时配置 |
| consumer_lag/oldest_message_age | 消费能力和用户影响 |
| dedup_hit_total | 重复投递实际规模 |
| consume_retry/dead | 消费错误健康度 |
| version_gap/out_of_order | 顺序和版本问题 |
| replay_rate/result | 人工恢复风险 |
指标必须按destination、eventType、consumerGroup、错误分类和必要业务域拆分,但避免把高基数eventId直接作为Metrics标签。
二十六、JDK 8 Demo:重复投递与消费结果复用
import java.util.LinkedHashMap;
import java.util.Map;
public class AtLeastOnceMessageDemo {
static final class Consumer {
private final Map<String, String> results =
new LinkedHashMap<String, String>();
private int businessExecutions;
String consume(String eventId) {
String old = results.get(eventId);
if (old != null) {
return old;
}
businessExecutions++;
String result = "handled-" + businessExecutions;
results.put(eventId, result);
return result;
}
}
public static void main(String[] args) {
Consumer consumer = new Consumer();
String first = consumer.consume("EVENT-9001");
String duplicate = consumer.consume("EVENT-9001");
System.out.println("first=" + first);
System.out.println("duplicate=" + duplicate);
System.out.println("businessExecutions=" + consumer.businessExecutions);
}
}输出:
first=handled-1
duplicate=handled-1
businessExecutions=1真实消费者用数据库唯一键和本地事务,而不是内存Map。Demo只证明重复到达和重复业务效果是两件不同的事。
二十七、常见错误与后果
| 错误 | 后果 | 正确方向 |
|---|---|---|
| DB提交后直接send无Outbox | 进程宕机造成永久漏消息 | 同事务持久化事件意图 |
| Outbox标SENT后再发 | 标记后宕机造成漏发 | 先发再标,允许重复并幂等消费 |
| 多投递器先查后发 | 同一事件被并发大量发送 | ownerToken条件抢占和租约 |
| 领取事务持锁等待MQ | 数据库锁和连接长期占用 | 短事务领取,提交后发送 |
| 消费者先查日志再执行业务 | 并发重复执行 | 唯一键原子占位,同事务写业务 |
| ACK在业务提交前 | 宕机后消息不再投递但业务没完成 | 业务事务提交后ACK |
| 外部HTTP成功前先标消费成功 | 外部失败时永久漏执行 | 幂等下游或Command Outbox |
| 重放生成新eventId | 绕过去重导致重复副作用 | 保留原eventId并记录replay元数据 |
| 死信无限自动回原Topic | 毒消息循环和下游雪崩 | 修复后小批限速重放 |
| 只靠Broker顺序不带版本 | 旧事件晚到覆盖新状态 | aggregateVersion条件更新 |
afterCommit当可靠消息 | 提交后进程宕机仍丢 | 持久Outbox或事务消息 |
二十八、面试标准回答
28.1 Outbox解决什么
Outbox把业务数据和待发事件写入同一个数据库本地事务,保证业务提交时一定留下可恢复事件意图,关闭“数据库成功但消息没记录”的双写窗口。独立投递器或CDC随后发送,发送确认未知时允许重复,因此消费者仍要幂等。
28.2 为什么Outbox仍然会重复消息
Broker已持久化消息后,投递器可能在把Outbox标SENT前宕机;恢复扫描看到未SENT会再次发送。同理ACK丢失也会让Broker重投。至少一次系统选择重复而不是静默丢失,业务通过稳定eventId和原子消费日志吸收重复。
28.3 消费日志怎样和业务保持一致
若消费日志和业务表同库,消费者在一个本地事务中先用consumerGroup + eventId唯一键原子占位,校验状态后更新业务并标记成功,提交后再ACK。业务提交后ACK丢失时,重复消息命中成功记录并重新ACK。
28.4 CDC Outbox和轮询Outbox怎么选
轮询实现直观、可控,但有扫描延迟、索引和多实例抢占成本;CDC读取数据库提交日志,延迟低且减少轮询,但需要Connector位点、Schema演进和运维能力。两者都不能取消重复和消费者幂等。
28.5 事务消息能保证消费者只执行一次吗
不能。事务消息主要协调生产者本地事务与消息是否可见;消费者仍可能因ACK丢失、重平衡和恢复而重复收到。消费者的数据库或外部副作用仍需稳定业务键、幂等、状态机和对账。
二十九、关联知识点
本章小结
可靠消息的本质是把每个不可靠确认窗口变成可恢复状态:生产端用Outbox或事务消息保存事件意图,投递端用租约和重试发送,消费端用唯一键和本地事务吸收重复,死信和对账负责长期失败。至少一次不是缺陷,而是网络不确定下“宁可重复、不能静默丢失”的工程选择。
