Outbox Pattern、CDC 与幂等消费:构建可靠事件驱动系统

在微服务和事件驱动架构中,一个非常常见的业务场景是:

  1. 修改数据库中的业务数据。
  2. 向 Kafka、RocketMQ 等消息中间件发送一条消息。
  3. 下游服务消费消息,并执行后续业务逻辑。

例如,用户提交订单后,订单服务需要:

  • 在数据库中创建订单;
  • 发送“订单已创建”事件;
  • 库存服务收到事件后扣减库存;
  • 积分服务收到事件后增加积分;
  • 通知服务收到事件后发送通知。

看起来只是一次数据库操作加一次消息发送,但这里隐藏着一个非常关键的问题:

如何保证数据库事务和消息发送的一致性?

如果数据库写入成功,但消息发送失败,下游服务就永远无法感知这次业务变化。

如果消息发送成功,但数据库事务回滚,下游服务却可能处理一条实际上并不存在的业务数据。

为了解决这些问题,工程中经常组合使用以下三种机制:

  • Outbox Pattern:保证业务数据和事件记录的一致性
  • CDC:可靠地捕获数据库中的事件变化
  • 幂等消费:解决消息重复投递和重复处理问题

这三者并不是互相替代的关系,而是分别解决可靠消息链路中的不同问题。

一、问题背景:数据库事务和消息发送为什么难以保持一致

假设订单服务需要创建订单并发送消息:

@Transactional
public void createOrder(CreateOrderCommand command) {
    Order order = new Order();
    order.setOrderNo(command.getOrderNo());
    order.setStatus("CREATED");

    orderRepository.save(order);

    kafkaTemplate.send(
            "order-created-topic",
            order.getOrderNo(),
            new OrderCreatedEvent(order.getId(), order.getOrderNo())
    );
}

这段代码看起来被 @Transactional 包裹,但需要注意:

Spring 数据库事务只能控制数据库操作,通常无法直接控制 Kafka 或 RocketMQ 的消息发送。

数据库事务和消息中间件属于两个独立系统,它们具有各自的事务边界。

因此可能出现以下异常情况。

1. 数据库提交成功,消息发送失败

执行顺序如下:

1. INSERT INTO orders 成功
2. 数据库事务提交成功
3. 发送 Kafka 消息失败

最终结果:

数据库:存在订单
Kafka:没有订单事件

下游库存、积分、通知等服务无法感知订单已经创建。

2. 消息发送成功,数据库事务回滚

执行顺序如下:

1. INSERT INTO orders 暂时成功
2. Kafka 消息发送成功
3. 数据库事务因为异常回滚

最终结果:

数据库:不存在订单
Kafka:存在订单事件

下游服务会基于一条无效消息执行后续逻辑。

3. 应用在临界点崩溃

即使采用“先提交数据库,再发送消息”的方式,也会存在临界点问题:

1. 数据库提交成功
2. 应用进程崩溃
3. 消息发送逻辑没有执行

这种问题并不能通过简单的重试完全解决,因为应用重启后,可能已经不知道哪些业务数据尚未发送消息。

二、为什么 Kafka ACK 不能解决数据库和消息的一致性问题

Kafka ACK 解决的是:

Kafka Producer 发送消息时,消息是否已经被 Kafka Broker 接收并持久化。

例如:

acks=all

表示 Producer 需要等待所有同步副本确认后,才认为发送成功。

ACK 可以提高 Kafka 内部的消息可靠性,但它无法解决数据库事务和 Kafka 事务之间的一致性。

原因是:

数据库事务 ≠ Kafka ACK

即使 Kafka 返回 ACK,也只能说明:

Kafka 已经成功保存消息

它无法说明:

数据库事务一定已经成功提交

同样,数据库事务提交成功,也无法保证 Kafka 一定收到消息。

因此:

Kafka ACK 解决的是消息在 Kafka 内部是否可靠写入的问题,而 Outbox Pattern 解决的是业务数据和事件记录是否原子提交的问题。

三、Outbox Pattern 是什么

Outbox Pattern,通常翻译为:

事务发件箱模式

其核心思想是:

不在业务事务中直接向消息中间件发送消息,而是把需要发送的事件先写入数据库中的 Outbox 表。

业务数据和 Outbox 事件记录放在同一个本地数据库事务中提交。

例如,创建订单时同时执行:

INSERT INTO orders (...);

INSERT INTO outbox_event (...);

由于这两条 SQL 位于同一个数据库事务中,因此它们具备原子性:

要么全部成功
要么全部失败

这样就能保证:

只要业务数据成功提交,就一定存在一条对应的待发送事件。

四、Outbox Pattern 的整体流程

完整流程如下:

客户端
  |
  v
订单服务
  |
  | 本地事务
  |-------------------------------|
  | 1. 写入 orders                |
  | 2. 写入 outbox_event          |
  |-------------------------------|
  |
  v
数据库
  |
  | CDC 或轮询程序读取 Outbox
  v
Kafka / RocketMQ
  |
  v
下游消费者
  |
  | 幂等处理
  v
库存服务 / 积分服务 / 通知服务

可以将这条可靠消息链路拆成三个阶段:

业务数据 -> Outbox 表 -> 消息中间件 -> 消费者

其中:

阶段 主要机制 解决的问题
业务数据写入 Outbox 本地数据库事务 业务数据与事件记录的一致性
Outbox 发送到 MQ CDC 或轮询 事件可靠发布
消费者处理消息 幂等消费 重复投递、重复处理

五、Outbox 表设计

一个典型的 Outbox 表可以设计为:

