目录
一、为什么需要消息队列?RocketMQ 的定位
二、快速实战:从 0 搭建并发送第一条消息
1. 环境准备
① 创建宿主机挂载目录
② 编写 Broker 配置文件
③ 启动容器(完整命令,含注释)
④ 开机自启设置
⑤ 验证安装
2. 生产者代码
3. 消费者代码
三、核心概念深度拆解
1. 整体架构图
2. NameServer:极简路由中心
3. Broker:消息存储与高可用核心
4. Topic、Tag、Message Key
5. Queue:并行度的最小单位
6. Producer Group 与 Consumer Group
7. 消息消费与负载均衡
四、顺序消息原理与代码
五、事务消息原理与代码
六、与 Kafka、RabbitMQ 的对比
七、总结
一、为什么需要消息队列?RocketMQ 的定位
在微服务架构中,服务间的同步调用会形成强依赖,一旦某个服务变慢或宕机,整个链路可能雪崩。消息队列可以:
-
解耦:订单服务只需发消息,库存服务、物流服务各自订阅,互不干扰。
-
削峰:秒杀瞬间的 10 万 QPS 先灌进 MQ,后端系统按自己的处理能力慢慢消费。
-
异步:发送短信、生成报表等非核心逻辑异步化,提高接口响应速度。
RocketMQ 的特点:
-
纯 Java 开发,社区活跃(已捐给 Apache)。
-
支持 10 万级 QPS 吞吐,毫秒级延迟。
-
原生支持事务消息、顺序消息、延时消息。
-
架构简单,NameServer 比 ZK 更轻量,运维成本低。
二、快速实战:从 0 搭建并发送第一条消息
1. 环境准备
以下使用 Docker 搭建 RocketMQ 单机环境,重点讲解 JVM 参数、数据持久化挂载 和 容器自启,让学习环境更接近生产要求。
① 创建宿主机挂载目录
在宿主机上创建目录,用于保存 RocketMQ 的数据、日志和配置文件,防止容器重启后数据丢失。
# 创建 NameServer 日志目录(NameServer 只存路由信息,无消息数据)
mkdir -p /opt/rocketmq/namesrv/logs
# 创建 Broker 数据目录(CommitLog、ConsumerQueue 等落盘文件)
mkdir -p /opt/rocketmq/broker/data
# 创建 Broker 日志目录
mkdir -p /opt/rocketmq/broker/logs
# 创建 Broker 配置文件目录(存放 broker.conf)
mkdir -p /opt/rocketmq/broker/conf
② 编写 Broker 配置文件
在宿主机 /opt/rocketmq/broker/conf/broker.conf 创建一个配置文件,挂载进容器,方便修改参数。
创建文件(用 cat 或 echo 均可):
cat > /opt/rocketmq/broker/conf/broker.conf << 'EOF'
# broker.conf 示例
# 指定 Broker 名称(集群内唯一)
brokerName = broker-a
# 指定所属集群名称
brokerClusterName = DefaultCluster
# Broker ID,0 表示 Master,>0 表示 Slave
brokerId = 0
# 消息存储路径(对应宿主机挂载点)
storePathRootDir = /home/rocketmq/store
# CommitLog 存储路径
storePathCommitLog = /home/rocketmq/store/commitlog
# ConsumerQueue 存储路径
storePathConsumeQueue = /home/rocketmq/store/consumequeue
# 是否自动创建 Topic(线上建议关闭,手动管理)
autoCreateTopicEnable = true
# 是否自动创建订阅组
autoCreateSubscriptionGroup = true
EOF
如果习惯用 vim,也可以 vim /opt/rocketmq/broker/conf/broker.conf 然后粘贴内容。
检查文件内容:
cat /opt/rocketmq/broker/conf/broker.conf
③ 启动容器(完整命令,含注释)
NameServer:
docker run -d \\
–name rmq-namesrv \\ # 容器名称
–restart=always \\ # 开机自启,守护进程挂了自动重启
-p 9876:9876 \\ # 暴露 NameServer 端口
-v /opt/rocketmq/namesrv/logs:/home/rocketmq/logs \\ # 挂载日志目录到宿主机
-e "MAX_POSSIBLE_HEAP=256m" \\ # 设置 JVM 最大堆内存(生产建议 2G,测试用 256m 足够)
-e "JAVA_OPT_EXT=-server -Xms256m -Xmx256m -Xmn128m" \\ # 精细控制 JVM 参数
apache/rocketmq:5.1.4 sh mqnamesrv # 镜像名 + 启动命令
参数解释:
-
MAX_POSSIBLE_HEAP:RocketMQ 容器内置的环境变量,用于自动计算 JVM 堆大小。生产环境建议设置为物理内存的一半,例如 2g。
-
JAVA_OPT_EXT:额外的 JVM 参数,这里显式指定了初始堆 -Xms、最大堆 -Xmx 和新生代大小 -Xmn,避免容器内存超限。
发现有注释的直接粘贴到Linux上会无法识别到末尾的 \\,增加一个无注释版本的:
docker run -d \\
–name rmq-namesrv \\
–restart=always \\
-p 9876:9876 \\
-v /opt/rocketmq/namesrv/logs:/home/rocketmq/logs \\
-e "MAX_POSSIBLE_HEAP=256m" \\
-e "JAVA_OPT_EXT=-server -Xms256m -Xmx256m -Xmn128m" \\
apache/rocketmq:5.1.4 sh mqnamesrv
Broker:
docker run -d \\
–name rmq-broker \\ # 容器名称
–restart=always \\ # 开机自启
-p 10911:10911 \\ # Broker 对外通信端口(生产者、消费者连接)
-p 10909:10909 \\ # Broker 内部管理端口(FastFailure 等)
-e "NAMESRV_ADDR=192.168.31.72:9876" \\ # 指定 NameServer 地址(改为自己内网 IP)
-e "MAX_POSSIBLE_HEAP=512m" \\ # Broker 堆内存,测试 512m,生产至少 2~4G
-e "JAVA_OPT_EXT=-server -Xms512m -Xmx512m -Xmn256m" \\ # 显式 JVM 参数
-v /opt/rocketmq/broker/data:/home/rocketmq/store \\ # 挂载消息存储目录,防止数据丢失
-v /opt/rocketmq/broker/logs:/home/rocketmq/logs \\ # 挂载日志目录
-v /opt/rocketmq/broker/conf/broker.conf:/home/rocketmq/rocketmq-5.1.4/conf/broker.conf \\ # 挂载自定义配置
apache/rocketmq:5.1.4 sh mqbroker -c /home/rocketmq/rocketmq-5.1.4/conf/broker.conf # 启动 Broker 并加载外部配置
无注释:
docker run -d \\
–name rmq-broker \\
–restart=always \\
-p 10911:10911 \\
-p 10909:10909 \\
-e "NAMESRV_ADDR=192.168.31.72:9876" \\
-e "MAX_POSSIBLE_HEAP=512m" \\
-e "JAVA_OPT_EXT=-server -Xms512m -Xmx512m -Xmn256m" \\
-v /opt/rocketmq/broker/data:/home/rocketmq/store \\
-v /opt/rocketmq/broker/logs:/home/rocketmq/logs \\
-v /opt/rocketmq/broker/conf/broker.conf:/home/rocketmq/rocketmq-5.1.4/conf/broker.conf \\
apache/rocketmq:5.1.4 sh mqbroker -c /home/rocketmq/rocketmq-5.1.4/conf/broker.conf
注意:192.168.31.72 换成你本机的局域网 IP,不能写 127.0.0.1,否则客户端连不上。
④ 开机自启设置
如果在 docker run 时没有加 –restart=always,可以用 docker update 补救:
# 设置 NameServer 和 Broker 容器开机自启
docker update –restart=always rmq-namesrv
docker update –restart=always rmq-broker
重启策略说明:
-
–restart=always:无论容器退出状态码如何,Docker 守护进程都会重启它(Docker 守护进程启动时也会拉起来)。
-
还可以使用 –restart=unless-stopped:如果手动执行 docker stop 停止了容器,不会自动重启,其他情况重启。
⑤ 验证安装
# 查看容器运行状态
docker ps | grep rmq
# 进入 Broker 容器,使用命令行测试发送消息
docker exec -it rmq-broker bash
cd rocketmq-5.1.4/bin
# 发送一条测试消息
sh tools.sh org.apache.rocketmq.example.quickstart.Producer
# 消费测试消息
sh tools.sh org.apache.rocketmq.example.quickstart.Consumer

