RocketMQ 架构:三个取舍决定了一切
一句话结论(30s)
RocketMQ 高吞吐的本质是「写入顺序化 + 读取索引化 + 路由弱一致化」三个取舍换来的。因为所有 Topic 共用一份 CommitLog 顺序追加,才能把随机写变成顺序写;因为引入 ConsumeQueue 固定 20 字节二级索引,才能用一次随机读的代价换回按序消费;因为 NameServer 无主、允许路由秒级不一致,才能换来零运维复杂度。
核心原理(2min)
所有 Topic 的消息统一追加到同一个 CommitLog 文件(每 1GB 滚动新建),把多 Topic 并发随机写收敛成一次顺序写;消费时先顺序读 ConsumeQueue(每个 Queue 独立文件,固定 20B = 8B offset + 4B size + 8B tagHash),拿到 CommitLog 物理偏移后再随机读一次定位完整消息——写全顺序、读仅一次随机。NameServer 各节点互不通信、各自维护全量 Topic→Broker 路由,Broker 每 30s 上报心跳,Producer/Consumer 轮询获取路由。延时消息不单独建系统,复用同一套 CommitLog + ConsumeQueue,写入 SCHEDULE_TOPIC_XXXX(18 个 Queue 对应 18 个延迟级别),由 ScheduleMessageService 定时扫描到期转投真实 Topic。
底层深入(5-10min)
取舍一:CommitLog 单文件顺序追加
所有 Topic 的所有消息写入同一个 CommitLog 文件:
Topic A, Queue 0, Msg1
Topic B, Queue 3, Msg2
Topic A, Queue 1, Msg3
↓ 全部顺序追加到同一个文件末尾
[Msg1][Msg2][Msg3][...] (CommitLog, 最大 1GB 后新建文件)
思考·内化:先别急着往下看答案。如果让你设计 RocketMQ 的存储,你会选「每个 Topic 各写一份文件」还是「所有 Topic 共用一份文件」?停下来想 30 秒,想想磁盘此刻到底在发生什么物理动作——答案就藏在这里。
为什么不每个 Topic 一个文件? 如果 100 个 Topic 各写一个文件,并发 = 100 次随机写(磁盘磁头跳到不同位置)。CommitLog 单文件串行追加 = 1 次顺序写。HDD 顺序写比随机写快约 4 倍,SSD 下也减少了 IO 调度开销。
思考·内化:顺序写到底「快」在哪?快不在「写」这个动作本身,而在省掉了磁盘最慢的部分——HDD 的磁头寻道与盘片旋转。随机写让磁头到处跳,顺序写让磁头一路往下铺;页缓存还能成块预读、操作系统能合并相邻写请求。想通这一点,你就明白为什么后面每个设计都在死保「顺序写」这条命脉。
代价:读取时不能直接按 Topic/Queue 顺序读到相关消息——需要 ConsumeQueue 做间接索引。
顺序追加的落地源码
「所有 Topic 共用一份文件顺序写」不是一句口号,而是写死在 CommitLog#asyncPutMessage 里的:入口只取队列最后一个 MappedFile,用一把 putMessageLock 把并发写串行化,再调用 mappedFile.appendMessage 追加到末尾。
// org.apache.rocketmq.store.CommitLog#asyncPutMessage(关键片段,已省略非核心行)
public CompletableFuture<PutMessageResult> asyncPutMessage(final MessageExtBrokerInner msg) {
// Set the storage time
...
MappedFile unlockMappedFile = null;
MappedFile mappedFile = this.mappedFileQueue.getLastMappedFile();
...
putMessageLock.lock(); //spin or ReentrantLock, depending on store config
try {
...
if (null == mappedFile || mappedFile.isFull()) {
mappedFile = this.mappedFileQueue.getLastMappedFile(0); // Mark: NewFile may be cause noise
...
}
if (null == mappedFile) {
log.error("create mapped file1 error, topic: {} clientAddr: {}", msg.getTopic(), msg.getBornHostString());
return CompletableFuture.completedFuture(new PutMessageResult(PutMessageStatus.CREATE_MAPPED_FILE_FAILED, null));
}
result = mappedFile.appendMessage(msg, this.appendMessageCallback, putMessageContext);
switch (result.getStatus()) {
case PUT_OK:
onCommitLogAppend(msg, result, mappedFile);
break;
case END_OF_FILE:
onCommitLogAppend(msg, result, mappedFile);
unlockMappedFile = mappedFile;
// Create a new file, re-write the message
mappedFile = this.mappedFileQueue.getLastMappedFile(0);
...
result = mappedFile.appendMessage(msg, this.appendMessageCallback, putMessageContext);
if (AppendMessageStatus.PUT_OK.equals(result.getStatus())) {
onCommitLogAppend(msg, result, mappedFile);
}
break;
case MESSAGE_SIZE_EXCEEDED:
case PROPERTIES_SIZE_EXCEEDED:
return CompletableFuture.completedFuture(new PutMessageResult(PutMessageStatus.MESSAGE_ILLEGAL, result));
case UNKNOWN_ERROR:
default:
return CompletableFuture.completedFuture(new PutMessageResult(PutMessageStatus.UNKNOWN_ERROR, result));
}
...
} finally {
beginTimeInLock = 0;
putMessageLock.unlock();
}
...
PutMessageResult putMessageResult = new PutMessageResult(PutMessageStatus.PUT_OK, result);
...
return handleDiskFlushAndHA(putMessageResult, msg, needAckNums, needHandleHA);
}
分析:putMessageLock.lock() 把「100 个 Topic 的并发随机写」强行收敛成「同一时刻只有一条消息在追加」,这是顺序写的前提。mappedFile.appendMessage 只往 mappedFileQueue.getLastMappedFile()(当前最后一个 1GB 文件)末尾写,只有写满返回 END_OF_FILE 才滚动新建文件重写一次。真正的字节写入(切 mappedByteBuffer 定位到 wrotePosition 再 put)发生在 MappedFile 内部,详见《RocketMQ 刷盘机制》一文。
取舍二:ConsumeQueue 二级索引
CommitLog:
[Msg_A_Q0][Msg_A_Q1][Msg_B_Q0][Msg_A_Q2][Msg_B_Q1]... (按到达顺序)
ConsumeQueue (Topic A, Queue 0):
[offset=0, size=128, tagHash][offset=256, size=128, tagHash]...
每条 ConsumeQueue Entry: 固定 20B = 8B offset + 4B size + 8B tagHash
ConsumeQueue 也是顺序追加写——每个 Queue 独立文件。消费者先顺序读 ConsumeQueue(快)→ 拿到 CommitLog offset → 随机读一次 CommitLog 定位消息(只有 1 次随机 IO)。
写全顺序,读仅一次随机——以 1 次随机 IO 的代价换取所有写入的顺序 IO。
思考·内化:为什么一定要有 ConsumeQueue 这层二级索引?因为 CommitLog 里消息是按「到达顺序」混排的,Topic A 的第 3 条可能夹在 Topic B、C 中间。没有索引,消费者想拿下一条自己的消息就只能全量顺序扫描 CommitLog——写有多快,读就有多慢。于是用一层「极轻目录」(每条仅 20 字节指针)换「一次随机读定位」,这就是二级索引诞生的逻辑。
取舍三:NameServer 无主设计
NameServer 各节点互不通信,每个节点维护全量 Topic→Broker 路由表。Broker 定时向所有 NameServer 上报心跳(30s),Producer/Consumer 轮询 NameServer 获取路由。
为什么不直接用 ZooKeeper? ZooKeeper 的 Raft 选主是强一致性的保证,但 NameServer 的路由信息允许秒级不一致(某个 Broker 宕机后最长 30s 路由才更新)。用”允许短暂不一致”换取”零运维复杂度”——NameServer 可以随时重启、动态扩缩,不需要协调选举。
思考·内化:为什么 NameServer 敢「无主」?关键是想清楚「路由信息到底需要多强的一致」。如果业务能容忍最长 30 秒才感知 Broker 宕机,那 Raft 选主这套强一致就纯属「用高复杂度换一个用不到的保证」。先问「能容忍多长的不一致」,再决定「付出多少复杂度」——这是所有分布式取舍的通用问法。
延时消息:复用机制
RocketMQ 不单独实现一个”延时消息系统”——而是复用已有的 SCHEDULE_TOPIC_XXXX(18 个 Queue,对应 18 个延迟级别)。延时消息写入此 Topic,ScheduleMessageService 定时扫描到期消息转到真实 Topic。
同一套 CommitLog + ConsumeQueue 体系同时服务正常消息和延时消息,统一的存储引擎,极简的设计。
总结
| 取舍 | 得到的 | 付出的 |
|---|---|---|
| CommitLog 单文件 | 顺序 IO,写吞吐最大化 | 需要 ConsumeQueue 间接定位 |
| ConsumeQueue | 读路径清晰,消费者按序消费 | 额外 20B/条 + 顺序读取一次 |
| NameServer 无主 | 零运维,秒级发现 | 路由信息允许秒级不一致 |
章末提问
1. 为什么 RocketMQ 所有 Topic 共用一份 CommitLog?代价是什么?
结论先行:为了把「多 Topic 并发随机写」收敛成「单一文件的顺序写」,最大化写吞吐。因为磁盘顺序写远快于随机写(HDD 约 4 倍),共用一份文件 + 一把 putMessageLock 串行化,能把 100 个 Topic 的随机跳点压成一路往下写。代价是读路径无法直接按 Topic/Queue 定位,必须引入 ConsumeQueue 二级索引补回来。
2. CommitLog 的「顺序写」具体是怎么在代码里保证的?
结论先行:靠「只取最后一个 MappedFile + putMessageLock 串行化 + appendMessage 追加末尾」三步保证。因为 asyncPutMessage 入口用 mappedFileQueue.getLastMappedFile() 锁定当前最后一个 1GB 文件,putMessageLock.lock() 把并发写串行成「同一时刻只有一条消息在追加」,appendMessage 只往末尾 wrotePosition 写;写满返回 END_OF_FILE 才滚动新建文件重写一次。
3. 为什么要有 ConsumeQueue 二级索引,不能直接扫 CommitLog 吗?
结论先行:不能,因为 CommitLog 按到达顺序混排,全量扫描是 O(全量数据) 的读放大,扛不住高吞吐消费。因为 ConsumeQueue 每条固定 20 字节、按 Queue 分离,消费时用 ConsumerOffset × 20 直接算落点 O(1) 定位,再随机读一次 CommitLog 取完整消息,用「一次随机读 + 20B/条磁盘」换回按序消费。
4. NameServer 为什么不选 ZooKeeper 这样的强一致协调?
结论先行:因为路由信息允许秒级不一致,不需要为强一致付出选主与运维复杂度。因为 Broker 宕机后最长 30s 内 Producer/Consumer 靠轮询就能感知,这段时间的不一致业务可容忍;而 ZooKeeper 的 Raft 选主在这里是「用高复杂度换一个用不到的保证」。