- 文档
- 教程
- 知识库
【免费下载链接】source-code-hunter
😱 从源码层面,剖析挖掘互联网行业主流技术的底层实现原理,为广大开发者 “提升技术深度” 提供便利。目前开放 Spring 全家桶,Mybatis、Netty、Dubbo 框架,及 Redis、Tomcat 中间件等
本文基于 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。这条链路共分四步:
- 检查消息存储状态(
checkStoreStatus):Broker 是否可用、角色是否为主、存储是否可写、pageCache 是否繁忙; - 检查消息合法性(
checkMessage):主题长度、属性长度是否超限; - 获取当前可写的 CommitLog 文件:通过
MappedFileQueue.getLastMappedFile定位或创建MappedFile; - 将消息写入 MappedFile:
CommitLog#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; }这里有两种触发新文件创建的场景:
- 首次写入:
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); }几个关键点:
- 写入位置:以
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)更新多队列派发信息。
topicQueueTable(ConcurrentMap<String, Long>)维护着"主题@队列"与逻辑偏移量的映射,它是后续构建 ConsumeQueue 的核心数据源。ConsumeQueue 中每条记录只存 20 字节(8 字节 CommitLog 物理偏移量 + 4 字节消息长度 + 8 字节 tag 哈希码),具体格式与查询逻辑见 RocketMQ ConsumeQueue 详解。
六、写入完成后的下行链路:提交、刷盘与派发
putMessage返回PutMessageResult后,存储流程并未结束,剩余工作由 CommitLog 内部及后台线程接力完成:
- 提交(commit)与刷盘(flush):根据刷盘策略(同步刷盘
SYNC_FLUSH/ 异步刷盘ASYNC_FLUSH),由FlushConsumeQueueService、CommitLog内的刷盘线程将writeBuffer(TransientStorePool 路径)或mappedByteBuffer中的数据落到磁盘。MappedFile.commit0会将committedPosition至wrotePosition之间的数据写入 FileChannel,isAbleToFlush/isAbleToCommit均以OS_PAGE_SIZE(默认 4KB)为粒度计算脏页数量,具体逻辑见 RocketMQ MappedFile 内存映射文件详解; - 逻辑队列派发(dispatch):
ReputMessageService定时从 CommitLog 消费新写入的消息,将 20 字节的索引条目追加到对应主题队列的 ConsumeQueue 文件,并更新 IndexFile(消息索引)。消息存储成功后"物理文件(CommitLog)+ 逻辑文件(ConsumeQueue)"双写一致的机制,是 RocketMQ 高吞吐读写分离设计的基石; - 主从同步:对于 SYNC_MASTER 模式,还会在刷盘完成后向从节点同步该条消息,等待从节点 ACK 后才向 Producer 返回发送成功。
七、全流程总结与故障排查对照
将四步主流程与返回状态码汇总如下,便于日常排障对照:
| 步骤 | 核心方法 | 检查点/动作 | 失败返回状态码 |
|---|---|---|---|
| 1. 检查存储状态 | DefaultMessageStore#checkStoreStatus | Broker 是否 shutdown | SERVICE_NOT_AVAILABLE |
| Broker 角色是否为 SLAVE | SERVICE_NOT_AVAILABLE | ||
| runningFlags 是否可写(磁盘满/逻辑队列写失败/索引写失败) | SERVICE_NOT_AVAILABLE | ||
| OS pageCache 是否繁忙 | OS_PAGECACHE_BUSY | ||
| 2. 检查消息 | DefaultMessageStore#checkMessage | 主题长度 > 127 | MESSAGE_ILLEGAL |
| 属性长度 > 32767 | MESSAGE_ILLEGAL | ||
| 3. 定位 CommitLog 文件 | MappedFileQueue#getLastMappedFile | 首写或文件写满时创建新文件 | CREATE_MAPEDFILE_FAILED |
| 4. 写入 MappedFile | MappedFile#appendMessagesInner+doAppend | 顺序写入 mmap/堆外缓冲,事务消息特殊处理 queueOffset | PUT_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 中间件等
相关推荐
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),仅供参考