2. 生产者代码
// 1. 创建生产者实例,参数为“生产者组名”,同一组内的生产者功能相同,常用于事务回查
DefaultMQProducer producer = new DefaultMQProducer("order_producer_group");
// 2. 设置 NameServer 地址,多个地址用分号分隔
producer.setNamesrvAddr("192.168.31.72:9876");
// 3. 启动生产者,底层会与 NameServer 建立 TCP 长连接,并定期更新 Topic 路由信息
producer.start();
// 4. 构建消息对象:指定 Topic(主题)、Tag(标签)和消息体字节数组
Message msg = new Message(
"OrderTopic", // 主题,消息的一级分类
"OrderPaid", // 标签,消息的二级分类,用于消费者过滤
"订单已支付,订单号1001".getBytes(RemotingHelper.DEFAULT_CHARSET) // 消息体,UTF-8 编码
);
// 5. 同步发送消息,等待 Broker 返回结果
SendResult sendResult = producer.send(msg);
// 打印发送结果:包含发送状态、消息ID、队列偏移量等
System.out.printf("发送结果: %s%n", sendResult);
// 6. 应用关闭时释放资源
producer.shutdown();
3. 消费者代码
// 1. 创建消费者实例,指定消费者组名,同一组内的消费者共享消费进度
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("order_consumer_group");
// 2. 设置 NameServer 地址
consumer.setNamesrvAddr("192.168.31.72:9876");
// 3. 订阅 Topic 和 Tag,* 表示订阅所有 Tag
consumer.subscribe("OrderTopic", "OrderPaid");
// 4. 注册并发消息监听器,收到消息后回调此方法
consumer.registerMessageListener(new MessageListenerConcurrently() {
@Override
public ConsumeConcurrentlyStatus consumeMessage(
List<MessageExt> msgs, // 一批消息(默认一批最多32条)
ConsumeConcurrentlyContext context) {
// 遍历这一批消息并打印消息体
for (MessageExt msg : msgs) {
System.out.printf("收到消息: %s %n", new String(msg.getBody()));
}
// 返回 CONSUME_SUCCESS 确认消费成功,Broker 会标记该消息已消费
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
});
// 5. 启动消费者,开始长轮询拉取消息
consumer.start();
System.out.println("消费者启动成功");
三、核心概念深度拆解
RocketMQ 的整体架构是基于四个核心角色:NameServer、Broker、Producer、Consumer。下面详细拆解各个组件的职责、交互流程和底层存储原理。
1. 整体架构图

2. NameServer:极简路由中心
职责:
-
维护 Topic → Broker 的映射关系,以及每个 Topic 有哪些 Queue。
-
为客户端(生产者、消费者)提供路由发现服务。
设计精髓:
-
无状态、对等:每个 NameServer 节点独立,不进行数据同步。所有路由数据由 Broker 通过心跳上报。
-
最终一致:Broker 每 30 秒向所有 NameServer 发送心跳,上报自己管理的 Topic 信息。NameServer 2 分钟未收到心跳就剔除该 Broker。因此短暂不一致不影响可用性。
-
轻量:不负责存储,只做简单查询,比 ZooKeeper 简单太多,部署两个即可(三台完全够用)。
路由发现流程:

3. Broker:消息存储与高可用核心
职责:
-
接收生产者消息,持久化存储,并在消费者拉取时转发。
-
管理消费进度(Offset),处理重试消息和死信。
主从机制:
-
一个 Broker 可以有 Master 和 Slave。Master 负责写入,Slave 同步数据并提供读服务(可分担部分读压力)。
-
同步方式:同步复制(消息写入 Master 后必须等待 Slave 确认)或异步复制(默认,性能优先)。同步复制可保证高可用下不丢消息,但延迟略高。
-
故障转移:当 Master 宕机,Slave 可以被提升为 Master(需配合 Controller 或手动切换)。
消息存储原理:
RocketMQ 的存储设计是它能达到高吞吐的核心,它不直接按 Queue 存文件,而是所有消息顺序写入一个 CommitLog 文件。

-
CommitLog:所有 Topic 的消息混在一起,顺序追加写,每个文件默认 1GB,满了自动创建新文件。顺序写入磁盘性能远高于随机写。
-
ConsumerQueue:按照 Topic-Queue 维度建立的索引文件,每个条目固定大小(20 字节),存储了该消息在 CommitLog 中的物理偏移量、消息长度和 Tag 的哈希值。消费者拉取消息时,先根据 Queue 和 Offset 找到 ConsumerQueue 中的索引,再定位到 CommitLog 读取完整消息体。
-
为什么这么设计?
-
顺序写 CommitLog 保证写入吞吐量。
-
ConsumerQueue 体积小,可全内存加载,消费时只需几次磁盘 IO 即可获取消息。
-
4. Topic、Tag、Message Key
-
Topic:消息的一级分类,逻辑上相当于数据库的表。一个 Topic 可以被多个生产者发送,多个消费者订阅。
-
Tag:同一 Topic 下的二级过滤,消费者可以在订阅时指定 Tag 表达式(如 "Pay" 或 "Pay || Order")。Broker 在消息投递时支持根据 Tag 进行过滤,减少不必要网络传输。
-
Message Key:业务自定义的唯一标识(如订单 ID),可以通过 Key 在 Broker 上查询消息,用于排查问题。
5. Queue:并行度的最小单位
-
一个 Topic 可以被划分为多个 Queue(默认 8 个,建议设置为 8 或 16),Queue 分布在不同的 Broker 上。
-
并行消费:同一个 Consumer Group 内的多个实例可以分别持有不同的 Queue,实现水平扩展。比如 Topic 有 8 个 Queue,最多可以有 8 个消费者实例同时消费(再多实例会空闲,因为一个 Queue 只能被组内一个实例消费)。
-
顺序消息:需要全局顺序的消息,必须将所有消息发送到同一个 Queue,并且消费者只能单线程消费该 Queue。如果只需要局部顺序(比如一个订单的操作有序),则按订单 ID 路由到固定 Queue 即可,不同订单之间并行。
6. Producer Group 与 Consumer Group
-
Producer Group:生产者的集合,通常用于事务消息的回查。普通消息可不分组。
-
Consumer Group:消费者的逻辑分组,最重要的概念。
-
集群消费(CLUSTERING):默认模式。一条消息只会被组内的一个消费者消费,消费进度保存在 Broker 上。
-
广播消费(BROADCASTING):一条消息会被组内所有消费者消费,消费进度保存在消费者本地。
-
消费进度(Offset):每个 Consumer Group 对每个 Topic-Queue 维护一个 Offset,表示已经消费到哪个位置。Broker 定期持久化 Offset,避免重启后重复消费。
-
7. 消息消费与负载均衡
消费流程:
Consumer 从 NameServer 获取 Topic 的 Queue 分布。
Consumer 组内所有实例进行 Rebalance,平均分配 Queue(如 3 个实例 6 个 Queue,每个实例分 2 个)。
实例对分配给自己的 Queue 发起长轮询拉取请求(默认拉取 32 条)。
Broker 根据请求中的 Offset 读取 ConsumerQueue,找到 CommitLog 位置,返回消息列表。
消费者处理消息后返回消费状态,Broker 更新 Offset。
Rebalance 示意图:
Consumer Group 有 3 个实例 (C1, C2, C3) Topic 有 6 个 Queue (Q0 ~ Q5)
分配结果: C1: [Q0, Q3] C2: [Q1, Q4] C3: [Q2, Q5]
每个实例消费自己负责的 Queue,互不干扰。
四、顺序消息原理与代码
RocketMQ 实现顺序消息的关键是:
-
发送端将需要有序的消息路由到同一个 Queue。
-
消费端对该 Queue 使用单线程顺序消费(MessageListenerOrderly)。
发送端代码:
// 发送多条消息,但保证同一订单 ID 进入同一个 Queue
for (int i = 0; i < 10; i++) {
// 构造消息,设置订单 ID 为业务标识
Message msg = new Message("OrderTopic", "OrderStep",
("订单步骤-" + i).getBytes());
// 同步发送,并传入自定义的选择器
SendResult result = producer.send(msg, new MessageQueueSelector() {
@Override
public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
long orderId = (long) arg; // 从 arg 中获取订单 ID
int index = (int) (orderId % mqs.size()); // 按订单 ID 取模选定队列
return mqs.get(index); // 返回选中的队列
}
}, orderId); // 传入订单 ID 作为选择依据
System.out.printf("消息发送到 Queue: %s%n", result.getMessageQueue().getQueueId());
}
消费端代码:
// 使用顺序监听器,保证同一 Queue 内的消息单线程消费
consumer.registerMessageListener(new MessageListenerOrderly() {
@Override
public ConsumeOrderlyStatus consumeMessage(List<MessageExt> msgs,
ConsumeOrderlyContext context) {
// 遍历消息,因为同一 Queue 的消息是顺序到来的,所以按顺序处理
for (MessageExt msg : msgs) {
System.out.printf("顺序消费: Queue=%d, 消息=%s%n",
msg.getQueueId(), new String(msg.getBody()));
}
return ConsumeOrderlyStatus.SUCCESS;
}
});
五、事务消息原理与代码
RocketMQ 的事务消息流程如下:

关键点:Half 消息是真正写入 CommitLog 的,但会打上事务标记,消费者无法拉取。只有当收到 commit 确认后,才会创建 ConsumerQueue 索引,消息变得可见。
代码示例:
// 1. 创建事务消息生产者
TransactionMQProducer producer = new TransactionMQProducer("tx_producer_group");
producer.setNamesrvAddr("192.168.31.72:9876");
// 2. 设置事务监听器,实现本地事务执行和回查逻辑
producer.setTransactionListener(new TransactionListener() {
// 执行本地事务
@Override
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
try {
// 模拟执行业务操作(如扣减库存)
System.out.println("执行本地事务: " + new String(msg.getBody()));
return LocalTransactionState.COMMIT_MESSAGE; // 返回提交状态
} catch (Exception e) {
return LocalTransactionState.ROLLBACK_MESSAGE; // 异常则回滚
}
}
// 回查本地事务状态(用于 Broker 未收到提交/回滚时主动询问)
@Override
public LocalTransactionState checkLocalTransaction(MessageExt msg) {
// 查询数据库或其他状态,确定本地事务是否成功
System.out.println("回查事务状态: " + new String(msg.getBody()));
return LocalTransactionState.COMMIT_MESSAGE;
}
});
// 3. 启动生产者
producer.start();
// 4. 发送事务消息
Message msg = new Message("OrderTopic", "OrderPaid", "事务消息测试".getBytes());
producer.sendMessageInTransaction(msg, null); // 第二个参数可传附加参数给 executeLocalTransaction
// 5. 生产者关闭
producer.shutdown();
六、与 Kafka、RabbitMQ 的对比
| 设计语言 | Java | Scala/Java | Erlang |
| 事务消息 | 原生支持,回查机制 | 需外部协调 | 不支持 |
| 顺序消息 | 支持,实现简单 | 支持,需单分区 | 不直接支持 |
| 延时消息 | 支持 18 个级别 | 不支持 | 插件支持 |
| 消息过滤 | Tag + SQL 表达式 | 不支持 | 不支持 |
| 吞吐量 | 单机 10 万 QPS+ | 百万级 QPS | 1~2 万 QPS |
| 存储机制 | CommitLog + ConsumerQueue | 分区顺序写 + 索引 | 基于队列,内存优先 |
| 集群协调 | NameServer(极简) | ZooKeeper / Kraft | 内置分布式 |
| 运维难度 | 低(无中心节点) | 中(依赖 ZK) | 中 |
七、总结
NameServer 是核心路由:它不存储消息,不需要集群同步,通过 Broker 定期心跳维持路由表,运维时要确保 Broker 正确配置了 NameServer 地址。
CommitLog 是高性能基石:所有消息顺序写到一个文件,再异步构建索引,实现了磁盘写入的极致优化。
Queue 决定了并行度:合理设置 Queue 数量(建议 8~16),太少无法并行,太多可能造成资源浪费和 Rebalance 开销。
Consumer Group 的 Rebalance 机制:理解它才知道消费实例数设置多少合适(实例数 ≤ Queue 总数,否则有实例空闲)。
事务消息是独特优势:帮助保证本地事务和消息发送的一致性,适合分布式事务场景,但使用中要注意回查逻辑的正确实现。
顺序消息牺牲了吞吐量:非必要不用,如果只需要局部顺序,按业务 ID 路由到同一 Queue 即可。







