做后端这些年,我被问得最多的一个问题就是:消息队列到底怎么保证消息不被重复消费?标准答案比较扎心:做不到。真实答案其实是:不要指望消息不重复,而是要把重复消息变成无副作用的操作。这个词就是幂等。刚带团队那会儿,我们做过一个百万级订单同步系统,上线第二周就撞上一次事故:一条库存扣减消息被重复消费了两次,两万多件商品的库存被多扣了两万多。排查了一整夜,最后定位到是消费者处理完业务后还没来得及提交 offset,进程就崩了,重启后消息被重新投递。从那天起我彻底想明白了:在分布式系统里,重复是常态,不重复才是巧合。我们要做的,不是去消灭重复,而是在架构上让重复变得无害。
这篇内容适合正在做消息队列、后端服务、支付对账、订单同步这类系统的同学参考。无论你用的是 Kafka、RocketMQ、RabbitMQ 还是别的 MQ,只要涉及“消费消息后写数据”的场景,幂等设计就是一道绕不过去的坎。我会把重复消息的来源、幂等方案的选型、一套完整的落地方案,以及我踩过的坑全部拆开讲,尽量做到看完了能直接参考。
1. 为什么消息一定会重复,这不是 Bug 是宿命
很多刚入行的同学会很自然地觉得:消息队列不应该保证消息不丢、不重吗?这个直觉没错,但现实世界的消息队列,在“不重”这件事上几乎没有谁能给出绝对承诺。这不是 MQ 厂商不给力,而是分布式系统的网络模型决定了:你无法精确区分“消息没发出去”和“消息发出去了但响应丢了”。
1.1 投递语义:At-Most-Once、At-Least-Once、Exactly-Once
要理解重复,先要理解消息队列的三种投递语义。
At-Most-Once(最多一次)模式下,消息可能丢,但不会重复。这种语义通常用于允许丢失数据的日志采集场景,一般不会用在对账、支付这类对准确性敏感的业务里。
At-Least-Once(至少一次)模式下,消息不会丢,但可能重复。这是目前绝大多数主流消息队列的默认行为。Kafka 默认是 At-Least-Once,RocketMQ 的普通消息也是,RabbitMQ 配合手动 ack 也是这个语义。
Exactly-Once(精确一次)听起来最理想,但实现代价极高,而且往往只适用于单一消息队列内部的某个特定场景。比如 Kafka 的 exactly-once 需要配合事务 API 和幂等生产者,但它解决的是“生产者到 broker”、“消费者到 broker”这一段的问题,跨系统的端到端精确一次几乎不可能做到。
所以结论很直接:你只要在用 MQ 做业务系统,就默认活在 At-Least-Once 的世界里。重复不是异常,是语义的一部分。
1.2 重复消息的三个典型来源
结合我排查过的线上问题,重复消息的来源基本能归成三类。
第一类是生产者重试。我们在业务代码里发消息时,经常会设置一个 retry 次数。假如发送超时了,生产者的第一反应是重试再发一次。但这里有个要命的细节:超时到底超在哪儿了?可能是网络抖动,消息根本没到 broker,也可能是消息已经到了 broker,只是响应包丢了。如果是后者,生产者重试就会造成 broker 里存了两条一模一样的消息。这是最隐蔽的来源,因为你无法从代码层面判断响应丢失。
第二类是消费者已经处理完了,但 offset/ack 没有提交成功。拿 Kafka 举例,消费者处理完业务逻辑后,需要提交 offset 才能告知 broker“这条消息我消费完了”。如果业务代码抛出异常导致 offset 没提交,这条消息会在下一次 poll 时被再次拉取;如果进程直接宕机,Kafka 会根据上次提交的 offset 重新分配分区,把那一批消息重投给新的消费者。这就是最经典的“业务成功了但没提交”场景,也是我第一次事故的根因。
第三类是消息队列自身的重投机制。比如 RocketMQ 的消费重试机制,业务抛异常后会进入重试队列,间隔一段时间重新投递。Kafka 在分区副本切换、消费者组 Rebalance 时也可能出现少量重复投递。这些都不是配置错了,而是分布式系统内在的容错机制,它们为了“不丢消息”选择了“允许重复”。
1.3 前端点两下和重复消费是一回事吗
热搜里有句话叫“前端点两次算是发两条消息吗”。从接口层看,前端双击按钮确实会触发两次 HTTP 请求,后端就会收到两个请求。很多团队只做了前端按钮置灰,就以为万事大吉,实际上这只能挡住“正常人”的操作,挡不住网络重放、超时重试、脚本并发。退一步说,就算前端只发了一条消息,后端的重复消费也一样会发生。所以前端防重和后端幂等不是二选一,而是两层防线。前端负责体验,后端负责兜底。
2. 幂等方案全景对比,从数据库唯一键到 Redis 锁
聊完了“为什么重复”,接下来进入核心:怎么让重复消息不产生副作用。这一节的本质,是给你一张方案地图,你对照自己的业务场景去选就行。
2.1 幂等的本质,不是防重复而是防副作用
幂等(Idempotency)这个概念来自数学,f(f(x)) = f(x) 就是幂等。翻译成人话:同一个操作不管执行一次还是执行一百次,结果完全一样。
读接口天然幂等,查十次和查一次结果没区别。但写接口就不一定了:创建订单、扣库存、加积分、转账,这些操作执行两次就会产生两份订单、扣两次库存、加两次积分、转两次账。所以幂等设计的核心对象,永远是有副作用的状态变更操作。
你要记住一个判断标准:你的业务操作是否天然幂等?如果是,比如“把订单状态置为已关闭”这种覆盖式更新,那不需要额外处理;如果不是,比如“账户余额增加 100”,那就必须引入幂等机制。
2.2 方案一:数据库唯一键约束,生产环境最推荐
在所有幂等方案里,我会首选数据库唯一索引,没有之一。它的思路很直白:在业务表或者专门的幂等表上,给“唯一业务标识”加唯一索引,第二次插入时数据库会直接报唯一键冲突,从根上截断重复流量。
举个例子。你有一张订单表,业务上允许用户对同一笔支付请求重复发起确认操作,那你可以把biz_id字段设为唯一索引。第一次插入成功,第二次插入抛 DuplicateKeyException,你捕获异常后直接返回成功就行。
这个方案最大的优势是:正确性完全由数据库保证,没有并发竞态。两个消费者同时拿到同一条消息,同时插入同一个唯一键,数据库只会让一个成功,另一个必然冲突。不需要分布式锁,不需要额外组件,简单可靠。
2.3 方案二:状态机幂等,跟着业务状态走
如果你的业务本身有明确的状态流转,可以用状态机来兜底。比如订单状态:待支付 -> 已支付 -> 已发货 -> 已完成。消费消息时,先查当前状态,判断目标状态是否合法。如果消息要求把“已支付”的订单改成“已支付”,那就说明是重复消息,直接忽略;如果要求从“已支付”流转到“已发货”,才继续处理。
这种方案的好处是直接贴合业务语义,代码可读性好。坏处是它只能防住“状态已经变化后”的重复消息,挡不住同一时间点的并发重复。所以严格来说,状态机幂等更适合作为辅助手段,和其他方案叠加使用。
2.4 方案三:Redis 缓存去重,高性能但有时间窗
用 Redis 做去重也很常见。流程是:消费消息时,先用SET key value NX EX 3600把消息 ID 写进 Redis,能写入说明是第一次来处理,继续业务逻辑;写入失败说明已经处理过了,直接跳过。
这套方案性能极高,适合 QPS 非常大的场景,但它有个很难绕开的缺陷:过期时间。你设置 3600 秒过期,那 3600 秒之后呢?如果消息延迟重投了,Redis 里已经没有记录,就会再处理一次。为了解决这个问题,很多人会把过期时间拉得很长,但拉太长又占内存。所以纯 Redis 方案只能做到“大概率不重复”,做不到“绝对不重复”。
2.5 方案四:分布式锁,防止并发但不解决重复
用分布式锁(比如 Redis 的 Redisson、ZooKeeper)也可以处理重复消息:消费前先加锁,拿到锁才处理业务,处理完释放锁。但要注意,分布式锁解决的是“并发冲突”问题,不是“重复消费”问题。如果消息 A 在 10 点处理完了并提交了 offset,之后又因为某种原因在 11 点被重新投递,此时锁早已经被释放,分布式锁根本拦不住这条迟到的重复消息。
所以分布式锁只能作为并发控制工具,不能单独当作幂等方案用。
2.6 方案五:Token 机制,接口幂等的经典套路
Token 机制是接口幂等设计里的老办法,对应热搜里的“接口幂等性设计”。核心思路是:前端在发起请求前,先向后端申请一个唯一的 token;后端把 token 存到 Redis 里;前端提交业务请求时带上这个 token;后端在执行业务前,先去 Redis 删除这个 token,能删掉就说明是第一次请求,删不掉就说明是重复请求,直接拒绝。
这个方案对“点击按钮提交订单”这类场景很合适,但它需要改造前端交互,而且对消息队列场景不太适用——因为真正的中台消费逻辑不会在每次执行前去申请 token。
2.7 方案对比与选型建议
| 方案 | 原理 | 优点 | 缺点 | 适用场景 |
|---|---|---|---|---|
| 数据库唯一键 | 唯一索引拦截重复插入 | 最可靠,无并发竞态 | 依赖数据库,需设计幂等表 | 绝大部分写业务,首推 |
| 状态机 | 业务状态校验 | 贴合业务语义 | 挡不住并发重复 | 订单、审批等状态流 |
| Redis 去重 | SETNX + 过期时间 | 性能高 | 过期后有重复窗口 | 高频请求,容忍低概率重复 |
| 分布式锁 | 并发互斥 | 防止并发冲突 | 不解决重投 | 并发控制辅助手段 |
| Token | 预发令牌校验 | 用户体验好 | 需改造前端 | 接口幂等、表单提交 |
我的建议很简单:能用数据库唯一键,就别整花活。先把唯一键方案落地,再根据性能压力去叠加 Redis 去重做前置拦截。
3. 一套可以直接落地的幂等消费链路
方案看再多,落地才是关键。这一节我拿出一个真实项目里跑过的模式,从表结构到消费代码到事务边界,一步步拆给你看。
3.1 整体链路设计:消息 ID 贯穿始终
先说整体链路。为了保证“百万条消息一条都不重复处理”,核心思路是在消息生命周期的每个环节都带上一个全局唯一的业务键。
生产端发送消息时,在消息体里带一个业务唯一 ID,比如bizId,它可以是订单号、支付单号、用户操作流水号。如果生产框架支持消息 key,也把bizId设置到 key 上。消费端拿到消息后,第一件事不是处理业务,而是拿bizId去幂等表里尝试插入。
整个链路就是:生产者生成业务唯一 ID -> 发送消息到 MQ -> 消费者收到消息 -> 幂等表登记 -> 业务处理 -> 提交 offset/ack。链路里最关键的一环,就是幂等表登记。
3.2 幂等表设计与消费代码实现
幂等表的设计很简单,核心就三个字段:自增主键、业务唯一键、创建时间。下面是我常用的建表语句:
CREATE TABLE `idempotent_record` ( `id` BIGINT NOT NULL AUTO_INCREMENT COMMENT '自增主键', `biz_id` VARCHAR(64) NOT NULL COMMENT '业务唯一键,对应消息里的业务ID', `biz_type` VARCHAR(32) NOT NULL COMMENT '业务类型,区分订单、支付、库存等', `created_at` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间', PRIMARY KEY (`id`), UNIQUE KEY `uk_biz_type_biz_id` (`biz_type`, `biz_id`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='幂等记录表';注意唯一索引要建成(biz_type, biz_id),因为不同业务类型下的同一个 ID 可能代表完全不同的含义,直接全局唯一容易误伤。
消费端的核心逻辑我拿 Java 伪代码写一下,思路同样适用于 Python、Go 等其他语言:
public void onMessage(Message msg) { String bizId = msg.getBizId(); String bizType = msg.getBizType(); try { idempotentRecordMapper.insert(bizType, bizId); } catch (DuplicateKeyException e) { log.warn("重复消息,直接跳过. bizId={}", bizId); return; } try { // 这里开始处理真正的业务逻辑:扣库存、更新订单、加积分等 doBusiness(msg); } catch (Exception e) { // 业务处理失败,抛出异常,触发 MQ 重试 throw e; } }这段代码最关键的一步是:永远不要用“先 select 再判断”的方式做幂等。两个消费者线程同时来处理同一条消息,如果都是先查幂等表,发现不存在,然后都往业务表插入数据,那就都进来了。唯一的“查重”操作必须由数据库唯一索引来判定,也就是直接 insert,靠冲突来判断重复。
3.3 事务边界:幂等记录和业务操作必须同生死
这是我在生产环境踩过最大的坑,没有之一。
你可能会想:先插入幂等记录,再执行业务操作,如果业务操作失败,幂等记录还在,后续重试就会直接跳过,那重试机制不就废了吗?对,如果幂等记录和业务操作不在同一个事务里,就会产生这个矛盾。
正确的做法是:幂等表插入和业务数据变更必须在同一个本地事务里。要么一起成功,要么一起回滚。业务失败时幂等记录也一起回滚,MQ 重新投递后还能再次尝试;业务成功时幂等记录跟着提交,后续重复消息再来了,唯一索引直接拦截。
用伪代码表示就是:
@Transactional public void handleMessage(Message msg) { // 在同一个事务里:先插入幂等记录,再执行业务 idempotentRecordMapper.insert(msg.getBizType(), msg.getBizId()); // 业务操作和幂等记录在同一事务内 orderMapper.updateStatus(msg.getOrderId(), "PAID"); stockMapper.deduct(msg.getProductId(), msg.getCount()); pointService.add(msg.getUserId(), msg.getPoint()); }这里有个很多人忽略的细节:insert幂等记录和updateStatus业务操作放在同一个@Transactional方法里,如果后续deduct或add抛异常,整个事务回滚,幂等记录也会消失。下次 MQ 重试时重新进入这个方法,再次尝试插入幂等记录,如果这次业务成功了,记录才真正留下。
这样就形成了一个闭环:重复消息被唯一索引挡住,失败消息可以重试,永不丢数据。
3.4 消费确认:手动提交 offset 才是安全的
除了事务边界,消费确认机制也直接影响幂等效果。拿 Kafka 举例,我强烈建议你关闭自动提交,改用手动提交。为什么?自动提交是周期性提交的,可能你业务还在处理中,offset 已经被提交了。这时候如果业务抛异常,消息重投了,那重复窗口就更难控制了。
手动提交的正确时机是:业务逻辑全部处理完成之后再提交。伪代码如下:
while (true) { ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100)); for (ConsumerRecord<String, String> record : records) { handleMessage(record); } // 这一批全部处理成功才提交 consumer.commitSync(); }注意,这里有一个 trade-off:如果你一条一条地提交 offset,性能会下降;如果一批提交,这一批里有一条失败,重试时整批都会拿回来。但没关系,反正幂等表会兜底,重复了也只是走一遍唯一索引拦截,消耗不大。这就是为什么幂等设计能让你的 offset 策略变得简单粗暴——不怕重复,才敢大批次提交。
4. 百万消息场景下的性能与一致性权衡
消息量一上来,“怎么做”只是第一步,“怎么做还扛得住”才是真正的考验。这一节说几个高并发下你必须想清楚的权衡。
4.1 幂等表会不会成为性能瓶颈
很多人听到“每条消息先插一条幂等记录”的第一反应是:这得多少额外写入量?会不会把数据库拖垮?
我的实测经验是:不会。幂等表本身字段少、索引简单,单条插入的成本极低。在普通 SSD 加 InnoDB 的配置下,单表每秒插入几千条是完全没问题的。如果消息量真的很大,比如每秒几万条,你可以按业务键做分表,比如idempotent_record_0到idempotent_record_15,用bizId.hashCode() % 16路由。因为幂等表只做插入和唯一键冲突检查,分表后逻辑完全不变。
真正需要关注的反而是另一个问题:幂等表会无限增长。百万、千万、上亿条记录攒下来,即使索引能扛住,磁盘空间和备份时间都是负担。建议定期归档,比如保留最近 90 天的记录,更早的迁移到归档表。因为幂等本质上是“近期重复消息”的拦截器,太老的重复消息一般不会出现,偶尔出现的也会被业务状态机拦住。
4.2 先处理业务还是先插入幂等表
有人会纠结:业务操作和幂等记录插入,顺序到底是先业务还是先幂等?我的建议是:幂等插入先做。理由很简单,幂等记录是“闸门”,先关上门再干活,如果后面失败了再回滚开门。如果你先做业务再插入幂等记录,两条并发消息同时进来,就可能出现双方都通过了业务校验、都去改数据的情况,幂等记录根本来不及拦。
而且从性能角度讲,先 insert 幂等记录能让你尽早拿到数据库的唯一索引仲裁结果,重复消息直接 return,后续的业务查询、状态判断全部省掉。
4.3 接口幂等性设计和 MQ 消费幂等是一回事吗
热搜词里高频出现的“接口幂等性设计”,其实和 MQ 重复消费问题底层原理完全一致:都是对同一个操作加“一次性令牌”。
差异在于入口不同。接口幂等的入口是 HTTP 请求,你可以要求调用方在 Header 里带一个幂等键,或者像前面说到的 Token 机制,后端用 Redis 校验一次。MQ 消费幂等的入口是消息中间件,你能拿到的“令牌”就是消息里带的bizId或者消息本身的msgId。
我的实践建议是:如果团队内部有统一的 RPC 或 HTTP 框架,可以直接在框架层面做一个幂等注解,配合 Redis 拦截重复请求;MQ 消费则一律走数据库幂等表,因为它是最终的存储事实,Redis 挂了不会影响正确性。
4.4 延迟消息和定时任务场景下的特殊处理
还有一种场景容易忽略:延迟消息。比如订单超时未支付关闭订单,这类消息可能延迟 15 分钟、30 分钟才投递。如果业务上有多个延迟任务针对同一个订单,就可能在时间窗口内产生重复。
这时候单靠唯一的bizId还不够,最好把“执行时机”也纳入幂等判断。比如幂等表里除了biz_id,再加一个execute_time字段,唯一索引变成(biz_type, biz_id, execute_time)。这样同一个订单的 30 分钟延迟消息和 60 分钟延迟消息可以各自执行一次,互不干扰。
5. 常见问题与排查技巧实录
方案讲完了,最后分享一些我真正在线上踩过的坑和排查思路。这些内容不一定出现在官方文档里,但实战时非常要命。
5.1 两个进程同时消费到同一条消息,为什么唯一索引没拦住
这个问题的答案在事务隔离级别上。很多团队的 MySQL 默认隔离级别是REPEATABLE READ,但这不影响唯一索引的唯一性判定。唯一索引的冲突检测是数据库存储引擎层干的活,和事务隔离级别无关,两个并发插入同一个唯一键,必然只有一个成功。
如果你真的遇到了“两边都成功了”,那就要检查你的幂等表唯一索引到底建没建上,或者你是不是先 select 后 insert 了。后者是最大的坑:查询判断不存在,然后两边同时插入,都没撞上唯一索引——因为唯一索引不存在呀。必须直接 insert,让数据库仲裁。
5.2 重复消息被跳过,但业务实际没处理成功
这是一个很隐蔽的坑。假设你用了“先插入幂等记录,再处理业务,同一个事务”的方案,那么事务回滚时幂等记录也会消失,这种场景不会出现。但如果你把幂等插入和业务处理分成了两个事务——比如幂等表单独一个事务先提交了,业务事务再提交——那业务失败回滚后,幂等记录已经留下了,下次重试直接跳过,这条消息就彻底丢了。
解决方案就是前面强调的:同一个本地事务。不要在幂等表插入成功后立刻提交,一定要让幂等记录和业务操作在同一个事务里同生共死。
5.3 Redis 去重的时间窗问题怎么解
我见过不少团队只用 Redis 做幂等,每次都纠结过期时间设多久。设短了怕重复消息漏进来,设长了怕 Redis 内存爆炸。
我的建议是:不要试图通过拉长 TTL 解决这个问题。更稳妥的做法是 Redis 挡住 99% 的重复流量,数据库唯一索引兜住最后 1%。也就是说,Redis 去重只是前置拦截,真正的可靠保证还是落在数据库上。这样 Redis 的 TTL 可以设得很短,比如 1 小时,过期了也无所谓,数据库会兜底。
5.4 排查重复消息的具体操作路径
如果线上真的出现了“业务重复处理”的事故,我的排查顺序是这样的:
先看消息日志。消费端在入口处记录msgId、bizId、消费时间、处理结果。通过日志检索同一个bizId出现了几次,能快速判断重复消息从哪来的。如果日志显示第一次处理成功后提交 offset 失败,那就是消费者提交时机问题;如果第一次处理还在超时状态,第二次就进来了,那就是并发场景没控制好。
再看幂等冲突日志。在捕获DuplicateKeyException的分支里,一定要打一条 WARN 日志,记录bizId和幂等表的冲突时间。平时这些日志可能没人看,但出事故时它们就是铁证。
最后看 MQ 的重试记录。RocketMQ 控制台能看到消费轨迹,Kafka 可以查 consumer lag 和提交失败的日志。把这三方面的信息拼起来,基本能还原重复消息的完整路径。
5.5 监控指标建议
建议给消费链路加上几个监控指标:消费 TPS、消费失败重试次数、幂等冲突次数、重复消息率。其中幂等冲突次数是最值得关注的,它不代表任何异常,但能客观反映系统的重复消息压力。如果这个指标突然飙升,八成是上游生产者或消费者出了问题,提前介入比等事故爆发强得多。
我个人在实际项目里的体会是:幂等设计这件事,方案本身并不难,难的是在每一个“看起来不会重复”的角落依然坚持加幂等。很多事故都发生在你以为不会重复的地方。所以做设计时别赌人性,别赌网络,哪怕多一行唯一索引,也比你事后对账救火轻松得多。最后再分享一个小细节:所有涉及幂等的核心表,都建议把唯一索引命名字段写得规矩一些,比如uk_biz_type_biz_id,这样排障时一眼就能认出哪个索引在兜底,省下的都是真金白银的时间。