MQ - RocketMQ 5.4 集群模式安装和使用

RocketMQ 集群安装说明

从单机走向集群,是任何分布式系统迈向工业生产线的必经之路。在生产环境中,单机模式就像在钢丝绳上跳舞,任何一次突发的断电、硬件故障、磁盘损坏,都可能会导致业务系统全线瘫痪。RocketMQ 集群模式最核心的设计点可以总结为两句话:

  • NameServer 无状态对等:NameServer 之间互不通信,各自分立。Broker 启动时会向所有 NameServer 报到。这种 “无状态” 设计意味着即使挂掉大半,只要剩下一台,整个集群的路由大脑就依然清醒。
  • Broker 主从双写与水平分片
    • 水平分片(分担压力):不同的 Master 节点(如 Master A 和 Master B)各自承载 Topic 的一部分队列,多 Master 并存可以物理叠加集群的写吞吐量。
    • 主从多副本(防死防丢):每个 Master 绑定一个 Slave。Master 负责读写,Slave 负责实时同步备份。当 Master 挂掉,消费者会自动平滑切换到 Slave 读取数据,确保业务不流断。

在生产落地时,通常有以下三种架构路线:

  • 多 Master 模式:全是 Master,没有 Slave。优点是配置简单,缺点是一旦某台机器炸了,未消费的消息在机器恢复前物理不可访问。
  • 多 Master 多 Slave 模式(异步复制):性能高,但 Master 瞬间断电,可能存在毫秒级的数据未同步丢失。
  • 多 Master 多 Slave 模式(同步双写 / DLedger 纯Raft):金融生产级标准。消息必须强行落盘到 Slave 成功后才返回给客户端,数据零丢失,且 Raft 协议支持 Master 挂掉后 Slave 自动竞选提拔为新 Master。

为了兼顾高性能与高可用,我们今天采用业界应用最广、最稳健的经典模型:2主2从(2M-2S)异步复制、同步刷盘集群架构。我们准备两台物理机(或者虚拟机),通过交叉部署 Master 和 Slave,在节约机器成本的同时,实现硬件级容灾。如果 192.168.1.7 整台机器彻底烧毁,由于 broker-a 的副本 broker-a-s 活着在 .8 上,broker-b 的主节点也活着在 .8 上,整个集群的业务数据依然全面可访问。


服务端集群的安装步骤

单机版的安装可以直接参考 RocketMQ 5.x 的安装和使用。这里我们直接切入最核心的集群配置文件编写。RocketMQ 官方在 conf/ 目录下默认提供了 2m-2s-async/ 模板,我们直接对其进行改造。


机器1的配置

在 .7 机器上,我们需要配置两个文件,一个负责 Master A,一个负责 Slave B。

文件一:配置 Master A。创建并编辑 conf/2m-2s-async/broker-a.properties:

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
# 集群所属名称
brokerClusterName=DefaultCluster
# 当前 Broker 的名称(同一个主从组合名字必须一致)
brokerName=broker-a
# 0 代表它是 Master,大于 0 代表是 Slave
brokerId=0
# 路由大脑 NameServer 地址列表,用分号隔开
namesrvAddr=192.168.1.7:9876;192.168.1.8:9876

# 核心端口配置
listenPort=10911
# 存储根路径(建议与单机版隔离)
storePathRootDir=/root/store/broker-a
storePathCommitLog=/root/store/broker-a/commitlog

# 生产级核心物理策略
deleteWhen=04
fileReservedTime=48
# 异步复制(Master 写入成功即返回,后台异步复制给 Slave)
brokerRole=ASYNC_MASTER
# 同步刷盘(数据必须安全落盘才返回,死守数据不丢)
flushDiskType=SYNC_FLUSH

# 允许自动创建 Topic(生产建议关闭,此处测试开启)!!!!!!!!!!!! 重要 !!!!!!!!!!!!!
autoCreateTopicEnable=true

文件二:配置 Slave B。创建并编辑 conf/2m-2s-async/broker-b-s.properties:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
brokerClusterName=DefaultCluster
# 注意:名字叫 broker-b,说明它是 .7 机器上 Master B 的亲兄弟
brokerName=broker-b
# 大于 0 代表是 Slave
brokerId=1
namesrvAddr=192.168.1.7:9876;192.168.1.8:9876

# 核心端口(一台机器起多个实例,端口必须错开!)
listenPort=11011
storePathRootDir=/root/store/broker-b-s
storePathCommitLog=/root/store/broker-b-s/commitlog

deleteWhen=04
fileReservedTime=48
# 角色声明为从节点
brokerRole=SLAVE
# 从节点无脑采取异步刷盘即可,提升备份效率
flushDiskType=ASYNC_FLUSH

autoCreateTopicEnable=true


机器2的配置

