RocketMQNotes

第 13 章:幂等设计

zjc 于 2026-01-13 发布

这是《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;
}

限制:

  1. Redis 可能丢失;
  2. 过期后仍会重复;
  3. 主从切换可能带来窗口;
  4. 不能替代数据库唯一约束;
  5. 只能作为性能优化。

最终裁决应落在业务数据库或账务表。

13.7 并发控制

同一消息可能被并发重试:

thread A -> insert event
thread B -> insert event

处理方式:

  1. 数据库唯一键;
  2. 分布式锁,但锁只用于减少冲突;
  3. 按聚合 ID 分片;
  4. 顺序消费单 key;
  5. 乐观锁版本;
  6. 条件更新。

不要依赖“重试不会并发”的假设。

13.8 幂等测试

必测场景:

  1. 同一 eventId 连续处理两次;
  2. 同一 eventId 并发处理;
  3. 处理成功但位点提交失败;
  4. 处理失败后重试;
  5. 死信重放;
  6. 补偿任务与消息并发;
  7. 版本旧事件晚到;
  8. Redis 幂等键过期。

断言示例:

handle(event);
handle(event);

assertEquals(1, inventoryRepository.findReserved(event.orderNo()).size());

13.9 常见错误

错误 后果
只判断内存 Set 重启后失效
只使用 Redis 缓存丢失后重复
事件表与业务不同事务 部分成功
无条件更新状态 新状态被旧事件覆盖
用消息 ID 做业务幂等 重发绕过
先业务后插事件且不同事务 异常时漏记

本章小结

幂等要从业务语义出发,选择稳定的事件 ID 和业务唯一键,并把最终一致性落在数据库唯一约束、事务和状态机上。Redis 可以降低冲突和数据库压力,但不能作为唯一依据。所有幂等逻辑都应针对重复、并发、重放和乱序进行测试。

思考题

  1. 为什么消息 ID 不适合作为业务幂等键?
  2. 事件表为什么要和业务处理同事务?
  3. Redis 幂等在什么情况下失效?
  4. 状态机版本如何防止旧事件覆盖新状态?
  5. 如何验证死信重放不会重复扣款?