CREATE TABLE outbox_event (
    id BIGINT PRIMARY KEY,
    event_id VARCHAR(64) NOT NULL,
    aggregate_type VARCHAR(64) NOT NULL,
    aggregate_id VARCHAR(64) NOT NULL,
    event_type VARCHAR(128) NOT NULL,
    topic VARCHAR(128) NOT NULL,
    event_key VARCHAR(128),
    payload JSON NOT NULL,
    status VARCHAR(32) NOT NULL DEFAULT 'PENDING',
    retry_count INT NOT NULL DEFAULT 0,
    next_retry_at DATETIME NULL,
    created_at DATETIME NOT NULL,
    published_at DATETIME NULL,
    updated_at DATETIME NOT NULL,

    UNIQUE KEY uk_event_id (event_id),
    KEY idx_status_retry (
        status,
        next_retry_at,
        created_at
    ),
    KEY idx_aggregate (
        aggregate_type,
        aggregate_id
    )
);

各字段含义如下:

字段 含义
id 数据库主键
event_id 全局唯一事件 ID
aggregate_type 聚合类型,例如 ORDER
aggregate_id 聚合根 ID,例如订单 ID
event_type 事件类型,例如 OrderCreated
topic 目标消息 Topic
event_key Kafka 消息 Key
payload 事件消息体
status 事件发送状态
retry_count 已重试次数
next_retry_at 下一次重试时间
created_at 事件创建时间
published_at 事件发布时间
updated_at 更新时间

1. event_id

event_id 是事件的全局唯一标识,通常使用:

  • UUID;
  • 雪花算法 ID;
  • 数据库分布式 ID;
  • ULID。

例如:

01JY1N5H2J3Q8A5W9M7K2D4P6R

消费者可以通过 event_id 判断事件是否已经处理过。

2. aggregate_type 和 aggregate_id

这两个字段用于标识事件属于哪个业务聚合。

例如:

aggregate_type = ORDER
aggregate_id   = 100001

表示这条事件属于订单 100001

它们可用于:

  • 排查同一个业务对象产生的所有事件;
  • 维护同一聚合内的消息顺序;
  • 构造 Kafka 消息 Key;
  • 做事件追踪和审计。

3. event_type

event_type 表示事件的业务类型,例如:

OrderCreated
OrderPaid
OrderCancelled
OrderCompleted

建议使用明确的过去式命名,因为领域事件通常表示:

某件事情已经发生。

推荐:

OrderCreated
PaymentSucceeded
InventoryDeducted

不推荐:

CreateOrder
ProcessPayment
DeductInventory

后者更像命令,而不是事件。

4. payload

payload 保存完整的事件内容,例如:

{
  "eventId": "01JY1N5H2J3Q8A5W9M7K2D4P6R",
  "eventType": "OrderCreated",
  "eventVersion": 1,
  "occurredAt": "2026-07-16T10:30:00Z",
  "data": {
    "orderId": 100001,
    "orderNo": "ORD202607160001",
    "userId": 20001,
    "amount": 199.00
  }
}

建议事件中至少包含:

eventId
eventType
eventVersion
occurredAt
data

必要时还可以增加:

traceId
tenantId
source
correlationId
causationId

六、Spring Boot 中实现 Outbox Pattern

1. 订单实体

public class Order {

    private Long id;

    private String orderNo;

    private Long userId;

    private BigDecimal amount;

    private String status;

    private LocalDateTime createdAt;
}

2. Outbox 事件实体

public class OutboxEvent {

    private Long id;

    private String eventId;

    private String aggregateType;

    private String aggregateId;

    private String eventType;

    private String topic;

    private String eventKey;

    private String payload;

    private String status;

    private Integer retryCount;

    private LocalDateTime nextRetryAt;

    private LocalDateTime createdAt;

    private LocalDateTime publishedAt;

    private LocalDateTime updatedAt;
}

3. 在同一个事务中写入业务数据和 Outbox

@Service
@RequiredArgsConstructor
public class OrderService {

    private final OrderMapper orderMapper;

    private final OutboxEventMapper outboxEventMapper;

    private final ObjectMapper objectMapper;

    @Transactional
    public Long createOrder(CreateOrderCommand command) {
        Order order = new Order();
        order.setOrderNo(command.orderNo());
        order.setUserId(command.userId());
        order.setAmount(command.amount());
        order.setStatus("CREATED");
        order.setCreatedAt(LocalDateTime.now());

        orderMapper.insert(order);

        String eventId = UUID.randomUUID().toString();

        OrderCreatedEvent event = new OrderCreatedEvent(
                eventId,
                order.getId(),
                order.getOrderNo(),
                order.getUserId(),
                order.getAmount(),
                Instant.now()
        );

        OutboxEvent outboxEvent = new OutboxEvent();
        outboxEvent.setId(generateId());
        outboxEvent.setEventId(eventId);
        outboxEvent.setAggregateType("ORDER");
        outboxEvent.setAggregateId(order.getId().toString());
        outboxEvent.setEventType("OrderCreated");
        outboxEvent.setTopic("order-created-topic");
        outboxEvent.setEventKey(order.getId().toString());
        outboxEvent.setPayload(serialize(event));
        outboxEvent.setStatus("PENDING");
        outboxEvent.setRetryCount(0);
        outboxEvent.setCreatedAt(LocalDateTime.now());
        outboxEvent.setUpdatedAt(LocalDateTime.now());

        outboxEventMapper.insert(outboxEvent);

        return order.getId();
    }

    private String serialize(Object value) {
        try {
            return objectMapper.writeValueAsString(value);
        } catch (JsonProcessingException exception) {
            throw new IllegalStateException(
                    "Failed to serialize outbox event",
                    exception
            );
        }
    }

    private long generateId() {
        return ThreadLocalRandom.current().nextLong(
                1,
                Long.MAX_VALUE
        );
    }
}

