基础消息模型实战:普通消息、同步/异步/单向消息发送
2026/9/6 7:52:40 网站建设 项目流程

基础消息模型实战:普通消息、同步/异步/单向消息发送

作者:黒漂技术佬
系列专栏:RocketMQ核心原理与无人售货柜项目实战

一、Maven依赖和初始化配置

1.1 引入依赖

Spring Boot项目引入RocketMQ Starter,版本和RocketMQ Server对应即可:

<dependencies><!-- RocketMQ Spring Boot Starter --><dependency><groupId>org.apache.rocketmq</groupId><artifactId>rocketmq-spring-boot-starter</artifactId><version>2.3.1</version></dependency><!-- Spring Boot Web(微服务接口) --><dependency><groupId>org.springframework.boot</groupId><artifactId>spring-boot-starter-web</artifactId><version>2.7.18</version></dependency></dependencies>

如果不用Starter,直接用原生Client也行:

<dependency><groupId>org.apache.rocketmq</groupId><artifactId>rocketmq-client</artifactId><version>5.3.1</version></dependency>

本文代码用Spring Boot Starter方式,更贴合微服务项目实战。

1.2 application.yml配置

rocketmq:# NameServer地址name-server:127.0.0.1:9876producer:# Producer组名(必须唯一)group:vending_producer_group# 发送超时时间(毫秒)send-message-timeout:3000# 同步发送失败重试次数retry-times-when-send-failed:2# 异步发送失败重试次数retry-times-when-send-async-failed:2# 消息体最大值(默认4MB)max-message-size:4194304server:port:8081

几个配置项解释一下:

  • group:Producer组名,用于标识一组Producer实例。同一组内的Producer发送同一类消息,Broker做故障隔离时按组管理。
  • send-message-timeout:发送超时,默认3秒。售货柜场景建议保持3秒,太短容易误判超时,太长拖慢接口。
  • retry-times-when-send-failed:同步发送失败重试次数,默认2次(加上首次共3次尝试)。

二、三种消息发送方式

RocketMQ的Producer支持三种发送方式,区别在于"发完之后等不等Broker确认"。

2.1 同步发送(Sync)

原理:Producer发送消息后阻塞等待,直到Broker返回确认(ACK)才继续往下走。

Producer --发送消息--> Broker Producer <--阻塞等待-- Broker(写入CommitLog+构建索引后返回ACK) Producer 继续执行后续代码

适用场景:重要消息、不能丢的消息。比如售货柜关门后的"订单创建"消息、支付成功消息。这些消息丢了要么用户白拿商品,要么扣了钱没出货,都是事故。

代码示例

