在微服务与事件驱动架构中,最隐蔽的故障往往不是业务逻辑错误,而是「数据库已提交、消息却丢失」的幽灵。事务发件箱(Transactional Outbox)用一个本地事务加一张业务表,就化解了双写不一致这个经典难题。本文从模式原理出发,对比轮询发布与事务日志监听两种实现,并结合 CDC 与幂等消费给出生产级落地细节。
前置基础可先阅读 Spring 集成消息队列 与 分布式事务:Seata、TCC 与 Saga 模式详解。
1. 双写一致性困境
1.1 一个订单引发的数据分裂
业务系统最常见的需求是「落库之后发事件」。订单表写入与消息发送是两个独立系统,任何一次失败都会造成二者不一致:
| 失败场景 | 结果 |
|---|---|
| 先发消息后写库,写库失败 | 下游收到事件但库里没有订单 |
| 先写库后发消息,发送失败 | 订单已存在但下游无感知 |
| 发送超时但实际已送达 | 重发造成重复消息 |
| 库已回滚但消息已发送 | 幽灵事件 |
// 典型但错误的双写
@Transactional
public Order createOrder(OrderRequest req) {
Order order = orderDao.save(req.toOrder()); // ① 写库
eventPublisher.publish(new OrderCreated(order)); // ② 发消息
// 若第 ② 步失败,① 已提交,事务无法回滚消息
return order;
}
1.2 为什么不用本地事务保护消息发送
本地事务只能保护同一资源(同一个数据库连接),而消息队列是外部资源。把 kafkaTemplate.send 放进 @Transactional 事务里,依赖的是消息中间件参与 JTA 分布式事务,成本和复杂度都极高,且大多数场景下根本不值得。
1.3 发件箱模式的核心思想
把「发消息」这个外部副作用,替换为「往本地表插一行」这个本地副作用。二者共享同一个数据库事务,要么都成功,要么都回滚:
业务方法(一个本地事务):
① 更新业务表(订单状态、余额扣减)
② 往 outbox 表插入事件记录
→ 提交
后台发布器:
③ 扫描 outbox 表中未发布的事件
④ 发送到消息队列
⑤ 标记为已发布
2. 发件箱表设计
2.1 表结构与 DDL
发件箱表是模式的载体,字段设计要兼顾发布器查询效率与消息体存储:
CREATE TABLE outbox_event (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
aggregate_id VARCHAR(64) NOT NULL COMMENT '聚合根 ID,如订单号',
aggregate_type VARCHAR(32) NOT NULL COMMENT '聚合类型,如 ORDER',
event_type VARCHAR(64) NOT NULL COMMENT '事件类型,如 ORDER_CREATED',
payload JSON NOT NULL COMMENT '事件体,序列化后的业务数据',
status TINYINT NOT NULL DEFAULT 0 COMMENT '0待发布 1已发布 2失败',
retry_count INT NOT NULL DEFAULT 0,
created_at DATETIME(3) NOT NULL DEFAULT CURRENT_TIMESTAMP(3),
published_at DATETIME(3) DEFAULT NULL,
KEY idx_ready (status, id)
) ENGINE=InnoDB;
2.2 事件写入与本地事务
业务方法与发件箱写入必须在同一个事务内,顺序通常是先更新业务表、后插入事件:
@Service
@Slf4j
public class OrderApplicationService {
private final OrderDao orderDao;
private final OutboxDao outboxDao;
@Transactional
public Order createOrder(OrderRequest req) {
Order order = orderDao.save(req.toOrder());
outboxDao.insert(OutboxEvent.builder()
.aggregateId(order.getId())
.aggregateType("ORDER")
.eventType("ORDER_CREATED")
.payload(JsonUtils.toJson(order))
.build());
return order;
}
}
2.3 索引与写入优化
idx_ready (status, id) 让发布器能用最少的回表找到待发布事件。写入侧要保持轻量:payload 存序列化后的最终状态即可,不要在其中嵌入大对象。对极高吞吐场景,可对 outbox_event 做按事件类型的分区或分表。
3. 发布方式一 轮询发布器
3.1 分批查询与状态机
轮询发布器是发件箱模式最朴素可靠的实现:定时扫描未发布记录,发送成功后更新状态。状态机要覆盖「成功、失败、重试、死信」四种情况。
3.2 发布器实现
@Component
@Slf4j
public class OutboxPoller {
private final OutboxDao outboxDao;
private final KafkaTemplate<String, String> kafkaTemplate;
@Scheduled(fixedDelay = 1000)
public void publishReadyEvents() {
List<OutboxEvent> batch = outboxDao.findReady(100); // status=0 按 id 升序
for (OutboxEvent event : batch) {
publishOne(event);
}
}
private void publishOne(OutboxEvent event) {
try {
kafkaTemplate.send(event.getEventType(), event.getPayload())
.get(5, TimeUnit.SECONDS);
outboxDao.markPublished(event.getId()); // 幂等标记
} catch (Exception e) {
outboxDao.markFailed(event.getId()); // 计失败
log.error("outbox publish failed, id={}", event.getId(), e);
}
}
}
3.3 延迟、批量与有序性
- 延迟:
fixedDelay与轮询批次共同决定延迟上限,秒级即可满足大多数事件驱动场景。 - 吞吐:批次大小可调,
findReady加LIMIT防止单次处理过多。 - 有序性:对同一
aggregate_id的事件,发布器必须按id升序处理,且保证同一聚合的事件不被并发发布器并发发送。
// 保证有序:查询按 (aggregate_id, id) 排序,发布器单实例
// 或用分区锁:同一 aggregate_id 路由到同一发布线程
4. 发布方式二 事务日志监听器
4.1 从日志感知变更
轮询需要额外的定时任务与查询,而事务日志监听器直接消费数据库的变更日志(如 MySQL binlog、PostgreSQL WAL),零侵入地发现「新插入的 outbox 行」,即 CDC 思路。
// 事务监听器方案(借助 Debezium 或 Canal 读取 binlog)
// 当 outbox_event 出现 INSERT 时,实时解析行数据并投递到消息队列
// 优点:毫秒级延迟、无轮询查询开销
// 缺点:需要维护 CDC 组件,架构复杂度上升
4.2 事务边界内的本地监听器
若发布事件本身不强依赖消息队列,也可以借助 Spring 的 TransactionSynchronization,在事务提交后再发送,避免事务未提交就发消息导致下游读到旧数据:
@Transactional
public void createOrderWithEvent(OrderRequest req) {
orderDao.save(req.toOrder());
TransactionSynchronizationManager.registerSynchronization(
new TransactionSynchronization() {
@Override
public void afterCommit() {
eventPublisher.publish(new OrderCreated(req)); // 提交后才发送
}
});
}
注意:afterCommit 只解决「提交后再发」的时序,并不解决「发送失败」的可靠性,生产环境仍应配合发件箱表兜底。
4.3 两种发布方式对比
| 维度 | 轮询发布器 | 事务日志监听器 |
|---|---|---|
| 延迟 | 秒级(取决于轮询间隔) | 毫秒级 |
| 复杂度 | 低,仅需定时任务 | 高,需 CDC 组件 |
| 对业务表侵入 | 无 | 无 |
| 数据库压力 | 定期扫描索引 | binlog 消费,几乎无压力 |
| 适合场景 | 中小流量、运维简单 | 高吞吐、低延迟要求 |
5. 与 CDC 联动
5.1 Debezium 发件箱事件路由器
Debezium 官方提供 outbox-event-router 单消息转换(SMT),专门从 outbox 表读取变更并重组为规范的事件消息,是「事务日志监听 + 发件箱」的黄金组合。
transforms: outbox
transforms.outbox.type: io.debezium.transforms.outbox.EventRouter
transforms.outbox.table.fields.additional.placement: aggregate_id:header:aggregateId
transforms.outbox.table.fields.additional.placement: event_type:header:eventType
transforms.outbox.route.by.field: event_type
transforms.outbox.route.topic.regex: (?<routed>.+)
5.2 从发件箱表到 Kafka 的链路
MySQL binlog → Debezium Connector → outbox SMT 重塑
→ Kafka Topic(按 event_type 路由,带聚合 ID 头)
→ 下游消费者按 header 去重
这组链路把「谁去读发件箱表」完全委托给数据库日志,业务代码不再关心发布细节。
5.3 水位线与健康度
无论哪种发布方式,都要监控发件箱积压。核心指标是「未发布事件的最早时间」与「每分钟新增量」,它们直接反映事件投递的健康程度。
6. 幂等消费与重复消息处理
6.1 消费端幂等是最后防线
发件箱保证「至少一次」投递,重复不可避免:发布后标记失败会重发、下游处理超时会重投。消费端必须自身幂等。
@Component
@Slf4j
public class OrderEventConsumer {
private final IdempotencyDao idemDao;
private final OrderProjectionService projection;
@KafkaListener(topics = "ORDER_CREATED", groupId = "order-projector")
public void onOrderCreated(ConsumerRecord<String, String> record) {
String eventId = record.headers().lastHeader("eventId").value() != null
? new String(record.headers().lastHeader("eventId").value())
: record.key();
if (!idemDao.tryAcquire("ORDER_CREATED", eventId)) {
return; // 已处理过,直接跳过
}
projection.buildProjection(record.value()); // 业务处理
idemDao.commit("ORDER_CREATED", eventId); // 标记完成
}
}
6.2 去重表与唯一键
CREATE TABLE event_idempotency (
event_key VARCHAR(128) PRIMARY KEY, -- 聚合类型 + 事件 ID
process_at DATETIME(3) NOT NULL,
status TINYINT NOT NULL DEFAULT 1 -- 1处理中 2完成
) ENGINE=InnoDB;
tryAcquire 用 INSERT ... ON DUPLICATE KEY UPDATE 原子占位,多个并发消费者只有一个能抢到处理权。
6.3 乱序与重试策略
- 乱序:同一聚合的事件必须保证按序消费,最简做法是发送时按
aggregate_id做 key,Kafka 同 key 进同分区即天然有序。 - 重试:消费失败不回滚占位、指数退避重试,超过阈值转人工。
- 业务幂等:消息幂等之外,业务本身也应可重入,例如状态机校验「已支付订单拒绝重复扣款」。
7. 生产实践
7.1 表膨胀与归档
发件箱表只保留「待发布」数据才有价值,已发布行应定期归档:
-- 归档策略:每小时删除已发布超过 24 小时的数据
DELETE FROM outbox_event
WHERE status = 1 AND published_at < NOW() - INTERVAL 1 DAY
LIMIT 5000;
对超大数据量,可直接把已发布行迁入归档表或按时间分区后 DROP PARTITION。
7.2 监控与告警
| 指标 | 告警阈值 | 说明 |
|---|---|---|
| 未发布事件最早时间 | 超过 5 分钟 | 发布器挂掉或队列积压 |
| 未发布事件数量 | 超过 1000 | 大量事件堆积 |
| 发布失败率 | 超过 5% | 消息队列异常 |
| 消费处理延迟 | 超过 1 分钟 | 下游消费缓慢 |
// 指标埋点示例
Gauge.builder("outbox.pending.max.age", pendingAgeSupplier)
.register(registry);
Counter.builder("outbox.published.total").register(registry);
7.3 与分布式事务框架的关系
发件箱本质是最终一致性方案,适用「容忍短暂延迟」的异步链路;而 Seata AT/TCC 解决的是需要强一致的同步场景。两者选型边界清晰:订单主链路需要强一致时选 Seata,旁路事件通知优先发件箱。此外发件箱表还可作为对账依据,为下游重放提供数据源。
8. 总结
| 主题 | 核心要点 |
|---|---|
| 双写困境 | 业务库与消息队列无法被同一本地事务保护 |
| 模式本质 | 把发消息替换为写本地发件箱表,共享事务 |
| 发布方式 | 轮询发布器低延迟保障,事务日志监听毫秒级 |
| CDC 联动 | Debezium outbox SMT 把发布职责交给 binlog |
| 幂等消费 | 去重表 + 唯一键,至少一次投递下保证精确处理 |
| 运维治理 | 归档已发布数据,监控积压与失败率 |
事务发件箱用最小的成本把「业务提交」与「消息投递」绑定在同一个本地事务里,换来的是可靠的最终一致。它不要求消息中间件参与分布式事务,也不引入复杂的协调器,配合 CDC 与幂等消费,就能构建一条吞吐高、延迟可控、可观测的事件发布链路。
延伸阅读
- Spring 集成消息队列 — 幂等消费与死信策略的发件箱下游视角
- 分布式事务:Seata、TCC 与 Saga 模式详解 — 本地消息表与发件箱的演进关系
- Java 缓存策略与 Redis 集成 — 事件驱动下的缓存一致性配合
- 数据库连接池、读写分离与分库分表 — 发件箱表大流量下的存储选型
- 日志框架、MDC 与分布式链路追踪 — 事件投递链路的可观测性基础
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。