
顺序消息原理与落地——售货柜单设备指令有序执行保障作者黒漂技术佬系列RocketMQ核心原理与无人售货柜项目实战第4篇一、为什么需要顺序消息先说个翻车事故。无人售货柜某次大促用户扫码开门取货结果设备先执行了关门指令再执行出货指令。用户被柜门夹了手货也没出来投诉电话直接打到了12315。排查日志发现开门、出货、关门三条指令通过RocketMQ发送消费者收到顺序是乱的——关门跑到了出货前面。这不是RocketMQ的Bug是我们自己的使用姿势不对。普通消息是不保证顺序的Broker把消息分散到不同队列消费者多线程并行拉取消费先到先得顺序天然不保证。要保证顺序得用顺序消息Ordered Message。二、两种顺序全局顺序 vs 分区顺序2.1 全局顺序所有消息严格按照发送顺序进入同一个队列消费者单线程依次消费。优点绝对的FIFO先进先出缺点吞吐量极低完全丧失了消息队列的并发优势类比超市只开一个收银台所有人排一条队绝对不会插队但等到天荒地老。2.2 分区顺序Partitioned Order消息按某个业务KeyShardingKey做Hash路由同一Key的消息进入同一队列不同Key可以分散到不同队列。消费者对每个队列单线程消费。同一ShardingKey的消息严格有序不同ShardingKey的消息无序但可以并行消费类比超市开了10个收银台按顾客姓氏分配窗口。姓张的永远去1号台姓李的永远去2号台。张家人内部有序李家人内部有序但张家和李家之间谁先结账无所谓。实际项目中99%用分区顺序全局顺序几乎没有使用场景。三、分区顺序消息的核心原理3.1 发送端ShardingKey路由生产者发送消息时通过MessageQueueSelector选择目标队列队列下标 hash(ShardingKey) % queueCount同一个ShardingKey的Hash值固定永远路由到同一个队列。来看源码核心逻辑// RocketMQ默认的队列选择器SelectMessageQueueByHashpublicMessageQueueselect(ListMessageQueuemqs,Messagemsg,Objectarg){intvaluearg.hashCode();// 对ShardingKey求HashintindexMath.abs(value)%mqs.size();// 取模得到队列下标if(index0){index0;}returnmqs.get(index);}关键点ShardingKey的选择决定了顺序粒度。3.2 存储端队列天然有序RocketMQ的CommitLog是所有消息顺序写入的大文件但每个队列ConsumeQueue内部维护了自己消息的逻辑偏移量。同一个队列里的消息存储顺序就是发送顺序这个是物理层面保证的。3.3 消费端MessageListenerOrderly消费者使用MessageListenerOrderly接口RocketMQ保证每个队列同一时刻只有一个线程在消费publicinterfaceMessageListenerOrderlyextendsMessageListener{ConsumeOrderlyStatusconsumeMessage(ListMessageExtmsgs,ConsumeOrderlyContextcontext);}和MessageListenerConcurrently并发消费的区别特性MessageListenerConcurrentlyMessageListenerOrderly消费线程多线程并发每队列单线程消费失败返回RECONSUME_LATER进重试队列返回SUSPEND_CURRENT_QUEUE_A_MOMENT稍后重试不跳过顺序保证无队列内严格有序吞吐量高相对较低注意消费失败的处理差异并发消费失败会丢进重试队列后续消息继续消费顺序消费失败会挂起当前队列等一会重试不会跳到下一条——否则顺序就乱了。四、SpringBoot实战售货柜设备指令顺序执行4.1 业务场景无人售货柜一次购物流程涉及多条设备指令开门指令解锁电磁锁用户开门取货出货指令如果涉及自动升降/推送机构执行出货动作关门指令延时关门确认用户取货完毕这三条指令必须按顺序到达设备否则就会出现前面说的夹手事故。ShardingKey的选择设备ID。同一台设备的指令路由到同一队列保证有序不同设备之间互不影响可以并行消费。4.2 定义指令消息模型DataBuilderpublicclassDeviceCommandMessageimplementsSerializable{privateStringdeviceId;// 设备IDShardingKeyprivateStringorderId;// 关联订单号privateIntegercommandType;// 指令类型1开门, 2出货, 3关门privateIntegersequence;// 指令序号1, 2, 3privateLongtimestamp;// 发送时间戳privateStringpayload;// 指令参数JSON}4.3 生产者按设备ID路由发送ServicepublicclassDeviceCommandProducer{ResourceprivateRocketMQTemplaterocketMQTemplate;/** * 发送顺序消息——同一个deviceId的路由到同一队列 */publicSendResultsendOrderly(DeviceCommandMessagecommand){MessageDeviceCommandMessagemessageMessageBuilder.withPayload(command).build();// syncSendOrderly: 第三个参数是hashKey用于选择队列// 这里传 deviceId保证同一设备的消息进同一队列returnrocketMQTemplate.syncSendOrderly(device-command-topic,message,command.getDeviceId()// ShardingKey 设备ID);}/** * 一次购物流程按序发送开门→出货→关门 */publicvoidsendShoppingFlow(StringdeviceId,StringorderId){// 1. 开门指令sendOrderly(DeviceCommandMessage.builder().deviceId(deviceId).orderId(orderId).commandType(1).sequence(1).timestamp(System.currentTimeMillis()).payload({\lockId\:\deviceId_L1\}).build());// 2. 出货指令sendOrderly(DeviceCommandMessage.builder().deviceId(deviceId).orderId(orderId).commandType(2).sequence(2).timestamp(System.currentTimeMillis()).payload({\motorId\:\M3\,\channel\:2}).build());// 3. 关门指令sendOrderly(DeviceCommandMessage.builder().deviceId(deviceId).orderId(orderId).commandType(3).sequence(3).timestamp(System.currentTimeMillis()).payload({\delaySec\:5}).build());}}核心就是syncSendOrderly方法底层通过hash(deviceId) % queueCount选择队列。三次调用用的是同一个deviceId所以三条指令进入同一个队列存储顺序就是发送顺序。4.4 消费者顺序消费Slf4jServiceRocketMQMessageListener(topicdevice-command-topic,consumerGroupdevice-command-consumer-group,consumeModeConsumeMode.ORDERLY,// 关键顺序消费模式maxReconsumeTimes5)publicclassDeviceCommandOrderlyConsumerimplementsRocketMQListenerMessageExt{ResourceprivateDeviceCommandExecutorcommandExecutor;OverridepublicvoidonMessage(MessageExtmessage){try{DeviceCommandMessagecommandJSON.parseObject(message.getBody(),DeviceCommandMessage.class);log.info(收到设备指令 | deviceId{} | type{} | seq{} | reconsumeTimes{},command.getDeviceId(),command.getCommandType(),command.getSequence(),message.getReconsumeTimes());// 按指令类型执行switch(command.getCommandType()){case1:commandExecutor.executeOpen(command);break;case2:commandExecutor.executeDispense(command);break;case3:commandExecutor.executeClose(command);break;default:log.warn(未知指令类型: {},command.getCommandType());}}catch(Exceptione){log.error(设备指令执行失败 | msgId{},message.getMsgId(),e);// 顺序消费抛异常框架会返回SUSPEND_CURRENT_QUEUE_A_MOMENT// 当前队列会暂停稍后重试这条消息不会跳到下一条thrownewRuntimeException(指令执行失败暂停当前队列,e);}}}关键配置consumeMode ConsumeMode.ORDERLY。框架底层会为每个队列创建一个ProcessQueue通过加锁保证同一队列同一时刻只有一个消费线程。所以同一设备的开门、出货、关门会依次被消费不会乱序。4.5 指令执行器Slf4jServicepublicclassDeviceCommandExecutor{ResourceprivateDeviceMqttClientmqttClient;// MQTT下发到设备ResourceprivateRedisTemplateString,StringredisTemplate;/** * 执行开门 */publicvoidexecuteOpen(DeviceCommandMessagecommand){// 幂等校验防止重复开门StringdedupKeycmd:exec:command.getOrderId():command.getSequence();BooleanabsentredisTemplate.opsForValue().setIfAbsent(dedupKey,1,10,TimeUnit.MINUTES);if(Boolean.FALSE.equals(absent)){log.info(指令已执行过跳过 | orderId{} | seq{},command.getOrderId(),command.getSequence());return;}// 通过MQTT下发开门指令到设备mqttClient.publish(command.getDeviceId(),CMD_OPEN,command.getPayload());log.info(开门指令已下发 | deviceId{},command.getDeviceId());}/** * 执行出货 */publicvoidexecuteDispense(DeviceCommandMessagecommand){StringdedupKeycmd:exec:command.getOrderId():command.getSequence();if(Boolean.FALSE.equals(redisTemplate.opsForValue().setIfAbsent(dedupKey,1,10,TimeUnit.MINUTES))){return;}mqttClient.publish(command.getDeviceId(),CMD_DISPENSE,command.getPayload());log.info(出货指令已下发 | deviceId{},command.getDeviceId());}/** * 执行关门 */publicvoidexecuteClose(DeviceCommandMessagecommand){StringdedupKeycmd:exec:command.getOrderId():command.getSequence();if(Boolean.FALSE.equals(redisTemplate.opsForValue().setIfAbsent(dedupKey,1,10,TimeUnit.MINUTES))){return;}mqttClient.publish(command.getDeviceId(),CMD_CLOSE,command.getPayload());log.info(关门指令已下发 | deviceId{},command.getDeviceId());}}这里加了基于Redis的幂等校验和上一篇讲的消息幂等性方案一脉相承。顺序消费虽然保证了顺序但消费失败会重试重试时必须做幂等处理。五、顺序消息的坑与避坑指南5.1 坑一消费阻塞顺序消费失败时当前队列会被挂起后续消息全部卡住。如果某条消息一直消费失败整个队列就堵死了。避坑方案设置maxReconsumeTimes超过重试次数后走兜底逻辑告警人工处理别让一条毒消息卡死整条队列。5.2 坑二ShardingKey设计不当如果ShardingKey粒度太粗比如用区域ID同一个区域内所有设备的消息挤进一个队列并发度大幅降低。如果粒度太细比如用订单ID同一设备的不同订单可能分到不同队列设备指令就乱序了。售货柜场景正确选择ShardingKey 设备ID。同一设备的指令有序不同设备并行消费并发度等于队列数。5.3 坑三Broker扩缩容导致队列数变化队列数变化后hash(ShardingKey) % queueCount的结果会变原来在队列A的消息可能路由到队列B。扩容期间发送的消息和旧消息可能不在同一队列顺序会断。避坑方案扩缩容期间做好灰度业务层做序号校验兜底。设备指令消息里那个sequence字段就是干这个的——消费时校验序号是否连续断了就告警。5.4 坑四不要在顺序消费里做耗时操作顺序消费是单队列单线程如果你在消费逻辑里调一个3秒的外部接口后面所有消息都得等3秒。吞吐量直接跳水。避坑方案耗时操作改成异步消费逻辑只做接收落库异步触发。如果必须同步等待设备响应建议用MQTT的QoS机制 本地ACK不要在MQ消费链路里同步等待。六、顺序消息消费失败的重试机制顺序消费和并发消费的重试策略不同这个容易搞混维度并发消费顺序消费失败返回值RECONSUME_LATERSUSPEND_CURRENT_QUEUE_A_MOMENT重试方式消息进入%RETRY%队列当前队列挂起稍后重试当前消息后续消息继续消费被阻塞不会跳过重试间隔递增10s/30s/1m…固定默认1s最大重试默认16次默认Integer.MAX_VALUE顺序消费默认无限重试这很危险。必须显式设置maxReconsumeTimes超过次数后让消息进入死信队列避免队列永久阻塞。七、总结要点说明全局顺序单队列单线程吞吐极低基本不用分区顺序按ShardingKey路由同Key有序生产可用ShardingKey选业务实体ID设备ID/订单ID/用户ID消费模式ConsumeMode.ORDERLY每队列单线程失败处理挂起当前队列不跳过需设重试上限幂等性顺序消息也会重试必须做幂等售货柜场景的核心链路ShardingKey设备ID → 同设备指令进同一队列 → 顺序消费保证开门→出货→关门依次执行。简单直接但生产环境一定要处理好消费失败重试和幂等兜底否则一条失败消息能把整台设备卡死。