一、为什么要用springboot整合rocketmq?
你在搞一套无人售货柜系统,用户扫码开门、拿走饮料、关门扣款。这个流程涉及好几个环节:订单创建、库存扣减、出货指令下发、支付回调。如果全用http接口同步调用,服务之间强耦合,一个环节慢了整条链路都卡住。
用rocketmq做消息中间件,服务之间通过消息异步通信,解耦又高效。而springboot + rocketmq-spring-boot-starter 让整合过程极其丝滑,几乎零样板代码。
二、环境准备
确保你已经启动了rocketmq的nameserver和broker。如果你还没装,快速用docker起一套:
# 启动 nameserver docker run -d --name rmqnamesrv -p 9876:9876 \ apache/rocketmq:5.3.0 sh mqnamesrv # 启动 broker docker run -d --name rmqbroker -p 10911:10911 -p 10909:10909 \ --link rmqnamesrv:namesrv \ -e "namesrv_addr=namesrv:9876" \ apache/rocketmq:5.3.0 sh mqbroker
本地访问 localhost:9876 是nameserver,localhost:10911 是broker。
三、pom.xml 引入依赖
创建一个springboot项目,在 pom.xml 中添加依赖:
<!-- springboot 父工程 -->
<parent>
<groupid>org.springframework.boot</groupid>
<artifactid>spring-boot-starter-parent</artifactid>
<version>2.7.18</version>
</parent>
<dependencies>
<!-- springboot web -->
<dependency>
<groupid>org.springframework.boot</groupid>
<artifactid>spring-boot-starter-web</artifactid>
</dependency>
<!-- rocketmq springboot starter -->
<dependency>
<groupid>org.apache.rocketmq</groupid>
<artifactid>rocketmq-spring-boot-starter</artifactid>
<version>2.2.3</version>
</dependency>
<!-- lombok(简化代码) -->
<dependency>
<groupid>org.projectlombok</groupid>
<artifactid>lombok</artifactid>
</dependency>
</dependencies>
版本兼容性提示:rocketmq-spring-boot-starter 2.2.x 对应 rocketmq 4.x;如果你用 rocketmq 5.x,需要用 2.3.0 以上版本。两者api基本一致,本文以4.x为例。
四、application.yml 配置
rocketmq:
# nameserver 地址,多个用分号隔开
name-server: 127.0.0.1:9876
# 生产者配置
producer:
# 生产者组名(必须唯一)
group: vending-producer-group
# 发送超时时间(毫秒)
send-message-timeout: 3000
# 同步发送失败重试次数
retry-times-when-send-failed: 3
# 异步发送失败重试次数
retry-times-when-send-async-failed: 3
server:
port: 8080
消费者配置稍后在注解里写,不需要放yml里。
五、生产者开发:三种发送方式
rocketmq提供三种发送模式,各有适用场景。
5.1 同步发送(syncsend)
生产者发送消息后阻塞等待broker返回确认,收到成功响应才继续。
@service
public class ordermessageproducer {
@autowired
private rocketmqtemplate rocketmqtemplate;
/**
* 同步发送订单消息
* 适用场景:重要消息,必须确认送达,如订单创建
*/
public sendresult sendordersync(orderdto order) {
message<orderdto> message = messagebuilder
.withpayload(order)
.setheader("keys", order.getorderid()) // 设置消息key用于查询
.build();
// topic:order_topic, payload:消息对象
sendresult result = rocketmqtemplate.syncsend("order_topic", message);
system.out.println("发送结果: " + result.getsendstatus());
return result;
}
}
- 可靠性高:发送失败会重试
- 性能低:要等响应,不适合高吞吐
- 场景:售货柜下单、支付通知
5.2 异步发送(asyncsend)
发送后不阻塞,通过回调函数处理发送结果。
/**
* 异步发送出货指令
* 适用场景:对响应时间敏感,但需要知道发送结果
*/
public void sendshipmentasync(string orderid) {
message<string> message = messagebuilder
.withpayload(orderid)
.build();
rocketmqtemplate.asyncsend("shipment_topic", message, new sendcallback() {
@override
public void onsuccess(sendresult sendresult) {
log.info("出货指令发送成功: {}", sendresult.getmsgid());
}
@override
public void onexception(throwable throwable) {
log.error("出货指令发送失败, orderid={}", orderid, throwable);
// todo: 降级处理,比如写入本地表稍后重发
}
});
}
- 性能高:不阻塞主流程
- 有回调:失败可知
- 场景:出货指令下发、设备控制指令
5.3 单向发送(sendoneway)
只负责发送,不等响应,不回调。
/**
* 单向发送设备心跳
* 适用场景:日志、心跳等大量不重要的消息
*/
public void sendheartbeatoneway(string deviceid, string status) {
map<string, string> payload = new hashmap<>();
payload.put("deviceid", deviceid);
payload.put("status", status);
payload.put("timestamp", string.valueof(system.currenttimemillis()));
rocketmqtemplate.sendoneway("device_heartbeat_topic", payload);
}
- 性能最高:fire and forget
- 无可靠性保证:可能丢
- 场景:心跳上报、日志埋点
三种方式对比
| 方式 | 可靠性 | 性能 | 响应 |
|---|---|---|---|
| syncsend | 高 | 低 | 阻塞等结果 |
| asyncsend | 中 | 高 | 回调通知 |
| sendoneway | 低 | 最高 | 不关心 |
六、消费者开发
消费者开发只需要两步:加注解 + 实现接口。
6.1 基本消费者
@slf4j
@component
@rocketmqmessagelistener(
topic = "order_topic", // 消费的topic
consumergroup = "order_consumer_group", // 消费者组名(唯一)
messagemodel = messagemodel.clustering // 集群模式:同一组下每条消息只被一个消费者消费
)
public class ordermessageconsumer implements rocketmqlistener<orderdto> {
@override
public void onmessage(orderdto order) {
log.info("收到订单消息: orderid={}, amount={}",
order.getorderid(), order.getamount());
// 业务逻辑:处理订单
processorder(order);
// 正常返回 = 消费成功
// 抛异常 = 消费失败,会触发重试
}
private void processorder(orderdto order) {
// 实际处理逻辑:扣减库存、记录订单等
log.info("处理订单完成: {}", order.getorderid());
}
}
小白疑问:什么是集群模式(clustering)和广播模式(broadcasting)?
- 集群模式:同一个consumergroup下,一条消息只被一个消费者实例消费。适合分布式部署。
- 广播模式:同一个consumergroup下,每条消息会被所有消费者实例都消费一遍。适合本地缓存刷新。
6.2 注解核心参数详解
@rocketmqmessagelistener(
topic = "order_topic",
consumergroup = "order_consumer_group",
messagemodel = messagemodel.clustering,
selectorexpression = "tag_a || tag_b", // tag过滤,*表示全部
consumemode = consumemode.concurrently, // 并发消费
maxreconsumetimes = 5, // 最大重试次数
consumetimeout = 30000l // 消费超时时间(ms)
)
- consumemode:
- concurrently:并发消费,多线程同时处理,速度快但不保证顺序
- orderly:顺序消费,单线程按队列顺序处理,适合有顺序要求的场景
6.3 消费失败与重试
消费者抛异常 → broker判定消费失败 → 延迟一段时间后重试。默认重试16次,重试间隔逐步增大(1s→5s→10s→30s→1m→...)。16次后还不成功,进入死信队列(dead letter queue)。
@override
public void onmessage(orderdto order) {
try {
// 业务处理
dobusiness(order);
} catch (exception e) {
log.error("处理订单失败: {}", order.getorderid(), e);
// 抛出异常 → 触发重试
throw new runtimeexception("消费失败", e);
}
}
实际生产中,建议对死信队列单独写一个消费者做人工干预处理。
七、消息序列化方案
7.1 默认json序列化
rocketmq-spring-boot-starter 默认用 rocketmqmessageconverter 把对象序列化为json字符串。你直接传对象就行,不用手动 json.tojsonstring()。
// 生产者直接发对象
rocketmqtemplate.syncsend("order_topic", orderdto);
// 消费者直接收对象
public void onmessage(orderdto order) { ... }
7.2 自定义messageconverter
如果默认json方案不满足需求(比如你想用protobuf),可以自定义:
@configuration
public class rocketmqconfig {
@bean
public rocketmqmessageconverter rocketmqmessageconverter() {
// 这里可以替换为你自己的序列化器
// 默认是 jackson json
return new rocketmqmessageconverter();
}
}
大多数场景下默认json就够了,不用折腾。
八、完整实战:售货柜订单消息
把生产者和消费者串起来,模拟售货柜完整订单链路。
8.1 消息dto
@data
@builder
@noargsconstructor
@allargsconstructor
public class orderdto implements serializable {
private string orderid; // 订单号
private string deviceid; // 售货柜设备id
private string userid; // 用户id
private string productid; // 商品id
private integer quantity; // 数量
private bigdecimal amount; // 金额
private long timestamp; // 创建时间
}
8.2 订单controller(生产者入口)
@restcontroller
@requestmapping("/api/order")
public class ordercontroller {
@autowired
private ordermessageproducer producer;
@postmapping("/create")
public string createorder(@requestbody orderdto order) {
order.setorderid(uuid.randomuuid().tostring());
order.settimestamp(system.currenttimemillis());
// 同步发送订单消息
sendresult result = producer.sendordersync(order);
return result.getsendstatus().equals(sendstatus.send_ok)
? "订单创建成功: " + order.getorderid()
: "订单创建失败";
}
}
8.3 订单消费者(处理出货逻辑)
@slf4j
@component
@rocketmqmessagelistener(
topic = "order_topic",
consumergroup = "order_consumer_group",
messagemodel = messagemodel.clustering,
maxreconsumetimes = 5
)
public class ordermessageconsumer implements rocketmqlistener<orderdto> {
@autowired
private inventoryservice inventoryservice;
@autowired
private shipmentservice shipmentservice;
@override
public void onmessage(orderdto order) {
log.info("处理订单: orderid={}, deviceid={}",
order.getorderid(), order.getdeviceid());
// 1. 扣减库存
boolean stockok = inventoryservice.deductstock(
order.getproductid(), order.getquantity());
if (!stockok) {
throw new runtimeexception("库存不足: " + order.getproductid());
}
// 2. 下发出货指令
shipmentservice.sendshipmentcommand(
order.getdeviceid(), order.getproductid(), order.getquantity());
log.info("订单处理完成: {}", order.getorderid());
}
}
8.4 出货指令消费者(设备端模拟)
@slf4j
@component
@rocketmqmessagelistener(
topic = "shipment_topic",
consumergroup = "shipment_consumer_group",
messagemodel = messagemodel.clustering
)
public class shipmentconsumer implements rocketmqlistener<string> {
@override
public void onmessage(string orderid) {
log.info("收到出货指令, 准备出货: orderid={}", orderid);
// 模拟设备出货电机转动
// 实际场景中这里会调用设备sdk控制硬件
}
}
九、常见整合问题排查
9.1 版本不兼容
| 现象 | 原因 | 解决 |
|---|---|---|
| 启动报 org.apache.rocketmq.common.message 相关类找不到 | starter版本与rocketmq server不匹配 | 4.x server用2.2.x starter;5.x server用2.3.0+ |
| 消费者收不到消息 | producer和consumer连的nameserver地址不一致 | 检查yml配置 |
| 消息体反序列化失败 | 生产者和消费者dto字段不一致 | 确保两端dto字段名、类型完全一致 |
9.2 消费者不消费
排查清单:
- topic和tag是否匹配:producer发的topic和consumer监听的topic要完全一致
- consumergroup是否被其他实例占用:同一个group+cluster模式只能有一个实例消费
- broker是否开启autocreatetopicenable:默认开启,topic不存在会自动创建;生产环境建议关闭手动创建
- 防火墙:确认10911端口可达
9.3 序列化异常
org.apache.rocketmq.spring.support.rocketmqmessageconverter ... cannot deserialize
常见原因:dto没有无参构造方法。加上 @noargsconstructor 即可。另外dto必须实现 serializable。
十、总结
springboot整合rocketmq的核心套路就是三步:
1. 引依赖 → 2.1.x/2.2.x/2.3.0+选对版本
2. 配yml → name-server + producer.group
3. 写代码 → 注入template发消息 / 加注解收消息
三种发送方式按场景选:
- 同步 → 重要消息(订单、支付)
- 异步 → 高吞吐且需回调(出货指令)
- 单向 → 不重要的海量消息(心跳、日志)
消费者记住 @rocketmqmessagelistener 注解那几个参数,配合 rocketmqlistener<t> 接口,基本能覆盖80%的业务场景。剩下20%的事务消息和顺序消息,后面的文章会专门讲。
到此这篇关于springboot整合rocketmq实现生产者、消费者快速搭建的文章就介绍到这了,更多相关springboot整合rocketmq内容请搜索代码网以前的文章或继续浏览下面的相关文章希望大家以后多多支持代码网!
发表评论