RocketMQ 是阿里开源的分布式消息中间件,具备高吞吐、高可用、低延迟、支持事务消息、延时消息等特性,广泛应用于微服务异步通信、流量削峰、分布式事务解耦、日志收集等业务场景。本文将从零带你实现 SpringBoot 整合 RocketMQ,涵盖环境准备、依赖配置、生产者发送消息、消费者监听消息、常见消息类型实战及踩坑总结。
一、前置环境准备
整合前需提前搭建好 RocketMQ 服务,本地或服务器部署均可,核心需要两个组件:NameServer(注册中心)、Broker(消息存储与转发)。
1. 环境版本说明
- SpringBoot:2.7.x(稳定通用版本)
- RocketMQ:4.9.5(兼容主流SpringBoot版本)
- JDK:1.8+
- Maven:3.6+
2. 服务启动验证
启动 RocketMQ 的 NameServer 和 Broker 后,确保服务正常运行,默认端口:
- NameServer:9876
- Broker:10909、10911
可通过控制台命令验证服务状态,确保后续项目连接无异常。
二、引入Maven依赖
SpringBoot 整合 RocketMQ 核心依赖为 rocketmq-spring-boot-starter,无需手动配置原生客户端,starter 已自动封装核心配置。
<!-- RocketMQ SpringBoot 启动器 -->
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
<version>2.2.2</version>
</dependency>
版本适配说明:2.2.x 版本适配 SpringBoot2.x,若使用 SpringBoot3.x 需升级至 2.3.x 及以上版本。
三、全局配置文件配置
在 application.yml 中配置 RocketMQ 核心连接参数,包括 NameServer 地址、生产者组名,统一管理配置,便于环境切换。
# RocketMQ 配置
rocketmq:
# NameServer 地址,集群环境用分号分隔
name-server: 127.0.0.1:9876
# 生产者配置
producer:
# 生产者组名,同一业务统一组名
group: springboot-producer-group
# 消息发送超时时间,默认3000ms
send-timeout: 3000
# 失败重试次数
retry-times-when-send-failed: 2
四、消息生产者实现
RocketMQ 封装了 RocketMQTemplate 模板类,内置各类消息发送方法,可快速发送普通消息、同步消息、异步消息、对象消息。
1. 生产者工具类
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Component;
@Component
public class RocketMQProducer {
// 注入RocketMQ模板类
@Autowired
private RocketMQTemplate rocketMQTemplate;
// 定义Topic主题,统一管理
private static final String TEST_TOPIC = "springboot-test-topic";
/**
* 发送普通同步消息
* 特点:发送后阻塞等待响应,可靠性高,适合重要业务场景
*/
public void sendSyncMessage(String message) {
rocketMQTemplate.syncSend(TEST_TOPIC, MessageBuilder.withPayload(message).build());
}
/**
* 发送对象消息(自动序列化)
*/
public <T> void sendObjectMessage(T data) {
rocketMQTemplate.syncSend(TEST_TOPIC, MessageBuilder.withPayload(data).build());
}
/**
* 发送延时消息
* RocketMQ预设延时等级:1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h
* level=3 代表10秒后消费
*/
public void sendDelayMessage(String message) {
rocketMQTemplate.syncSend(TEST_TOPIC, MessageBuilder.withPayload(message).build(), 3000, 3);
}
}
2. 核心方法说明
- syncSend:同步发送,阻塞等待消息投递结果,保证消息可靠送达,适用于支付、订单等核心业务
- asyncSend:异步发送,不阻塞主线程,通过回调接收结果,适合高吞吐非核心业务
- sendOneWay:单向发送,不等待响应、无回调,适用于日志收集、统计等无需确认的场景
五、消息消费者实现
消费者通过 @RocketMQMessageListener 注解监听指定 Topic,实现 RocketMQListener 接口处理消息,框架自动完成消息监听、消费、重试机制。
1. 普通消息消费者
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Service;
// 消费者监听注解:指定消费者组、监听的Topic
@RocketMQMessageListener(
consumerGroup = "springboot-consumer-group",
topic = "springboot-test-topic"
)
@Service
public class RocketMQConsumer implements RocketMQListener<String> {
/**
* 消息消费方法
* 消息监听成功后自动执行,参数为生产者发送的消息内容
*/
@Override
public void onMessage(String message) {
try {
// 业务消费逻辑
System.out.println("消费者成功接收消息:" + message);
} catch (Exception e) {
// 消费异常,RocketMQ会自动重试
e.printStackTrace();
throw new RuntimeException("消息消费失败");
}
}
}
2. 消费者核心注解参数说明
- consumerGroup:消费者组名,同一组消费者负载均衡消费消息,不同组可重复消费同一消息
- topic:监听的消息主题,必须与生产者发送的 Topic 一致
- messageModel:消费模式,默认 CLUSTERING(集群模式,一条消息仅被组内一个消费者消费),可选 BROADCASTING(广播模式,组内所有消费者均消费)
六、接口测试验证
编写测试接口,调用生产者发送消息,验证整个收发流程是否正常。
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
@RestController
@RequestMapping("/mq")
public class MqTestController {
@Autowired
private RocketMQProducer rocketMQProducer;
@GetMapping("/send")
public String sendMessage() {
// 发送普通消息
rocketMQProducer.sendSyncMessage("Hello SpringBoot RocketMQ 整合测试!");
return "消息发送成功";
}
@GetMapping("/send/delay")
public String sendDelayMessage() {
rocketMQProducer.sendDelayMessage("这是一条10秒延时消息");
return "延时消息发送成功";
}
}
测试步骤
- 启动 SpringBoot 项目,确保无报错,消费者自动监听 Topic
- 访问接口
http://localhost:8080/mq/send - 查看控制台,可打印出消费者接收的消息,证明整合成功
- 访问
/mq/send/delay,10秒后控制台打印延时消息,验证延时消息生效
七、对象消息收发实战
实际业务中常需要传递实体对象,RocketMQ starter 内置 Jackson 序列化,无需手动转换。
1. 定义实体类
import lombok.Data;
import java.io.Serializable;
@Data
public class User implements Serializable {
private Long id;
private String username;
private String phone;
}
2. 对象消息消费者
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Service;
@RocketMQMessageListener(
consumerGroup = "springboot-user-consumer-group",
topic = "springboot-user-topic"
)
@Service
public class UserMessageConsumer implements RocketMQListener<User> {
@Override
public void onMessage(User user) {
System.out.println("接收用户对象消息:" + user.toString());
}
}
3. 新增对象消息发送接口
@GetMapping("/send/user")
public String sendUserMessage() {
User user = new User();
user.setId(1L);
user.setUsername("测试用户");
user.setPhone("13800138000");
rocketMQProducer.sendObjectMessage(user);
return "用户对象消息发送成功";
}
八、常见问题与踩坑总结
1. 连接失败:Cannot connect to NameServer
- 原因:NameServer 地址配置错误、RocketMQ 服务未启动、防火墙端口未开放
- 解决:核对
name-server配置,开启9876、10911端口,重启服务
2. 消息重复消费
- 原因:消费者消费后未正常返回、业务异常触发重试、集群模式下节点切换
- 解决:业务实现幂等性(唯一ID去重、数据库幂等校验)
3. 消费者不监听消息
- 原因:Topic 名称前后端不一致、消费者组名重复、注解未生效
- 解决:统一生产者消费者 Topic,保证消费者类被 Spring 托管(添加@Service)
4. 序列化失败
- 原因:实体类无空参构造、缺少序列化依赖
- 解决:实体类添加无参构造,确保 lombok 依赖正常
九、总结
本文完成了 SpringBoot 与 RocketMQ 的完整整合,核心流程可概括为:引入starter依赖 → 配置NameServer和生产者参数 → 通过RocketMQTemplate发送消息 → 通过@RocketMQMessageListener监听消费消息。
同时实现了普通消息、延时消息、对象消息三种常用消息类型,覆盖大部分业务场景。RocketMQ 凭借高可靠、高吞吐的特性,可完美解决微服务异步解耦、流量削峰、延时任务等核心需求,后续可在此基础上拓展事务消息、顺序消息、消息过滤等高级功能。