news 2026/9/20 14:57:40

RocketMQ 消息发送存储流程源码剖析:从 DefaultMessageStore 到 CommitLog 的写入链路

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
RocketMQ 消息发送存储流程源码剖析:从 DefaultMessageStore 到 CommitLog 的写入链路
  • 文档
  • 教程
  • 知识库

【免费下载链接】source-code-hunter

😱 从源码层面,剖析挖掘互联网行业主流技术的底层实现原理,为广大开发者 “提升技术深度” 提供便利。目前开放 Spring 全家桶,Mybatis、Netty、Dubbo 框架,及 Redis、Tomcat 中间件等

项目地址:https://gitcode.com/doocs/source-code-hunter
点击查看免费下载

本文基于 RocketMQ 4.9.3 源码,以org.apache.rocketmq.store.DefaultMessageStore#putMessage为核心入口,完整剖析一条消息从 Broker 接收、状态校验、消息校验、CommitLog 文件定位到内存映射文件写入的完整存储链路。读完本文,你将掌握 Broker 端消息写入的四步主流程、PutMessageStatus各状态码的触发条件、CommitLog 文件的创建与滚动机制,以及事务消息在存储层的特殊处理逻辑,并能在生产环境中据此快速定位"消息写入失败"类问题的根因。

一、存储链路全景:消息进入 Broker 后的第一站

在 RocketMQ 消息发送流程 中,Producer 通过SendMessageProcessor将消息请求发送到 Broker 后,Broker 端真正落盘的核心入口是org.apache.rocketmq.store.DefaultMessageStore#putMessage。这条链路共分四步:

  1. 检查消息存储状态checkStoreStatus):Broker 是否可用、角色是否为主、存储是否可写、pageCache 是否繁忙;
  2. 检查消息合法性checkMessage):主题长度、属性长度是否超限;
  3. 获取当前可写的 CommitLog 文件:通过MappedFileQueue.getLastMappedFile定位或创建MappedFile
  4. 将消息写入 MappedFileCommitLog#asyncPutMessage内部调用MappedFile.appendMessagesInner,最终由DefaultAppendMessageCallback#doAppend完成字节级写入。

其中,步骤 3、4 是整个存储流程的性能核心:CommitLog 采用**顺序写 + 内存映射(mmap)**的方式,这是 RocketMQ 单机吞吐量远高于传统随机写数据库的关键设计。相关文件模型可参考 RocketMQ CommitLog 详解 与 RocketMQ MappedFile 内存映射文件详解。

二、第一步:检查消息存储状态(checkStoreStatus)

org.apache.rocketmq.store.DefaultMessageStore#checkStoreStatus在消息真正写入前做四道"闸门"检查,任何一道不通过都会直接拒绝本次写入。

1. 检查 Broker 是否已关闭

if (this.shutdown) { log.warn("message store has shutdown, so putMessage is forbidden"); return PutMessageStatus.SERVICE_NOT_AVAILABLE; }

shutdown为 true 表示 Broker 正在执行优雅停机流程,此时存储组件已停止对外服务,消息写入直接被拒绝。

2. 检查 Broker 的角色

if (BrokerRole.SLAVE == this.messageStoreConfig.getBrokerRole()) { long value = this.printTimes.getAndIncrement(); if ((value % 50000) == 0) { log.warn("broke role is slave, so putMessage is forbidden"); } return PutMessageStatus.SERVICE_NOT_AVAILABLE; }

如果当前 Broker 的角色是 SLAVE(从节点),则禁止写入消息。注意这里通过printTimes.getAndIncrement()配合value % 50000 == 0日志节流——只有每 5 万次拒绝才打印一次告警,避免高频拒绝场景下刷爆日志。写入只允许发生在主节点(SYNC_MASTER / ASYNC_MASTER),从节点仅负责同步主节点的数据并对外提供读取。

3. 检查 messageStore 是否可写

if (!this.runningFlags.isWriteable()) { long value = this.printTimes.getAndIncrement(); if ((value % 50000) == 0) { log.warn("the message store is not writable. It may be caused by one of the following reasons: " + "the broker's disk is full, write to logic queue error, write to index file error, etc"); } return PutMessageStatus.SERVICE_NOT_AVAILABLE; } else { this.printTimes.set(0); }