@RestController@RequestMapping("/order")publicclassOrderController{@ResourceprivateRocketMQTemplaterocketMQTemplate;/** * 售货柜关门,生成订单(同步发送) */@PostMapping("/closeDoor")publicStringcloseDoor(@RequestBodyCloseDoorRequestrequest){// 1. 本地业务:创建订单Orderorder=orderService.createOrder(request.getDeviceId(),request.getGoods());// 2. 同步发送消息到MQMessage<OrderMessage>message=MessageBuilder.withPayload(newOrderMessage(order.getId(),order.getDeviceId(),order.getAmount())).build();SendResultsendResult=rocketMQTemplate.syncSend("order_topic",message);// 3. 根据发送结果处理if(sendResult.getSendStatus()==SendStatus.SEND_OK){log.info("订单消息发送成功, msgId={}, orderId={}",sendResult.getMsgId(),order.getId());return"订单创建成功,请前往支付";}else{log.error("订单消息发送失败, status={}",sendResult.getSendStatus());// 发送失败的处理:记录日志、告警、本地重试等return"订单创建成功,但通知下游服务失败,已记录待补偿";}}}

SendResult返回的状态有四种:

状态含义处理方式
SEND_OK发送成功且Broker刷盘成功正常
FLUSH_DISK_TIMEOUT发送成功但刷盘超时消息已在Broker内存,大概率不丢
FLUSH_SLAVE_TIMEOUT发送成功但Slave同步超时Master有数据,Slave可能没同步上
SLAVE_NOT_AVAILABLESlave不可用只有Master有数据,存在单点风险

后三种在异步刷盘+无Slave配置下不会出现,开发环境不用太关注。

2.2 异步发送(Async)

原理:Producer发送消息后立即返回,不阻塞。Broker处理完成后通过回调函数通知Producer结果。

Producer --发送消息--> Broker Producer 继续执行(不等待) ...Broker处理后... Producer <--回调通知-- Broker(成功/失败)

适用场景:对响应时间敏感,但又需要知道发送结果的场景。比如售货柜的"用户行为上报"——用户拿商品的动作要快速记录,但不能因为MQ发送拖慢关门流程。

代码示例

@ServicepublicclassUserBehaviorService{@ResourceprivateRocketMQTemplaterocketMQTemplate;/** * 异步上报用户行为(异步发送) */publicvoidreportBehavior(StringdeviceId,StringuserId,Stringaction){BehaviorMessagebehavior=newBehaviorMessage(deviceId,userId,action,System.currentTimeMillis());// 异步发送,传入回调函数rocketMQTemplate.asyncSend("behavior_topic",behavior,newSendCallback(){@OverridepublicvoidonSuccess(SendResultsendResult){log.info("行为上报消息发送成功, msgId={}",sendResult.getMsgId());}@OverridepublicvoidonException(Throwablethrowable){// 发送失败的回调log.error("行为上报消息发送失败",throwable);// 失败处理:写本地日志表,定时任务补偿重发localFailLogService.save("behavior_topic",behavior,throwable.getMessage());}});// 发送调用立即返回,不阻塞log.debug("行为上报消息已提交异步发送");}}

注意点:异步发送的回调在Producer内部的线程池中执行,不是调用线程。如果回调里要操作共享数据,注意线程安全。另外,回调失败后不会自动重试(和同步发送不同),需要在onException中自己处理。

2.3 单向发送(Oneway)

原理:Producer只负责把消息发出去,不等Broker确认,也不注册回调。发完就忘(Fire and Forget)。

Producer --发送消息--> Broker Producer 立即返回(不等任何响应)

适用场景:对可靠性要求极低、对吞吐量要求极高的场景。比如日志收集——丢几条日志无所谓,但每秒要能发几万条。

代码示例

@ServicepublicclassDeviceLogService{@ResourceprivateRocketMQTemplaterocketMQTemplate;/** * 售货柜运行日志上报(单向发送) */publicvoidsendDeviceLog(StringdeviceId,StringlogLevel,Stringcontent){DeviceLogMessagelogMsg=newDeviceLogMessage(deviceId,logLevel,content,Instant.now());// 单向发送:不等结果,发了就完事rocketMQTemplate.sendOneWay("device_log_topic",logMsg);}}

Oneway没有返回值,没有回调,调用完就结束了。底层实现是Netty的单向请求,不等待响应直接释放。

2.4 三种方式对比

对比项同步发送异步发送单向发送
可靠性最高最低
发送速度最慢(阻塞等待)快(不阻塞)最快(不等响应)
返回结果SendResult回调通知
失败重试自动重试需在回调中手动处理不处理
吞吐量
适用场景订单、支付等重要消息行为上报、通知推送日志收集、埋点
售货柜场景订单创建、支付回调库存同步、出货上报设备日志、操作埋点

选型原则:能用同步就别用异步(简单可靠),对性能有要求才异步,日志类用单向。别为了"看起来高级"就用异步发送,异步带来的回调处理、失败补偿逻辑的复杂度不值得。

三、消息消费方式

3.1 Push模式(DefaultPushConsumer)

Push模式是开发中最常用的消费方式,用起来像Broker主动推消息过来。

Spring Boot中的实现

@Slf4j@Component@RocketMQMessageListener(topic="order_topic",// 消费的TopicconsumerGroup="order_consumer_group",// 消费者组consumeMode=ConsumeMode.CONCURRENTLY,// 并发消费messageModel=MessageModel.CLUSTERING// 集群模式)publicclassOrderMessageConsumerimplementsRocketMQListener<OrderMessage>{@ResourceprivateInventoryServiceinventoryService;@ResourceprivatePushServicepushService;@OverridepublicvoidonMessage(OrderMessagemessage){log.info("收到订单消息: orderId={}, deviceId={}, amount={}",message.getOrderId(),message.getDeviceId(),message.getAmount());try{// 1. 扣减库存inventoryService.deduct(message.getDeviceId(),message.getGoodsList());// 2. 推送支付通知给用户pushService.notifyPay(message.getOrderId(),message.getUserId(),message.getAmount());// 正常返回即表示消费成功}catch(Exceptione){log.error("消费订单消息失败, orderId={}",message.getOrderId(),e);// 抛出异常,MQ会自动重试thrownewRuntimeException("消费失败,触发重试",e);}}}

关键点

  • onMessage方法正常返回 = 消费成功,MQ标记消息已消费
  • onMessage方法抛异常 = 消费失败,MQ按延迟等级重试(10s→30s→1min→2min→…→最多16次)
  • 重试16次还失败 → 进入死信队列(DLQ),人工处理

3.2 Pull模式(DefaultPullConsumer)

Pull模式是消费者主动从Broker拉取消息,需要自己管理消费位点(Offset)、拉取频率等。

publicclassPullConsumerDemo{publicstaticvoidmain(String[]args)throwsException{DefaultLitePullConsumerconsumer=newDefaultLitePullConsumer("pull_consumer_group");consumer.setNamesrvAddr("127.0.0.1:9876");consumer.subscribe("device_log_topic","*");// *表示订阅所有Tagconsumer.start();try{while(true){// 拉取消息(阻塞等待,最多等5秒)List<MessageExt>messages=consumer.poll(Duration.ofSeconds(5));for(MessageExtmsg:messages){Stringbody=newString(msg.getBody(),StandardCharsets.UTF_8);System.out.println("拉取到消息: "+body);// 处理消息...// 注意:Pull模式需要自己管理Offset}// 手动提交消费位点consumer.commitSync();}}finally{consumer.shutdown();}}}

Pull模式更灵活但也更繁琐,需要自己处理位点管理、负载均衡、异常重试。除非有特殊需求(比如要精确控制拉取速率、批量处理),否则不推荐日常使用。

3.3 Push本质也是Pull

重点说清楚一个概念:RocketMQ的Push模式底层也是Pull

DefaultPushConsumer内部并不是Broker主动推送,而是Consumer用长轮询(Long Polling)机制拉取:

Consumer向Broker发Pull请求 ├── Broker有消息 → 立即返回消息 └── Broker没消息 → 请求挂起,Hold在Broker端(默认5秒) ├── 5秒内有新消息写入 → 唤醒挂起的请求,立即返回 └── 5秒后仍无消息 → 返回空,Consumer重新发起Pull请求

这种设计的好处:

  • 不会像纯Push那样,Consumer处理不过来时被Broker压垮
  • 也不会像纯Pull那样,没消息时空轮询浪费资源
  • Consumer可以控制拉取速率(通过pullBatchSize参数)

所以你用Push模式时,如果Consumer处理慢了,消息会在Broker端堆积,不会把Consumer搞崩。这也是为什么RocketMQ天然适合做"削峰"。

四、完整实战:售货柜订单消息发送与消费

把前面学的串起来,走一个完整的业务流程。

4.1 消息定义

@Data@AllArgsConstructor@NoArgsConstructorpublicclassOrderMessageimplementsSerializable{privateStringorderId;// 订单IDprivateStringdeviceId;// 设备IDprivateStringuserId;// 用户IDprivateList<String>goodsList;// 商品列表privateBigDecimalamount;// 订单金额privateLongtimestamp;// 创建时间戳}

注意:RocketMQ消息体必须是可序列化的对象(实现Serializable)或byte[]。推荐用JSON序列化,RocketMQ Spring Boot Starter默认用Jackson做JSON转换。

4.2 Producer端:关门后发送订单消息

@Slf4j@RestController@RequestMapping("/api/vending")publicclassVendingController{@ResourceprivateRocketMQTemplaterocketMQTemplate;@ResourceprivateOrderServiceorderService;/** * 售货柜关门接口 * 流程:识别商品 → 生成订单 → 同步发送MQ消息 → 返回给用户 */@PostMapping("/closeDoor")publicResult<String>closeDoor(@RequestBodyCloseDoorRequestreq){// 1. 识别商品(视觉/重力传感器)List<String>goods=deviceService.identifyGoods(req.getDeviceId());// 2. 生成订单Orderorder=orderService.createOrder(req.getDeviceId(),req.getUserId(),goods);// 3. 构建MQ消息OrderMessagemsg=newOrderMessage(order.getId(),req.getDeviceId(),req.getUserId(),goods,order.getAmount(),System.currentTimeMillis());// 4. 同步发送(订单消息不能丢)SendResultresult=rocketMQTemplate.syncSend("order_topic:order_created",// Topic:Tag格式MessageBuilder.withPayload(msg).build(),3000// 超时3秒);if(result.getSendStatus()!=SendStatus.SEND_OK){// 发送失败:订单已创建,记录待补偿log.error("MQ发送失败,记录补偿表 orderId={} status={}",order.getId(),result.getSendStatus());compensateService.record(order.getId(),"order_topic",msg);}// 5. 返回给用户returnResult.success(order.getId());}}

4.3 Consumer端:库存服务消费订单消息

@Slf4j@Component@RocketMQMessageListener(topic="order_topic",consumerGroup="inventory_consumer_group",consumeMode=ConsumeMode.CONCURRENTLY,consumeThreadMax=20// 最大消费线程数)publicclassInventoryConsumerimplementsRocketMQListener<OrderMessage>{@ResourceprivateInventoryServiceinventoryService;@OverridepublicvoidonMessage(OrderMessagemessage){log.info("库存服务消费订单消息: orderId={}, deviceId={}, goods={}",message.getOrderId(),message.getDeviceId(),message.getGoodsList());// 1. 扣减设备本地库存inventoryService.deductLocal(message.getDeviceId(),message.getGoodsList());// 2. 同步到总部ERP库存inventoryService.syncToErp(message.getDeviceId(),message.getGoodsList());// 正常返回 = 消费成功// 如果这里抛异常,MQ会自动重试}}

4.4 Consumer端:推送服务消费同一消息

@Slf4j@Component@RocketMQMessageListener(topic="order_topic",consumerGroup="push_consumer_group",consumeMode=ConsumeMode.CONCURRENTLY)publicclassPushConsumerimplementsRocketMQListener<OrderMessage>{@ResourceprivatePushServicepushService;@OverridepublicvoidonMessage(OrderMessagemessage){log.info("推送服务消费订单消息: orderId={}, userId={}",message.getOrderId(),message.getUserId());// 给用户推送"请前往支付"的通知pushService.sendPayNotification(message.getUserId(),message.getOrderId(),message.getAmount());}}

注意:两个Consumer的consumerGroup不同(inventory_consumer_groupvspush_consumer_group),所以它们各自消费全量消息。如果设成同一个Group,消息只会被其中一个消费——这是新手最容易踩的坑。

4.5 完整流程图

用户关门 │ ▼ VendingController.closeDoor() ├── 识别商品 ├── 创建订单(DB) ├── 同步发送MQ消息 → order_topic └── 返回订单ID给用户 order_topic消息 │ ├──→ InventoryConsumer (inventory_consumer_group) │ ├── 扣减本地库存 │ └── 同步ERP库存 │ └──→ PushConsumer (push_consumer_group) └── 推送支付通知给用户

五、消费幂等性:新手必知

最后一个重点:RocketMQ不保证消息不重复投递

什么意思?比如库存消费者处理完一条消息,正准备返回成功时,Consumer实例突然挂了,MQ没收到确认,重启后会重新投递这条消息。结果库存被扣了两次。

解决方案叫幂等性——同一条消息被消费多次,结果和消费一次一样。常见做法:

@OverridepublicvoidonMessage(OrderMessagemessage){// 用订单ID做幂等键StringdedupKey="order:"+message.getOrderId();// 1. 检查是否已处理过(Redis或DB唯一索引)if(redisTemplate.hasKey(dedupKey)){log.info("消息已处理过,跳过 orderId={}",message.getOrderId());return;}// 2. 执行业务逻辑inventoryService.deduct(message.getDeviceId(),message.getGoodsList());// 3. 标记已处理redisTemplate.opsForValue().set(dedupKey,"1",24,TimeUnit.HOURS);}

或者用数据库唯一索引:

-- 幂等记录表CREATETABLEmessage_idempotent(msg_idVARCHAR(64)PRIMARYKEY,consumer_groupVARCHAR(64),create_timeDATETIMEDEFAULTCURRENT_TIMESTAMP,UNIQUEKEYuk_msg_group(msg_id,consumer_group));

消费前先插入这条记录,插入成功说明是第一次消费,插入失败(违反唯一约束)说明已消费过,直接跳过。简单粗暴有效。

六、小结

这一篇我们实战了RocketMQ三种发送方式(同步/异步/单向),两种消费方式(Push/Pull),以及完整的售货柜订单消息发送和消费代码。最后强调了消费幂等性这个新手必知的关键点。下一篇深入Topic、Queue、Tag机制,设计售货柜多设备消息隔离方案。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询