深入 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 两个后台服务负责,清理的触发条件主要有三种:
是
否
是
是
否
否
是
否
文件清理检查
磁盘空间
是否超过阈值?
强制清理过期文件
(最先删除最老的文件)
当前文件
是否已过期?
过期时间
是否达到删除阈值?
删除过期文件
删除周期
是否到达?
跳过,保留文件
结束
具体策略:
文件过期删除:Broker 配置了 fileReservedTime(默认 72 小时),文件被创建后超过这个时间,且文件不再被写入(即不是当前活跃文件),就会被删除
磁盘空间强制删除:当磁盘使用率超过 diskMaxUsedSpaceRatio(默认 75%)时,会强制删除最老的文件,即使文件还未达到过期时间
手动触发:通过管理控制台或 API 手动触发清理
当前文件的保护:正在写入的 CommitLog 文件(当前活跃文件)不会被删除,即使它已经“过期”了。只有在文件滚动后,旧文件才会进入待删除队列。
消息的过期删除机制(文件时间戳与删除策略)
RocketMQ 的消息过期删除,本质上就是基于文件的删除,而不是基于单条消息的删除。这和很多其他消息中间件不同。
核心原理:
每个文件(CommitLog 文件或 ConsumeQueue 文件)都有一个最后修改时间戳
当后台清理线程扫描时,如果当前时间 - 文件最后修改时间 > fileReservedTime(默认 72 小时),该文件就会被删除
因为 CommitLog 中的消息是按时间顺序写入的,所以删除文件 = 删除该时间点以前的所有消息
为什么要这样设计?
基于文件的删除比基于消息的删除高效得多(直接删除文件 vs. 遍历并标记删除)
RocketMQ 假设消息在保存一段时间后,要么被消费了,要么不再需要了
这适用于大多数场景——消息是有“时效性”的,过期了就应该被清理
💡 小贴士:如果你的业务有“长期保存消息”的需求(比如审计场景),可以设置更长的 fileReservedTime,也可以将消息转存到其他存储系统(如 OSS 或 HDFS)进行长期归档。
消息的零拷贝技术原理(MMAP + FileChannel)
接下来我们聊聊 RocketMQ 的性能黑科技——零拷贝。
什么是零拷贝?
传统的文件读取 → 网络发送,数据要经历多次拷贝:
从磁盘读到内核空间(DMA 拷贝)
从内核空间读到用户空间(CPU 拷贝)
从用户空间写到 Socket 缓冲区(CPU 拷贝)
从 Socket 缓冲区写到网卡(DMA 拷贝)
—— 4 次拷贝,2 次 CPU 参与,CPU 被占用做数据搬运,效率低。
零拷贝技术是指通过操作系统的 sendfile 或 mmap 系统调用,减少数据在内核空间和用户空间之间的拷贝次数。
RocketMQ 在写入和读取两个场景分别使用了不同的零拷贝技术:
写入场景:MMAP(内存映射文件)
RocketMQ 使用 FileChannel 的 map() 方法,将 CommitLog 文件映射到操作系统的虚拟内存中(PageCache 的映射)。
RocketMQ 使用 MMAP
直接写入映射内存
自动
应用程序
内存映射区域
(用户态可见的 PageCache)
物理磁盘
零 CPU 拷贝
由 MMU 硬件完成映射
传统写入方式
写入
系统调用
异步
应用程序
用户空间缓冲区
内核空间缓冲区
PageCache
物理磁盘
需要 1 次 CPU 拷贝
通过 MMAP,应用程序可以直接操作 PageCache 中的数据,省去了从用户态拷贝到内核态的过程。写入 CommitLog 时,数据直接写入 PageCache,由操作系统异步刷盘。
读取场景:sendfile(零拷贝网络传输)
当 Consumer 拉取消息时,RocketMQ 使用 FileChannel 的 transferTo() 方法,利用操作系统的 sendfile 系统调用,直接将数据从 PageCache 发送到网卡:
零拷贝 sendfile
DMA
DMA 直接传输
仅传递文件描述符和偏移量
物理磁盘
内核空间
PageCache
网卡
0 次 CPU 拷贝
由 DMA 直接完成
传统读取 + 发送
DMA
CPU 拷贝
CPU 拷贝
DMA
物理磁盘
内核空间
用户空间
Socket 缓冲区
网卡
2 次 CPU 拷贝
RocketMQ 为什么选择零拷贝?
减少 CPU 占用:CPU 不需要做数据搬运,可以专注处理业务逻辑
提高吞吐量:省去了多余的数据拷贝,数据传输更快
更低的延迟:减少了数据在内核/用户态之间的切换
💡 小贴士:RocketMQ 的零拷贝是基于 PageCache 的。如果消息还在 PageCache 中,零拷贝直接从内存到网卡;如果消息已经被刷到磁盘,需要先读入 PageCache,再做零拷贝。这就是为什么刚生产的消息消费延迟极低(都在 PageCache 中)。
消息顺序写的性能优势(对比随机写)
我们经常听到“顺序写比随机写快”,但到底快多少?为什么快?我们用一组对比来看清楚:
顺序写入模式
初始寻道
顺序写入 4KB
顺序写入 4KB
顺序写入 4KB
顺序写入 4KB
硬盘磁头
起始位置
位置 +4KB
位置 +8KB
位置 +12KB
位置 +16KB
🚀 仅初始寻道一次
后续几乎是纯数据传输
速度可达 100-200MB/s
随机写入模式
寻道
移动磁头到位置 A写入 4KB
寻道
移动到位置 B写入 4KB
寻道
移动到位置 C写入 4KB
硬盘磁头
位置 A
位置 B
位置 C
⚡ 每次写入 = 寻道时间 + 旋转延迟 + 传输时间
💥 大量时间浪费在磁头移动上
性能差距到底有多大?
指标 随机写入 顺序写入 差距
机械硬盘吞吐量 ~1-2 MB/s ~100-200 MB/s 100 倍
SSD 吞吐量 ~50 MB/s ~500 MB/s 10 倍
IOPS(每秒操作数) ~100-200 ~10000+ 50-100 倍
RocketMQ 如何利用顺序写?
所有消息都追加到同一个 CommitLog 文件,不区分 Topic,不随机跳转
单线程顺序写入(由 putMessageLock 保证),避免多线程竞争导致乱序
批量聚合:RocketMQ 支持批量发送消息,一次写入多条,进一步放大顺序写的优势
这就是 RocketMQ 能达到 十万级 TPS 的核心原因之一。
异步刷盘与同步刷盘策略
刷盘策略解决的是 “消息什么时候写到物理磁盘” 的问题。
异步刷盘
流程:消息写入 PageCache(操作系统缓存)后立即返回 ACK,由后台线程异步将数据刷入磁盘
优点:写入延迟极低(微秒级),吞吐量极高
缺点:如果 Broker 宕机且 PageCache 中的数据尚未刷盘,数据可能丢失(但概率极低)
适用:对吞吐量要求高,容忍少量数据丢失的场景
同步刷盘
流程:消息写入 CommitLog 后,调用 fsync() 强制将数据刷入磁盘,等待刷盘完成才返回 ACK
优点:数据可靠性极高,Broker 宕机也不会丢消息
缺点:写入延迟增加(毫秒级),吞吐量下降
适用:金融、交易等一条都不能丢的场景
同步刷盘
消息写入
写入 PageCache
调用 fsync 强制刷盘
刷盘完成,返回 ACK
🛡️ 延迟:毫秒级
✅ 保证数据零丢失
异步刷盘
消息写入
写入 PageCache
立即返回 ACK
后台线程每隔 500ms
批量刷盘
物理磁盘
⚡ 延迟:微秒级
⚠️ 宕机可能丢失少量数据
配置方式:在 broker.conf 中配置 flushDiskType = ASYNC_FLUSH 或 SYNC_FLUSH。
同步复制与异步复制(主从数据同步机制)
注意:刷盘是“主节点写磁盘”,复制是“主节点往从节点同步数据”。这是两个不同的维度,可以组合使用。
异步复制
Master 写入成功就返回 ACK,Slave 异步从 Master 拉取数据
优点:延迟低,吞吐量高
缺点:Master 宕机时,Slave 可能缺少部分数据
同步复制
Master 写入后,等待 Slave 也写入成功才返回 ACK
优点:Master 宕机时,Slave 拥有完整数据
缺点:延迟增加,吞吐量降低
同步复制
同步等待
Slave 确认后
返回 ACK
Master 写入
Slave 写入
Producer
🛡️ 高可靠性
✅ 主从切换不丢数据
异步复制
数据同步
直接返回 ACK
不等 Slave
Master 写入
Slave
Producer
⚡ 低延迟
⚠️ 主从切换可能丢数据
四种组合方式
刷盘方式 复制方式 可靠性 吞吐量 适用场景
异步刷盘 异步复制 ⭐⭐ ⭐⭐⭐⭐⭐ 日志、监控,可丢少量数据
异步刷盘 同步复制 ⭐⭐⭐ ⭐⭐⭐⭐ 普通业务,主从切换不丢数据
同步刷盘 异步复制 ⭐⭐⭐⭐ ⭐⭐⭐ 核心业务,单机不丢数据
同步刷盘 同步复制 ⭐⭐⭐⭐⭐ ⭐⭐ 金融级,任何情况都不丢数据
配置方式:在 broker.conf 中配置 brokerRole = ASYNC_MASTER / SYNC_MASTER / SLAVE。
RocketMQ 高性能的四大杀手锏总结
最后,我们来总结一下 RocketMQ 高性能背后的四大核心技术:
🔪 杀手锏四:PageCache 利器
充分利用 OS 的
PageCache 缓存
热数据全在内存
读写近乎内存速度
OS 自动管理缓存
淘汰冷数据
🔪 杀手锏三:异步机制
异步刷盘
写入 PageCache 即返回
异步构建索引
ReputMessageService 后台处理
异步复制
主从同步不阻塞写入
🔪 杀手锏二:零拷贝技术
写入用 MMAP
省去用户态→内核态拷贝
读取用 sendfile
直接从 PageCache→网卡
CPU 不再做数据搬运
专注业务逻辑
🔪 杀手锏一:顺序写入
所有消息追加到
同一个 CommitLog
利用磁盘顺序写
比随机写快 100 倍
单线程写入保证
无竞争开销
一句话总结:
RocketMQ 通过 顺序写入 获得极致的写入速度,通过 MMAP + sendfile 零拷贝 获得高效的数据传输,通过 异步机制 降低链路上每一环的延迟,通过 PageCache 让热数据读写如同内存操作。四大杀手锏环环相扣,共同铸就了 RocketMQ 在双十一万亿级流量下的卓越表现。
小结
这篇我们非常硬核地深入了 RocketMQ 的存储与原理层,通过 8 张核心图,搞清楚了: