执行摘要
绝大多数人对消息中间件的理解停留在「发布订阅」和「削峰填谷」,但真正决定一个 MQ 能扛多少吞吐、延迟抖动有多大、堆积会不会把 broker 拖垮的,是它的存储引擎。Kafka 选择了「每个分区一个日志 + sendfile 零拷贝」,RocketMQ 选择了完全不同的另一条路:所有 Topic 共享一个 CommitLog 物理文件 + 轻量的逻辑队列索引 ConsumeQueue。
这两条路的分野,决定了它们在海量 Topic 场景下截然不同的表现。本文沿真实工程链路拆开 RocketMQ 的存储层:CommitLog 的顺序写与页缓存策略、ConsumeQueue 的 20 字节定长索引设计、IndexFile 的哈希链、延迟消息与重试队列的投递机制、事务消息的「半消息 + 回查」模型,以及 DLedger 如何用 Raft 复用 CommitLog 作为日志。每一节都附带可对照源码与生产参数。
一、先建立坐标:为什么要「所有 Topic 共享一个文件」
Kafka 的分区目录是 topic-partition/0000000000.log,一个 Broker 上若有 5000 个 Topic、每个 4 分区,就是 20000 个目录下的 20000 组文件。操作系统层面,这意味着20000 个独立的文件写句柄、20000 组页缓存竞争。当生产者的写入分散到这些分区时,原本顺序的磁盘写退化成近似随机写——机械盘上吞吐会掉一个数量级,SSD 上也会因为写放大与 GC 而抖动。
RocketMQ 的反直觉解法是:不管你有多少 Topic,Broker 上所有消息全部顺序追加到同一个 CommitLog。
CommitLog: [msg1(topicA)][msg2(topicB)][msg3(topicA)][msg4(topicC)]...
↑ 全局唯一物理offset(20位定长,左补零的文件名即起始offset)
这样磁盘写永远是纯粹的顺序写,分区数量与写性能解耦。代价是:消费侧必须有一种机制,快速从 CommitLog 里捞出「某个 Topic 某个队列的下一条消息」。这就是 ConsumeQueue。
二、CommitLog:顺序写、文件预分配与三段式落盘
2.1 MappedFileQueue 与预分配
CommitLog 由一组固定大小(默认 1GB)的 MappedFile 组成,文件名是该文件第一条消息的全局物理偏移量,用 20 位十进制左补零表示:00000000000000000000、00000000001073741824(1GB)……
// RocketMQ 中 MappedFile 的核心写入路径(简化)
public AppendMessageResult appendMessage(MessageExtBrokerInner msg, AppendMessageCallback cb) {
MappedFile file = this.mappedFileQueue.getLastMappedFile();
if (file == null || file.isFull()) {
file = this.mappedFileQueue.getLastMappedFile(msg.getTopic());
}
// 关键:写入只是往 mmap 出来的用户态地址 memcpy,不产生系统调用
return file.appendMessage(msg, cb);
}
写入热点是 mappedByteBuffer.put(...)。注意这行代码没有 syscall,它只是往页缓存映射的地址写内存,真正的落盘由内核的 writeback 线程异步完成。这是 RocketMQ 写入吞吐能到十万级 TPS 的根本原因。
但纯 mmap 写有一个隐蔽陷阱:文件刚创建时页缓存里还没有对应页,首次写入触发缺页中断,需要从磁盘读旧内容(即使是全零页也要建立映射),造成写入毛刺。RocketMQ 的处理是 warmMappedFile ——后台线程按页(4KB)逐字节写 0 预热,配合 mlock 把热文件锁在内存:
// MappedFile.warmMappedFile:预热避免缺页抖动
for (int i = 0; i < this.fileSize; i += OS_PAGE_SIZE) {
byteBuffer.put(i, (byte) 0);
}
this.mappedByteBuffer.force();
this.mlock(); // mlock 系统调用,锁定物理页不被换出
生产上这两个开关(warmMapedFileEnable、mlock)在延迟敏感场景建议打开,代价是启动期多花几十秒与常驻内存增加。
2.2 刷盘:三段式与 GroupCommit
页缓存是易失的,掉电即丢。RocketMQ 提供两种刷盘:
| 策略 | 机制 | 吞吐 | 宕机丢失 |
|---|---|---|---|
ASYNC_FLUSH | 后台线程定时 flush | 高 | 最多丢失一页缓存窗口 |
SYNC_FLUSH | 写入后同步等待刷盘 | 低 | 不丢(单节点) |
SYNC_FLUSH 的朴素实现是每条消息都 force(),那会把吞吐打到地板上。RocketMQ 的 GroupCommitService 做的是攒批:写入线程把请求放进 requestsWrite,交换后由刷盘线程一次性 force() 覆盖到最大 offset,然后统一唤醒所有等待者。这本质上是把随机 fsync 变成了顺序 fsync + 组提交,与 MySQL 的 binlog group commit 是同一个思想。
2.3 transientStorePoolEnable:把 OS 页缓存「旁路」掉
这是一个很多人开了反而变慢的参数。开启后,Broker 不再直接写 mmap 区域,而是先写堆外内存池(DirectByteBuffer),再由 CommitRealTimeService 定时 commit 到 FileChannel、由 FlushRealTimeService 刷盘:
写入线程 → 堆外内存池(write buffer) → commit → page cache → flush → disk
好处是读写隔离:消息写入不再污染页缓存,消费侧的冷读能独占页缓存,堆积消费时不会把写入打挂。坏处是多一次内存拷贝 + 多一次上下文切换。经验判断:只有在大量消费者滞后堆积、疯狂冷读历史消息的场景才值得开,正常追平消费的场景开了纯属自伤。
三、ConsumeQueue:20 字节定长索引的极致压缩
ConsumeQueue 是 CommitLog 的逻辑索引,每个 Topic@QueueId 一个目录,条目定长 20 字节:
struct ConsumeQueueEntry {
long commitLogOffset; // 8 字节:物理偏移量
int msgSize; // 4 字节:消息总长度
long tagsCode; // 8 字节:tag 的 hashCode(不是字符串!)
};
一个 1GB 的 ConsumeQueue 文件能存约 5368 万条索引,完全可以整体 mmap 进内存。消费时流程是:读 ConsumeQueue 拿到 offset+size → 拿 offset 去 CommitLog 随机读一条消息。
注意这个设计里的两个关键工程取舍:
- 不做真索引,只做顺序指针 + 定长条目。条目定长意味着可以用
offset / 20直接算出文件内位置,无需 B+ 树或哈希,O(1) 定位。 - 存 tag 的 hashCode 而不是 tag 字符串。Broker 端过滤只比对 8 字节哈希,速度极快;但哈希会冲突,所以客户端在收到消息后必须用原始 tag 字符串二次过滤。这是典型的空间/速度换正确性的权衡,也是很多人「为什么我过滤了 tag 却收到别的 tag 消息」困惑的来源。
// Broker 端 tag 过滤(ConsumeQueue 层)
public boolean isMatchedByConsumeQueue(Long tagsCode, SubscriptionData sub) {
if (tagsCode == null) return true; // 没有 tag,全通过
if (sub.getSubString().equals(SubscriptionData.SUB_ALL)) return true;
return sub.getCodeSet().contains(tagsCode); // 只比哈希码
}
// 客户端二次过滤:比对真实字符串,处理哈希冲突
public boolean isMatchedByCommitLog(ByteBuffer msgBuffer, SubscriptionData sub) {
String tag = decodeProperty(msgBuffer, MessageConst.PROPERTY_TAGS);
return sub.getSubString().equals(SUB_ALL) || sub.getTagsSet().contains(tag);
}
ConsumeQueue 是异步构建的:ReputMessageService 后台线程不断把新写入 CommitLog 的消息「派发」到各个 ConsumeQueue 与 IndexFile。这意味着 ConsumeQueue 落后于 CommitLog 是常态,Broker 重启时需要通过位点做恢复。
四、IndexFile:按 Key 查询的哈希链
除了按队列顺序消费,RocketMQ 还支持 queryMsgByKey。这依赖 IndexFile,结构是经典的槽位数组 + 链表:
Header(40B) | Slot Table(500万 × 4B) | Index Linked List(2000万 × 20B)
Index 条目 20 字节:4B KeyHash | 8B 物理offset | 4B 时间差 | 4B 上一个同槽条目序号。
slot[hash % 5000000] → 最新条目序号 → prev → prev → ... → 0(链表尾)
查询时取 slot 值沿 prev 链回溯,再按 key 原文与时间范围过滤。这是教科书式的开链哈希在磁盘上的实现,槽位数组常驻 mmap,单条查询最多几次随机读。
五、延迟消息与重试队列:不是定时器,是「分级 + 定时扫描」
RocketMQ 的延迟消息不支持任意时间,只有 18 个固定级别(1s 5s 10s 30s 1m 2m ... 2h)。实现机制很朴素:
- 生产者设置
delayLevel,Broker 不改内容,直接改写 Topic 为SCHEDULE_TOPIC_XXXX、把 queueId 改成 level-1,投递进中转 Topic; ScheduleMessageService为每个 level 起一个定时任务(1s 级每秒扫一次),扫描对应队列,到期的消息还原原始 Topic/QueueId 再重新投递一次到 CommitLog。
// 定时扫描到期延迟消息(简化)
for (int level : DELAY_LEVELS) {
long offset = offsetTable.get(level);
ConsumeQueue cq = findConsumeQueue(SCHEDULE_TOPIC, level - 1);
SelectMappedBufferResult buf = cq.getIndexBuffer(offset);
for (entry : buf) {
long deliverAt = entry.storeTimestamp + delayTime(level);
if (deliverAt <= now) {
MessageExtBrokerInner inner = rebuildFrom(entry); // 还原真实 topic
putMessage(inner); // 重新写入 CommitLog
} else { break; } // 队列按到期时间有序
}
offsetTable.put(level, newOffset);
}
消费失败重试同理:消息被投递到 %RETRY%<consumerGroup> 这个特殊的重试 Topic,按重试次数选择对应延迟级别,逐级拉长间隔(最多 16 次后进死信队列 %DLQ%<group>)。设计上非常干净——重试不需要任何额外组件,只是「换个 Topic 重新发一遍」。
六、事务消息:半消息 + 回查,而不是 2PC
RocketMQ 的事务消息常被误解为分布式事务,实际上它是「本地事务 + 两阶段提交语义 + 补偿回查」的最终一致方案:
1. Producer 发送 Half 消息(对消费者不可见)
→ Broker 改写 topic 为 RMQ_SYS_TRANS_HALF_TOPIC,落 CommitLog
2. Producer 执行本地事务(如扣库存写 DB)
3. 根据本地事务结果发 COMMIT / ROLLBACK
- COMMIT:Broker 找到 half 消息,还原 topic 投递到 ConsumeQueue(消费者可见)
并往 RMQ_SYS_TRANS_OP_HALF_TOPIC 写一条「已处理」标记
- ROLLBACK:只写 OP 标记,不投递(消息留在 CommitLog 里等过期删除)
4. 若 Producer 宕机 / 超时未决:Broker 启动回查线程,扫 HALF 队列,
按 transactionId 回调 Producer 的 checkLocalTransaction,最多回查 15 次
关键洞察:事务状态完全不修改 CommitLog 里的原始消息,只在 OP 队列里追加一条「这条 half 消息已被处理」的记录。half 与 op 两个队列做差集,剩下的就是待回查的悬挂事务。这是一个用「追加日志 + 差值」替代「原地更新状态」的典型设计——因为日志系统的原语只有 append。
// Producer 侧
TransactionListener listener = new TransactionListener() {
@Override
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
try { orderService.deductStock(arg); return LocalTransactionState.COMMIT_MESSAGE; }
catch (Exception e) { return LocalTransactionState.ROLLBACK_MESSAGE; }
}
@Override
public LocalTransactionState checkLocalTransaction(MessageExt msg) {
// 回查:必须幂等,且必须能查到「未知」的中间态
return orderService.queryTxState(msg.getTransactionId());
}
};
工程要点:回查接口必须幂等且必须处理「查不到」(返回 UNKNOW 而不是 ROLLBACK),否则本地事务已提交但查不到时会误回滚,造成数据不一致。
七、高可用:从主从复制到 DLedger 的 Raft 化
传统 RocketMQ 主从是「同步/异步双写 + haListenPort 长连接推送 CommitLog 字节」,主挂了从只能读不能自动切换,且可能丢数据。4.5 之后的 DLedger 模式把存储层换成了 Raft:
- CommitLog 直接作为 Raft log 复用:每条消息就是一条 Raft entry,主节点 append 本地后并行复制到多数派,多数派确认后 commit;
- Leader 选举后,用 Raft 的日志一致性检查做截断/补齐,天然解决了主从切换时的日志分叉;
- ConsumeQueue 不参与复制,切换后由新 Leader 从 CommitLog 重放重建(因为它是纯派生物)。
这一点非常优雅:只有一份权威数据(CommitLog),所有索引结构都是可重放的派生数据。Raft 化因此只需作用在 CommitLog 上,索引层代码零改动。
八、生产落地清单
| 旋钮 | 默认 | 生产建议 | 影响 |
|---|---|---|---|
flushDiskType | ASYNC_FLUSH | 金融级改 SYNC_FLUSH | 吞吐 vs 可靠性 |
transientStorePoolEnable | false | 仅堆积冷读场景开 | 读写隔离 vs 多一次拷贝 |
warmMapedFileEnable | false | 延迟敏感场景开 | 消除缺页毛刺 |
mappedFileSizeCommitLog | 1GB | 一般不动 | 太小则文件切换频繁 |
deleteWhen | 04 | 业务低峰 | 过期清理时机 |
fileReservedTime | 72h | 按堆积预算算 | 磁盘容量 |
useReentrantLockWhenPutMessage | false | 高冲突场景试开 | 自旋锁 vs 可重入锁 |
排障路径(按此顺序定位):
cat ~/store/config/delayOffset.json— 延迟消息是否卡在某个 level;grep "dispatch behind" storeerror.log— ConsumeQueue 构建是否跟不上 CommitLog;iostat -x 1看%util与await— 磁盘是否成为瓶颈(随机读放大来自冷消费);jstack看ReputMessageService与GroupCommitService线程状态;- 消费堆积时优先看客户端消费能力,而不是 Broker 写能力——90% 的「MQ 慢」是消费端慢。
九、结论
RocketMQ 存储层的设计可以用一句话概括:用一份全局顺序的权威日志,换取任意 Topic 规模下稳定的写入吞吐;所有索引都是可重放、可丢失、可重建的派生数据。
把这条链路串起来:
- CommitLog 共享单文件把「分区数」从写性能公式里剔除,代价是引入索引层;
- ConsumeQueue 定长 20 字节把索引压缩到可以整体 mmap,代价是 tag 只能用哈希、必须客户端二次过滤;
- 延迟消息与重试复用「改 Topic 重投」这一原语,不需要任何定时器组件;
- 事务消息用 half/op 两个队列做差集来表达「未决」,而不是原地修改状态;
- DLedger 只替换 CommitLog 的复制层,索引层零改动——因为它是纯派生。
真正值得带走的工程原则是:在一个只允许 append 的系统里,所有「状态」都应该表达为「另一个日志里的记录」,而不是「原地修改」。重试队列、事务回查、Raft 日志,本质上都是同一个模式的不同实例。理解了这一点,再看 Kafka 的 __consumer_offsets、Pulsar 的 cursor ledger、Flink 的 checkpoint manifest,会发现它们是同一套思想在不同坐标上的投影。

发表评论 取消回复