ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

RocketMQ 消息重复消费怎么办?四种去重方案对比与实战

RocketMQ 消息重复消费怎么办?四种去重方案对比与实战 1. 这事到底坑在哪消息重复从来不是“万一”做后端开发的只要系统里接了 RocketMQ早晚会被同一个问题找上门消费端收到了同一条消息两次而且业务被重复执行了。我最早遇到这事是在一个订单支付的场景用户支付成功之后回调通知走 MQ 触发积分发放结果某次 Broker 抖动消费端把同一条支付成功消息处理了两遍用户的积分凭空多了一倍。线上事故复盘的时候我的第一反应是“RocketMQ 出 Bug 了”后来翻了一堆文档和源码才发现重复消息是分布式消息队列的固有属性它不是异常而是常态。RocketMQ 官方定义的投递语义是At Least Once至少一次也就是说消息从 Producer 发送到 Broker再从 Broker 投递给 Consumer这个过程里任何一个环节的确认消息丢失都会触发重试和重投。正是这个“至少一次”的设计保证了消息不丢但也从根上决定了“重复”必然存在。这篇内容我会把消息重复的成因、去重的整体设计、代码层面的落地方式、生产环境的排查技巧全部串起来讲一遍。如果你是刚接触 RocketMQ 的开发者或者已经在生产上被重复消费坑过这篇文章可以帮你少走很多弯路。我会把我实际用过的方案贴出来包括 Redis 去重、数据库唯一键兜底、状态机去重等并说明它们各自的适用场景。先说一个核心结论RocketMQ 本身不提供消息去重能力去重是消费端必须自己解决的问题。去重的本质不是一个工具而是一套策略。你要根据业务容忍度、并发量、可用性要求来选择哪种策略而不是找到一个“万能钥匙”然后到处套用。2. 为什么消息会重复从发送端到消费端的全链路分析2.1 Producer 端重试发送确认丢了Producer 把消息发到 BrokerBroker 写入 CommitLog 之后会返回一个发送结果。问题在于如果 Broker 已经成功写入但返回给 Producer 的响应在网络传输中丢失了Producer 会认为发送失败从而触发重试把同一条消息再次发送给 Broker。这样 Broker 上就存在两条内容完全一样的消息它们的 Message ID 都不同。举个例子用户下单成功后订单服务向 MQ 发送一条“订单创建成功”的消息。Broker 收到了也返回了 ACK但 ACK 包在网络里丢了。订单服务的发送端超时后重试又发了一次。Broker 上就有了两条内容一样的消息。消费端如果不去重支付积分、短信通知、物流下单这些逻辑就会执行两次。这里有个容易忽略的细节Producer 重试发送时如果把setKeys或者业务主键带上两条消息的 keys 是一样的但 Message ID 是不同的。所以用 Message ID 做去重是无效的必须用业务唯一标识。2.2 Broker 端重投消费成功但 ACK 丢了RocketMQ 的消费模式分为集群消费和广播消费。集群消费模式下消费端拉取消息之后需要在处理完成后向 Broker 返回消费成功的确认。如果在处理完业务逻辑之后、ACK 返回之前消费端进程宕机了或者网络出现分区Broker 就认为这条消息还没被消费成功会把消息重新投递给其他消费者。这就是最经典的一个场景业务逻辑已经执行了但消息被重复消费了。而且这种情况比 Producer 端重试更常见因为它不要求发送端有任何异常纯粹是消费端的 ACK 机制导致的。在 RocketMQ 的 PushConsumer 实现里消息处理完成后会调用processQueue的removeMessage并返回ConsumeConcurrentlyStatus.CONSUME_SUCCESS。如果在这个返回之前 JVM 崩了Broker 过一会就会把消息重新投递出来。所以如果你用了 RocketMQ 的 PushConsumer你的消费逻辑必须天然具备幂等性或者你主动去重。2.3 Rebalance 带来的重复消费RocketMQ 集群模式下Consumer Group 内部会做负载均衡把 Topic 下的队列分配给每个消费者实例。当消费者实例发生变化——比如新实例上线、老实例宕机、实例所在机器负载过高触发心跳超时——Broker 会触发一次 Rebalance把队列重新分配。Rebalance 的瞬间某些消息可能已经被原消费者拉取到本地内存但还没处理完。重新分配之后这些消息会被分配给新的消费者实例新消费者会从 Broker 重新拉取这些消息进行消费于是同一个消息被两个消费者先后处理了一遍。即使你的消费者只有一个实例在发生 Rebalance 的时候也有可能把自己本地的消息“让”出去再拉回来造成自重复。我在实践中发现频繁的 Rebalance 往往是消费端 GC 停顿过长或者网络抖动引起的。每次全量 GC 超过几十秒心跳超时Broker 就把这个消费者踢出群组触发 Rebalance恢复后又触发一次。这个现象在 JVM 堆开得特别大的服务上特别明显。2.4 死信队列重投导致的重复当一条消息消费失败次数超过阈值默认 16 次RocketMQ 会把它投递到死信队列。如果后续你写了代码去消费死信队列里的消息并且消费成功了这条消息对应的原始业务也已经执行过了——于是又产生了重复。这种场景比较隐蔽因为死信队列里的消息常常是脏数据或者格式异常的消息很多人不会专门去消费它。但一旦你做了死信告警手动重发原始消息就会出现业务重复。所以处理死信消息之前务必先确认其业务状态再做补偿而不是盲目重发。2.5 延时消息到期后的重复投递RocketMQ 的定时/延时消息在到期后会由 Broker 把消息从延时队列重新写入原始 Topic。在这个“重新写入”的过程中如果 Broker 在写入后、返回结果前发生异常会导致同一个延时消息被写入两次。消费端就会收到两条相同的延时消息。这种场景不常见但也不是没有。我们曾经遇到过延迟关闭订单的任务重复触发排查了很久之后才发现是延时消息在 Broker 端重复写入导致的。所以不要以为延时消息就比普通消息更“安全”。3. 去重方案怎么选四种常见策略的对比与适用场景3.1 唯一键约束数据库层面的硬边界这是最可靠、最保底的方案。在订单、支付、流水这类核心业务表中设计一个唯一业务号字段比如biz_id然后建唯一索引。消费消息时先尝试插入一条记录如果插入成功说明这条消息是第一次处理继续执行后续业务如果插入时抛出唯一键冲突异常说明这条消息已经被处理过了直接跳过即可。这个方案的优点是彻底、绝不会漏因为它依赖数据库的强一致性约束不依赖 Redis 的过期策略也不依赖分布式锁的可靠性。缺点是有性能开销每条消息都要执行一次 INSERT 操作在高并发场景下会成为瓶颈。我的经验是唯一键约束应该作为兜底方案放在业务的核心链路上。比如用户下单、支付回调、发放权益这些直接涉及资金和资产的操作必须用数据库唯一键来做最终的幂等保障。因为 Redis 去重可能存在数据丢失或过期的问题而数据库的唯一索引是最底层的防线。建表的时候要注意一个细节唯一索引字段不能允许 NULL。MySQL 的 InnoDB 引擎下多个 NULL 值在唯一索引中是可以共存的也就是说如果biz_id字段是 NULL那么唯一索引失效重复数据会被插入。建表时务必把biz_id设置为NOT NULL并在业务代码里保证这个字段一定有值。3.2 Redis SetNX性能可控的去重利器Redis 的SETNX或者SET key value NX EX seconds命令是天然的去重工具。消费消息时用一个能代表业务唯一性的 key比如dedup:order:123456789执行SET key 1 EX 600 NX。如果返回 OK说明这是第一次处理继续执行业务如果返回 null说明之前已经处理过了直接返回消费成功。这里要注意Redis 去重的 key 必须设置过期时间。如果不设过期时间Redis 内存会被不断膨胀的 key 撑爆如果过期时间太短消息在业务处理期间还没完成 key 就过期了去重就会失效。我建议过期时间设置为至少比业务的峰值处理时间长 10 倍。比如业务处理最长 30 秒那 key 过期时间至少 300 秒。如果业务有跨天、跨月的维度就要用数据库的方案。还有一个常见的坑Redis 的SETNX必须跟EXPIRE在同一个命令里执行也就是使用SET key value NX EX seconds这种原子命令。如果你先执行SETNX再执行EXPIRE中间如果进程崩溃key 就变成了永不过期会一直占用内存。这一点我在代码评审里反复强调过。Redis 方案适合高频且对性能敏感的场景比如秒杀、风控、营销活动中的权益发放等。它的问题是如果 Redis 集群发生主从切换且未开启持久化极端情况下可能会丢失一部分 key 数据导致去重失效。所以 Redis 方案只能作为“快路径”核心链路仍然需要数据库兜底。3.3 本地布隆过滤器极低成本但要接受误判布隆过滤器可以在内存中判断一个元素“一定不存在”或者“可能存在”。用它来做消息去重就是把已经处理过的业务 ID 都加入过滤器新的消息来临时先判断是否存在如果不存在就执行业务并加入过滤器如果存在就认为重复跳过处理。布隆过滤器的优点是性能极高不需要访问 Redis 或数据库纯粹本地内存计算缺点是存在误判率。也就是说一个没有处理过的消息可能被误判为“已存在”从而被跳过导致业务漏处理。这个缺陷在做支付、转账这类业务时是绝对无法接受的。所以我一般不会把布隆过滤器作为唯一的去重方案而是把它放在真正需要去重的逻辑之前作为“前置过滤器”挡住绝大多数的重复消息减轻后续 Redis 和数据库的压力。比如一天之内同一用户的重复操作绝大多数都在布隆过滤器这一层被拦截了只有极少数“漏网之鱼”才会走到 Redis 和数据库层。必须注意布隆过滤器不支持删除。如果你只做“插入”和“判断”它就够用如果业务存在“失效”或“重置”的逻辑用它会很麻烦。另外 JVM 重启后布隆过滤器内存数据会全部丢失所以它只能用于短时间内的重复拦截不能作为长期去重依据。3.4 状态机去重适合长流程业务有些业务天然带有状态流转的属性比如订单的状态待支付、已支付、已发货、已完成。这种业务里消息里带上目标状态消费端根据当前状态判断是否可以流转到目标状态。如果当前状态已经是“已支付”又收到一条“支付成功”的消息就说明是重复消息直接忽略即可。这种方案不依赖额外的存储只依靠业务表自己的状态字段实现最简单。但它有一个前提业务必须是有状态的而且状态流转是单向的。如果重复消息携带的目标状态与当前状态一致或者落后就能被识别出来。如果业务本身无状态比如“记录一条日志”状态机去重就没法用。实际上在真实的订单系统里状态机去重往往和唯一键约束搭配使用状态字段用来判断“这条消息对应的状态是否已经达到”唯一键用来防止并发情况下“两条消息同时把订单从待支付改为已支付”的竞态问题。两者并不冲突。3.5 方案对比速查为了方便你决策我把四种方案整理成了一个表格方案可靠性性能实现成本适用场景数据库唯一键极高一般低订单、支付、权益发放等核心链路Redis SetNX较高高低高并发营销、秒杀、风控本地布隆过滤器一般有误判极高中前置过滤大量重复消息状态机较高高低有明确状态流转的业务我的建议是核心链路用数据库唯一键兜底前置再叠加 Redis 或布隆过滤器做性能优化。两者并不冲突反而能互补。4. 一个完整的实战案例支付回调消息的幂等消费4.1 需求背景与整体设计假设有一个订单系统用户支付成功后支付网关会发送一个 HTTP 回调订单服务收到回调后先把订单状态改为“已支付”然后向 RocketMQ 发送一条“支付成功”的消息。积分服务、短信服务、物流服务各自消费这条消息发放积分、发送通知、创建物流单。这条消息从发送到被消费的过程中可能存在的重复路径有支付网关的回调本身可能重复发送导致订单服务生成两条“支付成功”消息。Producer 发送时超时重试导致 Broker 中同时存在两条内容相同的消息。消费者处理完业务后 ACK 丢失Broker 重新投递。Rebalance 导致消息被多个消费者实例重复处理。针对以上重复路径我在消费端的处理策略是分两层第一层RedisSET NX EX快速去重挡掉大部分重复消息。 第二层数据库唯一键兜底确保即使在 Redis 失效的情况下也不会重复发积分、重复发短信。这个设计的核心思想是用 Redis 保性能用数据库保正确性。4.2 数据库表结构设计以积分发放为例建一张积分流水表CREATE TABLE points_record ( id bigint NOT NULL AUTO_INCREMENT, order_id varchar(64) NOT NULL COMMENT 订单号作为业务唯一标识, user_id bigint NOT NULL COMMENT 用户ID, points int NOT NULL COMMENT 积分数量, status tinyint NOT NULL DEFAULT 0 COMMENT 处理状态0-待入账1-已入账, create_time datetime NOT NULL DEFAULT CURRENT_TIMESTAMP, update_time datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, PRIMARY KEY (id), UNIQUE KEY uk_order_id (order_id) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4 COMMENT积分流水表;注意uk_order_id这个唯一索引它保证了同一个订单号只能插入一条积分流水记录。一旦插入成功说明积分发放这件事已经被处理过后续再收到相同的消息直接跳过即可。4.3 Redis 去重的代码实现下面是一个基于 Spring Boot RocketMQ 的消费端实现我简化了业务逻辑但保留核心的去重思路。Component public class PaySuccessMessageConsumer { private static final String DEDUP_KEY_PREFIX dedup:pay_success:; Autowired private StringRedisTemplate redisTemplate; Autowired private JdbcTemplate jdbcTemplate; Autowired private RocketMQListenerPaySuccessMessage rocketMQListener; // 手动确认消费模式 SneakyThrows public void onMessage(MessageExt messageExt) { String msgBody new String(messageExt.getBody(), StandardCharsets.UTF_8); PaySuccessMessage paySuccessMessage JSON.parseObject(msgBody, PaySuccessMessage.class); String orderId paySuccessMessage.getOrderId(); String dedupKey DEDUP_KEY_PREFIX orderId; // 第一层Redis 快速去重 Boolean firstVisit redisTemplate.opsForValue() .setIfAbsent(dedupKey, 1, Duration.ofMinutes(30)); if (Boolean.FALSE.equals(firstVisit)) { // Redis 中已存在说明这条消息之前已经处理过了直接消费成功 return; } try { processPaySuccess(paySuccessMessage); } catch (DuplicateKeyException e) { // 第二层数据库唯一键兜底 log.warn(duplicated message, orderId{}, orderId, e); } catch (Exception e) { // 业务处理失败需要回滚 Redis 标记让消息可以重试 redisTemplate.delete(dedupKey); throw e; } } }这段代码里有三个关键点需要特别说明第一setIfAbsent方法就是对应 Redis 的SET NX EX原子命令。Spring Data Redis 2.x 之后的版本已经不支持setIfAbsent(key, value, timeout)重载了吗其实支持但要注意传入的Duration不能为 null否则默认不设过期时间。第二业务处理抛出DuplicateKeyException时要单独捕获。这个异常表示数据库唯一键冲突说明消息实际上已经被处理过了只是 Redis 标记因为某种原因丢失了。这种情况下必须返回消费成功而不是让 RocketMQ 重试否则会陷入“重试-冲突-再重试”的死循环。第三业务处理抛出非重复的异常时必须把 Redis 标记删掉。因为消息消费失败后 RocketMQ 会重投如果 Redis 标记还在重投的消息会被直接跳过业务就永远无法处理成功了。这个细节非常关键我见过不少人在这个上面踩坑。4.4 核心业务逻辑数据库唯一键兜底processPaySuccess方法里做的事情是向积分流水表插入一条记录然后更新用户积分总余额。这里的关键是保证“插入流水”与“更新余额”的一致性。如果这两个操作不是原子的就可能在中间环节出问题。最简单的做法是先插入流水再更新余额。插入流水时如果遇到唯一键冲突说明已经处理过直接抛出DuplicateKeyException由上层捕获后返回消费成功。如果插入成功再更新用户积分总额。更新余额这一步即使失败也不会导致重复发放积分因为流水已经记录过了后续重试时只会从余额更新开始补做。private void processPaySuccess(PaySuccessMessage message) { String orderId message.getOrderId(); Long userId message.getUserId(); Integer points message.getPoints(); // 插入积分流水唯一键冲突说明重复 jdbcTemplate.update( INSERT INTO points_record (order_id, user_id, points) VALUES (?, ?, ?), orderId, userId, points ); // 更新用户积分余额 jdbcTemplate.update( UPDATE user_points SET total_points total_points ? WHERE user_id ?, points, userId ); }插入流水成功而更新余额失败的情况下消息会被 RocketMQ 重新投递重新执行processPaySuccess。这时插入流水会触发唯一键冲突并被上层捕获。上面这个代码就没有处理“插入成功更新余额失败后重试”的情况因为第二层的DuplicateKeyException直接返回消费成功了还没来得及更新余额。正确的做法应该是在捕获到DuplicateKeyException之后再单独执行“更新余额”的逻辑或者把“插入流水”和“更新余额”放进同一个数据库事务里。我推荐后者因为数据库事务能保证要么都成功要么都回滚不会出现中间状态。4.5 消息消费失败时的状态回滚消费消息时如果业务处理抛出异常说明这条消息消费失败RocketMQ 会按照重试策略重新投递。但我在 4.3 的代码里已经把 Redis 标记删掉了所以重投时setIfAbsent能重新加标记再次进入业务处理。这里有一个隐患如果消息一直消费失败Redis 标记会被反复“写入-删除-写入”造成不必要的开销。不过考虑到大多数消息最终都能在几次重试内成功这个开销是可以接受的。另一个更隐蔽的问题是如果 Redis 标记删除成功但 DB 事务已经提交了比如事务提交后、返回消费结果之前进程宕机此时 Redis 里没有标记数据库里有数据。消息重投后Redis 快速去重没有命中进入业务处理数据库唯一键冲突抛出DuplicateKeyException上层捕获后返回消费成功。这个链路是通的不会出问题。这也从侧面说明了数据库唯一键兜底的必要性。5. 常见问题与排查技巧实录5.1 消费端明明设置了去重为什么还是有重复这个问题我排查过很多次最后发现原因通常是去重 key 的生成不正确。比如用messageExt.getMsgId()做 Redis key前面说过Producer 重试发送同一条消息时两条消息的 Message ID 是不同的这样 Redis 根本拦不住重复。正确的做法是所有去重 key 必须基于业务唯一标识生成而不是基于消息本身的 ID。比如订单号、用户 ID 业务类型、上游请求的唯一流水号等。业界通常叫这个字段为 “bizKey” 或 “eventId”。在发送消息时把这个业务唯一标识放进消息体里消费时用它来生成去重 key。5.2 检测消息是否重复的标准是什么严格来说“重复”的定义有三种层级第一层消息 ID 相同。这是 RocketMQ 层面的重复如上所述不可靠。第二层业务 Key 相同。同一个事件比如“订单 12345 支付成功”被多次发送无论消息 ID 是否相同业务 Key 都相同。这是最常见、也最应该标准的判断依据。第三层消息内容相同。这种情况下业务 Key 可能也相同但内容相同的判定需要做 Json 序列化后比对成本较高一般不用。我在实践中统一使用“业务 Key”作为去重标准。在消息体里设计一个bizKey字段下游消费时以它为幂等键。5.3 如何利用 RocketMQ 的 Message Key 定位问题RocketMQ 的消息对象里有一个setKeys的概念。发送消息时可以设置一个业务唯一 ID比如订单号Message msg new Message(ORDER_TOPIC, , orderId, body); msg.setKeys(orderId);设置之后可以在 RocketMQ Dashboard 的“消息查询”页面按 Key 搜索到这条消息查看它的发送时间、消费状态、重试次数等信息。如果消费端发现问题让排查效率高很多。但注意Message Key 只是定位工具不是去重手段。因为 Message Key 在多个重复消息间可能是相同的也可能为空。我在代码评审时经常看到有人试图用 Message Key 做去重这是不对的。5.4 消费重复导致数据被覆盖的排查当你的数据库里出现了主键冲突或者某个订单的金额被多次累加基本可以断定重复消费真实发生了。排查顺序应该从下往上查看 RocketMQ Dashboard 中该 Topic 的消费组是否发生了 Rebalance。如果消费者实例频繁上下线重复率会升高。查看消费者日志中是否存在“重复消息已被跳过”之类的日志。如果日志没有输出说明去重逻辑没有被触发可能是 key 生成不对。查看消息轨迹里同一业务 Key 消息是否被多个消费者实例消费。如果消息轨迹显示两次消费的实例不同基本可以判断是 Rebalance 导致的重复。查看数据库记录是否出现重复并确认重复数据的插入时间。如果两条重复记录的插入时间相差很短毫秒级大概率是并发触发的如果相差几十秒甚至更久可能是 ACK 丢失后重投导致的。5.5 手动 ACK 与自动 ACK 对去重的影响RocketMQ 的 PushConsumer 有两种消费方式使用DefaultMQPushConsumer并注册MessageListenerConcurrently或者MessageListenerOrderly。它们在消息处理完成之后都会自动返回消费状态。但如果业务里用了 RocketMQ 的事务消息或者你在消费前做了异步处理就需要特别注意 ACK 的时机。我建议务必在业务逻辑全部执行完、并且业务结果确认安全之后再返回消费成功。不要在 Future 或异步线程里执行业务逻辑然后立即返回消费成功因为这时业务可能还没执行完进程一旦重启消息会重投就会出现重复执行。手动 ACK 模式下上述问题发生的概率更高。曾经有个项目在消费订单消息时把消息投递给本地线程池异步处理然后立刻返回CONSUME_SUCCESS结果线程池任务还没跑完进程就挂了消息重投后线程池里其实还有旧任务两边同时处理同一笔订单直接导致余额被扣了两次。5.6 消费堆积时如何防止重复爆发当某个 Topic 的消息消费速度跟不上生产速度Broker 上会积压大量消息。此时如果去新增消费者实例会触发 Rebalance更多的消息会被重新分配。在极端情况下积压的消息里如果有大量内容相同的重复消息去重层的压力会剧增。建议在监控大盘上对消费积压量设置告警比如积压超过 10 万条时立刻处理而不是等到业务被影响才发现。RocketMQ Dashboard 里可以直接看到 Consumer Group 的积压总数也可以接入 Prometheus 配合 Grafana 做趋势监控。5.7 关于消息幂等的最终建议去重方案要做的是“尽可能多地挡住重复”但没有任何方案能保证 100% 拦截。所以在核心业务上必须接受“最坏情况下业务可能被重复执行”的事实并用幂等设计来兜底。比如积分发放这个动作设计成“给用户增加积分”会重复累加但设计成“把积分流水的状态置为成功并且只有第一次置成功时增加余额”就能避免重复。每一条核心业务都应该问问自己这个操作天然幂等吗如果是“累加”“发送”“创建”这类非幂等操作就必须用数据库唯一键或者状态机来约束。6. 生产环境的心得与扩展去重这个话题表面上是代码层面加一个判断实际牵扯到消息语义、消费模型、存储选型、高并发设计多个维度。如果在架构设计阶段没有想清楚后期排查的成本会非常高。我自己的做法是在项目里维护一份“消息去重设计文档”里面记录每个 Topic 的消费场景、使用的去重策略、预期的重复率、监控指标和告警阈值。这样每次有新的消费场景先查文档再写代码不会因为方案选型不一致引入新的问题。还有一个容易被忽略的扩展点RocketMQ 的 SQL 过滤和 Tag 过滤可以帮助消费者只拉取自己关心的消息减少无效消息的消费次数间接降低去重层的工作量。如果某个消费组只关心少数几种消息建议在订阅时用 Tag 过滤。另外如果你的业务使用了 RocketMQ 的事务消息——也就是先执行本地事务再发消息的方案——那么事务消息的“回查”机制也可能导致消息在半成功状态下被重复投递。这种情况下消费者必须更加重视幂等设计不要假设“事务消息一定不会重复”。最后再分享一个我常用的工具技巧RocketMQ Dashboard 里自带的“消息轨迹”功能可以查看一条消息从生产到消费的完整路径包括每个节点的耗时。排查重复消费问题时我一般先看消息轨迹确认是哪条路径产生了重复再针对性地检查 Producer 还是 Consumer。这个习惯让我在几次线上故障中快速定位了根源。
返回列表