做RocketMQ源码阅读这几年,真正让我下定决心把“消息存储”和“Consumer”这两块吃透并写成文章的,是一次线上的消费积压事故。当时积压了三千多万条消息,Broker、NameServer的监控指标全部正常,客户端日志也没有一条报错,最后排查了三个多小时才定位到根因。事后复盘时我发现,问题恰好出在存储层和消费循环的“节奏配合”上——不读源码,光靠官方文档和监控面板,你连怀疑的方向都没有。
这篇文章我会沿着两条主线展开:一条是消息从Producer发出后,在Broker端如何落盘、如何构建消费索引;另一条是Consumer启动后,如何从Broker拉取消息、如何做重平衡、如何管理消费进度。涉及的核心类、方法名、关键源码分支我都会贴出来,同时把那些文档里不会写的取舍理由和排错经验一并讲清楚。适合已经在用RocketMQ、遇到积压或丢失问题,却不知道从哪下手的开发者。
1. 消息存储的架构骨架:CommitLog、ConsumeQueue、IndexFile各管哪一段
RocketMQ的存储目录下通常会看到三类东西:commitlog目录、consumequeue目录和index目录。很多人把它们统称为“存储文件”,但实际分工差异很大。要理解消息存储,第一步就是搞清楚这三者各自的边界。
1.1 CommitLog:所有消息唯一的落点
CommitLog是RocketMQ里唯一真正保存消息实体的地方,目录固定为${storePathRootDir}/commitlog。里面是一组大小固定为1GB的文件(MappedFile),文件名是一串固定20位的数字,比如00000000001073741824,这个数字代表该文件内第一条消息的全局物理偏移量。这种命名方式的妙处在于:通过offset / 1GB就能直接定位到文件,通过offset % 1GB就能定位到文件内偏移,不需要额外维护元数据。
所有Topic的所有Queue的消息,最终都追加到同一个CommitLog里。这与Kafka“每个分区独立日志文件”的设计形成鲜明对比。Kafka那样做的好处是分区内读写隔离,坏处是分区数量多了以后,文件句柄和磁盘寻道开销会急剧上升。RocketMQ之所以敢把所有消息写在一起,底气来自顺序写:在机械硬盘时代,顺序写能跑到上百MB/s,而随机写可能连10MB/s都到不了。这个“全局顺序追加”的模型,就是RocketMQ存储高吞吐的底层根基。
1.2 ConsumeQueue:按队列维度的轻量索引
消息全部堆在一个文件里,消费的时候却要按“Topic + QueueId”维度去取,如果每次消费都去扫CommitLog,性能是不可接受的。于是RocketMQ引入了ConsumeQueue,路径为${storeRootPath}/consumequeue/{topic}/{queueId},每个队列对应一个目录,目录内同样是一组固定大小的文件,每条索引记录固定20字节:
| 字段 | 长度 | 说明 |
|---|---|---|
| commitLogOffset | 8字节 | 消息在CommitLog中的物理偏移量 |
| msgSize | 4字节 | 消息序列化后的长度 |
| tagsCode | 8字节 | 消息Tag经过计算后的过滤码 |
这20字节的索引记录不包含消息内容,只记录了“消息在哪、多长、Tag是什么”。消费时先查ConsumeQueue拿到偏移量,再根据commitLogOffset去CommitLog里读完整消息。值得注意的是,严格来说消费路径不是“随机读”,而是先顺序读CQ,再按偏移读CommitLog,属于有规律的离散读,借助pageCache后性能依然可观,但相比写入还是有差距。
1.3 IndexFile:按Key查消息的哈希索引
IndexFile是第三类文件,路径为${storeRootPath}/index,用于根据消息Key做精确查询。它的结构是典型的哈希索引:文件头记录起始/结束时间戳、起始/结束偏移量、哈希槽数量;紧接着是哈希槽区;最后是索引条目链表区。每写入一条带Key的消息,会计算Key的哈希值,通过哈希槽定位到索引链表头,再头插一条索引项。
我自己的经验是:IndexFile平时存在感很低,因为绝大多数业务使用RocketMQ不会做按Key回溯查询。但一旦需要排查线上某条消息到底发出去没有,org.apache.rocketmq.store.index.QueryOffsetResult配合admin工具按Key检索,就比人工翻日志高效得多。这个文件可以容忍丢失或部分损坏,因为消息实体还在CommitLog里,只要CommitLog完好,最多是索引查不到,消息不会丢。
2. 消息写入链路:从Netty请求到CommitLog落盘,中间发生了什么事
理解了存储骨架,接下来进入最核心的写入链路。一条消息从客户端发出,到真正落盘,在Broker端至少经过四层:网络层、业务处理层、存储层、刷盘线程。每一层都有值得展开的源码细节。
2.1 从SendMessageProcessor到DefaultMessageStore
Broker收到Producer发来的SEND_MESSAGE请求后,NettyRemotingServer根据请求码路由到SendMessageProcessor。processRequest方法里会做一堆前置校验,比如Topic是否存在、延迟级别是否合法、是否有权限等。校验通过后,消息会被包装成MessageExtBrokerInner,随后调用DefaultMessageStore.putMessage。
putMessage这一步很关键,它不仅要校验当前Broker是否处于可写状态,还要调用CommitLog.asyncPutMessage完成真正的追加。源码里有一长串前置检查:
public PutMessageResult putMessage(MessageExtBrokerInner msg) { // 检查Broker是否可写、当前存储是否在运行、消息是否超长等 if (this.shutdown) { ... } if (!this.runningFlags.isWriteable()) { ... } ... return this.commitLog.asyncPutMessage(msg); }这里有一个容易被忽略的细节:putMessage入口处会先拿putMessageLock。RocketMQ默认使用自旋锁putMessageSpinLock(可通过useReentrantLockWhenPutMessage配置切换为可重入锁)。为什么用自旋锁?因为MappedFile.appendMessage的临界区代码非常短,锁竞争时间极短,自旋比线程挂起唤醒更划算。在超高并发写入场景下,这个锁的竞争策略直接影响吞吐。
2.2 MappedFile与mmap:为什么顺序写能这么快
asyncPutMessage里最核心的步骤是拿到当前可写的MappedFile,然后调用appendMessage。MappedFile封装了文件的内存映射逻辑,默认大小就是1GB。写入时先拿到fileChannel或writeBuffer,把消息体序列化到ByteBuffer,然后追加到缓冲区。
核心代码大概长这样:
public AppendMessageResult appendMessage(MessageExtBrokerInner msg, AppendMessageCallback cb) { assert msg != null; assert cb != null; long currentTime = System.currentTimeMillis(); return cb.doAppend(this.getFileFromOffset(), this.fileChannel, currentTime, this.mappedFile.getReadByteBuffer(), this.mappedFile.getWriteBuffer(), msg); }这里有个RocketMQ的高阶玩法:transientStorePool。默认情况下,消息直接写入MappedFile对应的MappedByteBuffer,相当于写入pageCache,后续由操作系统异步刷盘。如果开启transientStorePoolEnable=true,消息会先写入一个堆外内存池(writeBuffer,DirectByteBuffer),再由后台线程批量commit到pageCache,最后才刷到磁盘。这样做的本质是把“用户态到pageCache”的拷贝变成一个可控的批量动作,减少pageCache的锁竞争,适合追求更高写入吞吐且机器内存充足的场景。
需要特别提醒的是:开启transientStorePool后,还要正确设置堆外内存池大小(transientStorePoolSize默认5,每个池容量对应一个1GB文件)。如果内存不足,反而会频繁触发缺页或OOM,得不偿失。
2.3 文件翻转、刷盘策略与返回结果
当当前MappedFile剩余空间不足以容纳下一条消息时,CommitLog会先尝试复用剩余空间;如果剩余空间确实不够,则创建新的MappedFile。新文件的文件名就是当前已经写入的最大偏移量,也就是旧文件最后一条消息的末尾位置。这个“文件末尾占位”的逻辑在源码里也有体现:旧文件剩余的尾部空间会写入空的占位数据,保证消息不会跨文件存储。为什么要这么设计?因为如果一条消息跨文件,读取时需要同时打开两个文件拼接数据,复杂度和失败概率都会上升。
消息追加完成后,AppendMessageCallback.doAppend会返回AppendMessageResult,包含消息的物理偏移量、长度、写入时间等。消息ID就是根据Broker IP端口和偏移量生成的,这也是为什么同一台Broker上,消息ID的后面几位其实是文件偏移量的一部分。
刷盘策略决定了消息什么时候真正落到磁盘,RocketMQ的配置项是flushDiskType,分为两种:
| 刷盘策略 | 实现线程 | 行为 | 适用场景 |
|---|---|---|---|
| ASYNC_FLUSH | FlushRealTimeService | 每500ms或满足条件时批量刷盘 | 默认策略,吞吐高,存在极小概率丢消息 |
| SYNC_FLUSH | GroupCommitService | 每批消息提交后立即刷盘,返回成功前落盘 | 对可靠性要求极高,主备切换场景 |
同步刷盘不是每写一条就刷一次,而是把多个线程的写入请求合并成一次批量flush,由GroupCommitService统一调度。这样既保证了确认语义,又不会让IO完全退化。
3. 消费索引怎么建出来的:ReputMessageService异步分发
CommitLog是消息存储的本体,但Consumer消费时并不会直接去扫CommitLog,而是走ConsumeQueue。那ConsumeQueue里的20字节索引记录是怎么从CommitLog里“孵化”出来的?这就是ReputMessageService的工作。
3.1 消息落盘后,CQ到底是什么时候更新
ReputMessageService是DefaultMessageStore内部的一个后台服务线程,默认每1毫秒执行一次doDispatch。它维护了一个reputFromOffset指针,这个指针初始值是CommitLog在当前进程启动时的已写偏移量。启动后,它会循环执行:
public void run() { while (!this.isStopped()) { Thread.sleep(1); this.defaultMessageStore.doDispatch(this); } }doDispatch里会判断如果reputFromOffset < CommitLog的已写偏移,就把这段新增的字节读出来,逐条解析成消息,然后分发给各个Dispatcher。这就是RocketMQ里常说的“异步构建索引”。
写CommitLog和构建ConsumeQueue是两条完全不同的线程在做,意味着ConsumeQueue的更新一定滞后于CommitLog。滞后量通常只有几毫秒,很多人在排查消费延迟时会忽略这个因素。实际上在极端高并发下,如果Broker端IO被打满,ReputMessageService滞后几百毫秒甚至一秒都是正常的,这不代表消息丢失,只是消费侧会同步感知到延迟。
3.2 分发链:CommitLogDispatcherBuildConsumeQueue与IndexFile
DefaultMessageStore.doDispatch内部维护了一个DispatcherList,按注册顺序执行。核心分发器有两个:
CommitLogDispatcherBuildConsumeQueue:根据消息的topic和queueId找到对应的ConsumeQueue,调用putMessagePositionInfoWrapper把20字节索引写入。CommitLogDispatcherBuildIndex:如果消息带Key,就把Key的哈希写入IndexFile。
写ConsumeQueue时有一个和CommitLog类似的机制:CQ文件默认大小是30万个索引条目(约600KB),同样不允许索引记录跨文件,如果剩余空间不足,也会先补齐占用,再创建新文件。
这里必须强调一个恢复特性:ConsumeQueue是可以被“抛弃并重建”的。因为CQ本质上是CommitLog的一个可推导视图,只要CommitLog完好,即使整个consumequeue目录被误删,重启Broker后ReputMessageService可以按reputFromOffset重新扫描CommitLog,从头把CQ构建回来。这个特性在生产环境救援时非常有用,后面第6章我会讲具体操作。
4. Consumer消费循环:拉取请求的产生与完成
存储侧搞清楚了,接下来看消费侧。RocketMQ的消费模型看着简单——DefaultMQPushConsumer订阅Topic然后等回调,但背后PullMessageService和Broker之间的配合其实是一个永不停止的“拉取循环”。下面从源码角度拆开这个循环。
4.1 启动时铺好的两条线程线
DefaultMQPushConsumerImpl.start()里做了大量初始化,比如copySubscription复制订阅关系、初始化RebalanceImpl、PullAPIWrapper、OffsetStore等。最终会启动两个常驻线程:
RebalanceService:默认每20s执行一次重平衡,决定当前客户端应该消费哪些Queue。PullMessageService:从pullRequestQueue里不断取PullRequest,向Broker发拉取请求。
PullMessageService.run()的骨架代码非常直观:
public void run() { while (!this.isStopped()) { PullRequest pullRequest = this.pullRequestQueue.take(); if (pullRequest.isLockedFirst()) { ... } if (this.isConsumeGroupEmpty(pullRequest)) { ... continue; } this.pullMessage(pullRequest); } }pullRequestQueue是一个阻塞队列,如果拿不到拉取任务,线程就阻塞等待,不会空转。拉取完成后,无论有没有数据,都会把下一个PullRequest重新丢回队列里,形成循环。这条循环是消费侧的动力来源,只要线程活着,消费就不会停。
4.2 一次拉取从客户端到Broker再到回调
pullMessage方法里做的事情可以概括为四步:
- 计算本次要拉取的偏移量,来自
OffsetStore里记录的消费进度。 - 构建
PullMessageRequestHeader,包含消费者组、Topic、QueueId、nextOffset、maxMsgNums(默认32)等。 - 调用
PullAPIWrapper.pullKernelImpl发送PULL_MESSAGE请求给Broker。 - 注册
PullCallback,等待Broker响应。
Broker端收到请求后由PullMessageProcessor处理,先查ConsumeQueue拿到一批索引项,再根据commitLogOffset从CommitLog里读消息体,组装成PullResult返回。所以一次消费拉取的本质是:先读索引,再按索引读数据。返回的结果可能是FOUND、NO_NEW_MSG、NO_MATCHED_MSG或OFFSET_ILLEGAL。
PullCallback.onSuccess拿到结果后,客户端会做两件事:
boolean dispatchToConsume = processQueue.putMessage(pullResult.getMsgFoundList()); this.consumeMessageService.submitConsumeRequest(pullResult.getMsgFoundList(), processQueue, pullRequest.getMessageQueue(), dispatchToConsume);第一件事是把拉到的消息放进本地ProcessQueue,第二件事是提交给消费线程池去执行用户定义的MessageListener。这里有个容易忽略的点:processQueue.putMessage是先把消息放进树形Map,真正的消费线程执行完成后再从Map里删除,这个过程决定了消息“正在消费中”的状态。
4.3 ProcessQueue:消费端的内存缓冲与去重依据
ProcessQueue可以理解为每个Queue在客户端本地的内存缓冲,内部维护了一棵TreeMap<Long, MessageExt>,key是消息的偏移量,value是消息体。所有已经被拉取到本地、但还没被消费线程处理完的消息都存在这里面。
为什么用TreeMap而不是ArrayList?因为消息的偏移量天然有序,TreeMap可以高效支持:
- 按偏移量顺序取出消息,保证同一Queue消息的消费顺序被本地维护。
- 快速移除一批已消费消息。
- 快速计算
msgAccCount和pendingMsgCount,用于判断本地积压是否严重。
在顺序消费模式下,ProcessQueue还会持有一个锁,确保同一时间只有一个消费者线程在消费该队列;而在并发消费模式下,一个Queue的消息可以被多个消费线程并发处理,顺序不保证。如果ProcessQueue长时间没有消息消费,会有tryRemoveFromFifo之类的清理逻辑把它从processQueueTable里移除,这又是另一个源码细节,这里先不展开。
5. Rebalance与消费进度:并发消费背后的两个隐形引擎
如果说PullMessageService是消费侧的发动机,那Rebalance就是方向盘,OffsetStore则是仪表盘。三者配合不好,消费就会出现倾斜、重复或丢失。这一章把后两个讲透。
5.1 重平衡触发与队列分配
客户端启动后,RebalanceService会每隔20秒(rebalanceInterval)执行一次doRebalance,但也可以通过rebalanceImmediately()提前唤醒。触发重平衡的时机包括:Broker地址变化、消费者实例上下线、Topic路由信息变化。
重平衡的核心方法是RebalancePushImpl.rebalanceByTopic,流程如下:
- 从NameServer拿到该Topic下的所有MessageQueue列表(
mqSet)。 - 向Broker查询当前消费组内所有消费者客户端ID列表(
cidAll)。 - 用分配策略(默认
AllocateMessageQueueAveragely)计算当前客户端应该分到哪些Queue。
平均分配算法的源码思路值得单独看:
int index = cidAll.indexOf(currentCID); int mod = mqAll.size() % cidAll.size(); int averageSize = mqAll.size() <= cidAll.size() ? 1 : (mod > 0 && index < mod ? mqAll.size() / cidAll.size() + 1 : mqAll.size() / cidAll.size()); int startIndex = (mod > 0 && index < mod) ? index * averageSize : index * averageSize + mod; result.addAll(mqAll.subList(startIndex, startIndex + averageSize));这段逻辑简单但易出错:当队列数不能被消费者数整除时,多出来的队列会分配给前mod个消费者,每人多分一个。所以在重平衡后,某些实例负载会略高,这是正常现象。如果集群规模经常变化,可以考虑用AllocateMessageQueueConsistentHash做一致性哈希分配,减少抖动。
分配完成后,updateProcessQueueTableInRebalance会对新旧队列做diff:新增的Queue创建ProcessQueue和PullRequest;被移出的Queue则从本地表移除,对应队列的消费进度由Broker端的offset记录兜底。这里有个关键点:重平衡发生时,如果旧Queue还有消息正在本地ProcessQueue里排队未消费,这些消息不会被立即丢弃,而是会通过重平衡后的消费逻辑处理,具体取决于消费进度已经推进到哪里。
5.2 消费成功怎么记账:OffsetStore链路
消费进度管理是很多使用者最容易踩坑的地方。RocketMQ的OffsetStore有两个实现:
LocalFileOffsetStore:广播模式下使用,消费进度保存在本地文件。RemoteBrokerOffsetStore:集群模式下使用,消费进度定时上报给Broker。
并发消费的ConsumeMessageConcurrentlyService在处理结果时调用:
this.getConsumerGroupOffsetStore().updateOffset(messageQueue, offset, increaseOnly);updateOffset只更新本地内存里的偏移量,真正的持久化是定时任务做的。客户端默认每5秒(persistConsumerOffsetInterval)调用一次persistAll,把当前内存里的进度发给Broker;Broker端由ConsumerOffsetManager维护并写入磁盘。
这套“内存优先 + 定时上报”的模型带来一个重要推论:消费进度的持久化有最多几秒的延迟。如果消费者进程在内存offset还没来得及上报时异常退出,重启后会根据Broker上存的上一个进度重新拉取,导致部分已处理消息被重复消费。这也解释了为什么RocketMQ的语义是“至少一次”(At Least Once)而不是“恰好一次”。理解了这一点,你在设计消费端逻辑时就应该默认“消息可能重复”,把幂等放在首位,而不是寄希望于消息队列帮你扛。
5.3 批量消费参数怎么影响消费节奏
源码里有两个最直接影响消费节奏的参数:consumeMessageBatchMaxSize和pullBatchSize。
pullBatchSize默认32,决定一次拉取请求向Broker要多少条消息,过大容易瞬间把ProcessQueue撑大,增加客户端内存压力;过小则增加请求次数,拉取吞吐下降。consumeMessageBatchMaxSize默认1,决定每次提交给用户Listener的消息数量。需要注意:批量消费时,Listener返回的消费结果作用于整批消息,如果其中一条处理失败,整批会被判定为失败并走重试。
Spring Boot集成RocketMQ时,很多人习惯把consumeMessageBatchMaxSize调到很大来提升吞吐,但这会让单批消息处理时间变长,进而影响消费线程池的回收节奏,最终体现为消费不及时。我的建议是:不要盲目调大,批量消费的收益远小于风险,优先保证每条消息的处理时间和幂等性。
6. 存储与消费结合部的实战排错:几个容易被误判的问题
最后把近几年在存储和消费结合部遇到的高频问题做一个梳理,每一个都是我或团队同学踩过的真实场景。遇到类似现象,至少不用再从零开始怀疑。
6.1 消费积压但一切指标正常:先看本地ProcessQueue
回到开篇说的事故:消费积压3000万,Broker和客户端指标都正常。最后定位的过程是这样的——先看消费线程池,consumeThreadMin=20确实有20个线程在跑;再看线程Dump,发现所有消费线程都在执行用户逻辑,而且是慢IO操作。
真正的问题出在:拉取线程太快。每拉回一批消息,就立刻塞进ProcessQueue,消费线程处理不过来,ProcessQueue里堆积的消息越来越多。但监控面板查的都是Broker端消费位点和生产位点差,看不到客户端本地的未消费缓冲,所以指标显得“正常”。
这种问题的解决方法有两条路:一是消费线程池扩容,或者优化用户消费逻辑;二是通过pullInterval或者降低pullBatchSize,限制拉取速度,让ProcessQueue的积压保持在低位。RocketMQ没有提供太细粒度的消费限流阀,但通过这几个参数组合,完全可以把“拉取—消费”的节奏调到匹配。
6.2 “unable to read consumer identity”到底在说什么
这个报错出现时,客户端日志一般会有类似“no consumer ID list”或“unable to read consumer identity”的提示,很多同学误以为是权限或认证问题。实际上它发生在Rebalance阶段:当前消费者向Broker查询消费组内所有消费者ID列表时,拿到的结果为空或者异常。
常见触发原因有:
- 消费者实例启动时订阅关系还没注册完成,就被触发了一次重平衡,此时Broker端
ConsumerManager里的组信息不完整。 - 多个客户端使用相同的
clientId(默认是IP地址 + instanceName)。同一台机器部署多个实例时,如果instanceName重复,会互相覆盖注册信息。 - 客户端版本和Broker版本差异过大,导致
GetConsumerListByGroup响应反序列化失败。
排查思路很简单:先看Broker日志里该消费组的注册记录,确认ClientID是否唯一;再看订阅关系是否一致,同一个组内不同实例订阅不同Topic也会导致查询结果异常。运维上建议每个Container显式指定不同的instanceName,尤其是Spring Boot多实例部署时,这是最容易被忽略的坑。
6.3 控制拉取速率的几种有效手段
很多人问“RocketMQ能不能控制消费速率”,严格说它没有像消息队列那样内置的限速SDK,但可以从配置侧做到“软限速”:
| 参数 | 默认值 | 影响 |
|---|---|---|
| pullInterval | 0 | 两次拉取之间的间隔,大于0等于强制降速 |
| pullBatchSize | 32 | 每次拉取最大消息数,缩小可降低冲击 |
| consumeMessageBatchMaxSize | 1 | 单次回调消息数,不建议超过10 |
| consumeThreadMin/Max | 20 | 消费线程数,降低会影响并发消费能力 |
实际操作中,如果业务依赖外部接口且接口有QPS限制,我一般会在消费回调里按需Thread.sleep(),并把pullInterval设为一个较小值兜底。直接把pullInterval设得太大,会导致消费延迟变大,也不合理。
6.4 存储文件异常恢复与CQ重建
最后聊一个偏底层的救援操作:如果consumequeue目录下的CQ文件损坏,会出现一种奇妙的现象——Consumer消费不到消息,但CommitLog是完好的,消息也没丢。
RocketMQ本身在DefaultMessageStore.load阶段会对CommitLog和ConsumeQueue做加载校验;如果发现CQ损坏严重,最简单的恢复手段是:停Broker,备份(或删除)consumequeue目录,保留commitlog目录,然后重启Broker。ReputMessageService会从CommitLog的已有数据开始重新构建CQ。理论上耗时取决于CommitLog数据量,但这个过程是自动的。
删除CQ而不是整个存储目录,这是一个很多新手不敢做的动作,因为从直觉上看“索引文件”也是数据。但只要理解RocketMQ的存储架构,就会明白:CQ是可再生的派生数据,CommitLog才是必须保护的主数据。当然,操作前一定要备份,这是所有运维操作的底线。
拿我自己来说,那次线上故障之后,我做的最有价值的一件事,就是把RocketMQ存储和消费的源码调用链完整梳理了一遍,然后画在了团队的Wiki里。之后遇到任何消费异常,排查路径都清晰了很多。如果你也想真正掌握RocketMQ,建议不要停留在API使用层面,按这篇文章的顺序把CommitLog、ReputMessageService、PullMessageService、RebalanceService这几个核心类逐个读一遍,比读十篇博客都管用。源码不会骗人,文档才会。