零基础搞懂Topic、Queue、Producer、Consumer…配上一张全景图,再手敲代码发送第一条消息。


一、为什么先搞懂概念?

很多人刚接触RocketMQ,看文档时被一堆术语劝退:

“Topic?Tag?MessageQueue?ConsumerGroup?它们到底什么关系?”

打个比方你就懂了:

  • Topic = 某个业务的消息分类(比如“订单Topic”)

  • MessageQueue = Topic底下的具体队列,类似多个管道并行收发

  • Tag = 同一Topic下更细的标签(比如“订单已创建” vs “订单已支付”)

  • ConsumerGroup = 一组消费逻辑相同的消费者,共同分担队列里的消息

一句话总结

Producer往Topic里发消息,Broker把消息存进具体的MessageQueue,ConsumerGroup里的消费者从Queue里拉消息,Tag用来做过滤。


二、核心角色一图胜千言

重点理解

  • NameServer:所有Producer和Consumer都从它获取Broker地址(不负责存消息)

  • Broker:真正存储消息的服务器,一个集群可以有多个Master-Slave

  • MessageQueue:每个Topic下可以有多个Queue(默认4个),分布在不同的Broker上


三、核心概念逐个击破(附中英对照)

1. Topic(主题)

  • 作用:第一级消息分类,类似数据库里的“表”

  • 创建:自动创建(autoCreateTopicEnable=true)或手动创建

  • 注意:不同业务用不同Topic,不要混用

属性
最大数量 建议几百个以内,过多影响性能
存储 分布在多个Broker上,实现水平扩展

2. Tag(标签)

  • 作用:二级过滤,减少Consumer端无用消息

  • 场景:同一个Topic里区分“订单创建”、“订单支付”、“订单取消”

  • 使用:发送时指定,消费时通过subExpression过滤

// 发送时带Tag
Message msg = new Message("OrderTopic", "pay", orderJson.getBytes());
// 消费时只订阅pay标签
consumer.subscribe("OrderTopic", "pay");

3. MessageQueue(消息队列)

  • 本质:Topic下的最小存储单元,是一个顺序写入、顺序读取的FIFO队列

  • 分布:一个Topic的多个Queue可能分布在不同的Broker上(实现负载均衡)

  • 数量影响

    • 越多→并发越高(因为消费者数量受Queue数限制)

    • 但过多会增加Rebalance开销

4. Producer(生产者)

  • 职责:创建消息并发送到Broker

  • 发送方式:同步、异步、单向(后续文章详解)

5. Consumer(消费者)

  • 职责:从Broker拉取消息并处理

  • 两种模式

    • 集群模式(默认):同一条消息只被ConsumerGroup内的一个消费者消费

    • 广播模式:同一条消息被Group内所有消费者都消费

6. ConsumerGroup(消费者组)

  • 重要规则:一个ConsumerGroup内的所有消费者 订阅的Topic必须完全一致

  • 作用:实现负载均衡和水平扩展(增加Group内消费者数,自动分担Queue)

  • 消费进度:Group维度存储offset(即消费位置)

7. NameServer

  • 轻量设计:无状态、节点间不通信、不持久化任何数据

  • 核心功能

    • 管理路由表(Topic ↔ Broker映射)

    • 接收Broker心跳(每30秒)

    • 给Producer/Consumer提供Broker地址

  • 为什么不用ZooKeeper?

    • ZooKeeper太重,适合强一致性场景

    • NameServer更简单,允许短暂不一致(生产端会重试拉取)

8. Broker

  • 核心存储节点:接收消息、写入CommitLog、构建ConsumeQueue

  • 高可用:Master-Slave模式(一主一从或一主多从)

  • 5.x新增:Broker本身只负责存储,计算逻辑上移到Proxy层


四、一条消息的生命周期(从诞生到消费)

为了帮你把概念串起来,看一条消息的完整轨迹:

  1. Producer 从 NameServer 获取 Topic=Order 的路由信息(有哪些Broker、哪些Queue)

  2. Producer 选择其中一个Queue(默认轮询),构造消息(带Tag=pay)发送到对应的Broker Master

  3. Broker Master 收到消息后:

    • 先顺序写入 CommitLog(真实数据文件)

    • 然后异步/同步构建 ConsumeQueue(索引文件,记录消息在CommitLog中的物理偏移)

  4. Consumer 从NameServer拉取路由,找到该Topic下的所有Queue

  5. Consumer 向Broker(可读Slave或Master)发起拉取请求(长轮询)

  6. Broker 根据ConsumeQueue快速定位到CommitLog中的消息,返回给Consumer

  7. Consumer 处理消息,成功后向Broker发送ACK(确认),Broker移动消费进度(offset)

