Skip to content
Go back

RocketMQ 消息领域模型:Topic、Tag、MessageQueue、Group 与 Offset 的层次设计

RocketMQ 消息领域模型:五层架构的完整解析

一句话结论(30s)

RocketMQ 领域模型本质是五层从抽象到具象的层次——Topic/Tag 是逻辑抽象,MessageQueue/Offset 是物理实现,ConsumerGroup 是连接两者的桥梁,因为大多数 RocketMQ 使用 bug 的根源都是对这五层理解错了。关键设计是 MessageQueue 作为最小负载均衡单位:Consumer 实例数不能超过 Queue 数,Queue 数量直接决定消费并行度。核心权衡:Offset 用 64 位看似过度工程,但换取的是永不溢出、无需像 32 位计数器那样担心溢出。

核心原理(2min)

主流程:Topic 是命名空间(物理上只是一组 ConsumeQueue 目录的路径前缀,不是独立数据结构);Tag 是 Topic 内的二级分类,每条消息最多一个 Tag(其 hashCode 存在 ConsumeQueue 的 8 字节字段里,所以不能多标签);MessageQueue 是物理队列和负载均衡单位,写队列 >= 读队列、长度无限追加;ConsumerGroup 是消费者逻辑分组,集群模式一条消息只被一个实例消费且 Offset 存 Broker,广播模式所有实例全量消费且 Offset 存本地;Offset 有三个维度——CommitLog Offset(物理字节偏移)、ConsumeQueue Offset(每个 Group 独立维护的逻辑消费序号)、Queue Max Offset(用于算 Lag)。关键机制:consumer.subscribe("OrderTopic","PAID") 背后是 OrderTopic → 8 个 MessageQueue → 8 个 ConsumeQueue 目录 → 从各 Queue 当前 Offset 读 → Tag Hash 过滤 → 返回消息。

底层深入(5-10min)

为什么需要理解领域模型?

很多开发者在用 RocketMQ 时出现 bug,根源不是代码写错了,而是对领域模型理解错了。以为 Tag 能跨 Topic、以为 Offset 是全局的、以为 Queue 是虚拟的——这些误解直接导致消息丢失、重复消费、顺序错乱。

这篇文章从抽象到具象,逐层拆解 RocketMQ 的五层领域模型。

🤔 思考穿插:为什么很多 RocketMQ 的 bug 根源不是代码写错、而是领域模型理解错?—— 因为 Topic/Tag/Queue/Group/Offset 五层各司其职,一旦混用(以为 Tag 能跨 Topic、以为 Offset 是全局、以为 Queue 是虚拟的)就会直接导致丢消息、重复消费、顺序错乱。所以先把模型层次理顺,排查问题才有方向。

五层架构总览

Layer 1: Topic        ← 消息的逻辑分类(命名空间)
Layer 2: Tag          ← 消息的二级分类(同一 Topic 内)
Layer 3: MessageQueue ← 消息的物理队列(持久化存储 + 负载均衡单位)
Layer 4: Group        ← 消费者的逻辑分组(广播/集群模式、Offset 管理)
Layer 5: Offset       ← 消费进度的物理标记(64 位长整型)

Layer 1: Topic —— 命名空间,不是队列

Topic 是消息的逻辑命名空间,用于区分不同类型的消息。

Topic: OrderTopic    ← 订单相关消息
Topic: PaymentTopic  ← 支付相关消息
Topic: LogTopic      ← 日志消息

物理实现:一个 Topic = 一组 MessageQueue

创建 Topic "OrderTopic"(配置 8 个读队列 + 8 个写队列)

物理存储:
    CommitLog(一层)→ 包含该 Topic 的所有消息
    ConsumeQueue(二层)→ 每个 Queue 一个文件:
        /consumequeue/OrderTopic/0/...
        /consumequeue/OrderTopic/1/...
        /consumequeue/OrderTopic/2/...
        ...
        /consumequeue/OrderTopic/7/...

Topic 在物理层不是独立的数据结构——它只是一个路径前缀。

🤔 思考穿插:为什么 Topic 只是「路径前缀」而不是独立数据结构?—— 因为 RocketMQ 的写盘统一落进一个 CommitLog,不按 Topic 单独建存储;Topic 只在 ConsumeQueue 里以目录名形式存在,用于消费时定位。所以「Topic 是逻辑命名空间」,物理上它没有独立的存储单元。

Topic 命名规范

✅ 推荐:
order_paid_topic           ← 下划线分隔
OrderEventTopic            ← 大驼峰

❌ 不推荐:
order-topic                ← 包含连字符(可能被 URL 编码问题坑到)
orderTopicPaidMessageTopic ← 太长(增加元数据存储开销)