在 .8 机器上,同样配置两个文件,负责 Master B 和 Slave A。

文件一:配置 Master B。创建并编辑 conf/2m-2s-async/broker-b.properties:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
brokerClusterName=DefaultCluster
brokerName=broker-b
brokerId=0
namesrvAddr=192.168.1.7:9876;192.168.1.8:9876

listenPort=10911
storePathRootDir=/root/store/broker-b
storePathCommitLog=/root/store/broker-b/commitlog

deleteWhen=04
fileReservedTime=48
brokerRole=ASYNC_MASTER
flushDiskType=SYNC_FLUSH

autoCreateTopicEnable=true

文件二:配置 Slave A。创建并编辑 conf/2m-2s-async/broker-a-s.properties:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
brokerClusterName=DefaultCluster
# 名字叫 broker-a,负责死守 .5 机器上的 Master A
brokerName=broker-a
brokerId=1
namesrvAddr=192.168.1.7:9876;192.168.1.8:9876

listenPort=11011
storePathRootDir=/root/store/broker-a-s
storePathCommitLog=/root/store/broker-a-s/commitlog

deleteWhen=04
fileReservedTime=48
brokerRole=SLAVE
flushDiskType=ASYNC_FLUSH

autoCreateTopicEnable=true


集群启动与验证

配置完成后,严格按照以下时序在两台机器上启动对应的系统组件。

步骤 1:两台机器同时拉起 NameServer。在 192.168.1.7 和 192.168.1.8 上分别执行:

1
2
$ nohup sh bin/mqnamesrv &
$ tail -f ~/logs/rocketmqlogs/namesrv.log

步骤 2:在 192.168.1.7 上启动 Broker 矩阵

1
2
3
4
# 1. 启动 Master A
$ nohup sh bin/mqbroker -c conf/2m-2s-async/broker-a.properties > /dev/null 2>&1 &
# 2. 启动 Slave B
$ nohup sh bin/mqbroker -c conf/2m-2s-async/broker-b-s.properties > /dev/null 2>&1 &

步骤 3:在 192.168.1.8 上启动 Broker 矩阵

1
2
3
4
# 1. 启动 Master B
$ nohup sh bin/mqbroker -c conf/2m-2s-async/broker-b.properties > /dev/null 2>&1 &
# 2. 启动 Slave A
$ nohup sh bin/mqbroker -c conf/2m-2s-async/broker-a-s.properties > /dev/null 2>&1 &

步骤 4:用命令检验集群成果

1
2
3
4
5
6
sh bin/mqadmin clusterList -n 192.168.1.7:9876
#Cluster Name #Broker Name #BID #Addr #Version #InTPS(LOAD) #OutTPS(LOAD) #Timer(Progress) #PCWait(ms) #Hour #SPACE #ACTIVATED
DefaultCluster broker-a 0 192.168.1.7:10911 V5_4_0 0.00(0,0ms) 0.00(0,0ms|0,0ms) 0-0(0.0w, 0.0, 0.0) 0 1.16 0.0700 true
DefaultCluster broker-a 1 192.168.1.8:11011 V5_4_0 0.00(0,0ms) 0.00(0,0ms|0,0ms) 2-0(0.0w, 0.0, 0.0) 0 1.16 0.0700 false
DefaultCluster broker-b 0 192.168.1.8:10911 V5_4_0 0.00(0,0ms) 0.00(0,0ms|0,0ms) 0-0(0.0w, 0.0, 0.0) 0 0.58 0.0700 true
DefaultCluster broker-b 1 192.168.1.7:11011 V5_4_0 0.00(0,0ms) 0.00(0,0ms|0,0ms) 2-0(0.0w, 0.0, 0.0) 0 0.58 0.0700 false


客户端的使用

客户端代码

集群搭建完毕后,我们在之前单机版的 Spring Boot / Java 项目中接入时,代码改动极小。核心的物理变化是:NameServer 的地址可以改为传入全部的节点列表(当然也可以传一个或者一部分)。

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
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;

public class ClusterProducer {
public static void main(String[] args) throws Exception {
// 1. 实例化生产者
DefaultMQProducer producer = new DefaultMQProducer("GID_INDUSTRIAL_PRODUCER");

// 2. 核心:必须指定集群中所有健康的 NameServer 地址,用分号隔开
producer.setNamesrvAddr("192.168.1.7:9876");

producer.start();

// 3. 顺着双主双从集群发一条消息
for (int i = 0; i < 100; i++) {
Message msg = new Message("TOPIC_ORDER", "TagA", "ORDER_KEY_10086", "工业级集群测试数据".getBytes());
SendResult sendResult = producer.send(msg); ////
System.out.printf("消息集群投递成功,响应明细: %s%n", sendResult);
Thread.sleep(500);
}

producer.shutdown();
}
}


遇到的问题