这里最重要的不是具体代码,而是事务边界:

@Transactional
public Long createOrder(...) {
    // 写业务表
    orderMapper.insert(order);

    // 写 Outbox 表
    outboxEventMapper.insert(outboxEvent);
}

只要它们使用的是同一个数据库和同一个事务管理器,就可以保证:

订单成功 -> Outbox 事件一定存在
订单失败 -> Outbox 事件一定不存在

七、Outbox 事件如何发送到消息中间件

Outbox 表中的事件需要被转发到 Kafka、RocketMQ 等消息中间件。

常见方案有两种:

  1. 定时轮询 Outbox 表;
  2. 使用 CDC 监听数据库变更日志。

八、方案一:定时轮询 Outbox 表

轮询方案通常由应用中的定时任务完成:

1. 查询状态为 PENDING 的事件
2. 向消息中间件发送消息
3. 发送成功后更新为 PUBLISHED
4. 发送失败则增加 retry_count
5. 达到最大次数后进入 DEAD 状态

示例:

@Component
@RequiredArgsConstructor
public class OutboxPublisher {

    private final OutboxEventMapper outboxEventMapper;

    private final KafkaTemplate<String, String> kafkaTemplate;

    @Scheduled(fixedDelay = 1000)
    public void publish() {
        List<OutboxEvent> events =
                outboxEventMapper.selectPendingEvents(100);

        for (OutboxEvent event : events) {
            publishEvent(event);
        }
    }

    private void publishEvent(OutboxEvent event) {
        try {
            kafkaTemplate.send(
                    event.getTopic(),
                    event.getEventKey(),
                    event.getPayload()
            ).get();

            outboxEventMapper.markPublished(
                    event.getId(),
                    LocalDateTime.now()
            );
        } catch (Exception exception) {
            outboxEventMapper.markFailed(
                    event.getId(),
                    LocalDateTime.now()
            );
        }
    }
}

查询 SQL:

SELECT *
FROM outbox_event
WHERE status = 'PENDING'
  AND (
      next_retry_at IS NULL
      OR next_retry_at <= NOW()
  )
ORDER BY created_at ASC
LIMIT 100;

1. 轮询方案的优点

  • 实现简单;
  • 不依赖额外的 CDC 组件;
  • 应用代码可完全控制重试机制;
  • 容易理解和排查;
  • 适用于中小规模系统。

2. 轮询方案的缺点

  • 会周期性查询数据库;
  • 实时性取决于轮询周期;
  • 多实例并发时容易重复扫描;
  • 需要自行处理锁、重试、状态和死信;
  • 发送成功与状态更新之间仍然存在临界点。

最后一个问题尤其重要。

假设执行流程如下:

1. Kafka 消息发送成功
2. 应用进程崩溃
3. Outbox 状态没有更新为 PUBLISHED

应用重启后,会再次读取这条 PENDING 记录并重新发送消息。

因此:

即使使用 Outbox,也不能保证消息绝对只发送一次。

Outbox 通常提供的是:

At Least Once
至少一次投递

而不是:

Exactly Once
精确一次投递

九、多实例轮询时如何避免并发处理

如果服务部署了多个实例,多个实例可能同时查询到同一批 Outbox 事件。

例如:

实例 A 查询到事件 1001
实例 B 也查询到事件 1001

最终可能造成重复发送。

常见解决方案包括:

  • 悲观锁;
  • SELECT ... FOR UPDATE SKIP LOCKED
  • 抢占式状态更新;
  • 分片扫描;
  • 分布式锁。

1. 使用 SKIP LOCKED

MySQL 8 可以使用:

SELECT *
FROM outbox_event
WHERE status = 'PENDING'
ORDER BY created_at
LIMIT 100
FOR UPDATE SKIP LOCKED;

处理流程:

@Transactional
public List<OutboxEvent> claimEvents(int limit) {
    List<OutboxEvent> events =
            outboxEventMapper.selectForUpdateSkipLocked(limit);

    for (OutboxEvent event : events) {
        outboxEventMapper.markProcessing(event.getId());
    }

    return events;
}

不同实例会跳过已经被其他事务锁定的记录。

2. 使用抢占式更新

也可以先把事件从 PENDING 更新为 PROCESSING

UPDATE outbox_event
SET status = 'PROCESSING',
    updated_at = NOW()
WHERE id = ?
  AND status = 'PENDING';

只有更新行数为 1 的实例获得处理权。

int affectedRows =
        outboxEventMapper.claimEvent(event.getId());

if (affectedRows == 0) {
    return;
}

这种写法实际上使用数据库状态作为轻量级竞争锁。

3. PROCESSING 状态需要超时恢复

如果实例抢到任务后崩溃,事件可能永久停留在 PROCESSING

因此需要设置处理超时。

例如:

UPDATE outbox_event
SET status = 'PENDING',
    updated_at = NOW()
WHERE status = 'PROCESSING'
  AND updated_at < NOW() - INTERVAL 5 MINUTE;

或者在查询时直接把超时任务视为可重试任务。

十、CDC 是什么

CDC 是 Change Data Capture 的缩写,中文通常称为:

变更数据捕获

CDC 通过读取数据库的事务日志,捕获数据库中的数据变更。

不同数据库有不同的变更日志:

数据库 变更日志
MySQL Binlog
PostgreSQL WAL
SQL Server Transaction Log
Oracle Redo Log
MongoDB Oplog / Change Stream

以 MySQL 为例,当 Outbox 表插入一条记录时:

INSERT INTO outbox_event (...);

这次变更会被写入 MySQL Binlog。

CDC 工具读取 Binlog 后,就可以把这条数据库变更转换成消息,并发送到 Kafka。

常见 CDC 工具

