Skip to content

事件驱动微服务: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 消息。业务方只需要关注消费端逻辑。

代价

  1. 顺序:Debezium 保证单表变更顺序,但不能保证跨表顺序。如果订单表和支付表之间有因果关系,消费端需要自己排序或用状态机判断
  2. Schema 变更:业务表加字段,Debezium 会自动感知(因为 schema history topic 记录了变更),但下游消费端可能还没准备好处理新字段——你需要在消费端做 schema 兼容处理
  3. DDL 事件ALTER TABLE 本身也会被 Debezium 捕获,并且会阻断后续 binlog 事件的处理,直到 DDL 事件被消费。如果 DBA 在凌晨跑了一个 ALTER TABLE 加索引,而消费端没有对应的 DDL 处理逻辑,整个事件流会卡住

Outbox vs CDC 选型

维度OutboxCDC(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 SchemaProtobuf 做字段校验,而不是硬编码解析

CloudEvents 规范:CNCF 的 CloudEvents 定义了事件的标准格式,强制要求 sourcetypespecversionidtime 等元数据字段。用了这个规范,消费方至少知道「这是个什么事件」、「谁发的」、「什么时候发的」。

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
  • 美团技术博客:《事件驱动架构在美团履约系统的实践》

手撕 → 框架 → 生产化,一步步把 AI Agent 工程化搞透。
粤ICP备2026104257号-1