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/9/19 5:19:40 阅读更多 →
生命涌现的小龙虾技能之【Contactless Health Risk Screening Tool | 非接触式健康风险检测分析工具】简介

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

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

2026/9/19 3:32:25 阅读更多 →
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/9/21 20:35:21 阅读更多 →

最新新闻

3个坑教你手写实现装饰设计培训项目

3个坑教你手写实现装饰设计培训项目

3个坑教你手写实现装饰设计培训项目 版本升级后 API 全变了,昨天还能跑的装饰工程数据接口,今天全报 404。别急着骂娘,这其实是底层逻辑变了。很多从业者还在死记硬背旧版参数,结果被新版校验机制卡得死死的。与其天天查文档改参数,不如直接手…

2026/9/22 21:58:20 阅读更多 →
qq播放器下载源码拆解:3个实战项目级技巧

qq播放器下载源码拆解:3个实战项目级技巧

qq播放器下载源码拆解:3个实战项目级技巧 学会语法却不知怎么搭项目,是大多数开发者转行或进阶时的最大卡点。很多人背下了 Python 的类继承、Java 的并发包,甚至刷完了 LeetCode 的前 200 题,但面对一个真实的…

2026/9/22 21:58:20 阅读更多 →
2026最新条码查询价格接口源码拆解

2026最新条码查询价格接口源码拆解

2026最新条码查询价格接口源码拆解 配置环境就卡半天,这种痛谁懂?我见过太多开发者,为了接一个 条码查询价格 的功能,在依赖库里折腾一下午,结果连报错日志都看不清。别急,今天咱们不聊虚的,直接掀开底裤,看看2026年主流电商与供应链系统中…

2026/9/22 21:58:20 阅读更多 →
3步搞定火山石幼龙攻略:图解原理让你从入门到实战

3步搞定火山石幼龙攻略:图解原理让你从入门到实战

3步搞定火山石幼龙攻略:图解原理让你从入门到实战 学会语法却不知怎么搭项目?这大概是很多刚接触新工具或新框架的开发者最大的痛点。你背下了API,看懂了文档,但一动手写代码就卡壳,不知道模块怎么串联,数据流怎么走。别急,今天这篇火山石幼龙攻略…

2026/9/22 21:58:20 阅读更多 →
别背废话了!2868面试最佳实践,3分钟吃透核心考点

别背废话了!2868面试最佳实践,3分钟吃透核心考点

别背废话了!2868面试最佳实践,3分钟吃透核心考点 官方文档翻了三遍还是抓不住重点?别急,大厂面试官眼里,2868的核心逻辑其实只有三层。今天咱们直接撕开官方源码仓库的底层逻辑,用最佳实践帮你把这块硬骨头啃下来。…

2026/9/22 21:58:20 阅读更多 →
3步拆解为什么说双缝实验恐怖图解原理

3步拆解为什么说双缝实验恐怖图解原理

3步拆解为什么说双缝实验恐怖图解原理 版本升级后 API 全变了,代码跑不通,文档还跟不上。很多开发者在重构遗留系统时,常被这种“黑盒”逻辑卡死:输入输出明确,但中间过程完全不可观测,就像量子力学里的双缝实验一样令人抓狂。其实,这种“观测即…

2026/9/22 21:57:19 阅读更多 →

日新闻

3台商务办公笔记本实测:手写实现环境配置,告别卡半天

3台商务办公笔记本实测:手写实现环境配置,告别卡半天

3台商务办公笔记本实测:手写实现环境配置,告别卡半天 配置环境就卡半天?别怪机器慢,多半是你没选对工具链。在Java、Go或Python的项目现场, 手写实现…

2026/9/22 0:00:41 阅读更多 →
剑帝加点速查手册:3分钟搞懂核心逻辑

剑帝加点速查手册:3分钟搞懂核心逻辑

剑帝加点速查手册:3分钟搞懂核心逻辑 面试被问原理答不上来,是不是常态?别慌。很多开发者对着 GitHub 开源仓库里的代码发呆,看似简单实则暗藏玄机。今天这份【剑帝加点】速查手册,直接带你拆解核心实现,把面试必考的原理讲透。…

2026/9/22 0:00:41 阅读更多 →
手写实现图片压缩网站核心:搞定WebP转换与质量调优

手写实现图片压缩网站核心:搞定WebP转换与质量调优

手写实现图片压缩网站核心:搞定WebP转换与质量调优 复制来的代码跑不通不知道怎么调?别慌,这种“复制粘贴地狱”在开发圈太常见了。尤其是做 图片压缩网站…

2026/9/22 0:00:41 阅读更多 →

周新闻

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

直接铺开项目本身吧。这几个月我一直在折腾一件事:用Flutter给OpenHarmony做一款游戏集合类的App,说白了就是把若干小游戏塞进一个壳里,用统一入口分发。这个方向本身不算新鲜,真正让我花了不少心思的,是首页那堆游戏卡…

2026/9/22 4:32:41 阅读更多 →
Word表格编号全攻略:从列表编号到题注交叉引用

Word表格编号全攻略:从列表编号到题注交叉引用

写Word文档,最让人头疼的往往是那些“看起来不起眼”的小问题。比如表格编号这事:今天在表后面多加了两个空白行,明天给客户交稿前发现整个章节的编号全部错位,光是挨个改序号就能耗掉大半个下午。我前阵子帮人整理一份上百页的技…

2026/9/22 4:38:57 阅读更多 →
从第一个站到第二个站:独立开发者的静态网站选型与落地实践

从第一个站到第二个站:独立开发者的静态网站选型与落地实践

1. 项目概述1.1 核心需求解析做独立开发者这几年,说实话,第一个网站上线的那天晚上我兴奋得没睡着。但等它跑了半年,流量惨淡、功能臃肿、代码自己都懒得看第二遍之后,我才慢慢琢磨明白一个道理:第一个网站是练手&…

2026/9/22 8:51:04 阅读更多 →

月新闻

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能分类:[AI/大模型]细分主题:AI 增强型 CI/CD 流水线自动化与 GitOps 实践:Agent 工作流、工具调用与任务拆解:从原型到生产的验收清单很多团队在尝试用大…

2026/9/21 15:36:51 阅读更多 →
容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场分类:[工程技术]细分主题:Kubernetes 生产环境运维与排障实战:可复制的项目复盘模板与决策记录大部分团队的事故复盘报告,最后都变成了躺在 Confluence 或钉…

2026/9/21 15:36:51 阅读更多 →
容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步分类:[工程技术]细分主题:Docker 容器化技术与镜像安全管理:核心链路的逐步实现与关键代码取舍面对一个积累了五六年历史包袱的单体架构应用(包含 Web 接口、后台…

2026/9/22 2:43:42 阅读更多 →