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/7/22 18:15:15阅读更多 →
Linux 进程通信:管道、信号、共享内存、消息队列详解

Linux 进程通信:管道、信号、共享内存、消息队列详解

Linux 进程通信:管道、信号、共享内存、消息队列详解 一、引言 前面讲了多进程和多线程,但进程之间如何交换数据?总不能靠全局变量吧——进程的地址空间是隔离的,你在进程 A 里定义的变量,进程 B 压根看不到。 这就是 …

2026/7/22 18:13:14阅读更多 →
如何用Gemini快速入门加密货币回测?3步打造你的第一个交易策略

如何用Gemini快速入门加密货币回测?3步打造你的第一个交易策略

如何用Gemini快速入门加密货币回测?3步打造你的第一个交易策略 【免费下载链接】gemini Backtesting for sleepless cryptocurrency markets 项目地址: https://gitcode.com/gh_mirrors/gemini/gemini Gemini是一款专为加密货币市场设计的回测工具&#xff0…

2026/7/22 18:13:14阅读更多 →
EDMA3高级应用:从视频帧传输到乒乓缓冲与链式传输实战

EDMA3高级应用:从视频帧传输到乒乓缓冲与链式传输实战

1. 项目概述与EDMA3核心价值 在嵌入式系统开发,尤其是涉及音视频处理、高速数据采集或实时通信的领域,数据搬运的效率往往是整个系统性能的瓶颈。想象一下,一个640x480分辨率、每秒30帧的视频流,意味着CPU每秒钟需要处理超过900万…

2026/7/22 19:17:25阅读更多 →
EDMA3寄存器深度解析:从队列状态、事件管理到内存保护的实战指南

EDMA3寄存器深度解析:从队列状态、事件管理到内存保护的实战指南

1. 项目概述:从寄存器手册到实战理解的跨越如果你正在开发基于TI C6000系列DSP或类似SoC的嵌入式系统,并且性能瓶颈卡在了数据搬运上,那么你肯定绕不开EDMA3这个核心外设。手册里那几百页的寄存器描述,尤其是关于队列状态、事件管…

2026/7/22 19:17:25阅读更多 →
众阳公共卫生上报系统|医院公共卫生信息化、无纸化智能管理解决方案

众阳公共卫生上报系统|医院公共卫生信息化、无纸化智能管理解决方案

在传统医院公共卫生管理工作中,上报、查询、统计等工作长期依赖人工手工操作,流程繁琐、效率低下、易出现漏报、错报、数据滞后等问题。众阳公共卫生上报系统专为医疗机构公共卫生管理场景打造,全面替代传统手工模式,实现公共卫生…

2026/7/22 19:17:25阅读更多 →
FreeType 3.0路线图前瞻:下一代字体引擎的技术演进与功能预测

FreeType 3.0路线图前瞻:下一代字体引擎的技术演进与功能预测

FreeType 3.0路线图前瞻:下一代字体引擎的技术演进与功能预测 【免费下载链接】freetype Official mirror of https://gitlab.freedesktop.org/freetype/freetype 项目地址: https://gitcode.com/gh_mirrors/free/freetype FreeType作为一款广泛应用的开源字…

2026/7/22 19:17:25阅读更多 →
HarmonyOS应用开发实战:萌宠日记 - 健康分类数据隔离

HarmonyOS应用开发实战:萌宠日记 - 健康分类数据隔离

HarmonyOS应用开发实战:萌宠日记 - 健康分类数据隔离 前言 数据隔离 是多分类健康记录页的核心设计原则。在 萌宠日记 的 HealthRecordPage 中,5 个健康分类(体重、疫苗、驱虫、体检、其他)各自拥有独立的 数据源 和 展示视图&am…

2026/7/22 19:17:25阅读更多 →
高楼树木遮挡导致RTK失锁?惯导倾斜与激光测距组合方案

高楼树木遮挡导致RTK失锁?惯导倾斜与激光测距组合方案

测量员小王承接了城中村改造项目,GPS信号被密集房屋切割得支离破碎。传统RTK在树荫下、楼房间反复失锁,初始化时间比测量时间还长。倾斜测量和激光测距技术的组合,为这类复杂环境提供了有效解决方案。遮挡环境导致的问题本质是:卫…

2026/7/22 19:15:25阅读更多 →
Go语言静态资源打包方案对比与实践指南