当这行代码执行时,Producer 会随机挑一台 NameServer(如 .7)拉取 TOPIC_ORDER 的路由表。由于我们配置了 2 个 Master,NameServer 会返回告诉它:这个 Topic 目前在 broker-a 上有 4 个队列,在 broker-b 上也有 4 个队列。随后,你的客户端在发消息时,就会自动轮询向 192.168.1.5:10911 (Master A) 和 192.168.1.7:10911(Master B) 交叉砸入字节流。整个系统的单点瓶颈得到完美解决。好,我们来运行上述代码,看看消息发送时是不是真的轮询往 MasterA 和 MasterB 发呢?

1
2
3
4
5
6
7
8
9
10
消息集群投递成功,响应明细: SendResult [sendStatus=SEND_OK, msgId=24098A002455B220A0C02E99DF2717C139BE22D8CFE0241284B40000, offsetMsgId=C0A8010700002A9F0000000000038BF7, messageQueue=MessageQueue [topic=TOPIC_ORDER, brokerName=broker-a, queueId=0], queueOffset=0, recallHandle=null]
消息集群投递成功,响应明细: SendResult [sendStatus=SEND_OK, msgId=24098A002455B220A0C02E99DF2717C139BE22D8CFE0241284D40001, offsetMsgId=C0A8010700002A9F0000000000038D22, messageQueue=MessageQueue [topic=TOPIC_ORDER, brokerName=broker-a, queueId=1], queueOffset=0, recallHandle=null]
消息集群投递成功,响应明细: SendResult [sendStatus=SEND_OK, msgId=24098A002455B220A0C02E99DF2717C139BE22D8CFE0241284DA0002, offsetMsgId=C0A8010700002A9F0000000000038E4D, messageQueue=MessageQueue [topic=TOPIC_ORDER, brokerName=broker-a, queueId=2], queueOffset=0, recallHandle=null]
消息集群投递成功,响应明细: SendResult [sendStatus=SEND_OK, msgId=24098A002455B220A0C02E99DF2717C139BE22D8CFE0241284DF0003, offsetMsgId=C0A8010700002A9F0000000000038F78, messageQueue=MessageQueue [topic=TOPIC_ORDER, brokerName=broker-a, queueId=3], queueOffset=0, recallHandle=null]
消息集群投递成功,响应明细: SendResult [sendStatus=SEND_OK, msgId=24098A002455B220A0C02E99DF2717C139BE22D8CFE0241285180008, offsetMsgId=C0A8010700002A9F00000000000390A3, messageQueue=MessageQueue [topic=TOPIC_ORDER, brokerName=broker-a, queueId=0], queueOffset=1, recallHandle=null]
消息集群投递成功,响应明细: SendResult [sendStatus=SEND_OK, msgId=24098A002455B220A0C02E99DF2717C139BE22D8CFE02412851E0009, offsetMsgId=C0A8010700002A9F00000000000391CE, messageQueue=MessageQueue [topic=TOPIC_ORDER, brokerName=broker-a, queueId=1], queueOffset=1, recallHandle=null]
消息集群投递成功,响应明细: SendResult [sendStatus=SEND_OK, msgId=24098A002455B220A0C02E99DF2717C139BE22D8CFE024128524000A, offsetMsgId=C0A8010700002A9F00000000000392F9, messageQueue=MessageQueue [topic=TOPIC_ORDER, brokerName=broker-a, queueId=2], queueOffset=1, recallHandle=null]
消息集群投递成功,响应明细: SendResult [sendStatus=SEND_OK, msgId=24098A002455B220A0C02E99DF2717C139BE22D8CFE02412852B000B, offsetMsgId=C0A8010700002A9F0000000000039424, messageQueue=MessageQueue [topic=TOPIC_ORDER, brokerName=broker-a, queueId=3], queueOffset=1, recallHandle=null]
消息集群投递成功,响应明细: SendResult [sendStatus=SEND_OK, msgId=24098A002455B220A0C02E99DF2717C139BE22D8CFE024128532000C, offsetMsgId=C0A8010800002A9F000000000000E5A5, messageQueue=MessageQueue [topic=TOPIC_ORDER, brokerName=broker-b, queueId=0], queueOffset=1, recallHandle=null]
...

嗯?奇怪!怎么不往 broker-b 发呢?而且无论我们如何在客户端发消息,就是出不来 broker-b!😓


问题的解决

