Skip to content
Go back

RocketMQ架构——为什么CommitLog所有Topic共用一份文件

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 选主在这里是「用高复杂度换一个用不到的保证」。


Share this post on:

Previous Post
RocketMQ死信队列——16次重试后的最终兜底
Next Post
RocketMQ 延时消息原理:18 级延时与 SCHEDULE_TOPIC_XXXX 揭秘