常见的 CDC 工具有:

  • Debezium;
  • Canal;
  • Maxwell;
  • Flink CDC;
  • AWS DMS;
  • Kafka Connect JDBC Source。

在 Outbox 场景中,Debezium 是非常典型的实现方案。

十一、CDC 与 Outbox Pattern 的关系

Outbox Pattern 和 CDC 解决的是两个不同阶段的问题。

Outbox Pattern 解决:

如何保证业务数据和事件记录一起提交。

CDC 解决:

如何把已经提交到 Outbox 表中的事件可靠地发送到消息系统。

它们的关系可以表示为:

业务事务
   |
   | 同一个数据库事务
   v
业务表 + Outbox 表
   |
   | CDC 监听 Binlog
   v
Kafka

因此,CDC 并不是 Outbox Pattern 的替代品。

如果直接监听业务表,也可以产生消息,但会面临一些问题:

  • 很难区分业务变更意图;
  • 数据库字段和消息协议强耦合;
  • 一次业务操作可能涉及多张表;
  • UPDATE 很难表达具体领域事件;
  • 删除操作可能缺少完整业务上下文;
  • 数据修复 SQL 可能误触发业务事件。

例如,订单状态从:

CREATED -> PAID

确实可以推断为 OrderPaid

但如果管理员执行了一次历史数据修复:

UPDATE orders
SET status = 'PAID'
WHERE id = 100001;

CDC 无法天然判断这究竟是:

  • 正常支付;
  • 数据补偿;
  • 人工修复;
  • 测试操作。

因此更推荐:

由业务代码显式写入领域事件到 Outbox 表,CDC 只负责可靠搬运。

十二、基于 Debezium 的 Outbox 架构

典型架构如下:

Spring Boot
   |
   | 本地事务
   v
MySQL
   |
   | Binlog
   v
Debezium Connector
   |
   v
Kafka Connect
   |
   v
Kafka Topic
   |
   v
消费者

业务代码只需要:

写业务表
写 Outbox 表

不需要负责:

查询 PENDING 事件
发送 Kafka
更新发送状态
管理发送线程

Debezium 会读取数据库日志并将 Outbox 记录发送到 Kafka。

CDC 方案的优点

  • 基于数据库事务日志,数据库压力较小;
  • 不需要持续轮询业务表;
  • 实时性通常更高;
  • 不需要应用自行维护发送线程;
  • 能自然捕获所有已提交记录;
  • 适合中大型事件驱动架构。

CDC 方案的缺点

  • 需要部署 Kafka Connect、Debezium 等组件;
  • 运维复杂度更高;
  • 需要管理 Binlog、WAL 等日志;
  • 需要处理 Schema 变更;
  • 需要监控 CDC 位点和延迟;
  • 仍然可能产生重复消息。

需要再次强调:

CDC 也不能天然保证消费者只收到一次消息。

CDC Connector 发生重启、位点提交失败、Kafka 写入确认异常等情况时,都可能重新发送部分消息。

因此消费者仍然必须实现幂等。

十三、为什么消息重复几乎不可避免

在分布式系统中,很多操作存在这样的不确定状态:

操作可能已经成功,但调用方没有收到成功响应。

例如消费者处理 Kafka 消息:

1. 消费者执行数据库事务成功
2. 消费者进程崩溃
3. Kafka Offset 尚未提交
4. 消费者重启后重新消费

消息被再次投递。

再例如 Outbox Publisher:

1. 消息发送到 Kafka 成功
2. Publisher 没有收到 ACK
3. Publisher 认为发送失败
4. Publisher 再次发送

Kafka 中会出现两条相同事件。

因此可靠消息系统通常选择:

至少一次投递 + 幂等消费

而不是追求整个链路上的绝对精确一次。

十四、什么是幂等消费

幂等表示:

对同一个操作执行一次和执行多次,最终产生的业务结果相同。

数学形式可以表示为:

f(f(x)) = f(x)

在消息消费场景中:

同一条消息被消费一次

和:

同一条消息被消费多次

最终业务结果应当一致。

例如,“设置订单状态为已支付”通常具有一定幂等性:

UPDATE orders
SET status = 'PAID'
WHERE id = 100001;

执行一次和执行多次,最终状态都是 PAID

但“账户余额增加 100 元”不是天然幂等的:

UPDATE account
SET balance = balance + 100
WHERE user_id = 20001;

如果重复执行两次,余额会增加 200 元。

十五、幂等消费的常见实现方案

常见方案包括:

  1. 唯一约束;
  2. 消费记录表;
  3. 业务状态机;
  4. 乐观锁;
  5. Redis 去重;
  6. 使用业务唯一键;
  7. 天然幂等 SQL。

其中最可靠的通常是:

数据库唯一约束 + 本地事务。

十六、方案一:消费记录表

创建一张消息消费记录表:

CREATE TABLE consumed_event (
    id BIGINT PRIMARY KEY AUTO_INCREMENT,
    consumer_group VARCHAR(128) NOT NULL,
    event_id VARCHAR(64) NOT NULL,
    event_type VARCHAR(128) NOT NULL,
    consumed_at DATETIME NOT NULL,

    UNIQUE KEY uk_consumer_event (
        consumer_group,
        event_id
    )
);

联合唯一键:

UNIQUE KEY uk_consumer_event (
    consumer_group,
    event_id
)

表示:

同一个消费者组只能处理同一个事件一次。

之所以需要加入 consumer_group,是因为同一条事件可能需要被多个不同业务消费者分别处理。

例如:

OrderCreated

可能同时被以下消费者处理:

inventory-service
points-service
notification-service

这些消费者都应该处理一次,因此不能只对 event_id 建全局唯一约束。

消费代码示例