runningFlags是 Broker 存储层的运行状态位,当其被置为不可写时,通常意味着以下三类问题之一:

  • Broker 磁盘空间不足(写 CommitLog 物理文件失败);
  • 写逻辑队列(ConsumeQueue)出错
  • 写索引文件(IndexFile)出错

一旦状态恢复正常,printTimes会被重置为 0,日志节流计数也随之清零。

4. 检查 OS pageCache 是否繁忙

if (this.isOSPageCacheBusy()) { return PutMessageStatus.OS_PAGECACHE_BUSY; }

isOSPageCacheBusy()通过StoreStatsService统计的"写入数据量/耗时"计算写盘速度,若判定操作系统 pageCache 刷盘压力过大,则返回OS_PAGECACHE_BUSY,让上层(SendMessageProcessor)据此决定是否触发 Broker 端的快速失败或重试逻辑,从而避免内存写入速度被磁盘 IO 拖垮、导致消息大量堆积在堆外内存。

小结checkStoreStatus四个检查点中,前三个失败均返回SERVICE_NOT_AVAILABLE,只有 pageCache 繁忙返回独立的OS_PAGECACHE_BUSY状态码——这也是排查"消息发送超时/失败"问题时,Broker 端日志与返回状态码的首要对照依据。

三、第二步:检查消息合法性(checkMessage)

通过存储状态检查后,DefaultMessageStore#checkMessage对消息本身做两项合法性校验。

1. 校验主题长度不超过 127

if (msg.getTopic().length() > Byte.MAX_VALUE) { log.warn("putMessage message topic length too long " + msg.getTopic().length()); return PutMessageStatus.MESSAGE_ILLEGAL; }

主题长度上限为Byte.MAX_VALUE,即 127 个字符。这与客户端侧Validators.checkTopic的校验(见 RocketMQ 消息发送流程)遥相呼应,形成客户端预检 + 服务端兜底的双重防线。

2. 校验属性长度不超过 32767

if (msg.getPropertiesString() != null && msg.getPropertiesString().length() > Short.MAX_VALUE) { log.warn("putMessage message properties length too long " + msg.getPropertiesString().length()); return PutMessageStatus.MESSAGE_ILLEGAL; }

消息属性(properties 序列化后的字符串)长度上限为Short.MAX_VALUE,即 32767 字节。注意该上限针对的是属性字符串而非消息体——消息体的大小限制(默认 4MB,maxMessageSize)在客户端DefaultMQProducerImpl#checkMessage中校验。

两项校验失败均返回MESSAGE_ILLEGAL,Broker 端会拒绝该消息并记录告警日志。

四、第三步:获取当前可写入的 CommitLog 文件

通过校验后,CommitLog#asyncPutMessage开始定位本次写入的目标物理文件。CommitLog 文件的存储目录为${ROCKET_HOME}/store/commitlog,其中:

  • MappedFileQueue对应整个commitlog文件夹,负责管理该目录下的所有文件;
  • MappedFile对应文件夹下的单个文件,每个文件默认大小 1GB(可通过 broker 配置mappedFileSizeCommitLog调整),文件名即该文件的起始物理偏移量。

定位逻辑如下:

msg.setStoreTimestamp(beginLockTimestamp); 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: " + msg.getTopic() + " clientAddr: " + msg.getBornHostString()); return CompletableFuture.completedFuture(new PutMessageResult(PutMessageStatus.CREATE_MAPEDFILE_FAILED, null)); }

mappedFile.isFull()的判断标准是fileSize == wrotePosition.get(),即写入指针已触及文件末尾(详见 MappedFile 详解 中的isFull())。如果当前文件已写满,则通过getLastMappedFile滚动到下一个文件;如果创建失败,返回CREATE_MAPEDFILE_FAILED

文件滚动与创建:getLastMappedFile

public MappedFile getLastMappedFile(final long startOffset, boolean needCreate) { long createOffset = -1; MappedFile mappedFileLast = getLastMappedFile(); if (mappedFileLast == null) { createOffset = startOffset - (startOffset % this.mappedFileSize); } if (mappedFileLast != null && mappedFileLast.isFull()) { createOffset = mappedFileLast.getFileFromOffset() + this.mappedFileSize; } if (createOffset != -1 && needCreate) { return tryCreateMappedFile(createOffset); } return mappedFileLast; }

