Spring Boot 整合 RocketMQ 完全指南
rocketmq-spring-boot-starter 的版本选择与依赖引入在开始写代码之前我们面临第一个选择题用哪个版本的 Starter这看似是个小问题但在 Spring Boot 3.x 时代版本选不对项目可能连启动都起不来。版本选型的核心原则Spring Boot 版本 推荐 Starter 版本 说明Spring Boot 2.x 2.2.3 社区验证最稳生产案例最多Spring Boot 3.x 2.2.3 2.2.3 已支持 Jakarta EE兼容 Spring Boot 3需要 RocketMQ 5.x 新特性 2.3.x 可用但生产案例相对较少⚠️ 避坑提示2.2.0 以下版本使用 javax.* 包与 Spring Boot 3.x 的 jakarta.* 不兼容直接报错。Maven 依赖以最稳定的 2.2.3 为例org.apache.rocketmq rocketmq-spring-boot-starter 2.2.3 这个 Starter 已经传递依赖了 rocketmq-client所以你不需要再单独引入客户端依赖。但如果想精确控制客户端版本和服务端对齐可以额外声明 org.apache.rocketmq rocketmq-client 5.1.0 生产者的配置与使用 基础配置application.ymlrocketmq:name-server: 127.0.0.1:9876 # NameServer 地址多个用分号分隔producer:group: order-producer-group # 生产者组名send-message-timeout: 3000 # 发送超时时间毫秒retry-times-when-send-failed: 2 # 同步发送失败重试次数retry-next-server: true # 失败后是否换 Broker 重试compress-msg-body-over-how-much: 4096 # 超过多少字节压缩生产级配置建议NameServer 至少配置 2 个地址避免单点故障retry-next-server: true 开启后发送失败会自动换 Broker 重试提升可用性不要完全依赖自动重试解决所有问题业务层必须有兜底方案生产者代码import org.apache.rocketmq.client.producer.SendResult;import org.apache.rocketmq.spring.core.RocketMQTemplate;import org.apache.rocketmq.spring.support.RocketMQHeaders;import org.springframework.messaging.Message;import org.springframework.messaging.support.MessageBuilder;import org.springframework.stereotype.Service;Servicepublic class OrderProducer {private final RocketMQTemplate rocketMQTemplate; public OrderProducer(RocketMQTemplate rocketMQTemplate) { this.rocketMQTemplate rocketMQTemplate; } /** * 同步发送消息最常用 */ public SendResult sendOrder(String orderId, String content) { // destination 格式topic:tag String destination order-topic:order-create; MessageString message MessageBuilder .withPayload(content) .setHeader(RocketMQHeaders.KEYS, orderId) // 设置业务 Key用于查询和幂等 .build(); SendResult result rocketMQTemplate.syncSend(destination, message); // 生产环境需要检查 result.getSendStatus() return result; } /** * 异步发送消息 */ public void sendOrderAsync(String orderId, String content) { String destination order-topic:order-create; MessageString message MessageBuilder .withPayload(content) .setHeader(RocketMQHeaders.KEYS, orderId) .build(); rocketMQTemplate.asyncSend(destination, message, sendResult - { // 回调处理 if (sendResult.getSendStatus().name().equals(SEND_OK)) { System.out.println(异步发送成功 sendResult.getMsgId()); } }); } /** * 单向发送不关心结果最快 */ public void sendOrderOneway(String orderId, String content) { String destination order-topic:order-create; MessageString message MessageBuilder .withPayload(content) .setHeader(RocketMQHeaders.KEYS, orderId) .build(); rocketMQTemplate.sendOneWay(destination, message); }}KEY 的作用非常重要设置 RocketMQHeaders.KEYS 有三个核心用途消息查询在 Dashboard 中按业务 Key 快速定位消息事务回查事务消息回查时用于关联业务数据幂等控制消费者端用 Key 做去重判断消费者的配置与使用基础配置application.ymlrocketmq:name-server: 127.0.0.1:9876consumer:group: order-consumer-group # 消费者组名consume-mode: CLUSTERING # 消费模式CLUSTERING集群或 BROADCASTING广播consume-thread-min: 5 # 最小消费线程数consume-thread-max: 20 # 最大消费线程数consume-message-batch-max-size: 1 # 批量消费最大条数pull-batch-size: 32 # 批量拉取最大条数消费者代码使用 RocketMQMessageListener 注解import org.apache.rocketmq.spring.annotation.ConsumeMode;import org.apache.rocketmq.spring.annotation.MessageModel;import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;import org.apache.rocketmq.spring.core.RocketMQListener;import org.springframework.stereotype.Component;ComponentRocketMQMessageListener(topic “order-topic”,consumerGroup “order-consumer-group”,selectorExpression “order-create || order-pay”, // Tag 过滤* 表示全部consumeMode ConsumeMode.CONCURRENTLY, // 并发消费messageModel MessageModel.CLUSTERING, // 集群模式maxReconsumeTimes 16 // 最大重试次数-1 表示 16 次)public class OrderConsumer implements RocketMQListener {Override public void onMessage(String message) { // 1️⃣ 幂等校验最重要 // 2️⃣ 业务处理 System.out.println(消费订单消息 message); }}生产铁律一定要做幂等。幂等方式 适用场景数据库唯一键 订单、账务等有明确业务 ID 的场景Redis SETNX 高并发场景快速去重消息 KEY 通用方案配合业务状态判断事务消息的整合与实现事务消息是 RocketMQ 最有价值、也最容易用错的功能。它的核心是保证“本地事务”和“消息发送”要么一起成功要么一起失败。事务消息的完整流程本地数据库BrokerProducer业务应用本地数据库BrokerProducer业务应用Broker 未收到最终确认触发回查loop[事务回查默认每 60 秒]alt[本地事务成功][本地事务失败][事务状态未知异常/超时]发送事务消息发送半消息Half Message半消息持久化暂不可消费半消息发送成功回调执行本地事务执行本地事务如更新订单状态7a. 事务提交成功8a. 返回 COMMIT9a. 提交事务10a. 半消息→正式消息可消费7b. 事务回滚8b. 返回 ROLLBACK9b. 回滚事务10b. 删除半消息8c. 返回 UNKNOWN9c. 发起回查请求10c. 检查本地事务状态11c. 查询业务数据12c. 返回查询结果13c. 返回 COMMIT/ROLLBACK14c. 提交最终事务状态第一步定义事务监听器import org.apache.rocketmq.spring.annotation.RocketMQTransactionListener;import org.apache.rocketmq.spring.core.RocketMQLocalTransactionListener;import org.apache.rocketmq.spring.core.RocketMQLocalTransactionState;import org.springframework.messaging.Message;import org.springframework.stereotype.Service;ServiceRocketMQTransactionListener(txProducerGroup “order-tx-producer-group”) // 必须与发送方组名一致public class OrderTransactionListener implements RocketMQLocalTransactionListener {Autowired private OrderService orderService; /** * 执行本地事务 */ Override public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) { String orderId (String) msg.getHeaders().get(orderId); try { // 执行本地业务更新订单状态 boolean success orderService.updateOrderStatus(orderId, PAID); // 根据执行结果返回 COMMIT 或 ROLLBACK return success ? RocketMQLocalTransactionState.COMMIT : RocketMQLocalTransactionState.ROLLBACK; } catch (Exception e) { // 返回 UNKNOWN等待 Broker 回查 return RocketMQLocalTransactionState.UNKNOWN; } } /** * 事务回查方法 */ Override public RocketMQLocalTransactionState checkLocalTransaction(Message msg) { String orderId (String) msg.getHeaders().get(orderId); // 查询本地事务状态 String status orderService.getOrderStatus(orderId); if (PAID.equals(status)) { return RocketMQLocalTransactionState.COMMIT; } else if (CANCELLED.equals(status)) { return RocketMQLocalTransactionState.ROLLBACK; } // 状态仍未知继续等待下次回查 return RocketMQLocalTransactionState.UNKNOWN; }}第二步发送事务消息Servicepublic class OrderTransactionProducer {private final RocketMQTemplate rocketMQTemplate; public OrderTransactionProducer(RocketMQTemplate rocketMQTemplate) { this.rocketMQTemplate rocketMQTemplate; } public void createOrderWithTransaction(String orderId) { String destination order-tx-topic:order-create; MessageString message MessageBuilder .withPayload(订单创建 orderId) .setHeader(orderId, orderId) // 传递给事务监听器 .setHeader(RocketMQHeaders.KEYS, orderId) .build(); // 发送事务消息 rocketMQTemplate.sendMessageInTransaction(destination, message, null); }}消息监听器的多种用法RocketMQMessageListener 注解支持丰富的配置选项按 Tag 过滤RocketMQMessageListener(topic “order-topic”,consumerGroup “order-consumer-group”,selectorExpression “order-create || order-pay” // 只消费指定 Tag)2. 按 SQL92 表达式过滤RocketMQMessageListener(topic “order-topic”,consumerGroup “order-consumer-group”,selectorType SelectorType.SQL92, // 使用 SQL92 过滤selectorExpression “amount 1000 AND region ‘SH’” // SQL92 表达式)3. 顺序消费RocketMQMessageListener(topic “order-topic”,consumerGroup “order-consumer-group”,consumeMode ConsumeMode.ORDERLY // 顺序消费模式)public class OrderlyConsumer implements RocketMQListener {Overridepublic void onMessage(String message) {// 同一个 Queue 的消息会按顺序被消费}}4. 广播消费RocketMQMessageListener(topic “order-topic”,consumerGroup “order-consumer-group”,messageModel MessageModel.BROADCASTING // 广播模式)public class BroadcastConsumer implements RocketMQListener {Overridepublic void onMessage(String message) {// 每个消费者实例都会收到这条消息}}5. 接收原始 MessageExt获取更多元数据import org.apache.rocketmq.common.message.MessageExt;RocketMQMessageListener(topic “order-topic”,consumerGroup “order-consumer-group”)public class FullConsumer implements RocketMQListener {Overridepublic void onMessage(MessageExt message) {String msgId message.getMsgId();String body new String(message.getBody());String tags message.getTags();long bornTime message.getBornTimestamp();// 可以获取更丰富的消息元数据}}消费者线程池配置RocketMQMessageListener 中的线程池配置参数 默认值 说明consumeThreadNumber 20 消费线程数2.2.3 新参数推荐使用consumeThreadMax 64 已废弃5.x 不再推荐使用配置示例RocketMQMessageListener(topic “order-topic”,consumerGroup “order-consumer-group”,consumeThreadNumber 40 // 调大线程数提升并发消费能力)线程池调优建议消息处理逻辑轻量如简单计算→ 线程数可设大一些如 40-60消息处理逻辑重量如调用第三方 API、复杂数据库操作→ 线程数设小一些如 10-20避免资源争抢监控消费 TPS 和系统负载动态调整消息转换器的使用RocketMQ Spring Boot Starter 默认使用 RocketMQMessageConverter 进行消息序列化和反序列化。默认行为发送时对象 → JSON 字符串接收时JSON 字符串 → 目标类型自定义消息转换器import org.springframework.context.annotation.Bean;import org.springframework.context.annotation.Configuration;import org.springframework.messaging.converter.MessageConverter;Configurationpublic class RocketMQConfig {Bean public MessageConverter rocketMQMessageConverter() { // 自定义转换逻辑 return new CustomMessageConverter(); }}常见场景使用 Protobuf 替代 JSON提升序列化性能和减小消息体积使用 Kryo 等高性能序列化框架处理特殊的数据格式如二进制数据多环境配置与管理