其实在集群模式下,导致这种现象的原因只有一个:因为你在代码里,没有为 TOPIC_ORDER 提前创建路由元数据,导致 broker-b 根本没有向 NameServer 认领这个 Topic.。罪魁祸首就是我们在 Broker 的配置文件里都设置了 autoCreateTopicEnable=true

  • 当你的客户端发出第一条消息时,它去 NameServer 问:“TOPIC_ORDER 在哪?”
  • 此时因为是新 Topic,NameServer 账本里根本没有它!NameServer 就会把系统默认的一个兜底模板 Topic(TBW102)的路由丢给客户端。
  • TBW102 默认会指向当前第一个连接最快、最健康的 Broker(也就是刚才的 broker-a)。
  • 于是,客户端就在本地内存里为 TOPIC_ORDER 绑定了 broker-a 的 4 个队列(Queue 0, 1, 2, 3)。这就是为什么我们刚才的日志里只出现了 broker-a 的 2 -> 3 -> 0 -> 1 的轮询。
  • 既然客户端在内存里已经认定 TOPIC_ORDER 属于 broker-a,那么不管你 sleep 500ms 还是 5000ms,客户端每 30 秒去跟 NameServer 同步路由时,NameServer 返回的依然只有 broker-a。
  • 只有当有人真正往 broker-b 发送一次该 Topic 的消息,broker-b 才会动态在磁盘上为它创建文件夹并上报 NameServer。这就陷入了一个 “broker-b 不建立路由 $\rightarrow$ 客户端不往它发 $\rightarrow$ 它更不建立路由” 的死循环。

在真正的生产环境里是绝对禁止开启 autoCreateTopicEnable=true 的。所有的 Topic 必须在部署时由运维人员通过命令手动在所有 Master 上整齐划一地创建好。现在,我们只能在 centos10-01 上执行以下命令,强行打通 broker-b 的任督二脉:

1
2
3
4
5
6
7
8
# config/topics.json
# -c DefaultCluster:代表强行让该集群下的所有 Master(a和b)同时创建该 Topic
# -r 8:对应配置里的 readQueueNums=8(读队列数)。
# -w 8:对应配置里的 writeQueueNums=8(写队列数)。
$ sh bin/mqadmin updateTopic -n 192.168.1.7:9876 -c DefaultCluster -t TOPIC_ORDER -r 8 -w 8
create topic to 192.168.1.7:10911 success.
create topic to 192.168.1.8:10911 success.
TopicConfig [topicName=TOPIC_ORDER, readQueueNums=8, writeQueueNums=8, perm=RW-, topicFilterType=SINGLE_TAG, topicSysFlag=0, order=false, attributes={}]

再次重启你的 Java 业务程序,客户端这次再向 NameServer 询问。这一次 NameServer 会吐出一个完美的、包含了 broker-a(4个队列)和 broker-b(4个队列)共计 8 个队列的完整路由表。

