新闻详情

新闻详情

首页 / 资讯中心 / 详情

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

发布时间:2026/9/6 7:53:24来源:尧图网络
基础消息模型实战:普通消息、同步/异步/单向消息发送
基础消息模型实战普通消息、同步/异步/单向消息发送作者黒漂技术佬系列专栏RocketMQ核心原理与无人售货柜项目实战一、Maven依赖和初始化配置1.1 引入依赖Spring Boot项目引入RocketMQ Starter版本和RocketMQ Server对应即可dependencies!-- RocketMQ Spring Boot Starter --dependencygroupIdorg.apache.rocketmq/groupIdartifactIdrocketmq-spring-boot-starter/artifactIdversion2.3.1/version/dependency!-- Spring Boot Web微服务接口 --dependencygroupIdorg.springframework.boot/groupIdartifactIdspring-boot-starter-web/artifactIdversion2.7.18/version/dependency/dependencies如果不用Starter直接用原生Client也行dependencygroupIdorg.apache.rocketmq/groupIdartifactIdrocketmq-client/artifactIdversion5.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# 消息体最大值默认4MBmax-message-size:4194304server:port:8081几个配置项解释一下groupProducer组名用于标识一组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 继续执行后续代码适用场景重要消息、不能丢的消息。比如售货柜关门后的订单创建消息、支付成功消息。这些消息丢了要么用户白拿商品要么扣了钱没出货都是事故。代码示例RestControllerRequestMapping(/order)publicclassOrderController{ResourceprivateRocketMQTemplaterocketMQTemplate;/** * 售货柜关门生成订单同步发送 */PostMapping(/closeDoor)publicStringcloseDoor(RequestBodyCloseDoorRequestrequest){// 1. 本地业务创建订单OrderorderorderService.createOrder(request.getDeviceId(),request.getGoods());// 2. 同步发送消息到MQMessageOrderMessagemessageMessageBuilder.withPayload(newOrderMessage(order.getId(),order.getDeviceId(),order.getAmount())).build();SendResultsendResultrocketMQTemplate.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){BehaviorMessagebehaviornewBehaviorMessage(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){DeviceLogMessagelogMsgnewDeviceLogMessage(deviceId,logLevel,content,Instant.now());// 单向发送不等结果发了就完事rocketMQTemplate.sendOneWay(device_log_topic,logMsg);}}Oneway没有返回值没有回调调用完就结束了。底层实现是Netty的单向请求不等待响应直接释放。2.4 三种方式对比对比项同步发送异步发送单向发送可靠性最高高最低发送速度最慢阻塞等待快不阻塞最快不等响应返回结果SendResult回调通知无失败重试自动重试需在回调中手动处理不处理吞吐量低中高适用场景订单、支付等重要消息行为上报、通知推送日志收集、埋点售货柜场景订单创建、支付回调库存同步、出货上报设备日志、操作埋点选型原则能用同步就别用异步简单可靠对性能有要求才异步日志类用单向。别为了看起来高级就用异步发送异步带来的回调处理、失败补偿逻辑的复杂度不值得。三、消息消费方式3.1 Push模式DefaultPushConsumerPush模式是开发中最常用的消费方式用起来像Broker主动推消息过来。Spring Boot中的实现Slf4jComponentRocketMQMessageListener(topicorder_topic,// 消费的TopicconsumerGrouporder_consumer_group,// 消费者组consumeModeConsumeMode.CONCURRENTLY,// 并发消费messageModelMessageModel.CLUSTERING// 集群模式)publicclassOrderMessageConsumerimplementsRocketMQListenerOrderMessage{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模式DefaultPullConsumerPull模式是消费者主动从Broker拉取消息需要自己管理消费位点Offset、拉取频率等。publicclassPullConsumerDemo{publicstaticvoidmain(String[]args)throwsException{DefaultLitePullConsumerconsumernewDefaultLitePullConsumer(pull_consumer_group);consumer.setNamesrvAddr(127.0.0.1:9876);consumer.subscribe(device_log_topic,*);// *表示订阅所有Tagconsumer.start();try{while(true){// 拉取消息阻塞等待最多等5秒ListMessageExtmessagesconsumer.poll(Duration.ofSeconds(5));for(MessageExtmsg:messages){StringbodynewString(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 消息定义DataAllArgsConstructorNoArgsConstructorpublicclassOrderMessageimplementsSerializable{privateStringorderId;// 订单IDprivateStringdeviceId;// 设备IDprivateStringuserId;// 用户IDprivateListStringgoodsList;// 商品列表privateBigDecimalamount;// 订单金额privateLongtimestamp;// 创建时间戳}注意RocketMQ消息体必须是可序列化的对象实现Serializable或byte[]。推荐用JSON序列化RocketMQ Spring Boot Starter默认用Jackson做JSON转换。4.2 Producer端关门后发送订单消息Slf4jRestControllerRequestMapping(/api/vending)publicclassVendingController{ResourceprivateRocketMQTemplaterocketMQTemplate;ResourceprivateOrderServiceorderService;/** * 售货柜关门接口 * 流程识别商品 → 生成订单 → 同步发送MQ消息 → 返回给用户 */PostMapping(/closeDoor)publicResultStringcloseDoor(RequestBodyCloseDoorRequestreq){// 1. 识别商品视觉/重力传感器ListStringgoodsdeviceService.identifyGoods(req.getDeviceId());// 2. 生成订单OrderorderorderService.createOrder(req.getDeviceId(),req.getUserId(),goods);// 3. 构建MQ消息OrderMessagemsgnewOrderMessage(order.getId(),req.getDeviceId(),req.getUserId(),goods,order.getAmount(),System.currentTimeMillis());// 4. 同步发送订单消息不能丢SendResultresultrocketMQTemplate.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端库存服务消费订单消息Slf4jComponentRocketMQMessageListener(topicorder_topic,consumerGroupinventory_consumer_group,consumeModeConsumeMode.CONCURRENTLY,consumeThreadMax20// 最大消费线程数)publicclassInventoryConsumerimplementsRocketMQListenerOrderMessage{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端推送服务消费同一消息Slf4jComponentRocketMQMessageListener(topicorder_topic,consumerGrouppush_consumer_group,consumeModeConsumeMode.CONCURRENTLY)publicclassPushConsumerimplementsRocketMQListenerOrderMessage{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做幂等键StringdedupKeyorder: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机制设计售货柜多设备消息隔离方案。
网站建设高端定制企业官网
RELATED