相关新闻

Token 便宜,不等于 AI 便宜

Token 便宜,不等于 AI 便宜

过去两年,AI 行业最热闹的竞争几乎都围绕同一个问题:谁更强、谁更便宜、谁更快、谁更容易被大规模使用。任何新技术进入市场时,最先被比较的往往都是最表面的指标——价格、速度、参数、门槛。AI 也不例外。 模型刚开始从实验室走向现实世界时…

2026/8/17 1:52:23 阅读更多 →
生命涌现的小龙虾技能之【Contactless Health Risk Screening Tool | 非接触式健康风险检测分析工具】简介

生命涌现的小龙虾技能之【Contactless Health Risk Screening Tool | 非接触式健康风险检测分析工具】简介

🩺 Contactless Health Risk Screening Tool | 非接触式健康风险检测分析工具 智能健康/识别分析中枢 图片/视频智能分析 结构化报告 历史报告云端查询 🧭 技能概览 | Overview 模块内容🏷️ 技能名称非接触式健康风险检测分析工具&#…

2026/8/15 13:17:03 阅读更多 →
LongNet源码解析:DilatedAttention类的实现细节与设计思路

LongNet源码解析:DilatedAttention类的实现细节与设计思路

LongNet源码解析:DilatedAttention类的实现细节与设计思路 【免费下载链接】LongNet Implementation of plug in and play Attention from "LongNet: Scaling Transformers to 1,000,000,000 Tokens" 项目地址: https://gitcode.com/gh_mirrors/lo/Long…

