多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

RocketMQ 顺序消息:订单流程不乱套的秘密

RocketMQ 顺序消息:订单流程不乱套的秘密 作者鱼宵 实战驱动系列 · 第 5 篇完整课程与可运行源码已开源在 Giteehttps://gitee.com/j67mk2/rocketmq-journey 本文对应 lesson-05/你发了一条订单消息流创建订单 → 支付成功 → 商品发货 → 订单完成。结果消费者那边打印出来是这样的先收到商品发货再收到订单完成最后才慢悠悠收到创建订单。订单状态机直接乱套——用户还没支付你这边货都发出去了。这不是bug这是 RocketMQ 默认行为队列之间是并行的谁先谁后根本没保证。今天这篇就把这件事讲透怎么用MessageQueueSelector把同一个订单的消息摁到同一条队列里再用MessageListenerOrderly让消费者老老实实排队吃消息。讲完你会明白为什么生产上几乎没人做全局顺序而只做分区顺序。答案都在 lesson-05 的源码里clone 下来跑一遍乱序现场肉眼可见。一、为什么会乱把队列想成医院的 4 条挂号窗口先回忆第 1 课的常识一个 Topic 默认有 4 个队列MessageQueue。队列才是真正存消息、被消费的最小单位Topic 只是个名字。把它想成医院大厅Topic医院大厅 队列0 队列1 队列2 队列3 ← 4 条并行窗口 消息a 消息b 消息c 消息d默认发送时生产者是轮询的第 1 条丢窗口 0第 2 条丢窗口 1第 3 条丢窗口 2……而每条窗口背后都有一个独立的护士在并行处理。这就意味着窗口 2 上的商品发货可能先被处理完窗口 0 上的创建订单反而排在后面没轮到。队列与队列之间没有任何顺序保证。你发出的顺序是 创建→支付→发货→完成消费者看到的可能是 发货→完成→创建→支付。订单状态直接乱套。二、全局顺序 vs 分区顺序一张表记牢面试官最爱问这俩区别。先上结论表概念类比做法代价全局顺序医院只开一个窗口所有人排一条长队整个 Topic 只用 1 个队列所有消息串行处理完全牺牲并行度吞吐暴跌生产几乎不用分区顺序推荐同一个病人的号永远固定到同一个窗口办用业务 key订单号选队列同一 key 的消息进同一队列不同 key 散在不同队列并行只保证同一个订单内有序不同订单之间仍并行吞吐几乎不受影响面试标准答案一句话RocketMQ 能保证的是队列内有序也就是分区顺序全局顺序要靠单队列硬扛代价太大生产基本不做。为什么全局顺序代价大因为单队列 单线程消费并发度被砍到 1。原来 4 个队列并行每秒能吃 4000 条现在 1 个队列串行只能吃 1000 条吞吐直接掉到四分之一。订单量一大队列就堵成停车场。三、怎么做到分区有序两步两头都要配分区有序不是光改一个参数就行得发送端和消费端两头配合发送端同一个 orderId ──选队列──▶ 永远落到同一个队列 │ 队列内按写入顺序排好创建→支付→发货→完成 ▼ 消费端MessageListenerOrderly 加锁一个队列串行消费严格按队列内顺序处理第一步发送端调用三参数版producer.send(msg, selector, orderId)。选择器按 orderId 算队列下标同一个 orderId 永远选同一个队列。于是这个订单的 4 条消息全排进同一条队伍。第二步消费端用MessageListenerOrderly顺序监听器而不是默认的MessageListenerConcurrently并发监听器。顺序监听器会给每个队列加锁一条消息处理完、返回成功才放行下一条——队列内严格按顺序消费。⚠️ 最容易踩的坑光发送端选了队列消费端还用并发监听照样乱。两头都要配缺一头都不行。四、动手三个类对比着跑工程在lesson-05/三个主类OrderlyProducer顺序发送、UnOrderlyProducer对比组乱序发送、OrderConsumer顺序消费。4.1 OrderlyProducer按 orderId 选队列packagecom.example.lesson05;importorg.apache.rocketmq.client.producer.DefaultMQProducer;importorg.apache.rocketmq.client.producer.MessageQueueSelector;importorg.apache.rocketmq.client.producer.SendResult;importorg.apache.rocketmq.common.message.Message;importorg.apache.rocketmq.common.message.MessageQueue;importjava.nio.charset.StandardCharsets;importjava.util.List;/** * 顺序发送用 MessageQueueSelector 把同一个订单的消息固定到同一个队列。 */publicclassOrderlyProducer{publicstaticvoidmain(String[]args)throwsException{DefaultMQProducerproducernewDefaultMQProducer(lesson05_producer_group);producer.setNamesrvAddr(127.0.0.1:9876);producer.start();System.out.println(顺序生产者已启动……\n);StringorderIdORDER_1001;String[]states{创建订单,支付成功,商品发货,订单完成};for(intstep0;stepstates.length;step){Stringbody订单orderId状态states[step]第 (step1) 步;MessagemsgnewMessage(TopicLesson05,OrderFlow,body.getBytes(StandardCharsets.UTF_8));msg.setKeys(orderId);// 关键三参数 send第二个参数是队列选择器SendResultresultproducer.send(msg,newMessageQueueSelector(){OverridepublicMessageQueueselect(ListMessageQueuequeues,Messagem,Objectarg){Stringkey(String)arg;// 同一个 key - 同一个下标 - 同一个队列intindexMath.abs(key.hashCode())%queues.size();returnqueues.get(index);}},orderId);System.out.println(已发送body落到 queueIdresult.getMessageQueue().getQueueId());}producer.shutdown();}}这段代码最核心的就一行Math.abs(key.hashCode()) % queues.size()。把订单号哈希成队列下标同一个订单号哈希结果不变所以 4 条消息永远落到同一个队列不同订单散到不同队列照常并行。4.2 UnOrderlyProducer对比组乱序发送packagecom.example.lesson05;importorg.apache.rocketmq.client.producer.DefaultMQProducer;importorg.apache.rocketmq.client.producer.SendResult;importorg.apache.rocketmq.common.message.Message;importjava.nio.charset.StandardCharsets;/** * 乱序发送不指定队列选择器用默认轮询消息散到不同队列。 */publicclassUnOrderlyProducer{publicstaticvoidmain(String[]args)throwsException{DefaultMQProducerproducernewDefaultMQProducer(lesson05_producer_group);producer.setNamesrvAddr(127.0.0.1:9876);producer.start();System.out.println(乱序生产者已启动默认轮询不选队列发送订单 ORDER_2002 的 4 步状态……\n);StringorderIdORDER_2002;String[]states{创建订单,支付成功,商品发货,订单完成};for(intstep0;stepstates.length;step){Stringbody订单orderId状态states[step]第 (step1) 步;MessagemsgnewMessage(TopicLesson05,OrderFlow,body.getBytes(StandardCharsets.UTF_8));msg.setKeys(orderId);// 普通两参数 send不传选择器RocketMQ 自动轮询分到下一个队列SendResultresultproducer.send(msg);System.out.println(已发送body随机落到 queueIdresult.getMessageQueue().getQueueId());}System.out.println(\n4 步状态已发送散在不同队列。去 OrderConsumer 看是不是乱序了。);producer.shutdown();}}和 OrderlyProducer 的唯一区别调用的是producer.send(msg)两参数版本不传选择器。就这一点差别消费顺序天差地别。4.3 OrderConsumer顺序消费自带乱序检测packagecom.example.lesson05;importorg.apache.rocketmq.client.consumer.DefaultMQPushConsumer;importorg.apache.rocketmq.client.consumer.listener.ConsumeOrderlyContext;importorg.apache.rocketmq.client.consumer.listener.ConsumeOrderlyStatus;importorg.apache.rocketmq.client.consumer.listener.MessageListenerOrderly;importorg.apache.rocketmq.common.consumer.ConsumeFromWhere;importorg.apache.rocketmq.common.message.MessageExt;importjava.nio.charset.StandardCharsets;importjava.util.List;importjava.util.concurrent.ConcurrentHashMap;importjava.util.concurrent.atomic.AtomicLong;importjava.util.regex.*;/** * 顺序消费用 MessageListenerOrderly 按队列串行消费并打印顺序判断。 */publicclassOrderConsumer{publicstaticvoidmain(String[]args)throwsException{DefaultMQPushConsumerconsumernewDefaultMQPushConsumer(lesson05_consumer_group);consumer.setNamesrvAddr(127.0.0.1:9876);consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);consumer.subscribe(TopicLesson05,OrderFlow);// 每个订单已消费到第几步用来判断是否乱序ConcurrentHashMapString,IntegerorderStepMapnewConcurrentHashMap();AtomicLonglastMsgTimenewAtomicLong(System.currentTimeMillis());PatternstepPatternPattern.compile(第\\s*(\\d)\\s*步);// 注意这里注册的是 MessageListenerOrderly顺序监听器不是 Concurrentlyconsumer.registerMessageListener(newMessageListenerOrderly(){OverridepublicConsumeOrderlyStatusconsumeMessage(ListMessageExtmsgs,ConsumeOrderlyContextcontext){for(MessageExtmsg:msgs){StringbodynewString(msg.getBody(),StandardCharsets.UTF_8);StringorderIdmsg.getKeys();lastMsgTime.set(System.currentTimeMillis());MatchermstepPattern.matcher(body);intstepm.find()?Integer.parseInt(m.group(1)):-1;intexpectedorderStepMap.getOrDefault(orderId,0)1;if(stepexpected){System.out.println([顺序正确] body);orderStepMap.put(orderId,step);}elseif(stepexpected){System.out.println([重复/回退] body该订单已到第 (expected-1) 步跳过);}else{System.out.println([乱序告警] body本应先收第 expected 步却先收第 step 步);orderStepMap.put(orderId,step);}}returnConsumeOrderlyStatus.SUCCESS;}});consumer.start();System.out.println(顺序消费者已启动……连续 6 秒没新消息自动退出\n);longdeadlineSystem.currentTimeMillis()30_000;while(System.currentTimeMillis()deadline){Thread.sleep(500);if(System.currentTimeMillis()-lastMsgTime.get()6_000)break;}consumer.shutdown();System.out.println(观察结束消费者已关闭。);}}这段代码干了三件事注册MessageListenerOrderly顺序消费开关、用orderStepMap记录每个订单走到第几步、收到消息时对比期望步数 vs 实际步数自动打印[顺序正确]或[乱序告警]。五、实测输出选队列 vs 不选队列对比一目了然下面是同款环境Windows 11 Docker JDK 17下的真实运行输出直接贴给你对照。① OrderlyProducer 输出ORDER_1001 四步全部落到同一个队列 queueId1顺序生产者已启动按 orderId 选队列发送订单 ORDER_1001 的 4 步状态…… 已发送订单ORDER_1001状态创建订单第 1 步落到 queueId1 已发送订单ORDER_1001状态支付成功第 2 步落到 queueId1 已发送订单ORDER_1001状态商品发货第 3 步落到 queueId1 已发送订单ORDER_1001状态订单完成第 4 步落到 queueId1 4 步状态已全部发送。去 OrderConsumer 看消费顺序。② UnOrderlyProducer 输出ORDER_2002 四步被轮询散到 0/1/2/3 四个队列乱序生产者已启动默认轮询不选队列发送订单 ORDER_2002 的 4 步状态…… 已发送订单ORDER_2002状态创建订单第 1 步随机落到 queueId0 已发送订单ORDER_2002状态支付成功第 2 步随机落到 queueId1 已发送订单ORDER_2002状态商品发货第 3 步随机落到 queueId2 已发送订单ORDER_2002状态订单完成第 4 步随机落到 queueId3 4 步状态已发送散在不同队列。③ OrderConsumer 输出关键对比选队列的 ORDER_1001 全对没选队列的 ORDER_2002 乱了顺序消费者已启动MessageListenerOrderly队列内串行…… [顺序正确] 订单ORDER_2002状态创建订单第 1 步该订单此前已到第 0 步 [乱序告警] 订单ORDER_2002状态商品发货第 3 步按顺序本应先收到第 1 步却先收到了第 3 步 [乱序告警] 订单ORDER_2002状态订单完成第 4 步按顺序本应先收到第 1 步却先收到了第 4 步 [顺序正确] 订单ORDER_1001状态创建订单第 1 步该订单此前已到第 0 步 [顺序正确] 订单ORDER_1001状态支付成功第 2 步该订单此前已到第 1 步 [顺序正确] 订单ORDER_1001状态商品发货第 3 步该订单此前已到第 2 步 [顺序正确] 订单ORDER_1001状态订单完成第 4 步该订单此前已到第 3 步 [重复/回退] 订单ORDER_2002状态支付成功第 2 步该订单已到第 4 步这条更早跳过 观察结束消费者已关闭。这几行说明什么一张表看明白订单发送方式消费结果ORDER_1001按 orderId 选队列全在 queue1创建→支付→发货→完成4 条全是[顺序正确]ORDER_2002默认轮询散在 4 个队列实际收到 1→3→4→2两处[乱序告警]支付最后才到你本机跑 ORDER_2002 乱成什么样可能和上面略有不同取决于各个队列被消费的先后但只要它没按 1→2→3→4 来就证明不选队列会乱序。ORDER_1001 永远稳定有序。六、顺序消费的代价一条卡住全队堵死并发消费失败消息会被扔到重试队列不影响后面的消息。但顺序消费失败时RocketMQ 不会跳过这条——它会阻塞住这个队列反复重试这一条直到成功。好处是不会乱序代价是一条脏消息如果一直处理失败后面排队的所有消息都被它堵在这。所以顺序消费的业务逻辑要尽量稳对确实处理不了的消息业务上自己做兜底比如记死信、跳过别让它无限阻塞。常见坑汇总坑原因解法选了队列还是乱消费端还用MessageListenerConcurrently并发监听换成MessageListenerOrderly两头都配一条消息卡住全队不动顺序消费失败会阻塞当前队列反复重试业务逻辑做稳兜底记死信别无限阻塞乱序效果每次跑不一样不同队列被消费的先后本来就不确定乱序是概率性的正常现象出现乱序就达到教学目的了七、挑战题答案在仓库跑起来才知道⭐ 把 OrderlyProducer 里的orderId改成再发一个订单ORDER_1005同样四步观察它和 ORDER_1001 是否落在同一个队列、两个订单在消费端是否各自有序、互不影响。⭐⭐ 把 OrderConsumer 里的MessageListenerOrderly临时换回MessageListenerConcurrently注意返回值也要改成ConsumeConcurrentlyStatus再跑一遍看看连 ORDER_1001 是不是也变得时好时坏了体会消费端顺序监听的必要性。⭐⭐⭐ 思考如果你的 Topic 有 8 个队列按 orderId 选队列后某个热门大订单的消息全压在一个队列里这个队列成了性能瓶颈怎么办提示这是分区顺序在热点 key下的典型问题想想怎么拆分 key。八、面试回答模板面试官RocketMQ 怎么保证消息顺序RocketMQ 只能保证队列内有序分区顺序发送端用MessageQueueSelector按业务 key如 orderId选队列让同一个 key 的消息落到同一个队列消费端用MessageListenerOrderly加锁串行消费。两头配合缺一不可。追问全局顺序和分区顺序有什么区别为什么生产都用分区顺序全局顺序是整个 Topic 只用 1 个队列所有消息串行最严格但并发度被砍到 1吞吐暴跌分区顺序是按业务 key 分队列同 key 有序、不同 key 并行吞吐几乎不受影响。生产上订单、流水这类场景只需要保证同一个订单内有序不同订单之间本来就不需要互相排序所以分区顺序足够用。追问顺序消费失败了会怎样不会跳过这条而是SUSPEND_CURRENT_QUEUE_A_MOMENT阻塞当前队列反复重试直到成功。好处是保序代价是一条脏消息能堵死整个队列所以业务逻辑要稳、要做兜底。详细代码和实测对照见本文第四、五节完整可运行工程在 lesson-05/。关于这个系列本文是「Java 后端实战精通营」系列第 5 篇原则实战驱动、由浅到深、面试向每篇文章的结论都可以亲手验证。RocketMQ 实战精通营10 课https://gitee.com/j67mk2/rocketmq-journey本文对应源码位置lesson-05/三个主类OrderlyProducer 顺序发送、UnOrderlyProducer 对比组、OrderConsumer 顺序消费自带乱序检测打印系列文章一览按发布顺序篇主题1RocketMQ 入门Docker 一行起三件套跑通你的第一条消息2RocketMQ 发送方式同步异步批量单向消息都怎么发出去3RocketMQ 消费模式集群、广播与重试消息怎么被吃掉4RocketMQ 可靠性发送重试加幂等消息一条都不丢5RocketMQ 顺序消息订单流程不乱套的秘密6RocketMQ 事务消息订单与积分的最终一致7RocketMQ 延迟消息30 分钟未支付自动关单怎么做8RocketMQ 积压治理百万消息堵在队列怎么办9RocketMQ 集群高可用与过滤主从架构 Tag 精准投递10RocketMQ 面试冲刺高频考点一口气背完下一篇预告《RocketMQ 事务消息订单与积分的最终一致》——Half 半消息、本地事务执行、Broker 回查机制以及它和数据库 2PC 到底有什么区别。跑完有任何报错把终端输出发评论区一起排查。
返回列表