深入 RocketMQ 存储与原理(一)

发布时间:2026/7/20 15:10:19
深入 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 有几十 GBConsumeQueue 也只有它的几十分之一大小可以轻松加载到内存中消费者查找消息时几乎不会产生磁盘 I/O。 小贴士这就是 RocketMQ 的“空间换时间”策略——用一个轻量的索引文件让消息查找从 O(n) 变成了 O(1)。CommitLog 与 ConsumeQueue 的协作关系CommitLog 和 ConsumeQueue 是怎么配合工作的下面这张图展示了它们之间的完整协作关系消费端存储层生产端发送消息顺序写入消息写入成功返回偏移量异步构建索引提取消息信息提取消息信息提取消息信息根据 Queue 和偏移量返回 CommitLog 偏移量根据物理偏移量精准读取返回消息体ProducerBrokerCommitLog单一文件所有消息共用ReputMessageService后台线程ConsumeQueueTopicA-Queue0ConsumeQueueTopicA-Queue1ConsumeQueueTopicB-Queue0Consumer整个协作流程分为写入链路和消费链路两条线写入链路图中的 1→2→3→4Producer 发送消息到 BrokerBroker 将消息顺序写入 CommitLog这一步就返回 ACK 给 Producer后台线程 ReputMessageService 异步地从 CommitLog 中解析消息根据消息所属的 Topic 和 Queue将索引信息写入对应的 ConsumeQueue消费链路图中的 5→6→7→85. Consumer 根据自己的消费进度Queue 偏移量查询 ConsumeQueue6. ConsumeQueue 返回消息在 CommitLog 中的物理偏移量7. Consumer 根据物理偏移量直接从 CommitLog 读取消息内容8. CommitLog 返回完整的消息体注意写入和构建索引是异步解耦的——消息一旦写入 CommitLog 就返回成功索引的构建在后台慢慢追。这就是 RocketMQ 写入延迟极低的原因之一。消息写入 CommitLog 的完整流程一条消息从 Producer 发出到最终落盘到底经历了哪些步骤我们用一张详细的流程图来还原是否同步刷盘异步刷盘Producer 发送消息消息到达 Broker消息合法性校验Topic 是否存在 / 消息体大小是否超限消息内容准备生成消息 ID / 时间戳 / 计算 CRC是否配置了消息轨迹记录消息轨迹数据获取当前 CommitLog 文件的写入位置全局锁将消息按固定格式序列化写入 CommitLog 的追加位置更新 CommitLog 的写入指针位置刷盘策略强制将数据从 PageCache刷入物理磁盘返回写入结果给 Producer数据仅在 PageCache 中等待后台线程异步刷盘后台线程 ReputMessageService异步构建 ConsumeQueue 索引消费者后续可消费到该消息步骤解析消息到达 BrokerProducer 通过网络将消息发送到 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 的工作流程IndexFileConsumeQueueReputMessageServiceCommitLogProducerIndexFileConsumeQueueReputMessageServiceCommitLogProducer推动 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例如订单 IDBroker 在构建索引时对 Key 进行哈希运算存入 IndexFile用户通过控制台或 API 查询时输入 Key 值Broker 对 Key 进行同样的哈希运算在 IndexFile 中找到对应的消息位置再通过 CommitLog 偏移量读取完整消息 小贴士ConsumeQueue 和 IndexFile 的区别在于——ConsumeQueue 用于按队列顺序消费IndexFile 用于按业务 Key 精确查询。一个是“扫货架”一个是“查字典”。消息的物理文件布局与目录结构了解了各个文件的作用现在我们来看看它们在实际磁盘上是如何组织在一起的。RocketMQ 在 Broker 的存储目录下默认是 ~/store有这样一个完整的文件布局~/store/ Broker 存储根目录TopicA 目录下consumequeue 目录每个 queue 目录下00000000000000000000约 5.72MBconsumequeue/消息的队列索引TopicATopicBTopicCcommitlog/所有消息的物理存储commitlog 目录000000000000000000001GB000000000010737418241GB000000000021474836481GBqueue0queue1queue2queue3index/Key 查询索引checkpoint刷盘检查点文件abort异常关闭标识各文件/目录的作用目录/文件 作用commitlog/ 存储所有消息的原始数据每个文件 1GB文件名起始偏移量consumequeue/{topic}/{queueId}/ 存储 ConsumeQueue 索引每个文件约 5.72MB300000 条索引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 发送到网卡零拷贝 sendfileDMADMA 直接传输仅传递文件描述符和偏移量物理磁盘内核空间PageCache网卡0 次 CPU 拷贝由 DMA 直接完成传统读取 发送DMACPU 拷贝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 写入成功就返回 ACKSlave 异步从 Master 拉取数据优点延迟低吞吐量高缺点Master 宕机时Slave 可能缺少部分数据同步复制Master 写入后等待 Slave 也写入成功才返回 ACK优点Master 宕机时Slave 拥有完整数据缺点延迟增加吞吐量降低同步复制同步等待Slave 确认后返回 ACKMaster 写入Slave 写入Producer️ 高可靠性✅ 主从切换不丢数据异步复制数据同步直接返回 ACK不等 SlaveMaster 写入SlaveProducer⚡ 低延迟⚠️ 主从切换可能丢数据四种组合方式刷盘方式 复制方式 可靠性 吞吐量 适用场景异步刷盘 异步复制 ⭐⭐ ⭐⭐⭐⭐⭐ 日志、监控可丢少量数据异步刷盘 同步复制 ⭐⭐⭐ ⭐⭐⭐⭐ 普通业务主从切换不丢数据同步刷盘 异步复制 ⭐⭐⭐⭐ ⭐⭐⭐ 核心业务单机不丢数据同步刷盘 同步复制 ⭐⭐⭐⭐⭐ ⭐⭐ 金融级任何情况都不丢数据配置方式在 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 张核心图搞清楚了