主题
事件驱动微服务:Outbox 模式、CDC 数据捕获与事件编排
问题
下单成功后要发短信、扣库存、送积分、通知物流——同步调用链超过 3 跳,随便哪个下游慢半秒,整个接口就跟着崩。你加 MQ 收发了,但发现「消息发了但 DB 没更新」或者「DB 更新了但消息没发出去」,消息和 DB 永远不在一个事务里。
这是事件驱动架构最核心的矛盾:微服务之间通过事件解耦,但本地事务和消息发送不是原子操作。解决不了这个,事件驱动就是个笑话。
分析
同步调用链超过 3 跳就该考虑事件驱动
先画一条线:什么时候该从同步切成事件驱动?
- 1-2 跳(A → B):同步足够。B 挂了顶多降级,链路短好排查
- 3 跳以上(A → B → C → D):每一跳都增加 1 个故障点 + 1 段超时配置。一旦 C 超时,A 的线程池可能被占满
- 超过 5 跳:几乎不可能用同步做了。下单流程涉及 7-8 个服务,每个服务 50ms → 接口响应 400ms+,P99 更离谱
事件驱动本质上是把「同步依赖」变成「异步依赖」。A 只需要保证「我做的事已经被记录」,至于下游收到消息后做什么,A 不关心。代价是——你失去了「实时一致性」,必须接受「最终一致」。
Outbox 模式:业务表和消息表同事务写入
在这个问题:业务表更新了,但 MQ 消息因为网络抖动没发出去,下游永远不知道这条数据变了。反过来,消息发了但业务事务回滚了,下游收到一条脏数据。
Outbox 模式的解法很直接:业务操作和消息插入在同一个本地事务里。
java
@Service
@Transactional
public class OrderService {
private final OrderRepository orderRepo;
private final OutboxRepository outboxRepo;
public Order createOrder(CreateOrderRequest request) {
// 1. 业务操作:写入订单表
Order order = new Order(request.getUserId(), request.getAmount());
orderRepo.save(order);
// 2. 同一个事务写入 outbox 表
OutboxMessage message = OutboxMessage.create(
"order.created",
order.getId(),
objectMapper.writeValueAsString(order)
);
outboxRepo.save(message);
return order;
}
}关键设计:OrderRepository.save() 和 outboxRepo.save() 在同一个 @Transactional 里。要么都提交,要么都回滚。DB 不会出现「订单在了但消息没在」或者「消息在了但订单不在」的中间态。
但这就够了吗? 不够——outbox 消息只是躺在了表里,还没发出去。需要一个中继器把消息从 outbox 表投递到 MQ。
java
@Component
public class OutboxRelay {
private final OutboxRepository outboxRepo;
private final KafkaTemplate<String, String> kafka;
@Scheduled(fixedDelay = 1000)
@Transactional
public void relay() {
// 每次取一批未发送的消息,按创建时间排序
List<OutboxMessage> messages = outboxRepo.findTop100ByStatusOrderByCreatedAt(
MessageStatus.PENDING
);
for (OutboxMessage msg : messages) {
try {
// 投递到 MQ
kafka.send(msg.getTopic(), msg.getPayload()).get(3, TimeUnit.SECONDS);
// 更新状态
msg.markSent();
outboxRepo.save(msg);
} catch (Exception e) {
log.warn("Outbox 投递失败,下次重试: msgId={}", msg.getId());
// 不更新状态,下次循环继续尝试
msg.incrementRetryCount();
outboxRepo.save(msg);
}
}
}
}和本地消息表的关系:本地消息表本质上是 Outbox 模式的一个变体。区别在于,Outbox 的消息通常是领域事件("订单已创建"),本地消息表的消息通常是事务消息("请扣库存")。实现上没区别,都是同一个表定时扫、定时发。
坑:中继器投递了但 MQ 说没收到。MQ 返回 ACK 前网络断了,中继器认为消息没发出去,下轮重试又发一次。消费端必须幂等。消费者的去重键用消息 ID 或业务唯一键。
CDC 路线:Debezium 监听 binlog,业务零侵入
Outbox 方案需要业务代码里多写一行 outboxRepo.save()。如果业务方不配合,或者历史代码改不动,怎么办?
CDC(Change Data Capture)方案:不需要业务写任何消息代码。直接监听 DB 的 binlog,把业务表的变更转换成事件。
yaml
# Debezium MySQL Connector 配置
{
"name": "order-connector",
"config": {
"connector.class": "io.debezium.connector.mysql.MySqlConnector",
"database.hostname": "mysql-master.internal",
"database.port": "3306",
"database.user": "debezium",
"database.password": "***",
"database.server.id": "184054",
"database.server.name": "orders-db",
"database.include.list": "order_db",
"table.include.list": "order_db.orders",
"database.history.kafka.bootstrap.servers": "kafka:9092",
"database.history.kafka.topic": "schema-changes.order_db",
"transforms": "unwrap",
"transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
"transforms.unwrap.drop.tombstones": "false",
"snapshot.mode": "initial"
}
}Debezium 作为一个 Kafka Connect,监听 MySQL binlog 的 ROW 格式变更,把每一条 INSERT / UPDATE / DELETE 转换成 Kafka 消息。业务方只需要关注消费端逻辑。
代价:
- 顺序:Debezium 保证单表变更顺序,但不能保证跨表顺序。如果订单表和支付表之间有因果关系,消费端需要自己排序或用状态机判断
- Schema 变更:业务表加字段,Debezium 会自动感知(因为 schema history topic 记录了变更),但下游消费端可能还没准备好处理新字段——你需要在消费端做 schema 兼容处理
- DDL 事件:
ALTER TABLE本身也会被 Debezium 捕获,并且会阻断后续 binlog 事件的处理,直到 DDL 事件被消费。如果 DBA 在凌晨跑了一个ALTER TABLE加索引,而消费端没有对应的 DDL 处理逻辑,整个事件流会卡住
Outbox vs CDC 选型:
| 维度 | Outbox | CDC(Debezium) |
|---|---|---|
| 业务侵入 | 需要写 outbox 代码 | 零侵入,只配 connector |
| 事件语义 | 领域事件("订单已创建") | 数据变更("orders 表 INSERT") |
| 消息粒度 | 业务语义明确 | 裸数据行,需要做转换 |
| 历史数据 | 只记录 outbox 之后的数据 | 可以全量快照(snapshot) |
| 延迟 | 1 秒轮询 | 准实时(binlog 推送) |
| 运维复杂度 | 低,一个定时任务 | 中,需要维护 Kafka Connect 集群 |
| Schema 变更 | 不影响(消息是业务定义的) | 影响(会触发 schema 变更事件) |
实践建议:新系统用 Outbox 模式,业务语义更清晰;老系统改不动、或者需要全量数据同步的场景用 CDC。两者不冲突——可以 Outbox 发业务事件,CDC 做数据同步。
事件 schema 治理:版本兼容与 CloudEvents
随着事件数量增长,最头疼的问题不是事件的投递,而是事件的格式。订单服务发 {"id": "123", "status": "CREATED"},消费方解析 status 字段。下个版本订单服务改成了 {"orderId": "123", "state": "CREATED"},消费方全崩了。
版本兼容策略:
- 增字段不删字段(向后兼容)
- 不改已有字段的类型和语义
- 新版本加字段时给默认值
- 消费方用
JSON Schema或Protobuf做字段校验,而不是硬编码解析
CloudEvents 规范:CNCF 的 CloudEvents 定义了事件的标准格式,强制要求 source、type、specversion、id、time 等元数据字段。用了这个规范,消费方至少知道「这是个什么事件」、「谁发的」、「什么时候发的」。
json
{
"specversion": "1.0",
"type": "com.order.created",
"source": "/order-service/orders",
"id": "f3b1c4a2-7e8d-4c5a-9b6c-1d2e3f4a5b6c",
"time": "2026-09-26T10:00:00Z",
"datacontenttype": "application/json",
"data": {
"orderId": "ORD-20260926-001",
"userId": "user_abc123",
"amount": 99.00,
"items": [
{"sku": "SKU-001", "name": "商品A", "quantity": 1}
]
}
}消费者的 schema 兼容检查:消费者启动时,可以拉取事件 schema 注册表(如 Schema Registry),比对当前版本是否兼容。不兼容就直接报错,免得运行到一半才发现解析失败。
事件编排 vs 事件编舞:可观测性陷阱
这两个概念面试常问,但生产上更值钱的是它们的可观测性代价。
编舞(Choreography):每个服务独立监听事件、独立处理、独立发新事件。没有中心协调器。看起来松耦合,但排查问题时要看 5 套日志才能拼出一个完整链路。
text
订单服务发 "order.created"
→ 库存服务收到,扣库存,发 "inventory.deducted"
→ 物流服务收到,创建运单,发 "delivery.created"
→ 通知服务收到,发短信
某天库存扣了但运单没创建,查起来:
1. 订单日志:消息已发 ✓
2. 库存日志:消息已收 ✓,处理成功 ✓
3. 物流日志:没收到消息!编排(Orchestration):引入一个中心协调器,定义流程步骤,按顺序调用各服务。可观测性比编舞好,但引入了单点(协调器挂了整个流程卡住)。
跨异步边界的链路追踪:不管用哪种方式,事件驱动的追踪都比同步调用难——TraceId 怎么跨消息传递?
java
// 消费者:从消息头恢复 TraceId
@Component
public class OrderCreatedConsumer {
@KafkaListener(topics = "order.created")
public void onMessage(ConsumerRecord<String, String> record) {
// 从消息头拿 TraceId
String traceId = new String(record.headers()
.lastHeader("x-trace-id")
.value(), StandardCharsets.UTF_8);
// 设置到当前线程的 MDC
MDC.put("traceId", traceId);
try {
// 处理业务
inventoryService.deductStock(record.value());
} finally {
MDC.clear();
}
}
}生产者:发消息时把当前 TraceId 塞进消息头。
java
// 生产者:发送事件时附带 TraceId
@Service
public class OrderEventPublisher {
public void sendOrderCreated(Order order) {
String payLoad = objectMapper.writeValueAsString(order);
ProducerRecord<String, String> record = new ProducerRecord<>(
"order.created", order.getId(), payLoad
);
// 从当前 MDC 获取 TraceId
String traceId = MDC.get("traceId");
if (traceId != null) {
record.headers().add("x-trace-id", traceId.getBytes(StandardCharsets.UTF_8));
}
kafka.send(record);
}
}实践建议:3-5 个服务、流程简单、不太需要追踪 → 编舞。5-10 个服务、有合规或审计要求 → 编排。不管选哪个,TraceId 跨消息传递是必须实现的,否则事件驱动就是黑盒。
消费端幂等与事件重放
Kafka 在 exactly-once 语义下也会出现重复消息(生产者重试、消费者 rebalance 等)。消费端必须幂等。
java
@Component
public class IdempotentConsumer {
private final ProcessedEventRepo processedEventRepo;
private final InventoryService inventoryService;
@Transactional
public void handleOrderCreated(String eventId, OrderCreatedEvent event) {
// 幂等检查:去重表
if (processedEventRepo.existsById(eventId)) {
log.info("事件已处理,跳过: eventId={}", eventId);
return;
}
// 执行业务逻辑
inventoryService.deductStock(event.getSkuId(), event.getQuantity());
// 记录处理记录
processedEventRepo.save(new ProcessedEvent(eventId, LocalDateTime.now()));
}
}去重表设计:processed_events(event_id, processed_at),event_id 是唯一键。消费成功后再写入,消费失败不回滚(避免重复执行已成功的业务操作)。
事件重放:这是事件驱动架构的一个隐藏价值——出错时可以从 Kafka 的某个 offset 重新消费,弥补数据不一致。但需要确保重放不会副作用翻倍:
- 幂等检查到位 → 重放安全
- 幂等检查不到位 → 重放就是灾难
与 11/13 篇的关系
11 篇(数据一致性)讲了本地消息表和事务消息——那是「事务怎么保证」的层面。13 篇(Saga 模式)讲了补偿与隔离——那是「分布式事务怎么做」的层面。
本篇补的是「事件管道怎么建」——消息怎么从业务系统里出来、怎么保证不丢、怎么格式统一、跨异步边界怎么追踪。三篇合起来才构成事件驱动微服务的完整图景。
总结
事件驱动架构的核心不是「用 MQ 解耦」,而是让消息的可靠性和业务事务的可靠性一致。做不到这一点,事件驱动就是「消息丢了数据对不上」的灾难现场。
| 层面 | 关键设计 |
|---|---|
| 消息可靠 | Outbox 模式(业务表+消息表同事务),或 CDC(Debezium 监听 binlog) |
| 消息格式 | CloudEvents 规范 + Schema 版本兼容,增字段不删字段 |
| 消息追踪 | 生产方把 TraceId 塞消息头,消费方恢复 MDC |
| 消费幂等 | 去重表(event_id 唯一键),保证重放安全 |
| 选型判断 | 3 跳以上考虑事件驱动;新系统 Outbox,老系统 CDC;简单场景编舞,复杂场景编排 |
参考
- 微服务模块 11 篇:数据一致性(本地消息表与事务消息)
- 微服务模块 13 篇:Saga 模式(补偿与隔离性)
- 微服务模块 04 篇:服务通信与序列化选型
- Debezium 官方文档:《MySQL Connector 配置参考》
- CNCF CloudEvents 规范 1.0
- Chris Richardson. Microservices Patterns Chapter 5: Event-Driven Architecture
- 美团技术博客:《事件驱动架构在美团履约系统的实践》