RocketMQ 5.5 实践

RocketMQ 5.5 实践

之前看过阿里云微信公众号的几篇文章,觉得真的很有深度,也给了我很多启发。如 从 Loop 到 Graph Engineering 的演进思考与实战 里 Loop Engineering 的四个问题,我也在生产系统有遇到。

近期公众号发布了几篇关于 RocketMQ 的文章,如 Apache RocketMQ 面向 AI 演进:LiteTopic 支撑百万级多 Agent 会话协作、RocketMQ for AI,照着搭百炼、Qoder 同款异步通信架构 等。

我对几个主流的 MQ 都有过了解,也都有简单的使用过。MQ 简单来讲就是接收生产者发送的消息,并传递给消费者消费。而 MQ 架构组成(Topic、Partition or MessageQuene、ConsumerGroup)则各有各的实现(Rabbit MQ 还有 Exchange),消费者获取消息的渠道(push or pull)、序列化方式、持久化机制(如何保证消息不丢失)也是,而且很多内容在其他的中间件上也都有体现。

如果只是简单的使用,在了解 MQ 的组成部分后,SpringBoot 里使用对应的 Template 和 Listener 处理消息即可,所有的 MQ 几乎都可以满足需求。

那么,为什么选择 Rocket MQ?这里放一下 官网 的解释:

image-20260902161441404

无论是任意精度的延迟消息,还是轻量的事务消息都非常的吸引人,而目前的 5.x 版本又为多 Agent AI 系统提供了异步通信基础设施。我就又重新学习了下 Rocket MQ。

Docker 部署 Rocket MQ

官网里虽然也有 docker 的部署流程,但仅仅挂载了 broker 的配置文件,对于日志、持久化等没有考虑。

真正想要全面部署的话可能会出现 253 问题,github 上有相关的 issue。本质原因是:mq 容器内部创建了 ID 是 3000 的用户 rocketmq,服务也是该用户启动的,但挂载的 log、store 目录创建用户是 root,没有访问权限。

以下是我整理好的部署脚本,按需修改 HOST_IP、WORK_DIR 和端口号即可,其实就是把目录提前创建好了。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
#!/bin/bash
# ============================================================
# RocketMQ 5.5.0 单节点 + Dashboard 部署(UID=3000,权限加固)
# ============================================================

# -------- 可修改配置 --------
HOST_IP="192.168.41.190" # 宿主机 IP
WORK_DIR="/data/rocketmq" # 工作目录根路径
CONTAINER_UID=3000 # 容器内用户 UID
NAMESRV_PORT=25007 # 宿主机映射端口(容器内 9876)
BROKER_PORT=25011 # Broker 主端口
BROKER_HA_PORT=25009 # Broker HA 端口, Dashboard 要用
DASHBOARD_PORT=25015 # Dashboard 端口
DOCKER_NETWORK="rocketmq-net"
ROCKETMQ_IMAGE="rocketmq:5.5.0"
DASHBOARD_IMAGE="rocketmq-dashboard:2.0.0"

# -------- 脚本逻辑 --------
set -e

echo "=========================================="
echo "RocketMQ 单节点部署开始"
echo "宿主机 IP: $HOST_IP"
echo "容器内用户 UID: $CONTAINER_UID"
echo "Namesrv 宿主机端口: $NAMESRV_PORT (容器内 9876)"
echo "Broker 主端口: $BROKER_PORT"
echo "Dashboard 端口: $DASHBOARD_PORT"
echo "工作目录: $WORK_DIR"
echo "=========================================="

# 1. 清理旧容器和网络
docker rm -f mqnamesrv mqbroker rocketmq-dashboard 2>/dev/null || true
docker network rm $DOCKER_NETWORK 2>/dev/null || true

# 2. 创建 Docker 网络
docker network create $DOCKER_NETWORK

# 3. 预创建所有目录(先创建,后设权限)
mkdir -p $WORK_DIR/conf
mkdir -p $WORK_DIR/store/commitlog
mkdir -p $WORK_DIR/store/consumequeue
mkdir -p $WORK_DIR/store/rocksdbstore
mkdir -p $WORK_DIR/store/index
mkdir -p $WORK_DIR/logs/rocketmqlogs