2026/8/17 14:41:14 阅读更多 →

最新新闻

商照轨道灯口碑哪家强?这些厂商名声真响亮!

商照轨道灯口碑哪家强?这些厂商名声真响亮!

商照轨道灯口碑王:5年光衰不足5%,客户复购超70%在商业照明领域,轨道灯一直是商铺、展厅、办公空间的重点照明利器。面对市场上“雷士、三雄极光”等一众大牌,采购方往往陷入选择困境:商照轨道灯口碑哪家强?…

2026/8/18 1:51:46 阅读更多 →
基于PCA降维的高维数据可视化预处理工具——Python大数据分析实战

基于PCA降维的高维数据可视化预处理工具——Python大数据分析实战

摘要 在大数据时代,高维数据的可视化与分析面临“维度灾难”的严峻挑战。主成分分析(PCA)作为最经典的无监督线性降维方法,通过正交变换将原始高维特征投影至低维子空间,在保留最大方差信息的同时大幅降低数据维度,为后续可视化、聚类、分类等任务奠定基础。本文从PCA的…

2026/8/18 1:51:46 阅读更多 →
SRD-12VDC-SL-C继电器深度解析:从参数选型到驱动电路实战

SRD-12VDC-SL-C继电器深度解析:从参数选型到驱动电路实战