这里有两种触发新文件创建的场景:

  1. 首次写入mappedFileLast == null,此时createOffset = startOffset - (startOffset % mappedFileSize),即向下对齐到文件大小的整数倍——这正是 CommitLog 文件"以起始偏移量命名"的由来;
  2. 最新文件已写满createOffset = mappedFileLast.getFileFromOffset() + mappedFileSize,即上一个文件起始偏移量加上单个文件大小。

tryCreateMappedFile内部通过MappedFile.init完成文件初始化:以RandomAccessFile(file, "rw")打开文件通道,再调用fileChannel.map(MapMode.READ_WRITE, 0, fileSize)将整个文件映射到 JVM 虚拟内存,并将fileFromOffset解析为文件名对应的 long 值。该初始化过程(含ensureDirOK目录创建、TOTAL_MAPPED_VIRTUAL_MEMORY/TOTAL_MAPPED_FILES统计)的完整源码见 RocketMQ MappedFile 内存映射文件详解。

五、第四步:将消息写入 MappedFile

拿到目标MappedFile后,写入动作由MappedFile#appendMessagesInner完成:

public AppendMessageResult appendMessagesInner(final MessageExt messageExt, final AppendMessageCallback cb, PutMessageContext putMessageContext) { assert messageExt != null; assert cb != null; int currentPos = this.wrotePosition.get(); if (currentPos < this.fileSize) { ByteBuffer byteBuffer = writeBuffer != null ? writeBuffer.slice() : this.mappedByteBuffer.slice(); byteBuffer.position(currentPos); AppendMessageResult result; if (messageExt instanceof MessageExtBrokerInner) { result = cb.doAppend(this.getFileFromOffset(), byteBuffer, this.fileSize - currentPos, (MessageExtBrokerInner) messageExt, putMessageContext); } else if (messageExt instanceof MessageExtBatch) { result = cb.doAppend(this.getFileFromOffset(), byteBuffer, this.fileSize - currentPos, (MessageExtBatch) messageExt, putMessageContext); } else { return new AppendMessageResult(AppendMessageStatus.UNKNOWN_ERROR); } this.wrotePosition.addAndGet(result.getWroteBytes()); this.storeTimestamp = result.getStoreTimestamp(); return result; } log.error("MappedFile.appendMessage return null, wrotePosition: {} fileSize: {}", currentPos, this.fileSize); return new AppendMessageResult(AppendMessageStatus.UNKNOWN_ERROR); }

几个关键点:

  • 写入位置:以wrotePosition(AtomicInteger 写指针)作为本次写入的起始位置,保证并发安全;
  • 写入介质:若启用了TransientStorePool(堆外内存暂存池,即writeBuffer不为空),则先写入堆外writeBuffer的共享缓冲区切片,后续由 commit 阶段批量刷入 FileChannel;否则直接写入mappedByteBuffer(mmap 映射的直接内存)。两种路径的差异决定了刷盘模型(详见下文"刷盘与提交");
  • 消息类型分支MessageExtBrokerInner(单条消息)与MessageExtBatch(批量消息)走不同的doAppend重载;其他类型返回UNKNOWN_ERROR
  • 写入结果:追加完成后累加wrotePosition,并记录storeTimestamp

核心写入:DefaultAppendMessageCallback#doAppend

org.apache.rocketmq.store.CommitLog.DefaultAppendMessageCallback#doAppend是字节级写入的最终实现,包含三个关键动作。

动作一:计算写入偏移量

long wroteOffset = fileFromOffset + byteBuffer.position();

wroteOffset= 文件起始偏移量 + 文件内当前位置,这就是该消息在全局物理存储中的偏移量(committed offset 概念的基础)。

动作二:对事务消息做偏移量特殊处理

final int tranType = MessageSysFlag.getTransactionValue(msgInner.getSysFlag()); switch (tranType) { // Prepared and Rollback message is not consumed, will not enter the // consumer queue case MessageSysFlag.TRANSACTION_PREPARED_TYPE: case MessageSysFlag.TRANSACTION_ROLLBACK_TYPE: queueOffset = 0L; break; case MessageSysFlag.TRANSACTION_NOT_TYPE: case MessageSysFlag.TRANSACTION_COMMIT_TYPE: default: break; }

PREPARED(半消息)和 ROLLBACK(回滚)状态的消息不会被消费者消费,因此其在队列中的逻辑偏移量queueOffset直接置 0,后续也不会写入 ConsumeQueue(消费队列),从而避免半消息/回滚消息进入消费链路。只有 NOT(普通消息)和 COMMIT(已提交事务消息)才参与队列偏移量分配。

动作三:构造 AppendMessageResult

AppendMessageResult result = new AppendMessageResult(AppendMessageStatus.PUT_OK, wroteOffset, msgLen, msgIdSupplier, msgInner.getStoreTimestamp(), queueOffset, CommitLog.this.defaultMessageStore.now() - beginTimeMills);

AppendMessageResult封装了本次写入的完整结果信息:

字段含义
PUT_OK写入成功状态码
wroteOffset消息的全局物理偏移量
msgLen消息在 CommitLog 中占用的总长度(含头信息)
msgIdSupplier消息 ID(基于物理偏移量与进程地址生成)
storeTimestamp消息存储时间戳
queueOffset消息在逻辑队列中的偏移量
now() - beginTimeMills本次存储耗时

事务消息的队列偏移量更新

doAppend 的最后,根据事务类型决定是否推进主题队列的偏移量表:

switch (tranType) { case MessageSysFlag.TRANSACTION_PREPARED_TYPE: case MessageSysFlag.TRANSACTION_ROLLBACK_TYPE: break; case MessageSysFlag.TRANSACTION_NOT_TYPE: case MessageSysFlag.TRANSACTION_COMMIT_TYPE: // The next update ConsumeQueue information CommitLog.this.topicQueueTable.put(key, ++queueOffset); CommitLog.this.multiDispatch.updateMultiQueueOffset(msgInner); break; default: break; }
  • PREPARED / ROLLBACK:不更新topicQueueTable,不参与消费队列构建;
  • NOT / COMMIT:将topicQueueTable中该主题队列的偏移量自增,并调用multiDispatch.updateMultiQueueOffset(msgInner)更新多队列派发信息。

topicQueueTableConcurrentMap<String, Long>)维护着"主题@队列"与逻辑偏移量的映射,它是后续构建 ConsumeQueue 的核心数据源。ConsumeQueue 中每条记录只存 20 字节(8 字节 CommitLog 物理偏移量 + 4 字节消息长度 + 8 字节 tag 哈希码),具体格式与查询逻辑见 RocketMQ ConsumeQueue 详解。

六、写入完成后的下行链路:提交、刷盘与派发

putMessage返回PutMessageResult后,存储流程并未结束,剩余工作由 CommitLog 内部及后台线程接力完成:

  1. 提交(commit)与刷盘(flush):根据刷盘策略(同步刷盘SYNC_FLUSH/ 异步刷盘ASYNC_FLUSH),由FlushConsumeQueueServiceCommitLog内的刷盘线程将writeBuffer(TransientStorePool 路径)或mappedByteBuffer中的数据落到磁盘。MappedFile.commit0会将committedPositionwrotePosition之间的数据写入 FileChannel,isAbleToFlush/isAbleToCommit均以OS_PAGE_SIZE(默认 4KB)为粒度计算脏页数量,具体逻辑见 RocketMQ MappedFile 内存映射文件详解;
  2. 逻辑队列派发(dispatch)ReputMessageService定时从 CommitLog 消费新写入的消息,将 20 字节的索引条目追加到对应主题队列的 ConsumeQueue 文件,并更新 IndexFile(消息索引)。消息存储成功后"物理文件(CommitLog)+ 逻辑文件(ConsumeQueue)"双写一致的机制,是 RocketMQ 高吞吐读写分离设计的基石;
  3. 主从同步:对于 SYNC_MASTER 模式,还会在刷盘完成后向从节点同步该条消息,等待从节点 ACK 后才向 Producer 返回发送成功。

七、全流程总结与故障排查对照

将四步主流程与返回状态码汇总如下,便于日常排障对照:

步骤核心方法检查点/动作失败返回状态码
1. 检查存储状态DefaultMessageStore#checkStoreStatusBroker 是否 shutdownSERVICE_NOT_AVAILABLE
Broker 角色是否为 SLAVESERVICE_NOT_AVAILABLE
runningFlags 是否可写(磁盘满/逻辑队列写失败/索引写失败)SERVICE_NOT_AVAILABLE
OS pageCache 是否繁忙OS_PAGECACHE_BUSY
2. 检查消息DefaultMessageStore#checkMessage主题长度 > 127MESSAGE_ILLEGAL
属性长度 > 32767MESSAGE_ILLEGAL
3. 定位 CommitLog 文件MappedFileQueue#getLastMappedFile首写或文件写满时创建新文件CREATE_MAPEDFILE_FAILED
4. 写入 MappedFileMappedFile#appendMessagesInner+doAppend顺序写入 mmap/堆外缓冲,事务消息特殊处理 queueOffsetPUT_OK/UNKNOWN_ERROR

核心结论

  • CommitLog 是 RocketMQ Broker 端消息存储的唯一物理载体,采用顺序追加 + 内存映射写入,文件默认 1GB、以起始偏移量命名,写满即滚动创建新文件;
  • 事务消息的 PREPARED / ROLLBACK 状态在存储层直接"隐身"——queueOffset 置 0、不推进 topicQueueTable、不进入 ConsumeQueue,只有 COMMIT 后才能真正被消费者感知;
  • 消息写入的健壮性由checkStoreStatus的"四道闸门"和checkMessage的"双重校验"层层保障,任何一环节失败都会以明确的PutMessageStatus语义返回,开发者可结合 Broker 端 warn 日志(注意printTimes的 5 万次节流机制)快速定位根因。

对于想要继续深入源码的读者,建议按以下顺序阅读仓库中的关联文档:RocketMQ 消息发送流程(客户端侧完整发送链路)→ 本文(Broker 存储链路)→ RocketMQ CommitLog 详解(物理文件管理与恢复)→ RocketMQ MappedFile 内存映射文件详解(mmap 与刷盘机制)→ RocketMQ ConsumeQueue 详解(逻辑队列构建与检索)。

  • 文档
  • 教程
  • 知识库

【免费下载链接】source-code-hunter

😱 从源码层面,剖析挖掘互联网行业主流技术的底层实现原理,为广大开发者 “提升技术深度” 提供便利。目前开放 Spring 全家桶,Mybatis、Netty、Dubbo 框架,及 Redis、Tomcat 中间件等

项目地址:https://gitcode.com/doocs/source-code-hunter
点击查看免费下载

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/20 14:57:07

radare2 macOS 深度指南:代码签名、调试授权与安装打包全解析

radare2 macOS 深度指南&#xff1a;代码签名、调试授权与安装打包全解析 【免费下载链接】radare2 UNIX-like reverse engineering framework and command-line toolset 项目地址: https://gitcode.com/gh_mirrors/ra/radare2 导读 macOS 的代码签名&#xff08;Code …

作者头像 李华
网站建设 2026/9/20 14:54:57

爱思唯尔投稿必看:利益声明文件撰写与提交避坑指南

简介&#xff1a;这份爱思唯尔利益声明文件专为科研工作者与论文作者准备&#xff0c;用于学术出版前规范声明是否存在竞争性财务利益或个人关系&#xff0c;以维护研究工作的透明与公正。压缩包内含1个docx格式的标准模板&#xff0c;文件大小仅28KB&#xff0c;内容可直接参考…

作者头像 李华
网站建设 2026/9/20 14:52:59

空气调节用制冷技术习题集:压焓图与蒸气压缩循环核心考点精讲

简介&#xff1a;这份《空气调节用制冷技术习题》docx文档&#xff0c;面向建筑环境与能源应用工程、暖通空调及相关专业学生&#xff0c;用于巩固制冷原理、设备选型与系统运行等核心知识。文档涵盖填空题、单项选择题、判断题、简答题与综合题&#xff0c;知识点涉及蒸气压缩…

作者头像 李华
网站建设 2026/9/20 14:52:53

MATLAB仿真中EbN0与SNR换算全解析:让BER曲线不再偏移

先交代一个我见过太多人翻车的场景&#xff1a;在MATLAB通信仿真里折腾了一整天BER曲线&#xff0c;明明调制方式、信道编码、解调判决都检查过好几遍&#xff0c;曲线却始终和理论值差着那么一口气。有时候是整体右偏1dB&#xff0c;有时候是斜率对但横坐标平移了3dB&#xff…

作者头像 李华