# 4. 设置权限(容器用户 UID=3000)
echo "设置目录权限(UID=$CONTAINER_UID)..."
chown -R $CONTAINER_UID:$CONTAINER_UID $WORK_DIR
chmod -R 777 $WORK_DIR

if command -v chcon >/dev/null 2>&1; then
chcon -R -t svirt_sandbox_file_t $WORK_DIR 2>/dev/null || true
fi

# 验证权限(调试信息)
echo "目录权限确认:"
ls -ld $WORK_DIR/conf $WORK_DIR/store $WORK_DIR/logs

# 5. 生成 Broker 配置
cat > $WORK_DIR/conf/broker.conf << EOF
brokerClusterName = DefaultCluster
brokerName = broker-a
brokerId = 0
deleteWhen = 04
fileReservedTime = 48
brokerRole = ASYNC_MASTER
flushDiskType = ASYNC_FLUSH
brokerIP1 = $HOST_IP
listenPort = $BROKER_PORT
autoCreateTopicEnable = true
EOF
echo "✅ Broker 配置生成: $WORK_DIR/conf/broker.conf"

# 6. 启动 Namesrv
docker run -d --name mqnamesrv --network $DOCKER_NETWORK \
-p $NAMESRV_PORT:9876 \
$ROCKETMQ_IMAGE sh mqnamesrv
echo "等待 Namesrv 就绪..."
until docker logs mqnamesrv 2>&1 | grep -q "The Name Server boot success"; do
sleep 2
done
echo "✅ Namesrv 已就绪,宿主机端口: $NAMESRV_PORT → 容器 9876"

# 7. 启动前再次加固权限(防止其他进程干扰)
chown -R $CONTAINER_UID:$CONTAINER_UID $WORK_DIR
chmod -R 777 $WORK_DIR

# 8. 启动 Broker
docker run -d --name mqbroker --network $DOCKER_NETWORK \
-p $BROKER_PORT:$BROKER_PORT \
-p $BROKER_HA_PORT:$BROKER_HA_PORT \
-v $WORK_DIR/conf/broker.conf:/home/rocketmq/rocketmq-5.5.0/conf/broker.conf \
-v $WORK_DIR/store:/home/rocketmq/store \
-v $WORK_DIR/logs:/home/rocketmq/logs \
-e NAMESRV_ADDR="mqnamesrv:9876" \
$ROCKETMQ_IMAGE sh mqbroker -c /home/rocketmq/rocketmq-5.5.0/conf/broker.conf

sleep 5
if docker ps | grep -q mqbroker; then
echo "✅ Broker 启动成功"
else
echo "❌ Broker 启动失败,错误日志:"
docker logs mqbroker
exit 1
fi

# 9. 启动 Dashboard
docker run -d --name rocketmq-dashboard --network $DOCKER_NETWORK \
-p $DASHBOARD_PORT:8080 \
-e "JAVA_OPTS=-Drocketmq.namesrv.addr=mqnamesrv:9876" \
$DASHBOARD_IMAGE
echo "✅ Dashboard 启动成功"

# 10. 验证
echo "等待 15 秒后验证..."
sleep 15
echo "========== 集群状态 =========="
docker exec mqbroker sh mqadmin clusterList -n mqnamesrv:9876 2>/dev/null || echo "⚠️ 验证命令执行失败,可手动检查"
echo "=========================================="

echo "✅ 部署完成!"
echo ""
echo "========== 连接信息 =========="
echo "Namesrv 地址: $HOST_IP:$NAMESRV_PORT"
echo "Broker 地址: $HOST_IP:$BROKER_PORT"
echo "Dashboard 地址: http://$HOST_IP:$DASHBOARD_PORT"
echo ""
echo "日志文件位置(宿主机): $WORK_DIR/logs/rocketmqlogs/"
echo ""
echo "常用命令:"
echo " 查看 Broker 日志: docker logs mqbroker"
echo " 查看文件日志: tail -f $WORK_DIR/logs/rocketmqlogs/rocketmq_broker.log"
echo " 验证集群: docker exec mqbroker sh mqadmin clusterList -n mqnamesrv:9876"
echo "=========================================="

这个问题折腾了我很久,按照官网只挂载 conf 配置文件的话容器重启数据就会丢失,自己配置 log、store 又会踩坑。

