ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

消息队列重复消费怎么办?幂等设计原理与落地方案解析

消息队列重复消费怎么办?幂等设计原理与落地方案解析 做后端这些年我被问得最多的一个问题就是消息队列到底怎么保证消息不被重复消费标准答案比较扎心做不到。真实答案其实是不要指望消息不重复而是要把重复消息变成无副作用的操作。这个词就是幂等。刚带团队那会儿我们做过一个百万级订单同步系统上线第二周就撞上一次事故一条库存扣减消息被重复消费了两次两万多件商品的库存被多扣了两万多。排查了一整夜最后定位到是消费者处理完业务后还没来得及提交 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-OnceRocketMQ 的普通消息也是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) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4 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) { ConsumerRecordsString, String records consumer.poll(Duration.ofMillis(100)); for (ConsumerRecordString, 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这样排障时一眼就能认出哪个索引在兜底省下的都是真金白银的时间。
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进