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 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#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 主题长度校验:不能超过 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 消息发送流程 的发送前校验环节),Broker 侧的重复校验属于防御性编程,防止绕过客户端直接构造非法请求。
3.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,即用户设置的 keys、tags、自定义属性序列化后的字符串)长度被限制为Short.MAX_VALUE(32767)。其根因在于消息存储协议中属性长度字段以short类型编码(2 字节),超过该值将无法在消息体中正确表达。这与"消息体最大 4M"(maxMessageSize,由客户端DefaultMQProducer控制)是两类不同的限制:体重大小限制在客户端侧,主题/属性长度限制在 Broker 存储侧。
四、第三步:获取当前可写的 CommitLog 文件
通过合法性检查后,DefaultMessageStore#putMessage进入 CommitLog 写入路径CommitLog#asyncPutMessage,首先需要定位"当前应该往哪个文件里写"。
4.1 目录与文件组织
CommitLog 文件的存储目录为${ROCKET_HOME}/store/commitlog,其中:
MappedFileQueue对应整个 commitlog 文件夹,负责管理文件集合与滚动;MappedFile对应文件夹下的单个物理文件,每个文件默认 1G(mappedFileSizeCommitLog),文件名即为该文件的起始物理偏移量(左补零至 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; }两种需要创建新文件的场景:
- 首次写入:当前没有任何 MappedFile,
createOffset通过对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 的文件滚动机制是理解这三类文件的前提——它们的文件命名与滚动规则同源。
五、第四步:将消息写入 MappedFile
5.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(Collection<Message>))一次请求携带多条消息,存储层需要区分处理; - 双缓冲机制:
writeBuffer不为空(即启用了TransientStorePool堆外内存暂存,对应 broker 配置transientStorePoolEnable)时,先写入writeBuffer,再由后台线程 commit 到fileChannel;否则直接写入mappedByteBuffer。这两种路径下commit/flush的差异,详见 RocketMQ MappedFile 内存映射文件详解; - 追加成功后同步累加
wrotePosition并记录storeTimestamp,返回携带写入字节数等信息的AppendMessageResult。
5.2 DefaultAppendMessageCallback#doAppend
org.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环节)。
构造 AppendMessageResult
AppendMessageResult 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_TYPE),topicQueueTable(主题 + 队列维度维护的"下一条消息的队列偏移量"表)被推进,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_AVAILABLE | Broker 已 shutdown / Broker 为 SLAVE / 磁盘满、ConsumeQueue 或 IndexFile 写入失败导致不可写 | DefaultMessageStore#checkStoreStatus | 检查 Broker 生命周期状态、brokerRole配置、磁盘占用率与后台任务日志 |
OS_PAGECACHE_BUSY | PageCache 积压,刷盘速度跟不上写入 | DefaultMessageStore#isOSPageCacheBusy | 观察写入与刷盘速率差,必要时扩容或调整刷盘策略 |
MESSAGE_ILLEGAL | 主题长度 > 127 / 属性长度 > 32767 | DefaultMessageStore#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),仅供参考