Go语言静态资源打包方案对比与实践指南

1. 项目背景与核心需求在Go语言开发中,我们经常需要处理静态资源文件的打包问题。无论是Web应用的模板文件、前端资源,还是配置文件、证书等,都需要随程序一起分发。传统做法是将这些文件与编译后的二进制文件放在同一目录下,但这…

2026/7/22 0:53:59阅读更多 →
Go语言实现高性能LDAP认证服务的架构与实践

Go语言实现高性能LDAP认证服务的架构与实践

1. 项目背景与核心价值LDAP(轻量级目录访问协议)作为企业级身份认证的黄金标准,已经服务了超过80%的财富500强公司。我在金融科技领域实施统一认证体系时,发现传统Java方案存在启动慢、内存占用高等痛点。而Go语言凭借其协程并发模…

2026/7/22 0:53:59阅读更多 →
【AI面试官实战指南】:用ChatGPT模拟10类高频技术岗面试,3天提升应答精准度92%

【AI面试官实战指南】:用ChatGPT模拟10类高频技术岗面试,3天提升应答精准度92%

更多请点击: https://intelliparadigm.com 第一章:AI面试官实战指南的核心价值与适用场景 AI面试官并非替代人类HR的“黑箱工具”,而是以可解释、可审计、可迭代的方式,赋能招聘全链路的关键基础设施。其核心价值在于将主观经验沉…

2026/7/22 0:53:59阅读更多 →
中小企业小程序开发公司怎么选:预算、上手和售后避坑指南

中小企业小程序开发公司怎么选:预算、上手和售后避坑指南

中小企业做小程序,最常见的矛盾是预算有限,但又不希望功能太单薄;没有技术团队,但又希望后续能自己运营;想快速上线,又担心隐性收费和售后失联。选型时如果只看“低价套餐”或“案例数量”,很容…

2026/7/22 0:01:17阅读更多 →
GEO优化如何沉淀长期内容资产?广拓时代谈AI搜索时代的内容ROI

GEO优化如何沉淀长期内容资产?广拓时代谈AI搜索时代的内容ROI

企业做营销,最怕钱花完了,资产没有留下。 效果广告能带来一段时间的曝光,但预算停止后,流量往往也随之停止。短视频内容可能在几天内冲高,也可能很快沉下去。AI搜索时代,企业需要重新思考一个问题&#xff…

2026/7/22 0:01:17阅读更多 →
Agent 终态判定:何时该停止思考、给出最终回复

Agent 终态判定:何时该停止思考、给出最终回复

Agent 终态判定:何时该停止思考、给出最终回复 一、你的 Agent 在"再想想"的循环里绕了 12 轮,用户已经关窗口了 Agent 与人最大的区别是:人知道什么时候该停下来给答案,Agent 会一直"想"下去。你给 Agent 接…

2026/7/22 0:01:17阅读更多 →
YOLOv8推理性能优化:从1.2FPS到35FPS的全链路加速实践

YOLOv8推理性能优化:从1.2FPS到35FPS的全链路加速实践

如果你在部署 YOLOv8 时,发现推理速度只有可怜的 1-2 FPS,而别人的演示视频却能跑到 30 FPS 以上,那么问题很可能不在模型本身,而在于你的整个处理链路。很多开发者拿到一个训练好的 YOLOv8 模型后,会直接使用官方示例…

2026/7/21 22:53:50阅读更多 →
Coze与Dify对比指南:低代码AI应用开发从入门到实战

Coze与Dify对比指南:低代码AI应用开发从入门到实战

1. 从零到一:为什么你需要了解 Coze 和 Dify?如果你对 AI 应用开发感兴趣,但一看到“大模型”、“智能体”、“工作流”这些词就头疼,觉得门槛太高,那这篇文章就是为你准备的。很多开发者,包括我自己&#…

2026/7/22 18:55:50阅读更多 →
AI生图工具怎么选?2026年6月版实测对比

AI生图工具怎么选?2026年6月版实测对比

做自媒体的朋友应该都有体会:配图一直是个让人头疼的问题。2026年,AI生图工具已经非常成熟了,但工具太多反而不知道怎么选。以下是截至2026年6月我对主流AI生图工具的实测对比。Midjourney V8.1:速度之王2026年6月11日&#xff0c…

2026/7/22 18:55:50阅读更多 →