复盘积分翻倍事故
还记得王阿姨吗?积分服务消费超时,MQ 没等到确认,按「至少一次」的契约重投了一遍。站在 MQ 的角度,它没有任何错——错的是我们的消费代码裸奔:来一条加一条,来两条加两条。
这一篇把重复的根源挖到底,再把消费幂等的三板斧备齐。
重复不是意外,是出厂设置
先想清楚重复从哪来。上一章我们说过:ACK 可能丢。消息明明投到了,消费者处理完了,但确认消息在网络里迷了路——MQ 只能认定失败,重投一遍。这是「不丢」的代价,两个承诺(不丢、不重)在分布式里本质上是冲突的,MQ 集体选择了保「不丢」。
重复的来源清单:
- 生产端重试:ACK 超时,发送方重发,Broker 里出现两条相同消息;
- 消费确认丢失:处理成功但 ACK 没送达,Broker 重投;
- 消费者组 Rebalance:队列被重新分配,新消费者从上一次提交的位点接着干,上次处理到一半的会被再处理一遍。
那 Kafka 宣传的 Exactly Once 呢?它指的是生产端幂等(PID+序号防单分区内重复)加事务(跨分区原子写),主要面向流处理框架的「读-算-写」闭环。你的业务系统消费消息后写自己的数据库,这个环节的重复依然要靠业务幂等——恰好一次在这里帮不了你。
幂等三板斧之一:数据库唯一键
最朴素也最可靠的方案:让业务本身有唯一约束。比如「订单已支付」事件的消费逻辑是插一条积分流水,那么把(订单号+积分类型)建上唯一索引:
try {
pointLogMapper.insert(log); // 唯一索引撞上就抛 DuplicateKeyException
pointAccountMapper.addPoint(log.getUserId(), log.getPoint());
} catch (DuplicateKeyException e) {
log.warn("重复消息,幂等返回,msgId={}", msgId);
}重复消息来了,插入报错,捕获后当无事发生。数据库替你兜底,不引入新组件,缺点是要求消费动作里有「可唯一化的落库行为」。
幂等三板斧之二:去重表
消费动作本身不好唯一化(比如发短信、调外部接口),就单独建一张去重表:
consume_record (msg_key, consumer_group, status, create_time)
UNIQUE KEY (msg_key, consumer_group)消费逻辑变成:在业务事务里先 insert 去重记录,插入成功就干活,撞唯一键就说明处理过,直接返回。注意两点:去重表和业务操作必须在同一个事务里,否则插了记录没干活照样丢;按 consumer_group 隔离,让多个订阅组各自维护自己的去重账本。有人用 Redis setnx 做去重,快是快,但缓存和数据库不是原子提交,故障窗口照样重复——它可以当加速层,不能当唯一防线。
幂等三板斧之三:状态机
有状态流转的业务,用状态机防重最优雅。订单只允许从「待支付」到「已支付」单向流转:
int rows = orderMapper.markPaid(orderId);
// UPDATE orders SET status = 38 WHERE id = ? AND status = 20
if (rows == 0) {
return; // 状态已经不是待支付,重复消息或非法状态,直接放弃
}
pointService.addPoint(order.getUserId()); // 只有真正流转成功的才发积分重复消息第二次到来时,UPDATE 影响行数为 0,后面的动作一概不执行。状态机方案不需要额外表,还顺手挡住了乱序消息(第 8 篇会再见到它)。
三板斧怎么选
| 方案 | 适用前提 | 注意点 |
|---|---|---|
| 数据库唯一键 | 消费逻辑含落库,且能定出业务唯一键 | 唯一键设计要覆盖「同一次业务动作」的所有消息 |
| 去重表 | 动作无天然唯一键(调外部接口等) | 必须与业务同事务;Redis 只能做加速不能做兜底 |
| 状态机 | 业务有明确的状态流转 | 流转条件要写全,别漏掉中间态 |
最后校准一下预期:幂等保证的是「执行多次,效果如同执行一次」,而不是「只执行一次」。消息可能真的来了三遍,代码也真的跑了三遍,只是第二、三遍空转——这已经是分布式系统里能做到的最好了。
小结
重复是「至少一次」的影子,赶不走,只能免疫。三板斧记牢:能定唯一键就用唯一键,定不出就上同事务去重表,有状态流转就上状态机。王阿姨的积分事故,三板斧随便一斧都能挡住。
丢了不行,重了不行,还剩一个刺头:消息顺序。先扣库存再付钱和先付钱再扣库存,业务上是两回事。下一篇拆「全局有序」与「分区内有序」——以及为什么99%的场景只需要后者。
评论 (0)