思考:为什么要有ConsumeQueue?
如果没有它,每次消费都要扫描整个CommitLog,性能极差。ConsumeQueue相当于书的目录。


五、动手:发送第一条消息(体验这些概念)

环境准备(接第001篇末尾)

确保你已经启动了NameServer和Broker(5.0版本加上 --enable-proxy)。

1. Maven依赖

<dependency>
    <groupId>org.apache.rocketmq</groupId>
    <artifactId>rocketmq-client</artifactId>
    <version>5.0.0</version>
</dependency>

2. 生产者代码(发送带Tag的消息)

public class QuickStartProducer {
    public static void main(String[] args) throws Exception {
        // 1. 创建生产者,指定ConsumerGroup(这里Group只用于标识,不重要)
        DefaultMQProducer producer = new DefaultMQProducer("producer_group_demo");
        // 2. 设置NameServer地址
        producer.setNamesrvAddr("localhost:9876");
        // 3. 启动生产者
        producer.start();

        // 4. 创建消息(Topic=TestTopic,Tag=testTag)
        Message msg = new Message("TestTopic", "testTag", 
            "Hello RocketMQ".getBytes(StandardCharsets.UTF_8));
        
        // 5. 发送(同步)
        SendResult result = producer.send(msg);
        System.out.printf("发送结果: msgId=%s, sendStatus=%s, queueId=%d%n",
            result.getMsgId(), result.getSendStatus(), result.getMessageQueue().getQueueId());
        
        producer.shutdown();
    }
}

3. 消费者代码(订阅Topic和Tag)

public class QuickStartConsumer {
    public static void main(String[] args) throws Exception {
        // 1. 创建消费者,指定ConsumerGroup(重要:同一个Group内的消费者订阅必须一致)
        DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("consumer_group_demo");
        // 2. 设置NameServer
        consumer.setNamesrvAddr("localhost:9876");
        // 3. 订阅Topic和Tag(* 表示所有Tag)
        consumer.subscribe("TestTopic", "*");
        // 4. 注册消息监听器
        consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> {
            for (MessageExt msg : msgs) {
                System.out.printf("收到消息: topic=%s, tag=%s, body=%s%n",
                    msg.getTopic(), msg.getTags(), new String(msg.getBody()));
            }
            return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
        });
        // 5. 启动消费者
        consumer.start();
        System.out.println("消费者已启动");
        // 保持运行
        System.in.read();
    }
}

4. 预期输出

生产者打印类似:

发送结果: msgId=7F0000011C2018B4AAC2271A6D1C0000, sendStatus=SEND_OK, queueId=2

消费者打印:

收到消息: topic=TestTopic, tag=testTag, body=Hello RocketMQ

体验到了什么?

  • Topic: TestTopic

  • Tag: testTag

  • MessageQueue: queueId=2(说明发送到了该Topic下的第3个队列)

  • ConsumerGroup: consumer_group_demo(集群消费,如果你再启动一个相同Group的消费者,消息只会被其中一个拿到)


六、常见开发误区(新手必看)

误区 正解
以为NameServer是配置中心 NameServer只存Broker路由,不存任何业务配置
以为ConsumerGroup内的消费者订阅可以不同 订阅必须完全一致,否则会报错或导致消费混乱
以为Queue数量越多越好 Queue数应 ≤ 消费者数,否则部分Queue无消费者
以为Tag可以随意换 同一个消息的Tag在发送后不可修改,消费过滤仅对接收方有效
以为广播模式能保证顺序 广播模式下每个消费者都收到全部消息,并发消费可能导致顺序错乱

七、中英文术语速查表(收藏用)

中文 英文 说明
主题 Topic 消息分类
标签 Tag 二级过滤
消息队列 MessageQueue 实际存储单元
生产者 Producer 发消息
消费者 Consumer 收消息
消费者组 ConsumerGroup 一组消费逻辑相同的消费者
消息存储 CommitLog 真实数据文件
消费队列 ConsumeQueue 索引文件
偏移量 Offset 消费位置
重试队列 Retry Queue 消费失败后重试
死信队列 Dead Letter Queue 重试耗尽后的归宿
拉取模式 Pull Mode 消费者主动拉
推送模式 Push Mode Broker主动推(实际底层仍是拉)

八、下篇预告

下一篇我们将 手把手搭建完整的RocketMQ开发环境,包括Dashboard可视化监控、常见踩坑解决,让你本地环境丝滑运转。

思考题:
如果你设计一个聊天系统,用户A发送消息给用户B(B可能离线),应该用Topic还是Tag来区分不同会话?为什么?


本文基于 Apache RocketMQ 5.5.0,所有代码已在本地验证。
下一篇:[RocketMQ[第003篇]:手把手搭建RocketMQ开发环境(含Dashboard)]

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