@Service
@RequiredArgsConstructor
public class OrderCreatedConsumer {

    private final ConsumedEventMapper consumedEventMapper;

    private final InventoryService inventoryService;

    @KafkaListener(
            topics = "order-created-topic",
            groupId = "inventory-service"
    )
    @Transactional
    public void consume(String message) {
        OrderCreatedEvent event = deserialize(message);

        ConsumedEvent consumedEvent = new ConsumedEvent();
        consumedEvent.setConsumerGroup("inventory-service");
        consumedEvent.setEventId(event.eventId());
        consumedEvent.setEventType("OrderCreated");
        consumedEvent.setConsumedAt(LocalDateTime.now());

        try {
            consumedEventMapper.insert(consumedEvent);
        } catch (DuplicateKeyException exception) {
            return;
        }

        inventoryService.reserveInventory(
                event.orderId(),
                event.items()
        );
    }
}

执行流程如下:

1. 插入 consumed_event
2. 执行业务逻辑
3. 提交数据库事务

如果消息第一次消费:

consumed_event 插入成功
业务逻辑执行成功
事务提交

如果消息重复消费:

consumed_event 唯一键冲突
直接返回
不再执行业务逻辑

为什么消费记录和业务操作必须在同一个事务中

错误写法:

public void consume(OrderCreatedEvent event) {
    if (consumedEventMapper.exists(event.eventId())) {
        return;
    }

    inventoryService.reserveInventory(event);

    consumedEventMapper.insert(event.eventId());
}

这段代码存在多个问题。

问题一:先查后写存在并发竞争

两个消费线程可能同时执行:

线程 A:查询,不存在
线程 B:查询,不存在
线程 A:执行业务逻辑
线程 B:执行业务逻辑

最终业务重复执行。

因此不能只依赖:

SELECT 后再判断

必须依赖数据库唯一约束兜底。

问题二:业务成功,消费记录写入失败

执行顺序:

1. 扣减库存成功
2. 写入 consumed_event 失败
3. 消息重新投递
4. 再次扣减库存

因此消费记录和业务操作必须放入同一个本地事务:

@Transactional
public void consume(...) {
    insertConsumedEvent();
    executeBusinessLogic();
}

只有这样才能保证:

消费记录成功 + 业务处理成功

或者:

消费记录失败 + 业务处理回滚

十七、方案二:通过业务唯一键实现幂等

有些业务本身已经具有唯一业务标识。

例如,支付服务处理第三方支付回调时,可以使用:

payment_transaction_no

作为唯一键。

表结构:

CREATE TABLE payment_record (
    id BIGINT PRIMARY KEY,
    payment_transaction_no VARCHAR(128) NOT NULL,
    order_id BIGINT NOT NULL,
    amount DECIMAL(18, 2) NOT NULL,
    created_at DATETIME NOT NULL,

    UNIQUE KEY uk_payment_transaction_no (
        payment_transaction_no
    )
);

消费者直接尝试插入:

try {
    paymentRecordMapper.insert(paymentRecord);
} catch (DuplicateKeyException exception) {
    return;
}

如果支付流水号已经存在,就说明这笔支付已经处理过。

这种方案的优点是:

  • 无需额外消费记录表;
  • 幂等约束直接体现业务语义;
  • 数据库唯一索引强一致;
  • 适合订单号、支付流水号、退款单号等场景。

十八、方案三:基于状态机实现幂等

订单状态更新通常可以通过状态机约束。

例如,只允许订单从 CREATED 更新为 PAID

UPDATE orders
SET status = 'PAID',
    paid_at = NOW()
WHERE id = ?
  AND status = 'CREATED';

应用根据更新行数判断是否处理成功:

int affectedRows =
        orderMapper.markPaid(orderId);

if (affectedRows == 0) {
    return;
}

第一次消费:

CREATED -> PAID
affectedRows = 1

重复消费:

当前状态已经是 PAID
affectedRows = 0

这种方法不仅可以防止重复消费,还可以防止非法状态流转。

例如:

CANCELLED -> PAID

不应该被允许。

状态机方案的注意事项

只有当重复执行的业务操作可以通过状态流转限制时,状态机方案才适用。

例如:

订单支付
订单取消
订单完成

通常比较适合。

但以下业务不一定适合只通过状态机去重:

增加积分
发送优惠券
记录行为日志
发送短信

因为这些操作可能没有天然的唯一状态转换。

十九、方案四:乐观锁

可以给业务表增加版本号:

ALTER TABLE orders
ADD COLUMN version INT NOT NULL DEFAULT 0;

更新时带上版本号:

UPDATE orders
SET status = 'PAID',
    version = version + 1
WHERE id = ?
  AND status = 'CREATED'
  AND version = ?;

如果更新行数为 0,说明:

  • 订单状态已变化;
  • 数据已经被其他线程修改;
  • 当前消息可能已经被处理;
  • 版本号已经过期。

乐观锁适合解决并发更新,但需要注意:

乐观锁本身并不等价于完整的消息幂等机制。

它主要解决的是同一业务记录的并发修改问题,不一定能识别“这条具体事件是否已经处理过”。

二十、方案五:Redis 去重

也可以使用 Redis 保存已消费事件 ID:

String key =
        "consumed:event:inventory-service:" + event.eventId();

Boolean success = redisTemplate.opsForValue()
        .setIfAbsent(
                key,
                "1",
                Duration.ofDays(7)
        );

if (!Boolean.TRUE.equals(success)) {
    return;
}

inventoryService.reserveInventory(event);

对应 Redis 命令:

SET consumed:event:inventory-service:EVENT_ID 1 NX EX 604800

这种方案速度快,但存在一致性问题。

例如:

1. Redis SET NX 成功
2. 业务数据库处理失败
3. Redis 中已经标记为已消费
4. 消息重试时被直接跳过

最终业务消息丢失。

反过来也可能出现:

1. 业务处理成功
2. Redis 写入失败
3. 消息重试
4. 业务重复执行

因此 Redis 去重更适合:

  • 对强一致性要求不高的通知消息;
  • 日志、埋点等可容忍少量异常的场景;
  • 作为数据库幂等机制之前的快速过滤层;
  • 明确可接受有限时间窗口去重的场景。

对于支付、库存、账务等核心业务,不应只依赖 Redis。

二十一、数据库幂等和 Kafka Offset 的事务边界

消费者通常需要完成两件事:

1. 执行业务数据库事务
2. 提交 Kafka Offset

但这两个动作通常不属于同一个事务。

推荐处理顺序是:

1. 消费消息
2. 执行业务数据库事务
3. 数据库事务提交成功
4. 提交 Kafka Offset

如果在第 3 步和第 4 步之间崩溃:

数据库业务已经成功
Offset 没有提交

消息会再次被消费。

这正是为什么消费者必须幂等。

不推荐先提交 Offset 再处理业务:

1. 提交 Offset
2. 处理业务
3. 业务失败

此时 Kafka 会认为消息已经消费成功,消息可能永久丢失。

因此可靠消费的基本原则是:

先完成业务,再确认消息;允许消息重复,不允许消息丢失。

二十二、Outbox 同样可能重复发送消息

很多开发者误以为:

使用了 Outbox,就不会重复发送消息

这是错误的。

假设轮询程序执行:

1. 读取 Outbox 事件
2. 发送 Kafka 成功
3. 更新 Outbox 状态之前应用崩溃

Outbox 状态仍然是:

PENDING

应用恢复后会再次发送。

CDC 同样可能由于 Connector 重启、位点回退等原因重新发送。

因此完整的可靠链路必须接受这个事实:

消息可能重复

并通过幂等消费把重复影响消除。

二十三、消息顺序问题

Outbox Pattern、CDC 和幂等消费主要解决可靠性问题,但不自动保证消息顺序。

例如同一个订单产生两条事件:

OrderCreated
OrderCancelled

如果消费者先收到:

OrderCancelled

后收到:

OrderCreated

可能导致错误状态。

1. 使用相同的消息 Key

Kafka 可以保证:

同一个 Partition 内消息有序。

因此可以使用 aggregate_id 作为消息 Key:

kafkaTemplate.send(
        "order-event-topic",
        orderId.toString(),
        payload
);

同一个订单的事件会进入同一个 Partition。

2. 增加 aggregate_version

事件中增加聚合版本号:

{
  "eventId": "event-10002",
  "aggregateId": "order-10001",
  "aggregateVersion": 3,
  "eventType": "OrderCancelled"
}

消费者记录已经处理的最大版本号:

last_processed_version = 2

只有当:

event.aggregateVersion = last_processed_version + 1

时才处理。

对于小于等于当前版本的事件,可以认为是重复或过期消息。

3. 不要把全局有序作为默认目标

Kafka 只能在单个 Partition 内保证顺序。

如果为了全局有序而只使用一个 Partition,会严重限制吞吐量。

更合理的目标通常是:

同一个业务聚合内有序
不同业务聚合之间并行

例如:

同一个订单内有序
不同订单之间无需有序

二十四、消息版本兼容

事件一旦发送到消息中间件,就可能被多个消费者长期依赖。

因此事件结构需要考虑版本演进。

建议增加:

{
  "eventType": "OrderCreated",
  "eventVersion": 2
}

版本升级时尽量保持向后兼容。

例如,原事件:

{
  "orderId": 100001,
  "amount": 199.00
}

新增字段:

{
  "orderId": 100001,
  "amount": 199.00,
  "currency": "CNY"
}

通常是向后兼容的。

但删除字段、修改字段类型、改变字段语义可能造成消费者失败。

不推荐:

{
  "amount": "199.00 CNY"
}

因为原来的 amount 从数字变成了字符串。

更好的做法是新增字段:

{
  "amount": 199.00,
  "currency": "CNY"
}

二十五、失败重试和死信队列

消费者处理失败后,需要配置合理的重试策略。

推荐使用指数退避:

第 1 次:1 秒后
第 2 次:5 秒后
第 3 次:30 秒后
第 4 次:5 分钟后

不建议无限快速重试,否则可能导致:

  • CPU 占用过高;
  • 消息积压;
  • 数据库持续承压;
  • 日志大量刷屏;
  • 阻塞同一 Partition 后续消息。

超过最大重试次数后,可以进入死信队列:

order-created-topic
    |
    | 多次消费失败
    v
order-created-topic.DLT

死信消息中建议保留:

原始消息
eventId
eventType
消费组
失败原因
异常堆栈
重试次数
首次失败时间
最后失败时间

二十六、Outbox 数据清理

Outbox 表会不断增长,因此需要设计清理机制。

例如,删除 30 天前已经成功发送的事件:

DELETE FROM outbox_event
WHERE status = 'PUBLISHED'
  AND published_at < NOW() - INTERVAL 30 DAY
LIMIT 1000;

建议分批删除,不要一次删除大量数据:

每次删除 1000 条
循环执行
每轮适当休眠

否则可能造成:

  • 长事务;
  • 锁等待;
  • Binlog 急剧增长;
  • 主从复制延迟;
  • 数据库 IO 峰值。

如果事件具有审计价值,也可以迁移到历史表:

outbox_event
    |
    v
outbox_event_archive

或者存入低成本对象存储。

二十七、消费记录表清理

consumed_event 同样会持续增长。

