欢迎光临
我们一直在努力

【RocketMQ合集-01】RocketMQ 快速实战与核心概念学习:从安装到消息模型,搞懂阿里系分布式消息中间件

目录

一、为什么需要消息队列?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 的对比

    特性RocketMQKafkaRabbitMQ
    设计语言 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 即可。

  • 赞(0)
    未经允许不得转载:171主机测评 » 【RocketMQ合集-01】RocketMQ 快速实战与核心概念学习:从安装到消息模型,搞懂阿里系分布式消息中间件
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址