RocketMQTemplate 使用 RocketMQ

SpringBoot 里有很多的 Template,如 RedisTemplate、KafkaTemplate、JdbeTemplate 等。

pom 依赖

1
2
3
4
5
6
7
8
9
10
11
12
13
<properties>
<java.version>17</java.version>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<project.reporting.outputEncoding>UTF-8</project.reporting.outputEncoding>
<spring-boot.version>3.3.5</spring-boot.version>
<rocketmq-spring.version>2.3.1</rocketmq-spring.version>
</properties>

<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
<version>${rocketmq-spring.version}</version>
</dependency>

功能特性

  • 对于普通的消息,支持 sync、aysnc、oneway 方式发送(同步等结果、异步回调、单向发送)。消息可选携带 tag,消费者可以根据 tag 只处理想要的消息。
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
public void test() {
// sync
SendResult result = rocketMQTemplate.syncSend(TOPIC, "sync-msg-" + System.currentTimeMillis());
// async
rocketMQTemplate.asyncSend(TOPIC, "async-msg-" + System.currentTimeMillis(), new SendCallback() {
@Override
public void onSuccess(SendResult sendResult) {
System.out.println("[async] onSuccess, msgId=" + sendResult.getMsgId());
}

@Override
public void onException(Throwable e) {
System.out.println("[async] onException: " + e.getMessage());
}
});
// oneway
rocketMQTemplate.sendOneWay(TOPIC, "oneway-msg-" + System.currentTimeMillis());
// send to topic:tag
rocketMQTemplate.syncSend(TOPIC + ":" + tag, "tag-msg-" + tag);
}
  • 顺序消息,kafka 通过对发送消息的 key 路由进不同的 partition,rocketMQ 也类似,根据 key 路由到不同的 MessageQuene(相当于 partition)。
1
2
3
4
5
6
String[] steps = {"CREATED", "PAID", "SHIPPED"};
for (String step : steps) {
SendResult result = rocketMQTemplate.syncSendOrderly(TOPIC, userId + "-" + step, userId);
System.out.printf("[order-producer] %s -> queueId=%d%n", userId + "-" + step,
result.getMessageQueue().getQueueId());
}

基于分布式的理论,有 order 顺序要求的消息要不能并行发送,因为服务端是以消息到达的顺序排序的,客户端并行发送时,有可能后发送的消息先到 MQ。

  • 延时消息,RocketMQ 5.x 的延时消息使用 “时间轮” 实现,在分析 netty 时也有介绍过这个算法。
1
2
3
4
5
6
7
8
9
10
11
public void test(@RequestParam(name = "seconds", defaultValue = "10") int seconds) {
// 4.x version
Message<String> message = MessageBuilder.withPayload("legacy-delay-msg")
.setHeader(RocketMQHeaders.DELAY, 3)
.build();
rocketMQTemplate.syncSend(TOPIC, message);

// 5.x version
String payload = "check-order-" + System.currentTimeMillis();
SendResult result = rocketMQTemplate.syncSendDelayTimeMills(TOPIC, payload, seconds * 1000L);
}
  • 事务消息,这是 RocketMQ 非常受关注的一点。一般来讲,都是业务处理成功后(更新数据库)发送消息到 MQ,那么如果发送 MQ 失败怎么办?发送超时未知怎么办?RocketMQ 在普通消息基础上,支持二阶段的提交能力。将二阶段提交和本地事务绑定,实现全局提交结果的一致性。具体流程可以看官网的示意图:

事务消息

SpringCloudStream 使用 RocketMQ

如果系统中都使用 Template 的方式实现,那么换 MQ 就会非常麻烦,Spring 考虑了这种情况。Spring 中使用 SpringCloudStream(SCS)对各种 MQ 进行了适配,换 MQ 非常简单。

注意:虽然 SCS 中带了 SpringCloud,但这并不意味着 SCS 只能在 SpringCloud 微服务项目中使用,SpringBoot 中引入相关的依赖后也是可以使用的,示例配置如下:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
<properties>
<java.version>17</java.version>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<project.reporting.outputEncoding>UTF-8</project.reporting.outputEncoding>
<spring-boot.version>3.3.5</spring-boot.version>
<spring-cloud.version>2023.0.3</spring-cloud.version>
<spring-cloud-alibaba.version>2023.0.3.2</spring-cloud-alibaba.version>
</properties>

<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<dependency>
<groupId>com.alibaba.cloud</groupId>
<artifactId>spring-cloud-starter-stream-rocketmq</artifactId>
</dependency>
</dependencies>

SCS 模式下系统不再用 ***MQTemplate,而是 StreamBridge,如下所示:

生产者代码:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
@RestController
@RequestMapping("/scs")
public class ScsStreamController {

static final String OUT_BINDING = "scsDemo-out-0";

@Resource
private StreamBridge streamBridge;

/**
* StreamBridge 发送:参数是 binding 名;返回值反映消息是否递交成功。
*/
@GetMapping("/send")
public String send(@RequestParam(name = "msg", defaultValue = "scs-hello") String msg) {
boolean sent = streamBridge.send(OUT_BINDING, msg);
return sent ? "sent via SCS: " + msg : "send failed (binding unavailable): " + msg;
}

}

消费者代码:

1
2
3
4
5
6
7
8
9
10
11
12
13
@Configuration
public class ScsConsumerConfig {

/**
* 泛型 String 即消息体反序列化类型; SCS binder 负责收消息、反序列化后回调 accept。
*/
@Bean
public Consumer<String> scsDemoConsumer() {
return msg -> System.out.printf("[scs] thread=%s, msg=%s%n",
Thread.currentThread().getName(), msg);
}

}

那么最令人费解的地方来了:生产者和消费者的对应关系如何体现?我怎么知道我发送的 message 由哪个消费者处理?SCS 里的对应关系是通过 binding 维护的,配置文件里可以找到答案:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
spring:
cloud:
function:
# 圈定要绑定 MQ 的函数 bean, 单 bean 场景可省略
definition: scsDemoConsumer
stream:
rocketmq:
binder:
# rocketmq name-server 地址
name-server: ip:port
bindings:
# 消费端:bean 名 scsDemoConsumer + -in-0 派生而来
scsDemoConsumer-in-0:
destination: scs-demo-topic
group: scs-demo-group
# 发送端:无 Supplier bean,StreamBridge.send() 首次调用时按需创建
scsDemo-out-0:
destination: scs-demo-topic

我们来详细分析一下:

  1. scsDemoConsumer:这是消费者的 beanName(通过 @Bean 定义),实现方式为:Consumer<String> scsDemoConsumer(),这是 java.util.function 里的接口,用来声明消费者。该方式会被 SCS 派生 binding 信息 beanName-in-0,其中 in 的意思就像 netty 的 inboud 一样,代表消息的入口,也就是消费者。如果有多个这样的消费者,beanName 可以用,分隔。
  2. bindings:所有的消费者、生产者绑定信息。
  • scsDemoConsumer-in-0:消费端 binding,由 scsDemoConsumer 派生而来;
    • destination:真正的 topic,消费者所消费的信息就是从 destination 对应的 topic 来(如果是 RabbitMQ,则代表的是交换机);
    • group:ConsumerGroup 消费者组;
  • scsDemo-out-0:生产端 binding,生产者在调用 streamBridge.send() 方法时传递的第一个参数;
    • destination:真正的 topic,生产者消息发送的真正目的地;

可以总结为下面的流程:

1
2
3
4
Java bean 名              		binding 名(SCS 派生)           			destination(真正的 topic)
─────────────────────────────────────────────────────────────────────────────────────────────
Consumer scsDemoConsumer() → scsDemoConsumer-in-0 ──bindings 配置──→ scs-demo-topic
send("scsDemoConsumer-out-0") → scsDemoConsumer-out-0 ──bindings 配置──→ scs-demo-topic

SCS 的派生换来了 MQ 无关性,代价是 MQ 专属能力受限(顺序消息、延迟消息、事务消息),如果业务真的需要 MQ 的特色功能,官方 starter 才是首选。

虽然我在文章开始提了 MQ5.x 对于 Agent 系统的支持,奈何最近主要精力在系统业务,没有深入探索,后续学习~


RocketMQ 5.5 实践
https://zhuwenjie0716.github.io/2026/08/30/RocketMQ 5.5 实践/
作者
Wenjie Zhu
发布于
2026年8月30日
许可协议