动态创建 Topic

# 通过 Admin 工具创建
mqadmin updateTopic -n namesrv:9876 -t NewTopic -c DefaultCluster -r 8 -w 8

# 或生产/消费时自动创建(autoCreateTopicEnable=true,默认开启)

Layer 2: Tag —— 二级分类,不是多标签

Tag 是单个 Topic 内的消息二级分类。一条消息最多一个 Tag。

// Tag 的正确用法
Message msg = new Message("OrderTopic", "PAID", body);    // 只有 PAID
Message msg = new Message("OrderTopic", "SHIPPED", body);  // 只有 SHIPPED

// 错误:一条消息不能有多个 Tag
// msg.setTags("PAID||SHIPPED"); ← 这是 Tag 字符串为 "PAID||SHIPPED",不是两个 Tag!

Tag 的设计约束

维度说明
最多一个 Tag与 Kafka 的 key-value header 不同,RocketMQ 的 Tag 是单一值
作用域仅在所属 Topic 内有效
存储通过 hashCode 存在 ConsumeQueue 条目中(8 字节)
过滤Broker 端基于 hashCode 过滤

为什么只能有一个 Tag? 因为 Tag 的 hashCode 存在 ConsumeQueue 的 20 字节固定条目中(8 字节给 Tags Hash)。如果支持多 Tag,要么条目变大小(破坏固定长度设计),要么用 bitmap(8 字节只能表示 64 个 Tag,不够)。

🤔 思考穿插:为什么一条消息只能有一个 Tag,不能像标签一样贴多个?—— 因为 Tag 的 hashCode 要塞进 ConsumeQueue 的 20 字节固定条目里(8 字节给 Tag Hash),定长设计支撑了 O(1) 过滤;若支持多 Tag,条目要么变长破坏定长、要么用 bitmap(8 字节只能表示 64 个)不够用。所以单 Tag 是存储定长约束下的取舍。


Layer 3: MessageQueue —— 物理队列,最小负载均衡单位

MessageQueue 是消息落地的物理单位。 这是整个模型中最关键的概念。

Queue 的数量决定了并行度

writeQueueNums = 8, readQueueNums = 8

物理上创建 8 个 ConsumeQueue:
  /consumequeue/OrderTopic/0/
  /consumequeue/OrderTopic/1/
  ...
  /consumequeue/OrderTopic/7/

并行度:
  - 生产端:8 个 Queue,支持 8 个并发写通道(实际上 CommitLog 统一写,但路由过程可以并行)
  - 消费端:一个 Consumer Group 最多 8 个 Consumer 实例并行消费
           (第 9 个 Consumer 接不到 Queue,空转)

Queue 的三个核心约束

约束 1:写队列 >= 读队列
    writeQueueNums >= readQueueNums ← 永远成立
    原因:写队列是"完整的",读队列是写的"快照"
          如果读 > 写,有些 Queue 永远没消息

约束 2:读写数量可以不一致(扩容时)
    writeQueueNums = 8, readQueueNums = 4
    → 新消息写到 8 个 Queue,只有 4 个 Queue 被消费
    → 用于缩容过渡:先把读缩小,等消费者处理完旧消息,再缩小写

约束 3:Queue 长度理论无限
    Queue 不是 Ring Buffer,是无限追加的 ConsumeQueue 文件链。
    不设上限,只受磁盘容量限制。

Queue 的物理实现:ConsumeQueue 文件链

Queue-0 的 ConsumeQueue:
  /consumequeue/OrderTopic/0/00000000000000000000  (0-300,000 条消息的索引)
  /consumequeue/OrderTopic/0/00000000000006000000  (300,001-600,000 条消息的索引)
  /consumequeue/OrderTopic/0/00000000000012000000  (600,001-900,000 条消息的索引)
  ...

每个 ConsumeQueue 文件 300000 × 20 bytes ≈ 5.72 MB。Queue 的总大小 = 文件数量 × 5.72 MB,随着消息持续写入,文件不断增加。

Queue 选择与路由

// 生产者:选择 Queue
// 1. 不指定 → 默认轮询
producer.send(msg);  // 自动轮询:Q0 → Q1 → Q2 → ... → Q7 → Q0

// 2. 指定 MessageQueueSelector
producer.send(msg, new MessageQueueSelector() {
    @Override
    public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
        long orderId = (Long) arg;
        return mqs.get((int) (orderId % mqs.size()));  // 同订单发同 Queue
    }
}, orderId);

// 3. 指定 Queue
producer.send(msg, new MessageQueue("OrderTopic", "broker-a", 3));  // 直接指定 Queue-3

