深入 RocketMQ 存储与原理
前两篇我们搞清楚了 RocketMQ 的整体架构和核心概念——知道了消息怎么流转、组件怎么协作、Topic 和 Queue 怎么回事。但从现在开始,我们要换个视角了。
如果说前两篇是站在“使用者的角度”看 RocketMQ,那么这一篇,我们要钻进 RocketMQ 的肚子里面,看看它到底是怎么存储消息的,为什么能那么快,又为什么能那么可靠。
这就好比你开一辆车,前两篇你学会了挂挡、踩油门、看后视镜;今天我们要打开引擎盖,看看发动机和变速箱是怎么工作的——只有理解了原理,你才能在遇到性能瓶颈或数据丢失风险时,知道问题出在哪,怎么去优化。
这篇文章会比较硬核,但我会尽量用通俗的语言和大量的图示,把 RocketMQ 最核心的存储机制掰开揉碎了讲给你听。准备好了吗?我们开始。
五、消息存储机制
CommitLog 的设计与原理
我们先从最核心的 CommitLog 说起。
CommitLog 是什么?
简单说,CommitLog 是 RocketMQ 存储消息的“主文件”。每个 Broker 上只有一个 CommitLog 文件(严格说是有一组按大小滚动的文件),所有 Topic 的所有消息,都顺序写入到这个文件里。
你可以把 CommitLog 想象成一本巨大的流水账本——不管是谁的消息,来了就按顺序往本子上记,先来先记,后来后记,绝不跳着写。
为什么这么设计?
这里有个计算机存储的核心知识点:磁盘顺序写入的速度,比随机写入快几十甚至上百倍。
为什么?因为机械硬盘的读写依赖于磁头移动,随机写入意味着磁头要到处“跳”,每次跳跃都需要寻道时间(平均约 5-10ms)。而顺序写入,磁头几乎不需要移动,数据像流水一样连续不断地写到磁盘上。
即使是 SSD,顺序写入也能更好地利用带宽,减少写放大效应。
RocketMQ 正是利用了这个特性,让所有消息都顺序追加到 CommitLog,从而获得了极高的写入吞吐量。
CommitLog 的文件结构
CommitLog 在磁盘上是一个文件集合,每个文件默认 1GB,文件名就是起始偏移量(用 20 位数字表示,不足补 0):
~/store/commitlog/
├── 00000000000000000000 (第1个文件,偏移量 0 开始)
├── 00000000001073741824 (第2个文件,偏移量 1GB 开始)
├── 00000000002147483648 (第3个文件,偏移量 2GB 开始)
└── …
每个消息在 CommitLog 中的存储格式如下:
字段 长度 说明
消息总长度 4 字节 整个消息条目的字节数
消息序号 8 字节 消息的唯一递增序号
存储时间戳 8 字节 消息存储时的时间戳
消息体长度 4 字节 消息体的字节数
消息体 变长 实际的消息内容
扩展属性 变长 Topic、Tag、Key 等属性
… … 其他元数据
这样的设计意味着写入 CommitLog 完全不区分 Topic,所有消息一视同仁,顺序追加——这也是 RocketMQ 写入性能极高的根本原因。
ConsumeQueue 的设计与原理
如果 CommitLog 是所有消息的大杂烩,那消费者怎么快速找到自己想要的消息呢?这就是 ConsumeQueue 的用武之地。
ConsumeQueue 是什么?
ConsumeQueue 是消息的“索引文件”,每个 MessageQueue 对应一个 ConsumeQueue 文件。
如果说 CommitLog 是“流水账本”,那 ConsumeQueue 就是“分类目录”——它不存储消息本身,只存储每条消息在 CommitLog 中的物理位置(偏移量),以及消息的大小和 Tag 的哈希值。
ConsumeQueue 的存储格式
每个 ConsumeQueue 条目固定 20 个字节,非常轻量:
字段 长度 说明
CommitLog 偏移量 8 字节 消息在 CommitLog 中的物理位置
消息长度 4 字节 消息的字节数
Tag 哈希码 8 字节 Tag 的哈希值,用于消息过滤
这意味着,即使 CommitLog 有几十 GB,ConsumeQueue 也只有它的几十分之一大小,可以轻松加载到内存中,消费者查找消息时几乎不会产生磁盘 I/O。
💡 小贴士:这就是 RocketMQ 的“空间换时间”策略——用一个轻量的索引文件,让消息查找从 O(n) 变成了 O(1)。
CommitLog 与 ConsumeQueue 的协作关系
CommitLog 和 ConsumeQueue 是怎么配合工作的?下面这张图展示了它们之间的完整协作关系:
消费端
存储层
生产端
发送消息
-
顺序写入
-
消息写入成功
返回偏移量 -
异步构建索引
-
提取消息信息
-
提取消息信息
-
提取消息信息
-
根据 Queue 和偏移量
-
返回 CommitLog 偏移量
-
根据物理偏移量
精准读取 -
返回消息体
Producer
Broker
CommitLog
单一文件,所有消息共用
ReputMessageService
后台线程
ConsumeQueue
TopicA-Queue0
ConsumeQueue
TopicA-Queue1
ConsumeQueue
TopicB-Queue0
Consumer
整个协作流程分为写入链路和消费链路两条线:
写入链路(图中的 1→2→3→4):
Producer 发送消息到 Broker
Broker 将消息顺序写入 CommitLog(这一步就返回 ACK 给 Producer)
后台线程 ReputMessageService 异步地从 CommitLog 中解析消息
根据消息所属的 Topic 和 Queue,将索引信息写入对应的 ConsumeQueue
消费链路(图中的 5→6→7→8):
5. Consumer 根据自己的消费进度(Queue 偏移量),查询 ConsumeQueue
6. ConsumeQueue 返回消息在 CommitLog 中的物理偏移量
7. Consumer 根据物理偏移量,直接从 CommitLog 读取消息内容
8. CommitLog 返回完整的消息体
注意,写入和构建索引是异步解耦的——消息一旦写入 CommitLog 就返回成功,索引的构建在后台慢慢追。这就是 RocketMQ 写入延迟极低的原因之一。
消息写入 CommitLog 的完整流程
一条消息从 Producer 发出到最终落盘,到底经历了哪些步骤?我们用一张详细的流程图来还原:
是
否
同步刷盘
异步刷盘
Producer 发送消息
-
消息到达 Broker
-
消息合法性校验
Topic 是否存在 / 消息体大小是否超限 -
消息内容准备
生成消息 ID / 时间戳 / 计算 CRC -
是否配置了
消息轨迹?
记录消息轨迹数据
-
获取当前 CommitLog 文件的
写入位置(全局锁) -
将消息按固定格式序列化
写入 CommitLog 的追加位置 -
更新 CommitLog 的
写入指针位置 -
刷盘策略?
强制将数据从 PageCache
刷入物理磁盘
- 返回写入结果给 Producer
数据仅在 PageCache 中
等待后台线程异步刷盘
- 后台线程 ReputMessageService
异步构建 ConsumeQueue 索引
消费者后续可消费到该消息
步骤解析:
消息到达 Broker:Producer 通过网络将消息发送到 Broker 的指定端口
合法性校验:检查 Topic 是否存在、消息体是否超过 4MB(默认限制)等
消息内容准备:生成全局唯一的消息 ID,记录到达时间戳,计算 CRC 校验码
消息轨迹记录(可选):如果开启了消息轨迹功能,会记录消息的发送链路信息
获取写入位置:通过全局锁(putMessageLock)获取当前 CommitLog 文件的写入偏移量——注意这里是加锁的,但因为是顺序写,锁的持有时间极短,不影响并发
序列化写入:将消息按照固定格式(魔数、消息体大小、消息体、扩展属性等)序列化为字节数组,追加到 CommitLog 文件末尾
更新指针:更新内存中的写入位置指针,为下一条消息做准备
刷盘:根据配置的刷盘策略,决定是立即刷盘还是只写到操作系统缓存(PageCache)
返回结果:将写入状态(成功/失败)和消息 ID 返回给 Producer
异步构建索引:后台线程异步构建 ConsumeQueue 索引,不影响主流程的响应速度
ConsumeQueue 的异步构建机制(ReputMessageService)
上面多次提到了 ReputMessageService,这是 RocketMQ 中一个非常关键的后台服务。我们来深入了解一下它到底是怎么工作的。
ReputMessageService 是什么?
它是 RocketMQ 中负责异步构建 ConsumeQueue 索引的后台线程服务。它像一个勤劳的“搬运工”,不断从 CommitLog 中“搬运”消息的索引信息到对应的 ConsumeQueue 中。
为什么需要异步构建?
如果每写入一条消息,就同步去更新 ConsumeQueue,会有两个问题:
增加了写入链路的延迟(本来只需要写一次 CommitLog,现在要写两次)
无法保证 ConsumeQueue 的写入也是顺序的(不同 Topic/Queue 是分散的)
异步构建意味着:写入 CommitLog 就立即返回成功,索引在后台慢慢构建。
ReputMessageService 的工作流程:
IndexFile
ConsumeQueue
ReputMessageService
CommitLog
Producer
IndexFile
ConsumeQueue
ReputMessageService
CommitLog
Producer
推动 CommitLog 的消费进度指针 (reputFromOffset)
alt
[有新数据]
[无新数据]
loop
[每毫秒轮询]
- 发送消息,写入 CommitLog
- 返回写入成功(不等待索引)
- 检查 CommitLog 是否有新数据
- 解析新消息的物理位置和属性
- 计算 ConsumeQueue 偏移量,写入索引条目
- 如果配置了 IndexFile,构建哈希索引
短暂休眠(等待下一轮)
关键点:
ReputMessageService 会记录一个 reputFromOffset 指针,标记当前已构建到 CommitLog 的哪个位置
每次轮询时,从该位置读取新数据,构建索引,然后更新指针
如果某条消息的 Topic 或 Queue 不存在,ReputMessageService 会跳过并继续,同时记录错误日志
这种设计保证了即便索引构建失败,原始消息依然安全地存储在 CommitLog 中,数据不会丢失
消息的索引文件(IndexFile)的作用
除了 ConsumeQueue 这个“主索引”,RocketMQ 还有一个 IndexFile(索引文件),用于支持根据消息 Key 查询消息。
IndexFile 是什么?
IndexFile 是一个基于哈希索引的查询文件,用于快速定位消息在 CommitLog 中的位置。
IndexFile 的结构:
索引条目结构(每条 20 字节)
Key 的哈希值
4 字节
CommitLog 偏移量
8 字节
消息存储时间差
4 字节
同哈希槽的下一条索引
4 字节
IndexFile 文件结构
文件头部
(创建时间、消息总数、哈希槽数量等)
哈希槽数组
(默认 500 万个槽位)
索引条目数组
(默认 2000 万条)
Key 查询流程:
Producer 发送消息时指定 Key(例如订单 ID)
Broker 在构建索引时,对 Key 进行哈希运算,存入 IndexFile
用户通过控制台或 API 查询时,输入 Key 值
Broker 对 Key 进行同样的哈希运算,在 IndexFile 中找到对应的消息位置
再通过 CommitLog 偏移量读取完整消息
💡 小贴士:ConsumeQueue 和 IndexFile 的区别在于——ConsumeQueue 用于按队列顺序消费,IndexFile 用于按业务 Key 精确查询。一个是“扫货架”,一个是“查字典”。
消息的物理文件布局与目录结构
了解了各个文件的作用,现在我们来看看它们在实际磁盘上是如何组织在一起的。RocketMQ 在 Broker 的存储目录下(默认是 ~/store),有这样一个完整的文件布局:
~/store/ (Broker 存储根目录)
TopicA 目录下
consumequeue 目录
每个 queue 目录下
00000000000000000000
约 5.72MB
consumequeue/
(消息的队列索引)
TopicA
TopicB
TopicC
commitlog/
(所有消息的物理存储)
commitlog 目录
00000000000000000000
1GB
00000000001073741824
1GB
00000000002147483648
1GB
queue0
queue1
queue2
queue3
index/
(Key 查询索引)
checkpoint
(刷盘检查点文件)
abort
(异常关闭标识)
各文件/目录的作用:
目录/文件 作用
commitlog/ 存储所有消息的原始数据,每个文件 1GB,文件名=起始偏移量
consumequeue/{topic}/{queueId}/ 存储 ConsumeQueue 索引,每个文件约 5.72MB(300000 条索引)
index/ 存储基于 Key 的哈希索引文件,用于消息查询
checkpoint 记录最后一次刷盘的 CommitLog 和 ConsumeQueue 位置,用于宕机恢复
abort 文件存在表示 Broker 异常关闭,启动时需要做恢复检查
消息文件的滚动与清理策略
RocketMQ 的文件不是无限增长的,它有完善的滚动和清理机制。
文件滚动策略
CommitLog:每个文件固定 1GB,写满后自动创建新文件
ConsumeQueue:每个文件固定约 5.72MB(包含 300000 条索引),写满后自动创建新文件
文件名的设计非常巧妙——用文件的起始偏移量作为文件名。这样,通过任意一个偏移量,你可以立刻算出它属于哪个文件:
偏移量 15,000,000,000 → 文件名取整 → 00000000001500000000
文件清理策略
RocketMQ 的文件清理由 CleanCommitLogService 和 CleanConsumeQueueService 两个后台服务负责,清理的触发条件主要有三种:
更多推荐




所有评论(0)