
RocketMQ 消息发送存储流程源码剖析从 Broker 接收到 CommitLog 落盘的完整链路【免费下载链接】source-code-hunter 从源码层面剖析挖掘互联网行业主流技术的底层实现原理为广大开发者 “提升技术深度” 提供便利。目前开放 Spring 全家桶Mybatis、Netty、Dubbo 框架及 Redis、Tomcat 中间件等项目地址: https://gitcode.com/GitHub_Trending/so/source-code-hunter本文基于 RocketMQ 4.9.3 源码版本org.apache.rocketmq.store包从DefaultMessageStore#putMessage入口出发逐步拆解消息到达 Broker 后经历的存储状态检查、消息合法性校验、CommitLog 文件定位、消息序列化追加四个核心环节。读完本文你将理解PutMessageStatus各状态码的触发条件、CommitLog 顺序写盘与文件滚动机制、事务消息在存储层的特殊处理并能在排查发送失败/存储不可用类问题时精准定位代码路径。建议与仓库内 RocketMQ 消息发送流程客户端侧、RocketMQ MappedFile 内存映射文件详解、RocketMQ CommitLog 详解 配合阅读形成从客户端发消息到Broker 落盘的完整闭环。一、总览Broker 收到消息后发生了什么客户端Producer通过 Remoting 协议将消息发送到 Broker 后Broker 端由SendMessageProcessor处理请求最终调用存储层核心类org.apache.rocketmq.store.DefaultMessageStore#putMessage完成消息落盘。落盘过程可以划分为四个关键步骤检查消息存储状态checkStoreStatus确认 Broker 本身是否允许写入检查消息合法性checkMessage校验主题长度、属性长度是否越界获取当前可写的 CommitLog 文件通过MappedFileQueue定位或创建待写入的MappedFile将消息写入 MappedFile经由CommitLog#asyncPutMessage的appendMessagesInner与DefaultAppendMessageCallback#doAppend完成消息编码与追加。每一步的失败都会以PutMessageStatus枚举的形式返回最终封装成PutMessageResult回传给上层。下面逐层展开。二、第一步检查消息存储状态checkStoreStatus入口方法为org.apache.rocketmq.store.DefaultMessageStore#checkStoreStatus该方法依次执行四项检查任一不满足都会拒绝写入。值得注意的细节是前三项检查返回的都是SERVICE_NOT_AVAILABLE但它们的日志策略各不相同——shutdown 场景每次都打印告警而角色为从节点不可写场景每 50000 次才打印一次避免日志刷屏。2.1 检查 Broker 是否已关闭if (this.shutdown) { log.warn(message store has shutdown, so putMessage is forbidden); return PutMessageStatus.SERVICE_NOT_AVAILABLE; }当 Broker 正在执行优雅停机shutdown标志为 true时消息写入被直接拒绝。这也是为什么生产环境在重启 Broker 前需要先摘流量。2.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; }BrokerRole由 broker 配置文件中的brokerRole属性控制取值ASYNC_MASTER、SYNC_MASTER、SLAVE。从节点的职责是同步主节点数据并提供读服务默认不允许直接写入。这也解释了为什么读写分离场景下客户端必须把写流量路由到主节点。2.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 brokers 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#isWriteable是一个综合性的可写标志其状态由多个后台任务维护主要包括磁盘空间不足DiskCheckService周期性检查 CommitLog、ConsumeQueue、Index 等目录的磁盘占用率超过阈值diskMaxUsedSpaceRatio默认 75时置为不可写写入逻辑队列ConsumeQueue失败写入索引文件IndexFile失败。这三种原因在告警日志中均有明确提示。当检查通过后printTimes会被重置为 0保证下一次异常发生时能立即告警而不是等待取模命中。2.4 检查 PageCache 是否繁忙if (this.isOSPageCacheBusy()) { return PutMessageStatus.OS_PAGECACHE_BUSY; }isOSPageCacheBusy通过比较写入页数与刷盘页数的差值来判断操作系统 PageCache 是否积压过重。当差值超过阈值osPageCacheBusyTimeOutMills相关逻辑控制时说明刷盘速度跟不上写入速度此时返回OS_PAGECACHE_BUSY进行背压backpressure。这一点与客户端侧的行为直接相关客户端在同步发送模式下收到OS_PAGECACHE_BUSY后会走发送失败重试逻辑参考 RocketMQ 消息发送流程 中retryTimesWhenSendFailed的配置从而间接降低 Broker 的写入压力。生产环境中若频繁出现该状态码通常意味着需要扩容或调整刷盘策略。三、第二步检查消息合法性checkMessage入口方法为org.apache.rocketmq.store.DefaultMessageStore#checkMessage负责消息的静态合法性校验违反任一限制返回MESSAGE_ILLEGAL。3.1 主题长度校验不能超过 127if (msg.getTopic().length() Byte.MAX_VALUE) { log.warn(putMessage message topic length too long msg.getTopic().length()); return PutMessageStatus.MESSAGE_ILLEGAL; }主题长度被硬限制为Byte.MAX_VALUE127 字节。注意客户端侧的Validators#checkTopic也做了相同约束参见 RocketMQ 消息发送流程 的发送前校验环节Broker 侧的重复校验属于防御性编程防止绕过客户端直接构造非法请求。3.2 属性长度校验不能超过 32767if (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即用户设置的 keys、tags、自定义属性序列化后的字符串长度被限制为Short.MAX_VALUE32767。其根因在于消息存储协议中属性长度字段以short类型编码2 字节超过该值将无法在消息体中正确表达。这与消息体最大 4MmaxMessageSize由客户端DefaultMQProducer控制是两类不同的限制体重大小限制在客户端侧主题/属性长度限制在 Broker 存储侧。四、第三步获取当前可写的 CommitLog 文件通过合法性检查后DefaultMessageStore#putMessage进入 CommitLog 写入路径CommitLog#asyncPutMessage首先需要定位当前应该往哪个文件里写。4.1 目录与文件组织CommitLog 文件的存储目录为${ROCKET_HOME}/store/commitlog其中MappedFileQueue对应整个 commitlog 文件夹负责管理文件集合与滚动MappedFile对应文件夹下的单个物理文件每个文件默认 1GmappedFileSizeCommitLog文件名即为该文件的起始物理偏移量左补零至 20 位。这套组织方式的细节可进一步参考 RocketMQ CommitLog 详解文件命名、偏移量换算、findMappedFileByOffset的索引定位算法与 RocketMQ MappedFile 内存映射文件详解文件初始化、commit、刷盘、销毁。4.2 定位与创建文件的逻辑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 为空或该文件已写满mappedFile.isFull()即fileSize wrotePosition时调用getLastMappedFile(0)获取新的文件。源码注释Mark: NewFile may be cause noise提示这里可能在并发场景下创建新文件属于预期行为排查时不必惊慌。MappedFileQueue#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; }两种需要创建新文件的场景首次写入当前没有任何 MappedFilecreateOffset通过对startOffset按文件大小对齐取整得到保证文件名的偏移量是文件大小的整数倍最新文件已满createOffset 上一个文件起始偏移量 mappedFileSize即在上一个文件末尾顺延一个文件大小。若创建失败例如磁盘 IO 异常返回CREATE_MAPEDFILE_FAILED。4.3 与其他存储文件的呼应CommitLog 之外ConsumeQueue 和 IndexFile 同样依赖MappedFileQueue组织文件ConsumeQueue 的存储条目为定长 20 字节8 字节物理偏移量 4 字节消息长度 8 字节 tag 哈希码目录结构为主题/队列两级详见 RocketMQ ConsumeQueue 详解IndexFile 用于按消息 key 检索哈希槽与索引条目同样是定长结构详见 RocketMQ IndexFile 详解。理解 CommitLog 的文件滚动机制是理解这三类文件的前提——它们的文件命名与滚动规则同源。五、第四步将消息写入 MappedFile5.1 MappedFile#appendMessagesInner定位到目标 MappedFile 后调用org.apache.rocketmq.store.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通过 CAS 语义保证并发安全currentPos fileSize确保不会越界写入两种消息类型单条消息MessageExtBrokerInner与批量消息MessageExtBatch走不同的doAppend分支批量发送场景下客户端DefaultMQProducer的send(CollectionMessage)一次请求携带多条消息存储层需要区分处理双缓冲机制writeBuffer不为空即启用了TransientStorePool堆外内存暂存对应 broker 配置transientStorePoolEnable时先写入writeBuffer再由后台线程 commit 到fileChannel否则直接写入mappedByteBuffer。这两种路径下commit/flush的差异详见 RocketMQ MappedFile 内存映射文件详解追加成功后同步累加wrotePosition并记录storeTimestamp返回携带写入字节数等信息的AppendMessageResult。5.2 DefaultAppendMessageCallback#doAppendorg.apache.rocketmq.store.CommitLog.DefaultAppendMessageCallback#doAppend负责将MessageExtBrokerInner编码为 CommitLog 中存储的二进制格式并构造AppendMessageResult。计算写入偏移量long wroteOffset fileFromOffset byteBuffer.position();wroteOffset为本次消息在 CommitLog 中的全局物理偏移量等于文件起始偏移量加上文件内位置。该值将作为消息的物理位置写入 ConsumeQueue并被消费者用来定位消息。对事务消息做特殊处理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。sysFlag的取值在客户端发送侧设置——当消息携带PROPERTY_TRANSACTION_PREPARED属性时置位TRANSACTION_PREPARED_TYPE见 RocketMQ 消息发送流程 的sendKernelImpl环节。构造 AppendMessageResultAppendMessageResult result new AppendMessageResult(AppendMessageStatus.PUT_OK, wroteOffset, msgLen, msgIdSupplier, msgInner.getStoreTimestamp(), queueOffset, CommitLog.this.defaultMessageStore.now() - beginTimeMills);AppendMessageResult携带了本次追加的完整元信息写入状态PUT_OK、物理偏移量wroteOffset、消息长度msgLen、消息 ID、存储时间戳、队列偏移量以及本次追加的耗时now() - beginTimeMills用于统计存储延迟。更新主题-队列偏移量表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; }对于普通消息NOT_TYPE和提交后的消息COMMIT_TYPEtopicQueueTable主题 队列维度维护的下一条消息的队列偏移量表被推进multiDispatch.updateMultiQueueOffset同时更新多队列multiDispatch是 4.9.x 引入的 POP 消费 / 多队列投递支持的偏移量记录。这些偏移量随后在消息分发CommitLogDispatcherBuildConsumeQueue阶段被写入 ConsumeQueue供消费端按队列拉取衔接点可参考 RocketMQ ConsumeQueue 详解 的putMessagePositionInfo。六、与整体链路的衔接消息从客户端到磁盘要完整理解发送存储流程在 RocketMQ 中的地位需要把它放回整条消息链路上看客户端侧DefaultMQProducer#send经过消息校验主题、体大小 ≤ 4M、路由查找从 NameServer 拉取TopicPublishInfo、队列选择含故障延迟机制sendLatencyFaultEnable、sendKernelImpl组装SendMessageRequestHeader并发送详见 RocketMQ 消息发送流程路由维护Broker 启动后通过registerBrokerAll周期默认 30s可配置registerNameServerPeriod向 NameServer 注册路由NameServer 每 10s 扫描并剔除 120s 未活跃的 Broker保证生产者拿到的路由信息是新鲜的详见 RocketMQ NameServer 与 Broker 的通信存储层本文Broker 收到消息后走checkStoreStatus → checkMessage → 定位 MappedFile → doAppend四步落盘 CommitLog随后由分发线程构建 ConsumeQueue 与 IndexFile消费侧消费者通过 ConsumeQueue 定位消息的物理偏移量再从 CommitLog 按getMessage(offset, size)读取参考 RocketMQ 消息消费流程 与 RocketMQ CommitLog 详解getMessage与findMappedFileByOffset的查找算法。七、故障排查速查PutMessageStatus 一览PutMessageStatus触发场景关键代码位置排查建议SERVICE_NOT_AVAILABLEBroker 已 shutdown / Broker 为 SLAVE / 磁盘满、ConsumeQueue 或 IndexFile 写入失败导致不可写DefaultMessageStore#checkStoreStatus检查 Broker 生命周期状态、brokerRole配置、磁盘占用率与后台任务日志OS_PAGECACHE_BUSYPageCache 积压刷盘速度跟不上写入DefaultMessageStore#isOSPageCacheBusy观察写入与刷盘速率差必要时扩容或调整刷盘策略MESSAGE_ILLEGAL主题长度 127 / 属性长度 32767DefaultMessageStore#checkMessage检查消息主题与属性是否超限CREATE_MAPEDFILE_FAILED创建新 CommitLog 文件失败CommitLog#asyncPutMessage检查存储目录权限、磁盘空间与文件句柄数UNKNOWN_ERROR非MessageExtBrokerInner/MessageExtBatch的消息类型或文件已满仍尝试追加MappedFile#appendMessagesInner一般属于内部异常结合 Broker 错误日志分析八、小结RocketMQ 的消息存储流程是典型的多级校验 顺序写 内存映射设计多级校验保证非法消息、不可用状态下不会污染存储且每一类失败都有明确的PutMessageStatus语义便于上层客户端重试逻辑与运维侧日志告警快速定位CommitLog 顺序追加 MappedFile 内存映射将随机写转化为顺序写配合 PageCache 与可选的TransientStorePool双缓冲是 RocketMQ 高性能写入的基石事务消息与普通消息的差异化处理queueOffset是否推进、是否进入 ConsumeQueue体现了存储层对消息语义的完整建模。深入理解这条链路后再回头读 RocketMQ MappedFile 内存映射文件详解 中的 commit/flush 细节、RocketMQ CommitLog 详解 中的恢复与截断逻辑你会对整个 RocketMQ 存储子系统形成系统性的认知——这也正是本仓库从源码层面剖析主流技术底层实现的初衷。【免费下载链接】source-code-hunter 从源码层面剖析挖掘互联网行业主流技术的底层实现原理为广大开发者 “提升技术深度” 提供便利。目前开放 Spring 全家桶Mybatis、Netty、Dubbo 框架及 Redis、Tomcat 中间件等项目地址: https://gitcode.com/GitHub_Trending/so/source-code-hunter创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考