路由算法的核心约束:同一个业务标识必须路由到同一个 Queue(当需要顺序消费时)。


Layer 4: ConsumerGroup —— 消费的逻辑分组

ConsumerGroup 是一组 Consumer 实例的逻辑集合,用于负载均衡和消费进度管理。

集群模式(CLUSTERING)

ConsumerGroup: "OrderConsumerGroup"
  ├── Instance-A (IP: 10.0.1.1)
  ├── Instance-B (IP: 10.0.1.2)
  └── Instance-C (IP: 10.0.1.3)

Topic "OrderTopic" 有 8 个 Queue:
  Instance-A → Queue 0, 1, 2
  Instance-B → Queue 3, 4, 5
  Instance-C → Queue 6, 7

特点:
  - 每条消息只被 Group 中的一个实例消费
  - Rebalance 时重新分配
  - Offset 以 ConsumerGroup 为单位存储在 Broker 上

广播模式(BROADCASTING)

ConsumerGroup: "OrderNotifyGroup"
  ├── Instance-A
  ├── Instance-B
  └── Instance-C

Topic "OrderTopic" 有 8 个 Queue:
  Instance-A → Queue 0, 1, 2, 3, 4, 5, 6, 7  ← 全量
  Instance-B → Queue 0, 1, 2, 3, 4, 5, 6, 7  ← 全量
  Instance-C → Queue 0, 1, 2, 3, 4, 5, 6, 7  ← 全量

特点:
  - 每条消息被 Group 中的每个实例都消费到
  - 不需要 Rebalance
  - Offset 存储在消费者本地

订阅关系约束

同一个 ConsumerGroup 内,所有 Consumer 的订阅关系必须完全一致:
  ✅ Consumer1: Topic-A(Tag1 || Tag2), Topic-B(Tag1)
     Consumer2: Topic-A(Tag1 || Tag2), Topic-B(Tag1)
     → 完全一致,OK

  ❌ Consumer1: Topic-A(Tag1)
     Consumer2: Topic-A(Tag2)
     → 不一致!Rebalance 时会出错

Layer 5: Offset —— 64 位的”指针”

Offset 是 RocketMQ 中含义最多的概念,有三个不同的 Offset:

Offset 的三个维度

1. CommitLog Offset
   - 消息在 CommitLog 文件中的物理偏移量
   - 全局唯一,由 Broker 在写入时分配
   - 8 字节,单调递增(每条消息 +1 条消息的长度)
   - 物理含义:文件中的字节偏移

2. ConsumeQueue Offset
   - 消费者在 ConsumeQueue 中的逻辑偏移量
   - 每个 ConsumerGroup 独立维护
   - Offset = N → 表示"该 Group 已经消费了该 Queue 的前 N 条消息"
   - 逻辑含义:消费序号

3. Queue Max Offset
   - 当前 Queue(ConsumeQueue)中最大的逻辑序号
   - 表示该 Queue 当前有多少条消息
   - Lag = MaxOffset - ConsumerOffset

Consumer Offset 总结

消费进度 = ConsumerGroup + Topic + Queue + ConsumerOffset

存储位置:
  集群模式 → Broker 端(_consumer_offset.json / 内存)
  广播模式 → Consumer 本地文件

Offset 的意义:
  新建 Group 首次消费 → Offset = CONSUME_FROM_LAST_OFFSET(跳过历史消息)
                       或 CONSUME_FROM_FIRST_OFFSET(从第一条开始)
                       或 CONSUME_FROM_TIMESTAMP(从指定时间开始)

64 位 Offset 100 年不溢出

Offset 是 64 位长整型(long,Java 中 2^63-1 = 9,223,372,036,854,775,807)

假设每秒处理 100 万条消息(百万 TPS):
  9,223,372,036,854,775,807 / (1,000,000 × 86,400 × 365)
  ≈ 292,471 年

即使扩大到每秒 100 亿条:
  ≈ 29 年

结论:在任何现实场景下,Offset 永远不会溢出。

Offset 的 64 位设计”过度工程”了——但它保持了简单,不需要像 ZooKeeper znode 的 32 位计数器那样担心溢出。宁可浪费比特,不要引入溢出风险。 这是 RocketMQ 的一个隐式设计哲学。

🤔 思考穿插:为什么 Offset 用 64 位看起来「过度工程」,RocketMQ 却坚持?—— 因为 64 位长整型在任何现实负载下都不可能溢出(百万 TPS 也能撑近 30 万年),换来的是「永远不用处理溢出」的简单性;而 32 位计数器(如 ZK znode)迟早要面对溢出回绕。所以「宁可浪费比特,不要引入溢出风险」。


模型与物理存储的映射

领域模型                          物理存储
────────────────────────────────────────────────────
Topic "OrderTopic"       →       /consumequeue/OrderTopic/
  └─ Queue-0             →         └─ 0/
  └─ Queue-1             →         └─ 1/
  └─ Queue-2             →         └─ 2/
  └─ Queue-7             →         └─ 7/
                                     (每个目录 N 个 5.72MB 文件,文件链持续增长)

ConsumerGroup            →       Broker 内存中维护
"OrderConsumerGroup"               消费进度:
                                  OrderTopic/Q0 → Offset 150000
                                  OrderTopic/Q1 → Offset 148000
                                  ...

一条消息的完整路径:
  Producer → CommitLog (offset=12345) → ReputService → ConsumeQueue
          → OrderTopic/Q3/00000000 (索引条目 #5000)
          → Consumer 拉取 → Offset 5000 → 定位 CommitLog offset=12345 → 读消息

常见误解与纠正

误解正确理解
“Topic 就是队列”Topic 是命名空间,物理存储单位是 MessageQueue
“多个 Tag 可以像标签一样贴”每条消息只有一个 Tag,多条件用 SQL92
“Queue 长度有限”无限追加,只受磁盘限制
“Offset 会溢出”64 位,任何实际负载下都不会溢出
“Consumer 越多越好”Consumer 实例数不能超过 Queue 数量,多余的 Consumer 空转
“广播模式 Offset 在 Broker”广播模式 Offset 在 Consumer 本地文件
“Tag 可以跨 Topic”Tag 完全在 Topic 作用域内,不同 Topic 的相同 Tag 值没有任何关联

总结

RocketMQ 的五层领域模型是一个从抽象到具象的层次结构:

Topic(逻辑命名空间)
  └── Tag(消息分类标签)
       └── MessageQueue(物理存储单位,负载均衡粒度)
            └── ConsumerGroup(消费者逻辑分组)
                 └── Offset(消费进度 64 位指针)

理解这个模型的关键在于:Topic 和 Tag 是逻辑抽象(为了开发者理解),MessageQueue 和 Offset 是物理实现(为了存储和分发),ConsumerGroup 是这两者之间的桥梁(把物理的 Queue 分配给逻辑的消费者)。

当你写 consumer.subscribe("OrderTopic", "PAID") 时,RocketMQ 在背后做的映射是:OrderTopic → 8 个 MessageQueue → 定位到 8 个 ConsumeQueue 文件目录 → 从每个 Queue 的当前 Offset 开始读取 → 用 Tag Hash 过滤 → 返回消息。这个链路如果理解透了,排查任何 RocketMQ 问题都有方向。

章末提问

追问 1:为什么说「Topic 不是队列」?

结论先行:因为 Topic 是逻辑命名空间,物理存储单位是 MessageQueue(即 ConsumeQueue 文件链)。

因为:RocketMQ 所有消息统一写进一个 CommitLog,Topic 只是 ConsumeQueue 目录名的路径前缀、不是独立的数据结构;真正承载「队列」语义、负责负载均衡的是每个 Queue 对应的 ConsumeQueue 文件链。所以「Topic 是逻辑抽象,MessageQueue 才是物理队列」。

追问 2:为什么 Consumer 实例数不能超过 Queue 数?

结论先行:因为 MessageQueue 是最小负载均衡单位,Rebalance 时每个 Queue 最多分给一个实例,实例数超过 Queue 数时多出来的实例只能空转。

因为:消费并行度的上限 = Queue 数,每个 Queue 同一时刻只被一个实例消费;多余的实例拿不到队列、不产生吞吐。所以规划 Topic 时 Queue 数要 >= 预期的 Consumer 实例数,扩容 Consumer 前也要先确认 Queue 是否够。

追问 3:RocketMQ 里有哪几种 Offset?各自的含义是什么?

结论先行:有三种——CommitLog Offset(物理字节偏移)、ConsumeQueue Offset(每组独立的逻辑消费序号)、Queue Max Offset(队列最大序号,用于算 Lag)。

因为:CommitLog Offset 全局唯一、由 Broker 写入时按字节分配;ConsumeQueue Offset 是每个 ConsumerGroup 独立维护的消费进度、表示「已消费前 N 条」;Queue Max Offset 表示该队列当前总条数,Lag = MaxOffset - ConsumerOffset 是判断积压的关键指标。三者分别对应「存到哪、消费到哪、还剩多少」。


Share this post on:

Previous Post
RocketMQ 顺序消费:从队列路由到线程绑定的完整链路
Next Post
RocketMQ 消息过滤:Tag 过滤、SQL92 表达式与过滤原理