1
2
3
4
5
6
7
8
9
10
11
12
13
消息集群投递成功,响应明细: SendResult [sendStatus=SEND_OK, msgId=24098A002455B220A0C02E99DF2717C13B7D22D8CFE0241FBA6C0005, offsetMsgId=C0A8010700002A9F000000000004768A, messageQueue=MessageQueue [topic=TOPIC_ORDER, brokerName=broker-a, queueId=0], queueOffset=12, recallHandle=null]
消息集群投递成功,响应明细: SendResult [sendStatus=SEND_OK, msgId=24098A002455B220A0C02E99DF2717C13B7D22D8CFE0241FBA710006, offsetMsgId=C0A8010700002A9F00000000000477B5, messageQueue=MessageQueue [topic=TOPIC_ORDER, brokerName=broker-a, queueId=1], queueOffset=12, recallHandle=null]
消息集群投递成功,响应明细: SendResult [sendStatus=SEND_OK, msgId=24098A002455B220A0C02E99DF2717C13B7D22D8CFE0241FBA7C0007, offsetMsgId=C0A8010700002A9F00000000000478E0, messageQueue=MessageQueue [topic=TOPIC_ORDER, brokerName=broker-a, queueId=2], queueOffset=12, recallHandle=null]
消息集群投递成功,响应明细: SendResult [sendStatus=SEND_OK, msgId=24098A002455B220A0C02E99DF2717C13B7D22D8CFE0241FBA800008, offsetMsgId=C0A8010700002A9F0000000000047A0B, messageQueue=MessageQueue [topic=TOPIC_ORDER, brokerName=broker-a, queueId=3], queueOffset=14, recallHandle=null]
消息集群投递成功,响应明细: SendResult [sendStatus=SEND_OK, msgId=24098A002455B220A0C02E99DF2717C13B7D22D8CFE0241FBA8C0009, offsetMsgId=C0A8010800002A9F000000000001CF09, messageQueue=MessageQueue [topic=TOPIC_ORDER, brokerName=broker-b, queueId=0], queueOffset=14, recallHandle=null]
消息集群投递成功,响应明细: SendResult [sendStatus=SEND_OK, msgId=24098A002455B220A0C02E99DF2717C13B7D22D8CFE0241FBA93000A, offsetMsgId=C0A8010800002A9F000000000001D034, messageQueue=MessageQueue [topic=TOPIC_ORDER, brokerName=broker-b, queueId=1], queueOffset=14, recallHandle=null]
消息集群投递成功,响应明细: SendResult [sendStatus=SEND_OK, msgId=24098A002455B220A0C02E99DF2717C13B7D22D8CFE0241FBA99000B, offsetMsgId=C0A8010800002A9F000000000001D15F, messageQueue=MessageQueue [topic=TOPIC_ORDER, brokerName=broker-b, queueId=2], queueOffset=14, recallHandle=null]
消息集群投递成功,响应明细: SendResult [sendStatus=SEND_OK, msgId=24098A002455B220A0C02E99DF2717C13B7D22D8CFE0241FBA9F000C, offsetMsgId=C0A8010800002A9F000000000001D28A, messageQueue=MessageQueue [topic=TOPIC_ORDER, brokerName=broker-b, queueId=3], queueOffset=13, recallHandle=null]
消息集群投递成功,响应明细: SendResult [sendStatus=SEND_OK, msgId=24098A002455B220A0C02E99DF2717C13B7D22D8CFE0241FBAA5000D, offsetMsgId=C0A8010700002A9F0000000000047B36, messageQueue=MessageQueue [topic=TOPIC_ORDER, brokerName=broker-a, queueId=0], queueOffset=13, recallHandle=null]
消息集群投递成功,响应明细: SendResult [sendStatus=SEND_OK, msgId=24098A002455B220A0C02E99DF2717C13B7D22D8CFE0241FBAAA000E, offsetMsgId=C0A8010700002A9F0000000000047C61, messageQueue=MessageQueue [topic=TOPIC_ORDER, brokerName=broker-a, queueId=1], queueOffset=13, recallHandle=null]
消息集群投递成功,响应明细: SendResult [sendStatus=SEND_OK, msgId=24098A002455B220A0C02E99DF2717C13B7D22D8CFE0241FBAAF000F, offsetMsgId=C0A8010700002A9F0000000000047D8C, messageQueue=MessageQueue [topic=TOPIC_ORDER, brokerName=broker-a, queueId=2], queueOffset=13, recallHandle=null]
消息集群投递成功,响应明细: SendResult [sendStatus=SEND_OK, msgId=24098A002455B220A0C02E99DF2717C13B7D22D8CFE0241FBAB60010, offsetMsgId=C0A8010700002A9F0000000000047EB7, messageQueue=MessageQueue [topic=TOPIC_ORDER, brokerName=broker-a, queueId=3], queueOffset=15, recallHandle=null]
...


sendAsync 和 send 的执行

异步发送 sendAsync

异步发送利用了 TCP 通道复用。业务线程把字节流往 Netty 管道里一扔,瞬间返回(耗时接近 0ms)。当 Broker 的响应顺着网线飞回来时,Netty 读线程会根据 opaque 流水号去认亲,并执行对应的 SendCallback。通常用于高并发、大吞吐、允许轻微延迟的非核心主线业务(如日志收集、用户行为埋点、给用户发优惠券通知)。异步发送最突出作用就是能够释放网卡吞吐量。

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
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendCallback;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import java.util.concurrent.CountDownLatch;

public class AsyncSendDemo {
public static void main(String[] args) throws Exception {
DefaultMQProducer producer = new DefaultMQProducer("GID_ASYNC_DEMO");
producer.setNamesrvAddr("192.168.1.7:9876;192.168.1.8:9876");
producer.start();

// 搞一个门闩,仅仅是为了防止主线程退出了,看不了回调日志
CountDownLatch latch = new CountDownLatch(1);

// 1. 组装消息
Message msg = new Message("TOPIC_ORDER", "TagA", "KEY_ASYNC_002",
"高并发行为埋点数据".getBytes());

System.out.println("【异步发送】开始投递...");
long startTime = System.currentTimeMillis();

// 2. 调用异步发送,传入一个匿名内部类充当“认亲回调”
producer.send(msg, new SendCallback() {

@Override
public void onSuccess(SendResult sendResult) {
// 当 Broker 成功落盘并将 opaque 吐回来时,Netty 线程会顺藤摸瓜点燃这个方法
System.out.printf("【异步通知】收到 Broker 回音!消息投递成功。MsgId: %s, 耗时: %d ms%n",
sendResult.getMsgId(), (System.currentTimeMillis() - startTime));
latch.countDown();
}

@Override
public void onException(Throwable e) {
// 如果在异步重试后依然失败,或者超时,会点燃这个方法
System.err.println("【异步通知】由于网络抖动或 Broker 沦陷,该消息投递宣告失败!");
e.printStackTrace();
latch.countDown();
}
});

// 3. 证明业务线程根本没有被阻塞
System.out.printf("【异步发送】sendAsync 执行完毕!主业务线程耗时仅 %d ms,已重获自由去干别的大事了!%n",
(System.currentTimeMillis() - startTime));

// 死等回调执行完毕,防止 JVM 过早退出
latch.await();
producer.shutdown();
}
}
1
2
3
【异步发送】开始投递...
【异步发送】sendAsync 执行完毕!主业务线程耗时仅 1 ms,已重获自由去干别的大事了!
【异步通知】收到 Broker 回音!消息投递成功。MsgId: 24098A002455B220A0C02E99DF2717C17C3122D8CFE02640E4540000, 耗时: 477 ms

