这是《RocketMQ 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 消息系统通常提供至少一次投递语义,重复是常态而不是异常。幂等设计的目标是:同一条消息、同一个事件或同一个业务指令被处理多次,最终业务结果只生效一次。
13.1 为什么会重复
常见来源:
producer send timeout -> internal retry -> duplicate
consumer process success -> offset commit fail -> redeliver
consumer timeout -> broker retry -> duplicate
dead letter replay -> duplicate
compensation task -> duplicate
network retry -> duplicate
重复不能完全避免,必须在消费端和业务端处理。
13.2 幂等边界
先定义“什么算同一次”:
| 粒度 | 幂等键 |
|---|---|
| 事件 | eventId |
| 订单操作 | orderNo + eventType + version |
| 支付请求 | paymentId |
| 库存变更 | changeId |
| 用户任务 | taskId |
| 外部回调 | requestId |
不要用消息 ID 作为唯一业务幂等依据。消息 ID 标识一次物理消息,业务重发可能产生不同消息 ID,但代表同一业务事件。
13.3 数据库唯一约束
事件表:
CREATE TABLE processed_event (
event_id VARCHAR(64) PRIMARY KEY,
event_type VARCHAR(64) NOT NULL,
aggregate_id VARCHAR(64) NOT NULL,
consumed_at DATETIME(6) NOT NULL,
result_status VARCHAR(20) NOT NULL,
UNIQUE KEY uk_event_type_id (event_type, event_id)
);
处理:
@Transactional
public void handle(OrderCreatedEvent event) {
try {
processedEventRepository.insert(event.eventId(),
"OrderCreated", event.orderNo(), "SUCCESS");
orderService.create(event);
} catch (DuplicateKeyException e) {
log.info("duplicate event ignored, eventId={}", event.eventId());
}
}
事件插入和业务处理必须在同一个本地事务中。
13.4 状态机幂等
状态变更必须带前置条件:
UPDATE orders
SET status = 'PAID', paid_at = NOW(6), version = version + 1
WHERE order_no = ?
AND status = 'WAIT_PAY'
AND version = ?;
结果判断:
| 更新行数 | 含义 |
|---|---|
| 1 | 首次生效 |
| 0 | 已处理、状态不允许或版本不匹配 |
不能只执行“无条件 update”,否则重复事件可能覆盖较新的状态。
13.5 业务动作唯一键
扣减库存:
INSERT INTO inventory_change(
change_id, sku_id, order_no, quantity, created_at
) VALUES (?, ?, ?, ?, NOW(6));
用 change_id 唯一约束控制重复扣减,再通过触发器或应用事务更新库存汇总。
账户入账:
INSERT INTO account_transaction(
transaction_id, account_id, amount, direction, created_at
) VALUES (?, ?, ?, ?, NOW(6));
金融场景还应对账:
sum(account_transaction) = account.balance
13.6 Redis 幂等
Redis 适合做快速拦截:
Boolean first = redis.opsForValue()
.setIfAbsent("idempotent:" + event.eventId(), "PROCESSING", Duration.ofMinutes(10));
if (Boolean.FALSE.equals(first)) {
return;
}
try {
handleInDatabase(event);
} catch (Exception e) {
redis.delete("idempotent:" + event.eventId());
throw e;
}
限制:
- Redis 可能丢失;
- 过期后仍会重复;
- 主从切换可能带来窗口;
- 不能替代数据库唯一约束;
- 只能作为性能优化。
最终裁决应落在业务数据库或账务表。
13.7 并发控制
同一消息可能被并发重试:
thread A -> insert event
thread B -> insert event
处理方式:
- 数据库唯一键;
- 分布式锁,但锁只用于减少冲突;
- 按聚合 ID 分片;
- 顺序消费单 key;
- 乐观锁版本;
- 条件更新。
不要依赖“重试不会并发”的假设。
13.8 幂等测试
必测场景:
- 同一 eventId 连续处理两次;
- 同一 eventId 并发处理;
- 处理成功但位点提交失败;
- 处理失败后重试;
- 死信重放;
- 补偿任务与消息并发;
- 版本旧事件晚到;
- Redis 幂等键过期。
断言示例:
handle(event);
handle(event);
assertEquals(1, inventoryRepository.findReserved(event.orderNo()).size());
13.9 常见错误
| 错误 | 后果 |
|---|---|
| 只判断内存 Set | 重启后失效 |
| 只使用 Redis | 缓存丢失后重复 |
| 事件表与业务不同事务 | 部分成功 |
| 无条件更新状态 | 新状态被旧事件覆盖 |
| 用消息 ID 做业务幂等 | 重发绕过 |
| 先业务后插事件且不同事务 | 异常时漏记 |
本章小结
幂等要从业务语义出发,选择稳定的事件 ID 和业务唯一键,并把最终一致性落在数据库唯一约束、事务和状态机上。Redis 可以降低冲突和数据库压力,但不能作为唯一依据。所有幂等逻辑都应针对重复、并发、重放和乱序进行测试。
思考题
- 为什么消息 ID 不适合作为业务幂等键?
- 事件表为什么要和业务处理同事务?
- Redis 幂等在什么情况下失效?
- 状态机版本如何防止旧事件覆盖新状态?
- 如何验证死信重放不会重复扣款?