news 2026/7/22 18:15:00

Spring Boot 整合 RocketMQ 完全指南

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
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.yml):

rocketmq:
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;

@Service
public 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"; Message<String> 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"; Message<String> 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"; Message<String> message = MessageBuilder .withPayload(content) .setHeader(RocketMQHeaders.KEYS, orderId) .build(); rocketMQTemplate.sendOneWay(destination, message); }

}
KEY 的作用(非常重要):

设置 RocketMQHeaders.KEYS 有三个核心用途:

消息查询:在 Dashboard 中按业务 Key 快速定位消息
事务回查:事务消息回查时用于关联业务数据
幂等控制:消费者端用 Key 做去重判断
消费者的配置与使用
基础配置(application.yml):

rocketmq:
name-server: 127.0.0.1:9876
consumer:
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;

@Component
@RocketMQMessageListener(
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 最有价值、也最容易用错的功能。它的核心是保证“本地事务”和“消息发送”要么一起成功,要么一起失败。

事务消息的完整流程:

本地数据库
Broker
Producer
业务应用
本地数据库
Broker
Producer
业务应用
Broker 未收到最终确认,触发回查
loop
[事务回查(默认每 60 秒)]
alt
[本地事务成功]
[本地事务失败]
[事务状态未知(异常/超时)]

  1. 发送事务消息
  2. 发送半消息(Half Message)
  3. 半消息持久化,暂不可消费
  4. 半消息发送成功
  5. 回调执行本地事务
  6. 执行本地事务(如更新订单状态)
    7a. 事务提交成功
    8a. 返回 COMMIT
    9a. 提交事务
    10a. 半消息→正式消息,可消费
    7b. 事务回滚
    8b. 返回 ROLLBACK
    9b. 回滚事务
    10b. 删除半消息
    8c. 返回 UNKNOWN
    9c. 发起回查请求
    10c. 检查本地事务状态
    11c. 查询业务数据
    12c. 返回查询结果
    13c. 返回 COMMIT/ROLLBACK
    14c. 提交最终事务状态
    第一步:定义事务监听器:

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;

@Service
@RocketMQTransactionListener(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; }

}
第二步:发送事务消息:

@Service
public 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"; Message<String> message = MessageBuilder .withPayload("订单创建:" + orderId) .setHeader("orderId", orderId) // 传递给事务监听器 .setHeader(RocketMQHeaders.KEYS, orderId) .build(); // 发送事务消息 rocketMQTemplate.sendMessageInTransaction(destination, message, null); }

}
消息监听器的多种用法
@RocketMQMessageListener 注解支持丰富的配置选项:

  1. 按 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 {
@Override
public void onMessage(String message) {
// 同一个 Queue 的消息会按顺序被消费
}
}
4. 广播消费:

@RocketMQMessageListener(
topic = “order-topic”,
consumerGroup = “order-consumer-group”,
messageModel = MessageModel.BROADCASTING // 广播模式
)
public class BroadcastConsumer implements RocketMQListener {
@Override
public 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 {
@Override
public 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;

@Configuration
public class RocketMQConfig {

@Bean public MessageConverter rocketMQMessageConverter() { // 自定义转换逻辑 return new CustomMessageConverter(); }

}
常见场景:

使用 Protobuf 替代 JSON,提升序列化性能和减小消息体积
使用 Kryo 等高性能序列化框架
处理特殊的数据格式(如二进制数据)
多环境配置与管理

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/7/22 18:14:32

HighFive并行IO编程:基于MPI的HDF5高性能数据读写技巧

HighFive并行IO编程&#xff1a;基于MPI的HDF5高性能数据读写技巧 【免费下载链接】HighFive HighFive - Header-only C HDF5 interface 项目地址: https://gitcode.com/gh_mirrors/high/HighFive HighFive是一个Header-only的C HDF5接口&#xff0c;它简化了HDF5文件的…

作者头像 李华
网站建设 2026/7/22 18:12:40

Jellium Desktop色彩校准指南:让你的视频显示更准确

Jellium Desktop色彩校准指南&#xff1a;让你的视频显示更准确 【免费下载链接】jellium-desktop An unofficial desktop client for Jellyfin 项目地址: https://gitcode.com/GitHub_Trending/je/jellium-desktop Jellium Desktop是一款非官方的Jellyfin桌面客户端&am…

作者头像 李华
网站建设 2026/7/22 18:09:35

MetaGer元搜索引擎原理揭秘:聚合50+数据源的强大技术

MetaGer元搜索引擎原理揭秘&#xff1a;聚合50数据源的强大技术 【免费下载链接】MetaGer inofficial clone of the official MetaGer repository from https://gitlab.metager3.de/open-source/MetaGer.git 项目地址: https://gitcode.com/gh_mirrors/me/MetaGer MetaG…

作者头像 李华
网站建设 2026/7/22 18:08:33

前端焦虑?收藏!AI时代,前端如何转型成为AI产品经理或工程师?

文章针对前端开发者在AI时代的焦虑&#xff0c;提出前端非常适合转型AI应用层岗位&#xff0c;如AI工程师或AI产品经理。文章分析了前端业务逻辑简单、GitHub语料丰富的原因&#xff0c;并通过Cursor案例展示了AI在前端编程中的应用。文章还探讨了AI在前端开发流程中的实际作用…

作者头像 李华
网站建设 2026/7/22 18:08:04

【WorkBuddy从入门到精通实战教程】使用手册第 6 章 WorkBuddy的专家和专家团

专家和专家团与Skill的区别 WorkBuddy 本身是一个通用 Agent,什么任务都能接。但通用不意味着每个领域都应该用同一种方式处理。 比如, 同样是分析一份销售数据,普通 Agent 可能会读取数据、生成图表、总结趋势。数据分析专家会先理解业务目标,再确定核心指标,检查数据…

作者头像 李华