news 2026/9/12 4:18:32

MQ幂等性实战:重复消息产生的原理与四大去重方案

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
MQ幂等性实战:重复消息产生的原理与四大去重方案

凌晨一点被电话叫醒,线上报了一个"用户收到两条扣款通知"的问题。拉完流水之后定位到原因:订单表里同一个支付回调事件被消费端处理了两遍,第一遍正常入账,第二遍又把金额累加了一次。这不是网络抖动,也不是代码逻辑写错,而是消息队列在"保证不丢"的同时,天然会把同一条消息重复投递给消费者。这个问题的学名,就是消息幂等性

消息队列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配置了retries=3,第一次发送实际成功了,但响应包丢了,Producer重试又发了一次,同一条业务消息就出现了两份。

Consumer端消费成功但没来得及提交。这是最经典的重叠窗口。消费者处理完业务逻辑,正要提交offset或发送ack时,进程宕机、被重启、或者网络闪断。等消费者恢复后,Broker会把它当成"这条消息还没消费成功",重新投递一次。Kafka默认的enable.auto.commit=true,每5秒自动提交一次,这5秒窗口内挂掉,几乎必然产生重复消费。

Rebalance导致的重复。Kafka消费者组发生重平衡时,分区会重新分配,某些分区的最新offset可能回退到上一次提交的位置,之前已经消费过但还没提交offset的消息会被重新拉取。消费线程越多,重平衡越频繁,重复概率越高。

RocketMQ的重试队列。消费失败后,RocketMQ会把消息投递到重试队列,按延迟级别再次投递,默认可以重试16次。如果第一次消费时业务已经成功,但返回结果因为网络原因没送到Broker,这条消息还是会被重试投递。

以前我总以为"重复消费"只有在故障时才会出现,后来发现正常运行时也会因为超时、重平衡、并发变更产生重复。所以技术方案上必须把重复当成常态来设计。

1.3 为什么"MQ自带的MsgId"当不了幂等护身符

有人会问:RocketMQ每条消息都有msgId,拿它做去重不就完了?我最早也这么干过,后来发现这是个坑。

RocketMQ的msgId是Producer发送时生成的客户端ID,如果Producer重试发送同一条业务消息,msgId会重新生成,两个msgId不同,但业务内容完全一样,去重判断直接失效。Kafka则没有全局唯一的消息ID,offset只能唯一标识分区内的位置,消费者组变化后offset语义不稳定。RabbitMQ的deliveryTag是Channel级别的自增序号,消息重新入队后也会变化。

所以,MQ自带的ID只能用来查日志和排错,绝不能用来做业务幂等。真正靠谱的,是业务侧自己定义的唯一标识,比如orderIdpaymentIduserId + 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) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

判断是否重复时,不要用"先查再说",必须利用数据库唯一索引的原子性:

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:bizKey=order_12345 value=1 NX EX 300

NX表示键不存在时才设置成功,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) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

第三层,在同一个事务里先插入流水,再更新账户余额,余额更新成功依赖流水插入成功。这样即使消费端收到重复消息,第二次插入流水时就会触发DuplicateKeyException,事务回滚,余额不会被第二次累加。

有人问要不要用分布式锁,Redisson、ZooKeeper这些。我能想到的建议是:不要指望分布式锁单独扛幂等。锁有获取锁、释放锁、锁超时各种异常,一旦某次重复路径没拿到锁或锁提前释放,防线就穿了。分布式锁可以作为并发控制的手段,但幂等底线必须落在数据库的唯一索引和事务上。

4. 消费端并发与提交时机:别让保护层自己先破防

4.1 手动ack和自动commit,哪个更容易喂重复消息

幂等方案再好,如果消费端的提交时机不对,依然会放大重复消息的数量,甚至让合理请求被误杀。

Kafka的enable.auto.commit=true时,消费者每5秒自动提交一次拉取到的offset。如果业务处理耗时较长,消息已经被处理完,但offset还没等到自动提交就发生了Rebalance或宕机,重启后会从上次提交的offset重新消费,把处理过的消息再送一遍。关闭自动提交、改用手动commitSync,可以在业务处理结束后立刻提交,把重复窗口压缩到最小。

但要注意,手动提交也做不到真正的"业务+提交"原子性。如果业务写入成功了,commit之前消费者宕机,消息照样会重复投递。所以不要以为改成手动ack就万事大吉,它只是缩小了重复窗口,真正的防线依然是消费侧的幂等处理。

RocketMQ则是"消费成功后返回ConsumeConcurrentlyStatus.CONSUME_SUCCESS",Broker才会认为消息处理成功。如果返回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就能串起整个调用链。

幂等设计不是上线前临时补的,而是每个消息消费入口都必须回答的问题。我现在的习惯是,每次评审消费相关需求时,先问:"这条消息重复消费会产生什么后果?"如果答案是"后果很严重",就老老实实把去重表或乐观锁做上;如果答案是"重复了也没影响",至少要在代码里写清楚为什么可以不做。所有看似"没必要"的重复,迟早会在流量高峰期还给你一个惊喜。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/12 4:14:57

牙科就诊管理系统:SpringBoot+Vue3+MyBatis技术架构解析

1. 项目概述:牙科就诊管理系统的技术架构与核心价值这个牙科就诊管理系统采用了当前企业级开发中最主流的"前后端分离"架构方案。前端基于Vue3的Composition API实现响应式界面,后端采用SpringBoot快速构建RESTful API,数据持久层使…

作者头像 李华
网站建设 2026/9/12 4:14:51

电钢琴选购指南:从键盘到音源的全面解析

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/12 4:14:43

MFA安全新挑战:IDN同形攻击与零宽字符钓鱼防御

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/12 4:12:18

论文的方法论怎么选对?一篇讲透从问题到方法

论文方法论总选错,多半不是方法名不好听,而是它跟你的问句、你的资料、你后续要用的分析没接上。这里把三条对齐线与四类错配形态摊开,帮你判断该动哪一边。从问题怎么一步步推到方法、几类方法怎么选、两类方法怎么结合,这些各有…

作者头像 李华