MQ幂等性实战:重复消息产生的原理与四大去重方案 📅 发布时间:2026/9/12 4:18:33 👁 浏览次数: 凌晨一点被电话叫醒线上报了一个用户收到两条扣款通知的问题。拉完流水之后定位到原因订单表里同一个支付回调事件被消费端处理了两遍第一遍正常入账第二遍又把金额累加了一次。这不是网络抖动也不是代码逻辑写错而是消息队列在保证不丢的同时天然会把同一条消息重复投递给消费者。这个问题的学名就是消息幂等性。消息队列MQ如何保证消息的幂等性可以说是后端面试里出现频率最高的问题之一也是生产环境里最容易埋雷的地方。很多团队在引入Kafka、RocketMQ、RabbitMQ的时候只盯着吞吐量和削峰填谷忽略了消费端的重复风险结果流量一上来就爆出重复订单、重复发券、重复扣款。这篇文章我把自己的踩坑经历、方案选型和各种细节完整梳理一遍不管你是刚开始接触MQ还是已经在生产环境里摸爬滚打应该都能找到能直接落地的思路。1. MQ至少一次的承诺注定了重复消息躲不掉1.1 三种投递语义你最终绕不开至少一次很多人在设计消息消费逻辑时默认消息系统会刚刚好地把每条消息投递一次。但分布式系统里不存在这种理想情况。MQ官方文档里明确给出了三种投递语义At-most-once至多一次消息可能丢但不会重复。实时性要求高、允许丢弃的场景才会用。At-least-once至少一次消息不会丢但可能重复。这是Kafka、RocketMQ、RabbitMQ在绝大多数配置下的默认行为。Exactly-once精确一次不丢不重代价极高且通常只是某个环节内的精确一次不是端到端。Kafka有幂等Producer和事务API但它解决的是Producer到Broker这段路径的精确一次以及跨分区写入的原子性不是从Producer生成消息到Consumer完成业务写入的端到端精确一次。RocketMQ的事务消息能保证本地事务和消息发送的一致性但消息最终投递给消费者时依然可能重复。说白了只要你的业务和消息是两套系统重复就无法彻底消除必须在消费侧自己兜底。1.2 重复消息到底从哪几个环节冒出来我见过很多同学觉得重复投递是小概率事件直到线上出问题才去翻日志。下面这几条路径每一条都能真实地制造重复消息Producer端发送重试。发送消息时网络超时Producer会抛异常或触发重试但Broker可能已经写入成功。比如Kafka Producer配置了retries3第一次发送实际成功了但响应包丢了Producer重试又发了一次同一条业务消息就出现了两份。Consumer端消费成功但没来得及提交。这是最经典的重叠窗口。消费者处理完业务逻辑正要提交offset或发送ack时进程宕机、被重启、或者网络闪断。等消费者恢复后Broker会把它当成这条消息还没消费成功重新投递一次。Kafka默认的enable.auto.committrue每5秒自动提交一次这5秒窗口内挂掉几乎必然产生重复消费。Rebalance导致的重复。Kafka消费者组发生重平衡时分区会重新分配某些分区的最新offset可能回退到上一次提交的位置之前已经消费过但还没提交offset的消息会被重新拉取。消费线程越多重平衡越频繁重复概率越高。RocketMQ的重试队列。消费失败后RocketMQ会把消息投递到重试队列按延迟级别再次投递默认可以重试16次。如果第一次消费时业务已经成功但返回结果因为网络原因没送到Broker这条消息还是会被重试投递。以前我总以为重复消费只有在故障时才会出现后来发现正常运行时也会因为超时、重平衡、并发变更产生重复。所以技术方案上必须把重复当成常态来设计。1.3 为什么MQ自带的MsgId当不了幂等护身符有人会问RocketMQ每条消息都有msgId拿它做去重不就完了我最早也这么干过后来发现这是个坑。RocketMQ的msgId是Producer发送时生成的客户端ID如果Producer重试发送同一条业务消息msgId会重新生成两个msgId不同但业务内容完全一样去重判断直接失效。Kafka则没有全局唯一的消息IDoffset只能唯一标识分区内的位置消费者组变化后offset语义不稳定。RabbitMQ的deliveryTag是Channel级别的自增序号消息重新入队后也会变化。所以MQ自带的ID只能用来查日志和排错绝不能用来做业务幂等。真正靠谱的是业务侧自己定义的唯一标识比如orderId、paymentId、userId sourceType sourceId这类能唯一代表一次业务事件的值。这是整个幂等设计的地基地基如果歪了后面的方案全部白搭。2. 幂等方案怎么选去重表、版本号、Redis、状态机2.1 唯一键去重表代价最小兜底最稳我最推荐的方案是针对每一个需要幂等的业务事件建一张去重表把业务唯一键作为唯一索引。处理流程是收到消息后先往去重表里插入一条记录插入成功才继续执行真正的业务逻辑插入失败说明这条消息之前已经处理过直接向MQ提交消费确认什么都不做。CREATE TABLE idempotent_record ( id BIGINT AUTO_INCREMENT PRIMARY KEY, biz_key VARCHAR(128) NOT NULL COMMENT 业务幂等键, request_body TEXT COMMENT 原始消息内容便于排查, created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, UNIQUE KEY uk_biz_key (biz_key) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;判断是否重复时不要用先查再说必须利用数据库唯一索引的原子性try { idempotentRecordMapper.insert(record); } catch (DuplicateKeyException e) { // 唯一键冲突说明这条消息已经处理过直接返回 return; } // 插入成功继续执行业务逻辑 doBusiness(message);为什么强调数据库唯一索引而不是select count(*)后再insert因为两个消费者线程可能同时处理同一幂等键两个都查不到记录然后都去执行业务最终重复。唯一索引只有一条能插入成功另一个必然抛异常这才是真正防并发的手段。去重表和业务表的事务关系非常关键。如果业务逻辑本身有数据库操作强烈建议把去重表插入和业务更新放在同一个本地事务里。要么一起提交要么一起回滚。这样处理去重记录写成功、业务更新失败时事务回滚会把去重记录也回滚掉下一次重试还能重新执行不会因为已存在去重记录而跳过真正失败的业务。2.2 乐观锁/版本号并发更新场景下更新型的正确解法去重表适合新增型操作比如落订单流水、记录回调记录。但有些业务是更新型的比如修订单状态、扣库存、改账户余额。这类场景不能单纯靠插一条记录来表示完成需要给数据行加版本号用乐观锁保证更新只生效一次。UPDATE stock SET count count - #{count}, version version 1 WHERE id #{id} AND version #{expectVersion};执行结果只有两种影响行数为1说明这行数据在读取之后没有被别人改过扣减成功影响行数为0说明version已经变了消息被重复处理了或者有并发冲突直接丢弃本次操作。乐观锁的本质是利用同一份业务数据只能被推进一个状态的特性来拦截重复。它不需要额外建表业务表里加一列就行而且条件更新本身具备原子性并发下也不会出问题。不过要注意乐观锁不适用于金额累加这类可交换操作因为连续两次相同金额累加在业务上依然会产生重复影响这类操作后面单独讲。2.3 Redis SETNX/分布式锁高吞吐与性能的权衡去重表要写数据库高并发下会有压力所以很多人会用Redis做前置幂等。SET idempotent:bizKeyorder_12345 value1 NX EX 300NX表示键不存在时才设置成功EX设置过期时间。如果SET返回OK说明这是第一次处理如果返回空说明已经处理过直接跳过业务。Redis方案的优点很明显快抗压能力强而且天然有TTL不用像数据库去重表那样做历史清理。但它有一个隐患Redis的过期机制和持久化能力决定它只能当第一道闸门不能当唯一防线。如果业务处理时间超过TTL锁过期后再来一条重复消息就会穿透。加上Redis主从切换时可能存在少量数据丢失极端情况下幂等会失效。所以我的做法是Redis SETNX 数据库唯一约束双层结合。Redis扛住99%的重复流量数据库唯一索引兜底那1%的极端穿透。性能要好数据也要稳两边的好处都要。2.4 状态机校验让业务本身具备防重放能力订单、支付、退款这类的业务通常都有明确的状态流转比如待支付 - 已支付 - 已发货 - 已完成。如果每次更新订单状态时都带上当前状态必须等于期望状态的条件天然就能拦截重复消息。UPDATE order SET status PAID WHERE order_id #{orderId} AND status WAIT_PAY;第一条消息把状态从WAIT_PAY改到PAID影响行数为1第二条重复消息再执行时订单状态已经是PAIDWAIT_PAY条件不满足影响行数为0直接忽略。这比单纯靠版本号更直观因为状态机本身就把什么阶段能做什么操作讲清楚了。状态机方案适合有明确流程约束的业务但必须跟着业务规则走比如已支付不能再次支付已发货不能重复发货。如果业务本身允许退款后再支付这类状态回跳状态机要设计得更复杂不能简单靠前后置状态判断。四种方案没有银弹核心还是看业务操作的类型。我习惯用一张表来帮助团队做选择方案优点缺点适用场景唯一键去重表可靠性高支持事务可追溯多一次数据库写新增型操作回调记录、流水落库、发券乐观锁/版本号无需额外表并发安全不适合可交换累加操作更新型操作库存扣减、状态更新Redis SETNX性能高支持TTL有穿透和丢数据风险高并发入口配合数据库兜底状态机校验语义清晰契合业务流程状态流转设计复杂订单、支付、退款等强流程业务3. 三个高频业务场景的幂等落地细节3.1 支付回调幂等键宁可多拼不能少拼支付回调是MQ幂等问题的高发区。微信、支付宝回调我们的服务器我们把回调内容投递到MQ再由消费者更新订单状态。一个误判就可能造成重复发货或者重复入账。支付回调消息里字段不少但不是所有字段都适合做幂等键。只拿orderId做幂等键会出问题同一笔订单可能会有支付成功支付失败部分退款多个事件它们都属于同一个orderId如果在去重表里用orderId当唯一键后面到达的部分退款会被当成重复消息丢掉。正确的做法是把业务事件唯一ID做全out_trade_no trade_status或者orderId eventType transactionId。宁可多拼几个字段也不要为了省事导致不同事件互相误伤。有一点容易忽略transactionId是支付渠道侧的流水号如果支付渠道退款时重新生成了退款交易号那退款事件和支付事件之间用out_trade_no refund_flag区分会更稳妥。消费端的伪逻辑应该是String bizKey message.getOutTradeNo() : message.getTradeStatus(); // 1. 去重表插入失败说明处理过 if (!tryInsertIdempotentRecord(bizKey)) { ack(); // 重复消息直接提交 return; } // 2. 查询订单当前状态 Order order orderMapper.selectByOrderId(message.getOrderId()); if (order.getStatus() OrderStatus.PAID) { // 状态已经是终态按成功处理 ack(); return; } // 3. 更新订单状态并记录流水和去重表同一事务 paymentService.markOrderPaid(...);其实支付回调消息本身会携带动账金额如果回调里既有下单支付成功又有商家主动调价后的补差价支付幂等键里还得加上支付金额或者支付单号。核心原则就一句话唯一键要能区分出同一次业务动作而不是只区分同一个订单。3.2 库存扣减靠版本号和条件更新压住并发库存扣减是典型的更新型操作如果去抢购、秒杀场景同一商品的一条扣减消息被重复消费问题会被放大。我见过有同学直接用UPDATE stock SET count count - 1 WHERE id ?第一次执行把库存从100变99第二次重复执行变成98虽然库存足够但已经把别人的库存偷走了。一个可选方案是版本号乐观锁UPDATE stock SET count count - #{qty}, version version 1 WHERE sku_id #{skuId} AND version #{expectVersion};但这个消息重复场景下版本号方案并不完美。如果最终只有一个扣减单号重复消息携带的expectVersion是消费开始前读到的旧版本第一条执行成功会改变版本第二条执行必然影响0行——看起来没问题。问题在于如果两个扣减事件是同一个人同一订单的两次合法扣减比如买了两件不同尺码它们的expectVersion相同版本号会把第二次合法扣减也拦掉。副作用是误杀合法请求。所以我更推荐在扣减库存前先通过扣减流水表做一次唯一约束。每笔扣减生成一个deduct_id流水表加唯一索引插入成功才执行库存更新。这样重复消息到达时流水插入直接冲突不会碰库存数据。try { deductRecordMapper.insert(new DeductRecord(deductId, skuId, qty)); } catch (DuplicateKeyException e) { // 已处理过直接ack return; } int rows stockMapper.deduct(skuId, qty, requireEnough); if (rows 0) { throw new InsufficientStockException(); }流水表唯一约束保证了同一扣减单只处理一次库存条件更新保证了库存不足不会扣成负数。把两个机制配合起来才能既防止重复又保证业务正确。3.3 积分/余额加款查重、锁、事务三件套余额和积分加款属于可交换累加操作如果重复执行两次账户余额就是多出来的钱。这种场景有两个关键问题一是并发二是重复。第一层拦截消费端收到加款消息后先按userId sourceType sourceId去查流水表判断是否已加过。但这种先查后写在并发下不安全两个线程同时查到不存在然后都去执行加款。所以第二层必须用数据库唯一索引CREATE TABLE account_flow ( id BIGINT AUTO_INCREMENT PRIMARY KEY, user_id BIGINT NOT NULL, source_type VARCHAR(32) NOT NULL COMMENT 来源类型订单、活动、退款, source_id VARCHAR(64) NOT NULL COMMENT 来源单号, amount DECIMAL(12,2) NOT NULL, created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, UNIQUE KEY uk_source (user_id, source_type, source_id) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;第三层在同一个事务里先插入流水再更新账户余额余额更新成功依赖流水插入成功。这样即使消费端收到重复消息第二次插入流水时就会触发DuplicateKeyException事务回滚余额不会被第二次累加。有人问要不要用分布式锁Redisson、ZooKeeper这些。我能想到的建议是不要指望分布式锁单独扛幂等。锁有获取锁、释放锁、锁超时各种异常一旦某次重复路径没拿到锁或锁提前释放防线就穿了。分布式锁可以作为并发控制的手段但幂等底线必须落在数据库的唯一索引和事务上。4. 消费端并发与提交时机别让保护层自己先破防4.1 手动ack和自动commit哪个更容易喂重复消息幂等方案再好如果消费端的提交时机不对依然会放大重复消息的数量甚至让合理请求被误杀。Kafka的enable.auto.committrue时消费者每5秒自动提交一次拉取到的offset。如果业务处理耗时较长消息已经被处理完但offset还没等到自动提交就发生了Rebalance或宕机重启后会从上次提交的offset重新消费把处理过的消息再送一遍。关闭自动提交、改用手动commitSync可以在业务处理结束后立刻提交把重复窗口压缩到最小。但要注意手动提交也做不到真正的业务提交原子性。如果业务写入成功了commit之前消费者宕机消息照样会重复投递。所以不要以为改成手动ack就万事大吉它只是缩小了重复窗口真正的防线依然是消费侧的幂等处理。RocketMQ则是消费成功后返回ConsumeConcurrentlyStatus.CONSUME_SUCCESSBroker才会认为消息处理成功。如果返回RECONSUME_LATER消息会被送去重试。RabbitMQ的basicAck是消费者收到消息后显式确认。这三种机制本质上都承认一个事实消费端处理结果不可能与确认动作原子绑定重复是语义的一部分。4.2 多线程消费并发去重的竞态处理开启多线程消费之后同一分区或同一个消息组名下的消息会并行处理。并发场景下插去重表必须用数据库唯一索引拦截。我见过很多团队在并发测试正常、上线后突然出现重复数据原因就是去重逻辑写成了if (recordMapper.selectByBizKey(bizKey) null) { recordMapper.insert(record); doBusiness(); }两个线程同时走到selectByBizKey都返回null然后双双执行insert双双执行doBusiness。虽然业务数据在数据库里因为某些约束可能没炸但幂等逻辑已经名存实亡。正确写法就是前面说的直接insert靠唯一索引的DuplicateKeyException去判断不要在应用层自己做先查后插。多线程还带来一个问题同一业务消息的多个事件之间可能乱序。比如下单和支付两个事件同时被拉取支付事件先处理完下单事件才处理。如果幂等键设计不合理支付事件会把下单事件的去重记录提前占掉导致下单被误判为重复。解决方式要么是幂等键带上事件类型要么在业务逻辑里做状态校验允许事件不按顺序到达。4.3 幂等不等于不重复该怎么对业务方解释我经常被产品同事问你们不是用MQ保证不重复吗为什么我看到了重复的推送通知这里有个概念边界问题幂等性保证的是业务影响只发生一次不是消息只被消费一次。比如一条给用户发100积分的消息幂等方案保证用户只收到100积分不会收到200积分但MQ集群里这条消息可能被消费者取出来处理了两三次只是后几次被去重表挡住了。从用户和账务的视角来看结果是唯一的这就是幂等。面试时如果被问到MQ如何保证消息不重复我建议你先把这句话纠正过来不是保证消息不重复而是保证重复消息不会产生重复的业务结果。这个主动性表述往往比背方案更能体现对问题的理解深度。5. MQ幂等实战中那些最容易翻车的地方5.1 幂等键设计不稳一切白搭幂等方案选得再好幂等键选错整个防线就是纸糊的。我盘点一下最常见的几个错误只用主业务ID忽略了事件类型。同一条记录存在新建修改删除多个事件彼此会被误判为重复。用了时间戳或UUID当幂等键。每条消息都不同去重表形同虚设。依赖消息体里的某一字段但字段本身可能为空。比如退款单号在部分场景下为空去重直接失效。可空字段拼接出来的key不唯一。userId null null和另一个事件的key相同互相覆盖。设计幂等键时我的经验是先问三个问题这个业务事件最原子的标识是什么同一业务主键下会不会有多种事件类型消息重试时这个标识会不会变化三个问题都答清楚幂等键才算合格。5.2 去重表和业务表的写序错了有一个非常隐蔽的坑把去重记录插入放在业务执行之后。比如先更新余额成功后再插入去重表。如果更新余额后、插入去重表前进程崩溃这条消息就没有去重记录。消息被再次投递时余额会再更新一次。重复由此产生。正确顺序永远是把去重表插入和业务写入放到同一事务而且插入动作在业务更新之前。一旦插入成功事务提交后后续重复消息都会被挡掉。如果业务更新失败导致事务回滚去重记录也跟着回滚下一次重试还能重新处理。这里不需要考虑业务成功但去重没插入的问题因为它们在同一个数据库事务里不可能只发生一半。但如果去重表和业务表分属不同数据库就得引入分布式事务复杂度会飙升。我的建议是如果条件允许尽量让业务主库和幂等去重表在同一个实例里一次本地事务全部解决。跨库场景下先写Redis SETNX拦截再用对端MQ或定时任务做对账补偿这是后话。5.3 去重记录活得太短留下的窗口期数据库去重表一般不会无限增长团队都会做定时清理。清理策略如果没有计算好保留时间会留下一个危险的窗口期消息正因为某种原因在重试队列里延迟投递你去重表里的记录已经删掉了重试时消息直接穿透去重业务被重复执行。去重记录至少应该保留多久有一个简单的下限公式去重记录保留时间 消息最大重试间隔 最长业务处理时间 冗余Kafka默认的offset保留时间是7天如果消费者长时间离线后恢复会重新消费7天内的所有消息。RocketMQ的重试队列默认最多重试16次最大延迟级别可以到2小时左右。如果定时清理每天只留24小时遇到周末或长假积压的延迟消息重复概率会明显上升。我的实践是核心账务类去重表至少保留30天即使量大了要做归档也要把去重记录的TTL和归档策略分开防止归档后消息重试穿透。还有一个和TTL配套的细节清理去重记录时不要把biz_key一并删掉可以做逻辑过期——加一个expired_at字段定时任务只把记录标记为过期。万一穿透导致数据异常至少能通过保留的biz_key追踪到原始消息和消费日志。5.4 最后的小技巧把幂等键写进每次日志里排查幂等问题最难的一点是确认这条消息到底是重复的还是合法的新事件。所以我在所有涉及MQ消费的业务日志里强制打印三样东西消息唯一键MQ自带、业务幂等键自定义、消费结果首次/重复/忽略。日志格式固定、关键字统一线上出问题用一条grep就能串起整个调用链。幂等设计不是上线前临时补的而是每个消息消费入口都必须回答的问题。我现在的习惯是每次评审消费相关需求时先问这条消息重复消费会产生什么后果如果答案是后果很严重就老老实实把去重表或乐观锁做上如果答案是重复了也没影响至少要在代码里写清楚为什么可以不做。所有看似没必要的重复迟早会在流量高峰期还给你一个惊喜。