1. 项目缘起:一个被忽视的“开关”核心最近在折腾一个智能家居的改造项目,需要控制一个功率不小的交流电机。在选型控制电路的核心——继电器时,我遇到了一个老朋友:SRD-12VDC-SL-C。这个型号的继电器,在电子爱好者、工…

2026/8/18 1:51:46 阅读更多 →
基于随机森林的电信用户流失预测与特征重要性分析 —— 一个完整的Python大数据分析实战

基于随机森林的电信用户流失预测与特征重要性分析 —— 一个完整的Python大数据分析实战

摘要 在高度饱和的电信市场中,用户流失(Churn)是运营商面临的核心挑战之一。准确预测潜在流失客户并理解流失的关键驱动因素,能为精准挽留策略提供数据支撑。本文以电信客户数据集为对象,系统性地展示如何利用Python进行数据探索、预处理、特征工程、基于随机森林的流失预…

2026/8/18 1:51:46 阅读更多 →
很多人打 CTF 直接踩坑!赛事定义、核心考点、技术储备一次性讲透

很多人打 CTF 直接踩坑!赛事定义、核心考点、技术储备一次性讲透

在网络安全领域,CTF(Capture The Flag,夺旗赛)是检验技术实力的 “试金石”,也是白帽黑客成长的 “练兵场”。对于刚接触网络安全的新手来说,CTF 既神秘又充满吸引力 —— 它不像传统考试那样侧重理论&…

2026/8/18 1:51:45 阅读更多 →
电动垂直起降(eVTOL)技术解析与城市空中交通应用

电动垂直起降(eVTOL)技术解析与城市空中交通应用

1. 项目概述:低空经济时代的"空中出租车"创新实践 当全球主要城市都在为地面交通拥堵寻找解决方案时,上海东方枢纽出现的这抹亮色格外引人注目。御风未来带来的这款"空中出租车"并非科幻电影道具,而是已经完成适航取证、…

2026/8/18 1:50:45 阅读更多 →

日新闻

告别逐帧截图:用 extract-video-ppt 快速提取视频中的 PPT 并一键导出 PDF

告别逐帧截图:用 extract-video-ppt 快速提取视频中的 PPT 并一键导出 PDF

告别逐帧截图:用 extract-video-ppt 快速提取视频中的 PPT 并一键导出 PDF 【免费下载链接】extract-video-ppt extract the ppt in the video 项目地址: https://gitcode.com/gh_mirrors/ex/extract-video-ppt 如果你还停留在"看网课 不停暂停 截图 …