但清理消费记录需要谨慎,因为记录删除后,历史消息如果重新投递,可能再次被处理。

清理策略需要结合:

  • Kafka 消息保留时间;
  • 最大允许重放时间;
  • 业务数据生命周期;
  • 灾难恢复策略;
  • 消息补偿周期。

例如:

Kafka 保留消息 7 天
业务最多重放 15 天
consumed_event 保留 30 天

消费记录的保留时间应该明显大于消息可能被重新投递的时间窗口。

对于支付、账务等核心数据,也可以永久保留业务唯一流水,而不是依赖可清理的通用消费记录。

二十八、监控指标

可靠消息链路不能只实现功能,还必须具备可观测性。

Outbox 指标

建议监控:

PENDING 事件数量
PROCESSING 事件数量
FAILED 事件数量
DEAD 事件数量
最老 PENDING 事件等待时间
平均发布延迟
发送成功率
重试次数
每秒发布事件数

其中最重要的通常是:

最老未发送事件年龄

如果最老事件已经等待 30 分钟,即使 PENDING 数量并不多,也说明发布链路可能出现异常。

CDC 指标

建议监控:

CDC Connector 状态
Binlog/WAL 位点
源数据库与 Kafka 的延迟
Connector 重启次数
消息发送失败次数
Schema 解析失败次数
任务是否处于 RUNNING 状态

消费者指标

建议监控:

Consumer Lag
消费成功率
消费失败率
单条消息处理耗时
重复事件数量
幂等冲突数量
重试队列长度
死信队列数量
最老未消费消息年龄

二十九、完整代码流程示例

1. 生产端事务

@Transactional
public void createOrder(CreateOrderCommand command) {
    Order order = createOrderEntity(command);
    orderMapper.insert(order);

    OrderCreatedEvent event = new OrderCreatedEvent(
            UUID.randomUUID().toString(),
            order.getId(),
            order.getOrderNo(),
            Instant.now()
    );

    OutboxEvent outboxEvent = OutboxEvent.builder()
            .id(idGenerator.nextId())
            .eventId(event.eventId())
            .aggregateType("ORDER")
            .aggregateId(order.getId().toString())
            .eventType("OrderCreated")
            .topic("order-event-topic")
            .eventKey(order.getId().toString())
            .payload(toJson(event))
            .status("PENDING")
            .retryCount(0)
            .createdAt(LocalDateTime.now())
            .updatedAt(LocalDateTime.now())
            .build();

    outboxEventMapper.insert(outboxEvent);
}

事务提交后,数据库中一定同时存在:

orders 记录
outbox_event 记录

2. CDC 读取 Outbox

MySQL Binlog
    |
    v
Debezium
    |
    v
Kafka order-event-topic

消息 Key 使用:

aggregate_id

例如:

100001

这样同一个订单的事件可以进入同一个 Kafka Partition。

3. 消费端幂等处理

@Service
@RequiredArgsConstructor
public class InventoryOrderEventConsumer {

    private static final String CONSUMER_GROUP =
            "inventory-service";

    private final ConsumedEventMapper consumedEventMapper;

    private final InventoryService inventoryService;

    @KafkaListener(
            topics = "order-event-topic",
            groupId = CONSUMER_GROUP
    )
    @Transactional
    public void consume(OrderCreatedEvent event) {
        try {
            consumedEventMapper.insert(
                    new ConsumedEvent(
                            null,
                            CONSUMER_GROUP,
                            event.eventId(),
                            "OrderCreated",
                            LocalDateTime.now()
                    )
            );
        } catch (DuplicateKeyException exception) {
            return;
        }

        inventoryService.reserveByOrder(
                event.orderId()
        );
    }
}

执行结果:

第一次收到:
消费记录插入成功
库存预占成功
事务提交

第二次收到:
消费记录唯一键冲突
直接返回
库存不会重复预占

三十、常见错误设计

错误一:业务事务中直接发送 MQ

@Transactional
public void createOrder() {
    orderMapper.insert(order);
    kafkaTemplate.send(topic, message);
}

问题:

数据库事务和 Kafka 发送不属于同一个原子事务

错误二:事务提交后再发送,但没有补偿机制

orderService.createOrder();

kafkaTemplate.send(topic, message);

问题:

数据库提交后应用可能崩溃
消息没有发送,也无法自动发现

错误三:认为 Outbox 可以保证消息只发送一次

Outbox 只能保证:

业务数据和事件记录同时存在

它不能保证:

事件绝不重复发送

错误四:消费者先查再处理

if (exists(eventId)) {
    return;
}

executeBusiness();
insertConsumedEvent(eventId);

问题:

存在并发竞争窗口

必须使用:

唯一索引

作为最终防线。

错误五:消费记录和业务操作不在同一个事务

insertConsumedEvent();
executeBusiness();

如果第二步失败,但第一步已经提交,消息重试时会被错误跳过。

错误六:仅依赖 Redis 做核心业务幂等

Redis 和业务数据库之间缺少本地事务,容易出现状态不一致。

核心资金、库存、支付业务应优先使用数据库唯一约束或业务状态机。

错误七:把技术消息直接暴露为业务事件

直接把整张数据库表的 CDC 变更发送给所有消费者,会让下游系统与数据库结构强耦合。

例如:

{
  "before": {
    "status": 0
  },
  "after": {
    "status": 1
  },
  "op": "u"
}

下游必须知道:

status = 0 表示什么
status = 1 表示什么

数据库字段一旦变化,所有消费者都可能受到影响。

更好的方式是发送明确的业务事件:

{
  "eventType": "OrderPaid",
  "eventVersion": 1,
  "data": {
    "orderId": 100001,
    "paidAt": "2026-07-16T10:30:00Z"
  }
}