相关资讯

更多精彩内容,欢迎继续阅读

较早相关资讯

最新相关资讯

MLP图像生成:从基础原理到睡衣小马风格实践 2026/9/6 11:26:57

MLP图像生成:从基础原理到睡衣小马风格实践

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

阅读更多 →
Zephyr线程模型深度解析:Main Thread与Idle Thread的调度机制 2026/9/6 11:26:57

Zephyr线程模型深度解析:Main Thread与Idle Thread的调度机制

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

阅读更多 →
PCB布线进阶指南:从信号完整性与电源路径到高速电路设计 2026/9/6 11:26:57

PCB布线进阶指南:从信号完整性与电源路径到高速电路设计

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

阅读更多 →
Zotero附件爆满?用PanAttach Sync同步百度网盘实战指南 2026/9/6 11:26:57

Zotero附件爆满?用PanAttach Sync同步百度网盘实战指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

阅读更多 →
分布式采集系统丢数排查指南:从RS485布线到软件补传的完整解决思路 2026/9/6 11:26:57

分布式采集系统丢数排查指南:从RS485布线到软件补传的完整解决思路

做过工业数据采集的朋友,估计都有过这种经历:现场几十上百个分散测点,用分布式采集模块和RS485总线把数据汇聚到上位机,平时跑得好好的,一到生产高峰期或者天气恶劣的时候,采集软件就开始丢数。轻则某几个测…

阅读更多 →
便携式音频与ADC驱动中的OPA365AIDBVR:低失真高速轨到轨运放应用案例解析 2026/9/6 11:23:57

便携式音频与ADC驱动中的OPA365AIDBVR:低失真高速轨到轨运放应用案例解析

OPA365AIDBVR:50MHz零交叉轨到轨CMOS运算放大器深度解析在便携式音频设备、高精度数据采集系统和传感器信号链设计中,运算放大器的选型往往需要在带宽、噪声、失真和功耗之间做出精细权衡。传统的轨到轨运算放大器在输入级切换时存在交叉失真&#xff0c…

阅读更多 →

今日资讯

本周资讯

本月资讯

看完文章仍有疑问?

联系尧图顾问,获取一对一建站咨询

立即免费咨询 📞 400-888-8888
📞