在微服务和事件驱动架构中,一个非常常见的业务场景是:
- 修改数据库中的业务数据。
- 向 Kafka、RocketMQ 等消息中间件发送一条消息。
- 下游服务消费消息,并执行后续业务逻辑。
例如,用户提交订单后,订单服务需要:
- 在数据库中创建订单;
- 发送“订单已创建”事件;
- 库存服务收到事件后扣减库存;
- 积分服务收到事件后增加积分;
- 通知服务收到事件后发送通知。
看起来只是一次数据库操作加一次消息发送,但这里隐藏着一个非常关键的问题:
如何保证数据库事务和消息发送的一致性?
如果数据库写入成功,但消息发送失败,下游服务就永远无法感知这次业务变化。
如果消息发送成功,但数据库事务回滚,下游服务却可能处理一条实际上并不存在的业务数据。
为了解决这些问题,工程中经常组合使用以下三种机制:
- 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 等消息中间件。
常见方案有两种:
- 定时轮询 Outbox 表;
- 使用 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 元。
十五、幂等消费的常见实现方案
常见方案包括:
- 唯一约束;
- 消费记录表;
- 业务状态机;
- 乐观锁;
- Redis 去重;
- 使用业务唯一键;
- 天然幂等 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