三十一、轮询 Outbox 和 CDC 应该如何选择

可以根据系统规模和基础设施选择。

适合使用轮询的场景

  • 系统规模较小;
  • 消息量不高;
  • 不希望维护 Kafka Connect;
  • 对秒级延迟可以接受;
  • 团队希望先快速落地;
  • 已经有成熟的定时任务平台。

适合使用 CDC 的场景

  • 消息量较大;
  • 对实时性要求较高;
  • 已经使用 Kafka Connect;
  • 团队具备 CDC 运维能力;
  • 需要标准化事件发布链路;
  • 多个服务都需要 Outbox 能力;
  • 希望减少应用内消息发布代码。

推荐演进路线

很多系统可以按照以下路线演进:

第一阶段:
Outbox + 定时轮询 + 幂等消费

第二阶段:
Outbox + 独立发布服务 + 幂等消费

第三阶段:
Outbox + Debezium CDC + Kafka + 幂等消费

Outbox 表的数据模型可以保持不变,只替换发布方式。

三十二、三者分别解决什么问题

可以用一张表总结:

机制 解决的问题 不能解决的问题
Outbox Pattern 业务数据和事件记录的一致性 消息重复、消费者重复处理
CDC 可靠捕获和发布数据库变更 消息绝不重复、消费幂等
幂等消费 消除重复消息带来的业务影响 生产端事件缺失
Kafka ACK Kafka Broker 内部写入可靠性 数据库和 Kafka 的原子一致性
Kafka Offset 记录消费者读取进度 业务处理是否重复
重试机制 临时故障恢复 永久性业务错误
死信队列 隔离持续失败的消息 自动修复业务数据

三十三、推荐的企业级实现原则

在实际项目中,可以遵循以下原则。

1. 业务数据和 Outbox 使用同一个本地事务

业务表写入
Outbox 表写入
必须原子提交

2. 每条事件必须有全局唯一 event_id

event_id

需要从生产端一直透传到消费端。

3. 采用至少一次投递模型

不要假设消息不会重复。

At Least Once + Idempotent Consumer

通常是更现实、更可靠的工程方案。

4. 消费者必须默认实现幂等

不要等到生产环境出现重复消费后再补幂等。

重复消息不是异常现象,而是可靠消息系统的正常边界情况。

5. 数据库唯一约束必须作为最终防线

应用代码判断只能作为优化,不能代替数据库唯一约束。

应用判断:提高可读性或减少异常
数据库唯一键:保证并发正确性

6. 幂等记录和业务操作必须处于同一个事务

否则会出现:

标记成功但业务失败

或者:

业务成功但标记失败

7. 使用 aggregate_id 作为 Kafka Key

保证同一个业务聚合内的事件进入同一个 Partition。

orderId
userId
paymentId

都可以作为消息 Key,具体取决于聚合边界。

8. 设计重试、死信和人工补偿机制

可靠消息系统不能只考虑成功路径。

需要明确处理:

临时失败
永久失败
毒消息
消费超时
死信重放
人工补偿
事件重放

9. 事件协议必须版本化

推荐字段:

{
  "eventId": "...",
  "eventType": "...",
  "eventVersion": 1,
  "aggregateId": "...",
  "occurredAt": "...",
  "traceId": "...",
  "data": {}
}

10. 为整个链路建立可观测性

应该能够通过 eventId 查询:

事件何时产生
是否写入 Outbox
CDC 是否捕获
写入哪个 Topic
消费者是否收到
是否重复消费
是否处理成功
是否进入死信队列

三十四、最终架构

一个完整、可靠的事件驱动架构可以表示为:

┌──────────────────────┐
│     Order Service    │
│                      │
│  1. 写入 orders      │
│  2. 写入 outbox      │
└──────────┬───────────┘
           │ 同一个本地事务
           v
┌──────────────────────┐
│        MySQL         │
│                      │
│ orders               │
│ outbox_event         │
└──────────┬───────────┘
           │ Binlog
           v
┌──────────────────────┐
│ Debezium / CDC       │
└──────────┬───────────┘
           v
┌──────────────────────┐
│        Kafka         │
│ order-event-topic    │
└──────────┬───────────┘
           v
┌──────────────────────┐
│ Inventory Consumer   │
│                      │
│ consumed_event       │
│ inventory_business   │
│ 同一个本地事务       │
└──────────────────────┘

这套架构中:

Outbox 保证事件不会因为业务事务边界而缺失
CDC 保证已提交事件可以被可靠捕获
Kafka 提供持久化和分发能力
幂等消费解决重复投递问题
重试和死信机制解决消费失败问题
监控和追踪保证系统可以被运维

总结

Outbox Pattern、CDC 和幂等消费是构建可靠事件驱动系统的三个核心组成部分。

Outbox Pattern 的核心是:

将业务数据和事件记录写入同一个本地数据库事务。

CDC 的核心是:

读取数据库事务日志,将已经提交的 Outbox 事件可靠地发送到消息中间件。

幂等消费的核心是:

即使同一条消息被重复投递,也只产生一次有效业务结果。

三者组合后的完整语义是:

业务数据和事件记录原子提交
        +
事件至少一次可靠发布
        +
消费者幂等处理
        =
最终一致的事件驱动系统

需要特别理解的是:

Outbox 不是为了实现整个分布式链路的精确一次,而是把无法控制的“双写问题”,转换成可恢复、可重试、可观测的单库事务问题。

工程实践中,不应该把目标定义为:

消息永远不重复

更合理的目标是:

消息可以重复,但业务结果不能重复

这也是绝大多数可靠消息系统最终采用的设计:

At Least Once Delivery
        +
Idempotent Consumer
        =
Effectively Once
使用 Hugo 构建
主题 StackJimmy 设计