拓冰建站拓冰建站
首页 / 资讯中心 / 正文

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。这条链路共分四步检查消息存储状态checkStoreStatusBroker 是否可用、角色是否为主、存储是否可写、pageCache 是否繁忙检查消息合法性checkMessage主题长度、属性长度是否超限获取当前可写的 CommitLog 文件通过MappedFileQueue.getLastMappedFile定位或创建MappedFile将消息写入 MappedFileCommitLog#asyncPutMessage内部调用MappedFile.appendMessagesInner最终由DefaultAppendMessageCallback#doAppend完成字节级写入。其中步骤 3、4 是整个存储流程的性能核心CommitLog 采用**顺序写 内存映射mmap**的方式这是 RocketMQ 单机吞吐量远高于传统随机写数据库的关键设计。相关文件模型可参考 RocketMQ CommitLog 详解 与 RocketMQ MappedFile 内存映射文件详解。二、第一步检查消息存储状态checkStoreStatusorg.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 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是 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. 校验主题长度不超过 127if (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. 校验属性长度不超过 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 序列化后的字符串长度上限为Short.MAX_VALUE即 32767 字节。注意该上限针对的是属性字符串而非消息体——消息体的大小限制默认 4MBmaxMessageSize在客户端DefaultMQProducerImpl#checkMessage中校验。两项校验失败均返回MESSAGE_ILLEGALBroker 端会拒绝该消息并记录告警日志。四、第三步获取当前可写入的 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。文件滚动与创建getLastMappedFilepublic 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; }这里有两种触发新文件创建的场景首次写入mappedFileLast null此时createOffset startOffset - (startOffset % mappedFileSize)即向下对齐到文件大小的整数倍——这正是 CommitLog 文件以起始偏移量命名的由来最新文件已写满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); }几个关键点写入位置以wrotePositionAtomicInteger 写指针作为本次写入的起始位置保证并发安全写入介质若启用了TransientStorePool堆外内存暂存池即writeBuffer不为空则先写入堆外writeBuffer的共享缓冲区切片后续由 commit 阶段批量刷入 FileChannel否则直接写入mappedByteBuffermmap 映射的直接内存。两种路径的差异决定了刷盘模型详见下文刷盘与提交消息类型分支MessageExtBrokerInner单条消息与MessageExtBatch批量消息走不同的doAppend重载其他类型返回UNKNOWN_ERROR写入结果追加完成后累加wrotePosition并记录storeTimestamp。核心写入DefaultAppendMessageCallback#doAppendorg.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已提交事务消息才参与队列偏移量分配。动作三构造 AppendMessageResultAppendMessageResult 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)更新多队列派发信息。topicQueueTableConcurrentMapString, Long维护着主题队列与逻辑偏移量的映射它是后续构建 ConsumeQueue 的核心数据源。ConsumeQueue 中每条记录只存 20 字节8 字节 CommitLog 物理偏移量 4 字节消息长度 8 字节 tag 哈希码具体格式与查询逻辑见 RocketMQ ConsumeQueue 详解。六、写入完成后的下行链路提交、刷盘与派发putMessage返回PutMessageResult后存储流程并未结束剩余工作由 CommitLog 内部及后台线程接力完成提交commit与刷盘flush根据刷盘策略同步刷盘SYNC_FLUSH/ 异步刷盘ASYNC_FLUSH由FlushConsumeQueueService、CommitLog内的刷盘线程将writeBufferTransientStorePool 路径或mappedByteBuffer中的数据落到磁盘。MappedFile.commit0会将committedPosition至wrotePosition之间的数据写入 FileChannelisAbleToFlush/isAbleToCommit均以OS_PAGE_SIZE默认 4KB为粒度计算脏页数量具体逻辑见 RocketMQ MappedFile 内存映射文件详解逻辑队列派发dispatchReputMessageService定时从 CommitLog 消费新写入的消息将 20 字节的索引条目追加到对应主题队列的 ConsumeQueue 文件并更新 IndexFile消息索引。消息存储成功后物理文件CommitLog 逻辑文件ConsumeQueue双写一致的机制是 RocketMQ 高吞吐读写分离设计的基石主从同步对于 SYNC_MASTER 模式还会在刷盘完成后向从节点同步该条消息等待从节点 ACK 后才向 Producer 返回发送成功。七、全流程总结与故障排查对照将四步主流程与返回状态码汇总如下便于日常排障对照步骤核心方法检查点/动作失败返回状态码1. 检查存储状态DefaultMessageStore#checkStoreStatusBroker 是否 shutdownSERVICE_NOT_AVAILABLEBroker 角色是否为 SLAVESERVICE_NOT_AVAILABLErunningFlags 是否可写磁盘满/逻辑队列写失败/索引写失败SERVICE_NOT_AVAILABLEOS pageCache 是否繁忙OS_PAGECACHE_BUSY2. 检查消息DefaultMessageStore#checkMessage主题长度 127MESSAGE_ILLEGAL属性长度 32767MESSAGE_ILLEGAL3. 定位 CommitLog 文件MappedFileQueue#getLastMappedFile首写或文件写满时创建新文件CREATE_MAPEDFILE_FAILED4. 写入 MappedFileMappedFile#appendMessagesInnerdoAppend顺序写入 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点击查看免费下载相关推荐RocketMQ 消息发送存储流程源码剖析从 Broker 接收到 CommitLog 落盘的完整链路RocketMQ 消息发送存储流程源码剖析从 Broker 接收到 CommitLog 落盘的完整链路 本文基于 RocketMQ 4.9.3 源码版本 o文档教程技术博客知识库RocketMQ 消息发送全流程源码解析从生产者到 Broker 的完整链路RocketMQ 消息发送全流程源码解析从生产者到 Broker 的完整链路 本文基于 RocketMQ 4.9.3 源码围绕 rocketmq send文档教程技术博客知识库RocketMQ 消息拉取流程源码解析从 PullMessageService 到 Broker 拉取的完整链路RocketMQ 消息拉取流程源码解析从 PullMessageService 到 Broker 拉取的完整链路 本文基于 doocs/source code文档教程知识库创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门