在 RocketMQ 的物理世界里,客户端和 broker-a 之间,通常只有一条(或者极少数几条)长连接 TCP 管道。当 100 个业务线程并发调用 sendAsync 时,这些消息就像是 100 颗不同颜色的子弹,被无锁(或者极轻量锁)地倾泻进同一个 TCP 管道中。由于网线是串行的,这 100 颗子弹在网线上是相互交错、前后挨着的。 Broker 收到后,由于 netty 采用了 SendMessageExecutor 业务线程池并发写盘,写盘有快有慢。 这就导致 Broker 吐回来的响应顺序,跟客户端发过去的请求顺序,是完全对不上的!比如客户端先发了请求 1,再发了请求 2。但 Broker 可能会先返回请求 2 的结果,再返回请求 1 的结果。既然网络全乱了,客户端怎么知道哪一个响应属于哪一个业务回调? 答案就是核心组件:opaque(请求唯一标识)与 ResponseFuture。

在 RocketMQ 的通讯内核 rocketmq-remoting 中,有一个核心类叫 ResponseFuture。每次你调用异步发送,底层都会在内存中凭空组装出这样一个精密的状态机:

1
2
3
4
5
6
7
8
9
public class ResponseFuture {
private final int opaque; // 请求的唯一流水号(快递单号)
private final Channel channel; // 对应的 Netty 管道
private final long timeoutMillis; // 超时时间
private final InvokeCallback invokeCallback; // 你在业务层传进来的 SendCallback 隐形斗篷
private volatile RemotingCommand responseCommand; // 存放最终飞回来的响应字节
private final CountDownLatch countDownLatch = new CountDownLatch(1); // 物理挂起锚点
// ...
}

我们来看底层的执行代码到底是怎么走的:

第一,业务线程挂载账本(不阻塞):

  • 你的业务线程调用 sendAsync。
  • RocketMQ 底层通过一个内部原子计数器(AtomicInteger),生成一个全局唯一的 opaque(比如:9527)。
  • 把你的 SendCallback 业务回调函数包装进 ResponseFuture 对象中。
  • 关键一步是把这个对象塞进一个客户端全局的内存哈希表里:responseTable.put(9527, future)。这个表就是“认亲账本”。
  • Netty 异步把消息通过 TCP 网线推出去。业务线程完全不等待,立刻挥一挥衣袖返回。线程重获自由,可以去处理下一个请求。

第二,网线与 Broker 端,疯狂并网:

  • 请求夹带着 opaque = 9527 的标记划过网线。
  • Broker 收到后,高并发写盘。处理完毕后,组装响应报文,并原封不动地把 opaque = 9527 塞回响应头里,顺着网线吐回给客户端。

第三,Netty 读线程(Selector),顺藤摸瓜,惊醒回调:

  • 客户端的网卡收到了二进制字节流,Netty 的 Selector 线程(I/O 线程)抓起这串字节,经过编解码器翻译,还原成一个 RemotingCommand 响应对象。
  • Netty 读线程一眼看到了这个响应对象里面的 opaque = 9527。
  • 它立刻冲向刚才那个 “认亲账本”:ResponseFuture future = responseTable.remove(9527)。
  • 拿到这个 future 后,Netty 读线程(或者丢给公用的异步回调线程池)直接调用 future.getInvokeCallback().operationComplete(future)
  • 顺着这个入口,最终点燃了你在 Java 业务层里写下的 onSuccess(SendResult sendResult) 代码块!

RocketMQ 压榨网卡的终极魔法,就是用一个基于全局唯一流水号(opaque)的内存认亲表(responseTable),把串行的 TCP 网线强行扭置成了并发的立交桥。网络 I/O 线程只管收发,业务线程只管扔消息,彼此通过 ResponseFuture 隔空对齐。


同步发送 send

同步发送会物理挂起当前业务线程(利用 CountDownLatch),直到 Broker 返回 SendResult。通常用于对数据正确性要求极高、不能容忍任何丢失的场景(如钱包扣款、订单创建)。

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
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;

