RocketMQ 存储架构:CommitLog 与 ConsumeQueue
以 RocketMQ 4.9.8 经典存储为边界,沿消息追加、索引分发、按队列读取和恢复清理解释 CommitLog、ConsumeQueue、IndexFile、内存映射与持久确认的分工。
业务看到的是 Topic 和队列,磁盘上不一定为每条队列单独保存一份消息正文。RocketMQ 的经典存储把消息正文集中追加到 CommitLog,再为 Topic 下的各个队列建立轻量索引。消费者先找到队列索引,再定位正文。这两层组织方式解释了它的写入路径,也解释了“发送成功但暂时拉不到”为什么不能直接判成消息丢失。
本文固定分析 Apache RocketMQ rocketmq-all-4.9.8 的经典存储路径,不覆盖 DLedger,也不把它等同于所有 5.x 存储实现。文中的 20 字节 ConsumeQueue 条目、追加锁和分发流程均属于这个版本边界。较新版本的批量队列、其他索引实现和复制控制模式,需要另外核对,不能沿用同一套文件细节。
正文按一条订单消息从写入到读取展开,代码使用源码结构的简化表达,具体依据链接到相应类。应用级的不丢失、重复消费和顺序恢复已经在本专栏前文讨论,这里关注它们下面的存储机制。
一、为什么正文日志与消费队列要分开
假设一个 Broker 承载订单、支付、物流等 Topic,每个 Topic 又有多条队列。消费者按队列位置读取,生产者却同时向这些队列写入。如果每条队列都维护完整正文文件,写入流量会分散到很多文件尾部,需要管理更多独立写入位置和文件生命周期。
经典 RocketMQ 选择在一个 Broker 的 CommitLog 中混合保存不同 Topic、队列的消息,形成集中追加路径。每条消息仍记录自己属于哪个 Topic 和 queueId,随后建立对应队列的索引。消费者的逻辑顺序由各自 ConsumeQueue 表达,正文存储顺序由 CommitLog 表达。
这种设计没有消灭索引写入。它让大块正文集中存储,逻辑队列只追加小条目;代价是消费者要多走一次定位,异步分发也多出可见性和恢复进度。判断架构优劣要看整个读写负载,不能只数用了几个文件。
这里的“一个 CommitLog”描述一个存储实例的经典日志组织,不是整个集群只有一条全局日志。多个 Broker 有各自存储,文件也会分段;一个机器上的追加顺序,更不能推导出整个业务的全局因果顺序。
如果两个生产者把同一订单事件无序发来,日志只会按实际收到并追加的次序保存。它不知道哪个业务事件应该先发生。业务顺序仍需要稳定路由、发送协调和消费端状态条件,详见 顺序消费。
二、先辨认四种位置
存储排障最容易把不同 offset 混在一起。CommitLog 物理位置按字节计量;ConsumeQueue 的逻辑位置按条目计量;消费者提交进度表示某个消费组在队列上恢复的位置;业务版本则属于订单自己的状态协议。
| 位置 | 作用范围 | 计量与用途 |
|---|---|---|
| CommitLog 物理 offset | 一个 Broker 的正文日志 | 字节地址,定位消息正文 |
| queueOffset | 某 Topic、queueId 的逻辑队列 | 消息逻辑位置,定位队列条目 |
| 消费组进度 | 消费组与相应队列 | 恢复依据,不是正文删除指令 |
| 订单业务版本 | 一个业务对象 | 校验状态转移和重复、乱序 |
例如订单队列的第 42 条消息,对应 CommitLog 字节位置 900000,正文长度 300。不能用 42 直接当正文地址,也不能用 900000 判断订单已经执行了九十万次。两个字段之间的映射由 ConsumeQueue 保存。
不同 Topic 或不同队列的 queueOffset 不能直接比较。队列 A 的 42 和队列 B 的 42 是两个独立位置。集群多个 Broker 的物理 offset 也各自有作用域,需要连同 Broker、Topic 和 queueId 一起解释。
拉取请求中的进度与真实业务完成也不同。SDK 可以已经把消息拉入本地缓存,业务尚未提交;业务已经提交,进度却未持久更新。底层存储告诉我们记录在哪里,不能证明外部数据库已经完成什么。
三、CommitLog 怎样追加一条完整消息
Broker 接收请求后进行相应检查,把消息编码成存储记录,包含正文及恢复、路由和检索所需的元信息。不能把 CommitLog 理解为只保存裸 JSON:读取和恢复还需要知道记录长度、所属队列及其他字段。
在 4.9.8 的 CommitLog 中,可以沿 asyncPutMessage、追加回调和刷盘处理查看实际路径。追加部分在相应锁保护下选择当前文件并写入;文件剩余空间不足时切换下一段,再追加消息。
经典单条追加路径的简化形状:
消息检查与编码
→ 进入追加临界区
→ 获取当前 MappedFile
→ 写入记录;若当前段结束,切换下一段
→ 更新追加结果与位置
→ 离开临界区
→ 按配置等待刷盘、复制等结果
锁保护共享的日志追加状态,并不意味着生产者请求所有阶段都只能单线程执行。网络接收、编码、后续等待和异步流程有各自执行路径;也不能因为看到一把追加锁,就直接推断某个固定吞吐上限。
分段文件让日志不必无限增长为一个文件。物理位置可以拆成所在段与段内偏移;前面的段保持原有记录,当前段继续追加。文件大小与预分配策略属于配置和实现细节,不应该把常见默认值当成无法更改的物理定律。
顺序追加有利于形成稳定写入路径和批量持久化,但也会受到磁盘延迟、文件分配、页缓存压力和写入竞争影响。消息越大,编码、复制、网络和磁盘占用都越大;消息条数相同,不代表存储压力相同。
追加结果中的成功也要结合后续条件理解。记录进入写入缓冲路径之后,按配置可能还需要等刷盘或副本进度,才能形成发送方看到的最终结果。把内部追加成功直接当成全部可靠性条件成立,会漏掉一层确认边界。
四、ConsumeQueue 的 20 字节里放了什么
4.9.8 经典 ConsumeQueue 定义每个条目 20 字节:8 字节正文物理 offset,4 字节正文大小,8 字节 tagsCode。它不保存一份完整消息体,而是保存读取正文所需的定位信息以及部分过滤信息。
经典 ConsumeQueue 条目,固定 20 字节:
[ physicalOffset: 8 ][ messageSize: 4 ][ tagsCode: 8 ]
逻辑 queueOffset = q
索引逻辑字节地址 = q × 20
条目给出正文物理位置和长度,再读取 CommitLog
固定长度的好处是按逻辑位置快速计算索引地址。索引自己也会分段,因此还要定位所在索引文件和段内位置;不是每次从第一条扫描到第 q 条。大量消息的索引仍占空间,但与大正文相比通常轻得多。
tagsCode 不是完整 Tag 字符串。在相应路径中可以表示 Tag 的哈希,也可能指向 ConsumeQueueExt 的扩展信息。不能把一个 8 字节字段当成所有过滤语义,更不能认为哈希永远不会碰撞。过滤最终还需要符合具体实现的再次校验。
queueOffset 持续增长,前面过期记录清理后,并不从零重新编号。读取小于最小有效位置的请求,需要处理位置过小的状态;仅有一个数值位置,不代表那段正文仍在保存。
消费组 A 和消费组 B 通常共享同一份队列数据索引,各自维护消费进度。ConsumeQueue 不是“某个消费者的私人待办表”,消费组完成一条消息不会立刻把这个共享条目和正文都删除。否则另一个消费组和历史重放就没有数据可读。
五、Reput 把正文变成可消费的索引
消息写入正文日志后,存储分发服务解析新增记录,根据 Topic、queueId 和记录中的位置建立相应索引。4.9.8 的 ReputMessageService 位于 DefaultMessageStore,可以从它的推进位置看出分发过程与追加过程有独立进度。
它按正文记录推进,正常记录分发到对应处理器,文件尾部则按实现规则进入下一段。索引不是消费者从所有正文中临时筛出来的;Broker 已经把各条逻辑队列的地址序列整理好。
因此存在一个窗口:正文已经追加,相关 ConsumeQueue 条目还未建立。消费者此刻可能暂时看不到新消息。这个窗口是否很短,要看分发速度和当前负载;它不能单靠生产者已经收到确认来推断为零。
需要观察追加端、分发端和消费端三个进度。CommitLog 一直增长,分发位置长时间停住,是索引分发问题;索引正常增长,消费组位置停住,是消费或进度问题;两者都正常而业务表缺数据,要查业务过滤、回调和提交。
分发滞后可以按字节衡量,但字节差不是消息数。大消息和小消息导致的数值差异很大;还应结合最老不可见消息的年龄理解业务影响。只看总日志大小,无法区分正常保留和异常积压。
异步派生索引让恢复时可以利用正文补齐缺失索引,但不是随意删索引文件就一定无损。恢复依赖有效正文、已有进度与正确操作流程;正文已经清理的历史区间,无法通过重新分发凭空恢复。
六、按队列拉取怎样读到正文
消费者请求某 Topic、queueId 的某个逻辑位置。Broker 先检查队列范围,再找到对应索引缓冲区,逐条取出物理位置、大小和过滤信息,通过初步过滤后读取 CommitLog 的相应区域,按返回数量、大小等限制组装结果。
相应源码可以沿 DefaultMessageStore.getMessage 查看。它区分位置过小、位置越界、没有匹配消息和正文正在移除等情况;客户端不能把所有“这次没有正文”都解释成同一种故障。
假设请求从 queueOffset=42 开始,42 的 Tag 不匹配,43 匹配,44 的正文又太大导致本批次达到限制。返回的是符合本次条件的一批记录和后续读取依据,不一定恰好是“请求位置加返回消息数”。过滤过的队列位置也要计入正确推进。
过滤分成索引阶段和正文阶段,有助于先减少不必要正文读取,但可用信息与过滤实现有关。Tag 哈希、SQL 属性过滤和扩展位图不能混成一个算法。过滤能够降低某些读成本,却不会让磁盘只保存当前订阅者需要的消息。
读取正文通常可以命中 OS 页缓存;消费历史积压时,数据可能已不在缓存,需要磁盘读取。索引访问顺序较清晰,正文却混有不同队列记录,具体磁盘访问和预读效果取决于负载。不能从“日志顺序写”推出“所有消费者都是完全连续读”。
读路径结束仍然不代表业务完成。Broker 返回字节,SDK 分发回调,应用提交数据库,再按照消费模型更新进度。每个阶段都可能失败,存储中的消息存在只能证明其中一段事实。
七、IndexFile 用于检索,不是消费位点
根据业务 Key 查某条消息,和根据队列位置连续消费,是不同访问需求。经典 IndexFile 为 Key 检索保存哈希索引及正文位置,方便定位问题;ConsumeQueue 则维护队列内的地址序列。两者不能互相替代。
IndexFile 的结构包含哈希槽与索引条目,同槽记录通过索引关联。检索得到候选正文位置后,还要校验完整消息中的实际条件,不能把哈希相等当成原始业务 Key 必然相等。
这个索引没有数据库唯一约束的语义。两条消息使用相同订单 Key,可能都是合法事件,也可能是同一事件重复发送;存储都可以接受。生产者设置 Key 不能自动实现去重,也不能阻止两次扣款。
时间范围和索引保留会影响查找。索引查不到一条记录,可能是条件不匹配、索引范围或建立进度问题,也可能是历史清理;需要与发送结果、正文和业务证据交叉判断。它是定位工具,不应成为判断退款唯一性的权威来源。
业务 Key 最好便于关联但不暴露敏感信息。事件身份和业务动作身份可以单独保存,日志记录它们的关联。不要为了方便查消息,把完整凭证或用户隐私塞进索引 Key。
八、内存映射、页缓存和刷盘分别做什么
4.9.8 的 MappedFile 使用文件映射,并维护写入、提交、刷盘位置。映射让程序通过缓冲区访问文件对应的内存页,但这些页的管理与落盘涉及操作系统,不等于把所有日志复制到 Java 堆里。
页缓存提高最近数据的读写效率。写入缓存和持久写到设备之间有不同状态;进程重启、操作系统故障、整机掉电及磁盘损坏,影响也不同。只说“写到内存了”或“已经是文件了”,无法描述可靠性边界。
开启 transient store pool 的路径还会使用额外写缓冲,提交步骤将其内容推进到文件写入路径,之后才涉及刷盘进度。不要把这里的 committedPosition 当成消息事务已经提交,也不要与消费者提交 offset 混同。名字相似,含义不同。
同步刷盘让发送确认等待相应持久化条件,异步刷盘可以按策略批量推进,通常改变延迟和故障窗口。具体结果要结合存储状态和错误处理,不是只要配置字符串写成 SYNC 就能在所有故障下绝对不丢。
系统会同时受到 Java 内存与 OS 缓存的约束。把 JVM 堆设得过大,可能挤压页缓存和系统余量;堆设得太小,又可能增加 GC 或限制消息缓存。配置需要按实际对象占用、文件访问和机器资源判断,不能给所有 Broker 一个固定比例。
热数据命中率下降后,磁盘读会与追加、刷盘、复制竞争。历史回溯读取可以成为明显的资源变化,即使业务新流量没增加。排障时要把缓存命中、磁盘等待、系统内存和读写流量放在一起看。
九、复制、刷盘与可见性是独立边界
本地写入、复制到副本、构建索引和业务消费,是不同进度。一个副本接收到数据,不代表该副本也完成了本地刷盘;本地刷盘完成,不代表其他节点已有可用副本;正文存在,不代表索引已经派生。
经典主从模式的复制机制与同步或异步确认条件,会影响故障后的数据范围。DLedger、Controller 等模式有自己的选主与安全条件,不应把经典主从的配置和恢复逻辑一概套上去。本文只把这些列为需要确认的边界,不展开其他模式的完整协议。
从 producer 看,一次超时可能发生在正文追加之后,也可能发生在等待刷盘或复制的时候。仅凭超时不能证明消息没存入;重发仍可能产生重复。相反,仅凭成功也不能保证任何机器、任何副本、任何外部业务都安全完成。消息不丢失讨论了端到端的责任链。
把这几个进度分开,才容易理解性能取舍。加快索引分发解决可见滞后,不一定减少刷盘等待;增加消费并发解决应用积压,不会修复复制停滞;扩容读流量也不能让丢失的正文重新出现。
十、重启恢复与文件清理怎样配合
日志恢复需要识别有效记录边界,处理尾部未完整写入的数据,并让相应逻辑索引与正文范围一致。校验字段和记录长度提供恢复依据,但不能把任何磁盘损坏都当成能够自动修复。恢复可能需要副本、备份或业务事件来源。
正常关闭与异常退出也有区别。后台进度、检查点和尾部状态可能不同,恢复入口会根据相应条件决定扫描与派生范围。应用不应通过手工修改文件长度或索引位置来制造看似一致的状态,操作必须遵循对应版本的维护流程。
消息清理不是“所有消费组都确认才删除”。RocketMQ 官方存储与清理说明强调保留时长与消费状态的独立关系。经典存储按文件组织管理历史数据,慢消费者可能落到已经清理的范围之前。
CommitLog 某些段删除后,相关逻辑索引的最小有效位置也需要校正。消费组仍然记着旧 offset,不表示旧正文还可读。系统返回位置过小,需要业务选择追赶可用范围、重新建立投影或从其他持久来源恢复,不能默默跳过后宣称没有丢业务。
保留期限要大于允许的消费者故障与恢复窗口,同时评估磁盘容量和清理压力。空间接近阈值时,清理和拒绝写入策略可能改变实际可读范围。不能只计算“每秒消息数乘三天”,还要考虑消息平均大小、副本、索引和资源余量。
消费完成不立即删除正文,支持多个消费组与一定范围内回溯;同时也意味着日志不是永久业务档案。重要事件的长期归档与账务事实应由相应业务系统负责,不能把 MQ 默认保留当成永久审计存储。
十一、沿一条订单消息看正常与故障路径
订单服务可靠发送 eventId=e-42,路由到 Broker B 的订单队列 3。Broker 编码后在 CommitLog 的当前位置追加,记录包含 Topic 和 queueId;相应确认条件满足后,producer 得到结果。分发服务稍后解析记录,为订单队列 3 写入新的 ConsumeQueue 条目。
消费者请求队列 3 的下一个逻辑位置,Broker 找索引,取正文,过滤后返回。应用按 eventId 与订单版本提交本地状态,再让 SDK 推进消费进度。整个过程中,物理 offset、queueOffset、消费组位置和业务版本分别描述各自阶段。
如果追加成功但 producer 响应丢失,发送器重发同一事件,正文可能多一条记录;消费者仍须幂等。如果分发服务滞后,消息暂时不可见,先看正文与分发进度,而不是立即重复补发大量新事件。如果应用数据库提交后进度丢失,恢复后重读,这是消费确认窗口。
如果消费者停机太久,原正文已经清理,重新启动只能知道旧位置不可用。此时单靠不断拉取不能恢复缺失数据,需要从事件归档、业务状态或其他可靠来源恢复。恢复投影与重做付款副作用又是不同动作,不能用一套“重放所有”处理。
这条路径也说明为什么存储问题不能只用消费者日志排查。生产者确认、Broker 存储进度、队列索引、消费组进度和业务结果,分别提供不同证据,缺一段就可能把可见滞后误诊成丢失,或者把提前确认误诊成 Broker 故障。
十二、怎样用这个模型排查问题
写入变慢,先拆网络、编码与追加等待、文件切换、刷盘和副本等待;不要先假设一定是 GC。拉取变慢,要看读取的是近期热数据还是历史冷数据、索引过滤效率和磁盘压力。消息不可见,要确认正文是否已经存在以及分发是否跟上。
积压则看队列最大可见位置与消费组持久位置,并结合消息年龄。追赶时还应估计消费数据的字节量与下游处理成本;相同积压条数可以对应完全不同的恢复时间。更多消费者只有在队列分配、应用并发和下游容量允许时才有效。
维护中不要混用不同版本的内部文件假设。先确定 Broker 版本、存储实现、复制模式和客户端消费模型,再选择对应工具与恢复流程。经典 ConsumeQueue 的 20 字节布局是理解这条路径的好入口,却不是所有 RocketMQ 存储的统一接口。
我会用两个问题检验是否理解了这套设计:消息正文和队列位置怎样对应;故障后哪一层可以由哪一层恢复。CommitLog 提供有效正文,ConsumeQueue 提供队列视图,检索索引帮助定位,消费进度决定重读范围,业务记录决定效果是否完成。它们相互配合,也各自保留明确边界。
参考资料
- RocketMQ 4.9.8 CommitLog:追加和确认路径。
- ConsumeQueue:经典条目与逻辑位置。
- DefaultMessageStore:分发、读取与恢复。
- MappedFile与IndexFile:映射缓冲及检索索引。
- RocketMQ 存储与清理:保留与消费状态的区别;版本化源码是本文文件布局的依据。
如果这篇文章对你有帮助