128 lines
6.7 KiB
Markdown
128 lines
6.7 KiB
Markdown
|
|
---
|
|||
|
|
tags: [MQ, RocketMQ, 存储引擎, CommitLog]
|
|||
|
|
create time: 2026-05-24 19:52
|
|||
|
|
---
|
|||
|
|
|
|||
|
|
# RocketMQ CommitLog 三层存储模型
|
|||
|
|
|
|||
|
|
## 概述
|
|||
|
|
|
|||
|
|
RocketMQ 采用 CommitLog + ConsumeQueue + IndexFile 三层存储架构,将所有 Topic 的消息统一追加写入同一个 CommitLog 文件,再通过异步构建的索引文件实现高效消费和查询。这种设计牺牲了一定的读取局部性,但换来了极致的写入吞吐和灵活的消息查询能力。
|
|||
|
|
|
|||
|
|
## 正文
|
|||
|
|
|
|||
|
|
### 1. 三层存储模型总览
|
|||
|
|
|
|||
|
|
RocketMQ 的存储核心是三个文件的协作:
|
|||
|
|
|
|||
|
|
- **CommitLog**:消息的物理存储,所有 Topic 的消息统一追加写入。
|
|||
|
|
- **ConsumeQueue**:消费队列索引,每个 Topic-Queue 维护一个轻量级索引文件,存储消息在 CommitLog 中的位置。
|
|||
|
|
- **IndexFile**:哈希索引,支持按 MessageKey 或时间范围查询消息。
|
|||
|
|
|
|||
|
|
```mermaid
|
|||
|
|
graph TD
|
|||
|
|
Producer["Producer 发送消息"] --> CommitLog["CommitLog: 顺序追加写入"]
|
|||
|
|
CommitLog -->|"异步构建索引"| CQ["ConsumeQueue: 按 Topic-Queue 组织"]
|
|||
|
|
CommitLog -->|"异步构建索引"| Index["IndexFile: 按 MessageKey 哈希"]
|
|||
|
|
CQ -->|"顺序读取索引"| Consumer["Consumer 消费消息"]
|
|||
|
|
Index -->|"按 Key 查询"| Query["消息回溯/查询"]
|
|||
|
|
|
|||
|
|
style CommitLog fill:#4A90D9,color:#fff
|
|||
|
|
style CQ fill:#F5A623,color:#fff
|
|||
|
|
style Index fill:#6EC1E0,color:#fff
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
### 2. CommitLog:万物皆追加
|
|||
|
|
|
|||
|
|
CommitLog 是 RocketMQ 存储的灵魂。所有 Topic、所有 Queue 的消息,不分青红皂白,全部顺序追加到同一个文件(或一组分片文件)中。
|
|||
|
|
|
|||
|
|
每个 CommitLog 文件默认 1GB,写满后自动切换到下一个文件。单条消息的存储格式包括:消息长度、魔数(MagicCode)、CRC 校验、Body、Properties 等字段。
|
|||
|
|
|
|||
|
|
> [!question] 为什么 RocketMQ 要把所有 Topic 写到同一个 CommitLog,而不是像 Kafka 那样按 Partition 分开存储?
|
|||
|
|
>
|
|||
|
|
> 核心原因是**磁盘顺序写**。机械硬盘的顺序写性能接近 SSD,但随机写性能极差。Kafka 按 Partition 分文件,当 Topic 数量很多时(几千甚至上万),每个 Partition 都要维护独立的文件句柄,大量小文件同时写入会导致磁盘随机 I/O 激增。RocketMQ 把所有消息塞进同一个 CommitLog,不管有多少 Topic,写入永远是**单一文件的顺序追加**,在 Topic 数量爆炸的场景下优势明显。
|
|||
|
|
>
|
|||
|
|
> 代价是什么?Consumer 读取时需要先查 ConsumeQueue 拿到 offset,再去 CommitLog 随机读——多了一次寻址。但读操作通常由 Page Cache 命中,实际影响不大。
|
|||
|
|
|
|||
|
|
### 3. ConsumeQueue:消费的桥梁
|
|||
|
|
|
|||
|
|
ConsumeQueue 是 CommitLog 和 Consumer 之间的桥梁。每个 Topic 的每个 Queue 对应一个 ConsumeQueue 文件,每条记录固定 20 字节:
|
|||
|
|
|
|||
|
|
| 字段 | 大小 | 说明 |
|
|||
|
|
|------|------|------|
|
|||
|
|
| CommitLog Offset | 8 字节 | 消息在 CommitLog 中的物理偏移 |
|
|||
|
|
| Size | 4 字节 | 消息长度 |
|
|||
|
|
| Tag HashCode | 8 字节 | 消息 Tag 的哈希值,用于过滤 |
|
|||
|
|
|
|||
|
|
Consumer 拉取消息时,先读 ConsumeQueue 获取 offset 列表,再根据 offset 去 CommitLog 读取完整消息。由于 ConsumeQueue 的条目是定长的,可以直接通过下标计算偏移量进行随机读取,效率很高。
|
|||
|
|
|
|||
|
|
Tag 过滤也是在 ConsumeQueue 层完成的:Consumer 携带订阅的 Tag 哈希,Broker 遍历 ConsumeQueue 条目时先比对 Tag HashCode,不匹配的直接跳过,避免了读取 CommitLog 的开销。
|
|||
|
|
|
|||
|
|
### 4. IndexFile:按 Key 查消息
|
|||
|
|
|
|||
|
|
IndexFile 是可选的哈希索引文件,结构类似 HashMap:通过 MessageKey 的哈希值定位到槽位(Slot),每个槽位指向一个链表,链表节点存储消息的 CommitLog Offset 和时间戳。
|
|||
|
|
|
|||
|
|
这使得 RocketMQ 支持按 MessageKey 精确查询消息,也支持按时间范围回溯消息——在排查问题、重放历史消息时非常有用。
|
|||
|
|
|
|||
|
|
### 5. MappedFile 与 mmap
|
|||
|
|
|
|||
|
|
RocketMQ 使用 `MappedFile` 抽象封装了 mmap(内存映射文件)操作。每个 CommitLog / ConsumeQueue / IndexFile 文件都对应一个 MappedFile 实例。
|
|||
|
|
|
|||
|
|
mmap 的核心优势:将磁盘文件映射到进程虚拟内存空间,写入时直接操作内存,由操作系统内核负责异步刷盘。相比 `write()` 系统调用,mmap 减少了一次内核态到用户态的数据拷贝,写入性能更高。
|
|||
|
|
|
|||
|
|
```go
|
|||
|
|
// 简化版 CommitLog 写入逻辑
|
|||
|
|
func (cl *CommitLog) PutMessage(msg *Message) error {
|
|||
|
|
// 1. 序列化消息为字节数组
|
|||
|
|
data := serialize(msg)
|
|||
|
|
|
|||
|
|
// 2. 在 MappedFile 当前写入位置追加数据(内存映射写入)
|
|||
|
|
mappedFile := cl.mappedFileQueue.GetLastMappedFile()
|
|||
|
|
offset := mappedFile.GetCurrentWritePosition()
|
|||
|
|
mappedFile.AppendMessage(data)
|
|||
|
|
|
|||
|
|
// 3. 根据刷盘策略决定何时落盘
|
|||
|
|
if cl.flushPolicy == SyncFlush {
|
|||
|
|
mappedFile.Flush() // 同步刷盘:阻塞等待数据写入磁盘
|
|||
|
|
}
|
|||
|
|
// 异步刷盘则由后台线程定时执行,不阻塞写入
|
|||
|
|
|
|||
|
|
// 4. 异步构建 ConsumeQueue 和 IndexFile 索引
|
|||
|
|
cl.reputService.BuildIndex(msg, offset)
|
|||
|
|
|
|||
|
|
return nil
|
|||
|
|
}
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
这段代码体现了 CommitLog 写入的核心流程:序列化 -> 追加到 MappedFile -> 刷盘 -> 异步建索引。整个过程是单文件顺序写,没有锁竞争,吞吐极高。
|
|||
|
|
|
|||
|
|
### 6. 同步/异步刷盘与主从同步
|
|||
|
|
|
|||
|
|
**刷盘策略**决定了消息写入 MappedFile(内存)后何时真正持久化到磁盘:
|
|||
|
|
|
|||
|
|
- **同步刷盘(SYNC_FLUSH)**:写入后立即调用 `fsync`,数据安全性高,但吞吐较低。适合对可靠性要求极高的场景(如金融交易)。
|
|||
|
|
- **异步刷盘(ASYNC_FLUSH)**:由后台 FlushCommitLogService 线程定时刷盘,吞吐高但宕机时可能丢失少量消息。适合大多数业务场景。
|
|||
|
|
|
|||
|
|
**主从同步**方面,RocketMQ 4.x 引入了基于 Raft 协议的 **DLedger** 模式替代传统的 Master-Slave 异步复制。DLedger 要求消息写入多数节点后才算成功,从根本上解决了主从异步复制时 Master 宕机丢数据的问题。
|
|||
|
|
|
|||
|
|
```mermaid
|
|||
|
|
graph LR
|
|||
|
|
Producer["Producer"] -->|"发送消息"| Master["Master Broker"]
|
|||
|
|
Master -->|"同步/异步刷盘"| Disk1["磁盘"]
|
|||
|
|
Master -->|"DLedger Raft 复制"| Slave1["Follower 1"]
|
|||
|
|
Master -->|"DLedger Raft 复制"| Slave2["Follower 2"]
|
|||
|
|
Slave1 -->|"本地刷盘"| Disk2["磁盘"]
|
|||
|
|
Slave2 -->|"本地刷盘"| Disk3["磁盘"]
|
|||
|
|
|
|||
|
|
style Master fill:#4A90D9,color:#fff
|
|||
|
|
style Slave1 fill:#6EC1E0,color:#fff
|
|||
|
|
style Slave2 fill:#6EC1E0,color:#fff
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
## 关联笔记
|
|||
|
|
|
|||
|
|
- [[04-存储引擎/8-MQ-存储引擎设计|MQ 存储引擎设计]]
|
|||
|
|
- [[07-主流MQ对比/24-RocketMQ|RocketMQ]]
|
|||
|
|
- [[05-可靠性保障/12-MQ-消息确认与持久化|MQ 消息确认与持久化]]
|
|||
|
|
- [[04-存储引擎/9-Kafka-存储设计|Kafka 存储设计]]
|