public class SyncSendDemo {
public static void main(String[] args) {
DefaultMQProducer producer = new DefaultMQProducer("GID_SYNC_DEMO");
producer.setNamesrvAddr("192.168.1.7:9876;192.168.1.8:9876");

try {
producer.start();

// 1. 组装消息
Message msg = new Message("TOPIC_ORDER", "TagA", "KEY_SYNC_001",
"核心扣款业务数据".getBytes());

System.out.println("【同步发送】准备砸入网络,当前业务线程即将挂起休眠...");

// 2. 调用同步 send。此时业务线程原地“死等”网线对面的回音
SendResult sendResult = producer.send(msg);

// 3. 只有收到 Broker 的 ACK,下面这行代码才会被惊醒并执行
System.out.printf("【同步发送】成功被唤醒!Broker 响应明细: %s%n", sendResult);

} catch (Exception e) {
// 同步发送一旦抛出异常,说明网络彻底断开或 Broker 宕机,必须立刻进行业务补偿
System.err.println("【同步发送】捕获到物理崩溃,启动紧急容灾预案!");
e.printStackTrace();
} finally {
producer.shutdown();
}
}
}
1
2
【同步发送】准备砸入网络,当前业务线程即将挂起休眠...
【同步发送】成功被唤醒!Broker 响应明细: SendResult [sendStatus=SEND_OK, msgId=24098A002455B220A0C02E99DF2717C17CB022D8CFE026447A1C0000, offsetMsgId=C0A8010800002A9F00000000000205EE, messageQueue=MessageQueue [topic=TOPIC_ORDER, brokerName=broker-b, queueId=2], queueOffset=21, recallHandle=null]

如果你调用的是同步的 producer.send(msg),它的底层网络架构和 sendAsync 一模一样,用的也是同一个 TCP 管道,也生成了 opaque。唯一的物理区别在于:

  • 同步发送时,业务线程把消息交给 Netty 后,没有立刻返回,而是抓起 ResponseFuture 内部的那个 countDownLatch,调用了 future.countDownLatch.await(timeout, TimeUnit.MILLISECONDS)

  • 业务线程原地物理挂起(休眠)。

  • 当 Netty 读线程根据 opaque = 9527 捞出 future 时,它不会执行回调,而是执行

    future.countDownLatch.countDown()

  • 这一脚油门下去,把正在休眠的业务线程瞬间唤醒,业务线程从 future.responseCommand 里拿走结果,假装自己完成了一次 “同步阻塞” 调用。


消息接收端代码注意事项

在 RocketMQ 的物理世界里,消息接收端(Consumer)通常是整个微服务分布式架构中最脆弱、最容易沦陷的一环。因为发送端(Producer)往往只负责高频地往外抛字节,而接收端却要实打实地承载执行业务逻辑、写数据库、调第三方接口等沉重的肉身。如果接收端代码写得不规范,极易引发内存溢出(OOM)、消息积压、顺序错乱、以及死锁等生产事故。


必须死守幂等性

消息接收端必须死守“幂等性”,绝不盲目信任中间件。由于网络抖动、重试机制或消费者扩容重平衡,RocketMQ 只能保证 “At least once”(至少投递一次),绝对无法避免重复投递。

❌ 错误示范:拿到消息直接做业务

1
2
3
4
5
6
7
8
9
// 错误:一旦网络抖动,这条消息被投递了两次,用户就会被扣款两次!
consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> {
for (MessageExt msg : msgs) {
String orderId = msg.getKeys();
int amount = Integer.parseInt(new String(msg.getBody()));
walletService.deductMoney(orderId, amount); // 直接扣款,极度危险!
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
});

✅ 正确示范:利用数据库唯一键或 Redis 去重表实施物理拦截

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> {
for (MessageExt msg : msgs) {
String msgId = msg.getMsgId(); // 获取全局唯一的消息ID(或者业务的订单号)

// 1. 利用 Redis 的 setNX 建立防线
boolean isFirstConsuming = redisTemplate.opsForValue()
.setIfAbsent("LOCK_CONSUME:" + msgId, "Y", Duration.ofHours(1));

if (!isFirstConsuming) {
System.err.printf("【防线拦截】检测到重复投递的消息: %s,直接宣告成功,拒绝重复执行业务!%n", msgId);
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; // 优雅放行
}

try {
// 2. 执行真正的扣款业务
walletService.deductMoney(msg.getKeys(), Integer.parseInt(new String(msg.getBody())));
} catch (Exception e) {
redisTemplate.delete("LOCK_CONSUME:" + msgId); // 业务失败,释放防线以便重试
return ConsumeConcurrentlyStatus.RECONSUME_LATER; // 丢进重试队列
}
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
});


区分并发消费与顺序消费

RocketMQ 提供了两种监听器:

  • MessageListenerConcurrently(并发消费):多线程同时吃不同的消息,速度极快,但无法保证顺序。
  • MessageListenerOrderly(顺序消费):单线程锁定同一个 Queue,严格按照物理写入顺序一条一条吃。如果你的业务(如订单状态流转:创建 $\rightarrow$ 支付 $\rightarrow$ 发货)要求绝对不能乱,就必须使用 Orderly。