2026/8/18 0:00:57 阅读更多 →
思源宋体TTF一站式上手:7个字重免费商用,从下载到上线的完整走查

思源宋体TTF一站式上手:7个字重免费商用,从下载到上线的完整走查

思源宋体TTF一站式上手:7个字重免费商用,从下载到上线的完整走查 【免费下载链接】source-han-serif-ttf Source Han Serif TTF 项目地址: https://gitcode.com/gh_mirrors/so/source-han-serif-ttf 你是不是也经历过这种时刻:设计稿里…

2026/8/18 0:00:58 阅读更多 →
华硕笔记本控制权回收指南:GHelper 如何用一个 10MB 文件替代 Armoury Crate

华硕笔记本控制权回收指南:GHelper 如何用一个 10MB 文件替代 Armoury Crate

华硕笔记本控制权回收指南:GHelper 如何用一个 10MB 文件替代 Armoury Crate 【免费下载链接】g-helper Lightweight Armoury Crate alternative for Asus laptops with nearly the same functionality. Works with ROG Zephyrus, Flow, TUF, Strix, Scar, ProArt, …

2026/8/18 0:00:59 阅读更多 →

周新闻

基于阿里云与通义千问(Qwen)构建AI应用:从模型调用到生产部署的完整实践指南

基于阿里云与通义千问(Qwen)构建AI应用:从模型调用到生产部署的完整实践指南

如果你是一名开发者,最近可能已经感受到了AI大模型正在从“玩具”变成“生产力工具”的强烈信号。从代码补全到智能Agent,从本地部署到云端API,我们正处在一个技术栈快速重构的节点。然而,面对层出不穷的模型、框架和工具&#xf…

2026/8/17 2:58:27 阅读更多 →
工业通信系统底层逻辑:04 反射——高频能量撞墙之后会发生什么?

工业通信系统底层逻辑:04 反射——高频能量撞墙之后会发生什么?

第四篇:反射——高频能量撞墙之后会发生什么? —— 你以为信号已经过去了,其实它正在回来打你 老Q的现场笔记 第五季,我们正式进入工业神经系统层。这里不再是单个设备的战斗,而是整个工厂“经脉”层面的秩序之战。从这一篇开始,你将第一次看清:看似简单的信号传播,背…

2026/8/17 2:58:30 阅读更多 →
【文章复现】非线性值迭代自适应动态规划(ADP):离散时间非线性系统的策略迭代自适应动态规划算法研究附Matlab代码

【文章复现】非线性值迭代自适应动态规划(ADP):离散时间非线性系统的策略迭代自适应动态规划算法研究附Matlab代码

✅作者简介:热爱科研的Matlab仿真开发者,擅长毕业设计辅导、数学建模、数据处理、建模仿真、程序设计、完整代码获取、论文复现及科研仿真。🍎 往期回顾关注个人主页:Matlab科研工作室👇 关注我领取海量matlab电子书和…

2026/8/17 2:58:32 阅读更多 →

月新闻

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南 【免费下载链接】BaiduNetdiskPlugin-macOS For macOS.百度网盘 破解SVIP、下载速度限制~ 项目地址: https://gitcode.com/gh_mirrors/ba/BaiduNetdiskPlugin-macOS 还在为百度网盘macOS版的龟速下…

2026/8/17 18:54:37 阅读更多 →
终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换

终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换

终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换 【免费下载链接】ncmdump 项目地址: https://gitcode.com/gh_mirrors/ncmd/ncmdump 还在为网易云音乐下载的NCM格式文件无法在其他播放器播放而烦恼吗?ncmdump解密工具帮你轻松解决这个困…

2026/8/17 18:55:16 阅读更多 →
HarmonyOS 应用开发《掌上英语》第81篇: 智能体卡片:为英语学习 App 打造桌面级学习助手

HarmonyOS 应用开发《掌上英语》第81篇: 智能体卡片:为英语学习 App 打造桌面级学习助手

AgentCard 智能体卡片:为英语学习 App 打造桌面级学习助手适用平台:HarmonyOS 7.0 (API 26 Beta)一、引言 HarmonyOS 7.0(API 26 Beta)新增了 AgentCard 智能体卡片能力,这是继 HMAF(鸿蒙智能体框架&#x…

2026/8/17 18:55:55 阅读更多 →