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
[本地事务成功]
[本地事务失败]
[事务状态未知(异常/超时)]
- 发送事务消息
- 发送半消息(Half Message)
- 半消息持久化,暂不可消费
- 半消息发送成功
- 回调执行本地事务
- 执行本地事务(如更新订单状态)
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 注解支持丰富的配置选项:
- 按 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 等高性能序列化框架
处理特殊的数据格式(如二进制数据)
多环境配置与管理