❌ 错误示范:用并发监听器去处理顺序业务

1
2
3
4
5
6
// 错误:虽然你用 MessageQueueSelector 投递到了同一个 Queue,但在这里被多线程并发消费,
// “发货”线程跑得比“支付”快,状态机瞬间崩塌!
consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> {
System.out.printf("当前线程 %s 正在乱序消费消息...%n", Thread.currentThread().getName());
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
});

✅ 正确示范:使用 MessageListenerOrderly 死守单通道

1
2
3
4
5
6
7
8
9
10
// 正确:RocketMQ 底层会自动用线程锁、内存锁把该队列死死锁住,强迫在单线程下顺次消费
consumer.registerMessageListener((MessageListenerOrderly) (msgs, context) -> {
for (MessageExt msg : msgs) {
// 此时,同一个订单号的消息,绝对是“创建”完,再“支付”,再“发货”
System.out.printf("【顺序消费】当前线程: %s, 正在处理队列 [%d] 的消息: %s%n",
Thread.currentThread().getName(), context.getMessageQueue().getQueueId(), new String(msg.getBody()));
}
// 注意:顺序消费返回的是 SUSPEND_CURRENT_QUEUE_A_MOMENT(挂起当前队列一会再试)
return ConsumeOrderlyStatus.SUCCESS;
});


严格禁止在监听器内将消息抛给自定义线程池处理

我们已经知道,消费点位(Offset)的右移,完全取决于消费监听器返回 CONSUME_SUCCESS 的时机。如果你在监听器内部自己搞了一个 new ThreadPoolExecutor 把消息丢进去异步执行,然后监听器立刻返回了 SUCCESS,这会导致 RocketMQ 的自愈链路完全失效!

❌ 错误示范:自作聪明的“异步消費”

1
2
3
4
5
6
7
8
9
10
// 致命错误:消息刚丢给自定义线程池,监听器就返回 SUCCESS 了,Broker 端认为消费成功,Offset 直接右移。
// 结果几毫秒后,你的自定义线程池因为内存爆了或者报错挂了,这条消息在物理上彻底丢失!
consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> {
for (MessageExt msg : msgs) {
myCustomThreadPool.submit(() -> {
doHeavyBusiness(msg); // 真正干活,但生死无人感知
});
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
});

✅ 正确示范:老老实实利用 RocketMQ 自带的线程参数调优。如果你觉得消费不够快,正确的姿势是去调大 Consumer 自带的底层线程池参数,而不是自己写线程池:

1
2
3
4
5
6
7
8
9
10
11
12
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("GID_INDUSTRIAL_CONSUMER");

// 正确姿势:通过配置直接调节底层 Netty 消费线程池的并发水位
consumer.setConsumeThreadMin(20); // 最小 20 个线程
consumer.setConsumeThreadMax(64); // 最大 64 个线程
consumer.setConsumeMessageBatchMaxSize(1); // 每次拉取只处理1条,防止大批消息一起卡死

consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> {
// 就在 RocketMQ 给你的线程里规规矩矩干活,执行完再返回,让点位安全右移
doHeavyBusiness(msgs.get(0));
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
});


重试和异常情况的处理

并发消费下:捕获一切异常,优雅返回 RECONSUME_LATER。在 Concurrently 模式下,不要尝试在代码里自己写 while(true) 死循环去重试。应该主动返回 RECONSUME_LATER,让消息进入我们第三章拆解的 %RETRY% 梯度阻尼重试队列,让出 CPU 算力。

1
2
3
4
5
6
7
8
9
10
consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> {
try {
executeRemoteHttpCall(msgs.get(0));
} catch (Throwable e) {
// 捕获所有包括 Runtime 的异常,交还给 RocketMQ 走 16 次重试自愈链路
System.err.println("远程调用失败,触发 RocketMQ 梯度退避重试机制");
return ConsumeConcurrentlyStatus.RECONSUME_LATER;
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
});

顺序消费下:千万不能返回 SUSPEND 太久。在 Orderly 顺序模式下,如果返回了 SUSPEND_CURRENT_QUEUE_A_MOMENT(挂起一会再试),RocketMQ 默认会每隔 1 秒钟疯狂重新投递这条失败的消息,直到成功。

  • 如果你的底层数据库彻底瘫痪了,代码一直返回 SUSPEND,这条消息会把整个 Queue 文件夹死死卡住几个小时,后面的订单消息全部积压。
  • 在顺序消费时,代码内部自己记录重试次数,如果超过 3 次依然失败,记录日志,强行返回 SUCCESS 释放通道,转为人工补偿,切记!