消息队列重复消费和顺序消费该怎么处理

juming
juming 正式会员超兽战士
发布于 2026-10-07 12:56 ·1 浏览 ·0 回复

重复消费只能靠消费端做幂等来解决,顺序消费只能靠「同一业务键进同一队列/分区 + 单线程按序消费」来解决——这两件事都别指望消息队列自动帮你搞定。幂等是兜底,顺序是约束,两者在订单、支付这类场景里必须同时上。

消息队列为什么会重复消费?

结论:因为主流消息队列(Kafka、RocketMQ、RabbitMQ)默认提供的都是「至少一次」(at-least-once)投递语义,也就是消息宁可重发也不能丢。

具体触发场景有三个。第一,消费者业务处理成功,但 ACK 确认丢失或超时,Broker 认为你没消费成功,重新投递。第二,消费者进程崩溃触发 rebalance(重平衡),分区被分给别的消费者,而 offset(消费位点)没提交成功,这批消息会被重新拉取。第三,生产端发送超时后重试,同一条消息被写入了两次——Kafka 要开启 enable.idempotence=true 才能在生产端消除这种重复,RocketMQ 默认不保证。

结论:重复消费是设计前提,不是异常情况,任何消费逻辑都必须假设「这条消息我可能收到 2 次以上」。

重复消费怎么处理?核心就是幂等

结论:重复消费没有服务端解法,唯一的解法是让消费逻辑本身幂等——同一条消息消费 1 次和消费 10 次,业务结果完全一致。

最可靠的方案是数据库唯一索引。给业务表加上业务唯一键(比如订单表的 uk_order_no),重复插入会直接抛 DuplicateKeyException(MySQL 错误码 1062),代码捕获后当成消费成功处理即可。这个方案不依赖任何中间件,并发下也不会失效。

其次是去重表 + 本地事务。建一张 msg_dedup 表,字段是 msg_id(主键)、status、create_time,把「业务数据写入」和「插入去重记录」放在同一个数据库本地事务里,保证原子性。这里有个高频踩坑点:先 SELECT 判断再 INSERT 在并发下必然失效,两个线程会同时查不到再同时插入——必须用唯一索引兜底,而不是靠查询。

Redis 的 SET msgId 1 NX EX 86400 也能做去重,返回 1 表示首次消费。但 Redis 和数据库不在同一个事务里,如果 Redis 写入成功、业务执行失败,这条消息就永远不会被重试了。所以它只适合通知、埋点、统计这类能容忍少量丢失的场景,不适合订单、资金。

最后是状态机和乐观锁。UPDATE orders SET status=2 WHERE id=? AND status=1,靠返回的 affected rows 是否为 1 来判断本次是否真正生效。它天然只允许状态单向流转,重复消息打到已完成的订单上会自动被忽略。去重表记得按 create_time 建索引并定期清理,通常保留 7 天,用定时任务每天删一次。

顺序消费怎么处理?

结论:消息队列只在「同一队列或同一分区内」保证顺序,跨队列、跨分区没有任何顺序保证。所以顺序消费就三件事:同业务键路由到同队列、单线程按序消费、失败不能跳过。

生产端要按业务键(订单号、用户 ID)选队列。Kafka 在发送时指定 key,Kafka 用 key 的 hash 决定分区,同一个 key 必然进同一个分区,这是最省事的做法。RocketMQ 要用 MessageQueueSelector 按 orderId 取模选队列,或者直接用顺序消息 API。RabbitMQ 则是单队列配单消费者,或者用一致性哈希交换机。

消费端要保证一个分区只被一个线程消费。Kafka 里同一 consumer group 内一个分区只会分配给一个消费者线程,天然满足;RocketMQ 要用 MessageListenerOrderly(内部对每个 MessageQueue 加锁,同一队列串行消费),千万别用 MessageListenerConcurrently,那个是并发消费,直接乱序。

失败处理是顺序消费最容易出事的地方。顺序消费失败不能「跳过继续」,否则后面的消息全部错位。必须返回失败让消息队列重试同一队列。RocketMQ 的顺序消费失败会阻塞当前队列,所以必须配置最大重试次数(默认 16 次)并接死信队列,否则这个队列会被永久卡住、后面消息全部堵死。Kafka 可以 seek 回退 offset 重试,但要防止无限循环,建议失败超过 3 次就写入独立的 DLQ topic 人工介入。

想并发消费又要保序,怎么做?

结论:用「内存队列按 key 分流」模式——拉取线程只负责按分区顺序取消息,按 key 的 hash 投递到内存队列,每个内存队列后面挂一个固定线程串行处理。

这个模式在 Kafka 高吞吐场景里很常见。有三个必须注意的点:内存队列要设容量上限,满了就调用 pause() 暂停对应分区的拉取,绝不能丢消息;提交 offset 时只能提交「所有内存队列都已处理完成」的那个最小 offset,提交早了就会丢消息;消费者线程数等于分区数就够了,不要再额外套线程池,否则分区内的顺序会被打乱。

顺序和幂等必须一起做

结论:顺序消费和幂等不是二选一,顺序消费的场景反而更需要幂等。

原因是顺序消费失败重试时,很可能是「业务已经写库了、只是 ACK 丢了」,重试就会重复执行。以订单状态流转为例:创建→支付→发货→完成,用状态机加唯一索引兜底,既能保证顺序推进,也能保证重复消息不会把已完成的状态改回去。

总结一下:重复消费靠数据库唯一索引或本地事务去重表做幂等,别用「先查后插」;顺序消费靠业务键路由到同一分区/队列、单消费线程、失败重试不跳过,并且必须配死信队列防止队列卡死;两者同时使用才是生产环境的正确姿势。

版权声明:本文来自 GJ站长论坛《消息队列重复消费和顺序消费该怎么处理》
原文链接:https://www.gj0.com/thread-292.html
转载请注明出处并保留本声明;内容仅代表作者观点,与本站立场无关。若本文涉嫌侵权,请联系本站处理。

全部回复 0

还没有回复,来抢沙发~