RocketMQ消息队列:业务异步解耦、设备消息分发 RocketMQ消息队列业务异步解耦、设备消息分发订单创建后要发邮件、发短信、扣库存、记日志——如果全串行执行用户等得花都谢了。消息队列就是来解这种一个操作触发一堆事情的耦的。今天主角是阿里系扛把子 RocketMQ。一、消息队列三兄弟怎么选先看一张对比表帮你快速定位场景特性RocketMQRabbitMQKafka开发语言JavaErlangJava/Scala吞吐量10万级TPS万级TPS百万级TPS延迟毫秒级微秒级毫秒级顺序消息原生支持需额外处理分区内有序事务消息原生支持插件支持幂等去重延迟消息18个级别插件死信实现不支持原生集群依赖NameServer无自管理ZooKeeper适用场景电商交易、工控业务解耦日志采集、流计算选型建议RocketMQ阿里血统、事务消息原生支持、延迟消息现成可用适合电商和工控场景RabbitMQ轻量灵活、社区成熟适合一般业务解耦Kafka吞吐怪兽、持久化强悍适合日志、埋点、大数据管道二、RocketMQ 核心概念理解这几个概念RocketMQ 就入门了一半┌─────────────────────────────────────────────────┐ │ RocketMQ 架视图 │ │ │ │ ┌──────────┐ ← 路由信息 ┌──────────────────┐ │ │ │NameServer│ │ Broker │ │ │ │ 注册中心 │ │ ┌──────┬──────┐ │ │ │ └──────────┘ │ │Queue0│Queue1│ │ │ │ ↑ │ ├──────┼──────┤ │ │ │ ┌────┴────┐ │ │Queue2│Queue3│ │ │ │ │ │ │ └──────┴──────┘ │ │ │ ▼ ▼ └──────────────────┘ │ │ ┌──────┐ ┌──────┐ ↑ ↓ │ │ │Producer│Consumer│ │ │ │ │ │ 生产者 │ 消费者 │ │ │ │ │ └──────┘ └──────┘ │ │ │ │ 发送消息 拉取/推送消息 ←──────┘ └────── │ └─────────────────────────────────────────────────┘概念比喻说明Topic新闻频道消息的逻辑分类比如order-topicTag子栏目Topic 下的二级分类如order-create、order-payMessageQueue队列Topic 的物理分区一个 Topic 可配多个队列实现并行Producer Group记者团队同一类生产者集合事务消息中用于回溯Consumer Group读者俱乐部同一类消费者集合同一条消息只被组内一个消费者处理集群模式NameServer电话簿记录 Broker 和 Topic 的路由信息轻量级无状态其中 Tag 的设计很巧妙——你可以一个 Topic 承载订单全流程用 Tag 区分创建、支付、取消等不同事件消费者按 Tag 过滤。三、Docker 部署version:3services:namesrv:image:apache/rocketmq:5.1.0container_name:rmq-namesrvports:-9876:9876command:sh mqnamesrvbroker:image:apache/rocketmq:5.1.0container_name:rmq-brokerports:-10911:10911-10909:10909environment:-NAMESRV_ADDRnamesrv:9876command:sh mqbroker-c /home/rocketmq/conf/broker.confdepends_on:-namesrv启动后NameServer 监听 9876 端口Broker 监听 10911客户端通信端口。四、SpringBoot 整合 RocketMQ依赖dependencygroupIdorg.apache.rocketmq/groupIdartifactIdrocketmq-spring-boot-starter/artifactIdversion2.2.3/version/dependency配置rocketmq:name-server:127.0.0.1:9876producer:group:order-producer-groupsend-message-timeout:3000retry-times-when-send-failed:24.1 发送消息三种模式ServicepublicclassOrderMessageProducer{AutowiredprivateRocketMQTemplaterocketMQTemplate;// 同步发送等待 Broker 确认可靠性最高publicvoidsendSync(LongorderId){rocketMQTemplate.syncSend(order-topic:order-create,orderId);}// 异步发送不阻塞回调处理结果publicvoidsendAsync(LongorderId){rocketMQTemplate.asyncSend(order-topic:order-create,orderId,newSendCallback(){OverridepublicvoidonSuccess(SendResultresult){log.info(消息发送成功: {},result.getMsgId());}OverridepublicvoidonException(Throwablee){log.error(消息发送失败,e);// 补偿逻辑写表、重试等}});}// 单向发送不等结果性能最高但没保障publicvoidsendOneway(LongorderId){rocketMQTemplate.sendOneWay(order-topic:order-create,orderId);}}三种模式的选择同步适合核心业务要结果、异步适合非核心但要知道成败、单向适合日志采集丢了也无所谓。4.2 消费消息ComponentRocketMQMessageListener(topicorder-topic,consumerGroupstock-consumer-group,selectorExpressionorder-create,// 只消费 order-create 标签consumeModeConsumeMode.CONCURRENTLY// 并发消费)publicclassStockDeductConsumerimplementsRocketMQListenerOrderMessage{OverridepublicvoidonMessage(OrderMessagemessage){log.info(收到订单消息: {},message.getOrderId());// 处理库存扣减stockService.deduct(message.getProductId(),message.getQuantity());}}消费模式对比模式行为场景集群 (CLUSTERING)同组内一条消息只被一个实例消费服务多实例任务不重复执行广播 (BROADCASTING)同组内所有实例都收到同一条消息配置刷新、缓存更新通知五、消息类型5.1 顺序消息有些业务对顺序有强要求——比如订单的创建→支付→发货这三个消息不能乱序处理。RocketMQ 保证同一队列内的消息严格有序// 发送时指定发送到同一个消息队列rocketMQTemplate.syncSendOrderly(order-topic,orderMsg,orderId.toString());// ↑ 按此hash选择队列// 消费端改为顺序消费RocketMQMessageListener(topicorder-topic,consumerGrouporder-process-group,consumeModeConsumeMode.ORDERLY// 顺序消费模式)原理相同hashKey的消息落到同一个 Queue消费者在同一个 Queue 内串行处理。5.2 延迟消息下单 30 分钟未支付自动取消——这是经典的延迟消息场景// 延迟级别1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h// 级别 16 30 分钟MessageStringmsgMessageBuilder.withPayload(orderId.toString()).build();rocketMQTemplate.syncSend(order-topic:order-delay,msg,3000,16);RocketMQ 不支持任意时间的延迟只能选 18 个预设级别开源版 18 个商业版支持任意时间。如果 30 分钟不够精确可以用定时任务扫描做补充。5.3 事务消息这是 RocketMQ 的杀手锏——保证本地事务和消息发送的原子性Producer Broker Local DB │ │ │ │──①发送半消息────────→│ │ │←──②半消息确认────────│ │ │ │ │ │──③执行本地事务──────────────────────────→│ │←──④本地事务结果────────────────────────│ │ │ │ │──⑤提交/回滚────────→│ │ │ (COMMIT or ROLLBACK) │TransactionalpublicvoidcreateOrderWithTransaction(OrderDTOdto){// 发送事务消息rocketMQTemplate.sendMessageInTransaction(order-topic:order-create,MessageBuilder.withPayload(dto).build(),dto.getOrderId()// 事务参数);}// 事务监听器RocketMQTransactionListenerpublicclassOrderTransactionListenerimplementsRocketMQLocalTransactionListener{OverridepublicRocketMQLocalTransactionStateexecuteLocalTransaction(Messagemsg,Objectarg){try{// 执行本地事务写订单表OrderDTOdto(OrderDTO)((Message)msg).getPayload();orderMapper.insert(dto);returnRocketMQLocalTransactionState.COMMIT;}catch(Exceptione){returnRocketMQLocalTransactionState.ROLLBACK;}}OverridepublicRocketMQLocalTransactionStatecheckLocalTransaction(Messagemsg){// Broker 回调检查本地事务是否真的成功了StringorderId(String)((Message)msg).getPayload();OrderorderorderMapper.selectById(orderId);returnorder!null?RocketMQLocalTransactionState.COMMIT:RocketMQLocalTransactionState.ROLLBACK;}}事务消息的关键Broker 发送半消息后如果迟迟收不到生产者的 COMMIT/ROLLBACK会主动回调checkLocalTransaction方法确认结果。这就是回查机制——防止生产者挂了导致消息一直悬在半空。六、工控设备消息分发场景我司有个智慧农业项目几百个物联网设备温湿度传感器、水肥机、通风设备实时上报数据。用 RocketMQ 做消息分发设备层 消息队列层 消费者层 ┌─────────┐ ┌──────────────┐ ┌──────────────┐ │ 温度传感器 │──MQTT──→│ │ ←Tag──│ 实时报警服务 │ │ 湿度传感器 │──MQTT──→│ device-data │ ←Tag──│ 时序数据库写入 │ │ 水肥机 │──MQTT──→│ Topic │ ←Tag──│ 设备状态监控 │ │ 通风设备 │──MQTT──→│ │ ←Tag──│ 大屏数据推送 │ └─────────┘ └──────────────┘ └──────────────┘每个设备上报的消息带 Tagtemp、humidity、fertilizer不同消费者按 Tag 订阅自己关注的数据。设备数据量很大但单条不重要用单向发送模式压到最低延迟。七、常见踩坑消息重复消费幂等性网络抖动可能导致同一条消息被消费多次。解决方案——消费者端去重// 用 Redis 记录已处理的消息IDStringkeymsg:consumed:message.getMsgId();BooleansuccessredisTemplate.opsForValue().setIfAbsent(key,1,10,TimeUnit.MINUTES);if(Boolean.FALSE.equals(success)){log.warn(消息重复跳过: {},message.getMsgId());return;}// 正常处理业务...消息顺序性保证顺序消息意味着同一 Queue 串行处理吞吐量会下降。不要滥用只在确实要求顺序的场景使用。总结RocketMQ 在阿里内部分布式电商场景久经考验原生的顺序消息、事务消息、延迟消息能力让它特别适合需要复杂消息语义的场景。从业务解耦下单异步通知到工控数据分发设备数据管道RocketMQ 都能打出漂亮仗。记住同步送核心、异步送非核心、单向送日志幂等防重复。