尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

RocketMQ源码精读:从路由到存储的核心链路解析

RocketMQ源码精读:从路由到存储的核心链路解析 1. 源码阅读前的思路准备1.1 为什么这么多人卡在RocketMQ源码上RocketMQ的源码在消息中间件里属于比较“耐读”的那一类——整体不算特别难但架不住量大。我见过不少朋友打开GitHub仓库之后第一反应都是“我该从哪看起”。这很正常RocketMQ的代码工程有几十个模块光Namesrv、Broker、Client、Remoting、Store这几个核心包就有几十万行如果没有一条清晰的主线很容易掉进细节里出不来读了两天还在看日志工具类。我的建议是不要想着“全部读懂”而是先建立一张底层代码的阅读地图。读源码和逛街一样你得先知道自己要去哪个商圈、走哪条主干道而不是一进门就钻进某个小店出不来。这张地图的关键就是几条主链路路由发现链路、消息发送链路、消息存储链路、消息消费链路外加一个贯穿始终的通信模块。还有一个容易被忽略的点——RocketMQ的源码版本差异。现在网上能找到的源码分析文章大多基于4.x版本而官方主推的5.x版本在模块划分和部分实现上已经有了不小变化比如引入了Proxy、gRPC通道等。我建议你以4.9.x的稳定版本作为第一遍阅读对象因为这个版本的代码结构最经典社区讨论最多遇到问题也最容易找到对应的资料。5.x可以在理解4.x之后再去看增量部分。1.2 带着问题去读而不是带着“读完”的目标很多人读源码失败是因为把目标定成了“把每一行都看懂”。其实这是最大的误区。读源码的正确姿势是带着问题去读比如生产者发送一条消息后Broker 是怎么把消息落盘的Consumer 从哪里拉消息拉到的消息存在哪里一个 Topic 有多个队列Consumer 是怎么分配队列的NameServer 和 Broker 之间是怎么保持心跳的这些问题就是一条条很清晰的阅读线索。你不需要关心每个类的每个方法只需要顺着线索把相关的类和方法串起来就能在比较短的时间里理解RocketMQ的核心骨架。读代码的过程很像做拼图先找边角料把框架定住再慢慢填充内部细节。2. 整体架构与模块划分2.1 源码工程的模块地图先从GitHub上把代码拉下来我用的是4.9.4版本。整个工程的核心模块可以列成一张表模块作用关键类remoting底层通信框架基于Netty封装NettyRemotingServer/Client, RemotingCommandcommon公共类、常量、协议定义MessageConst, TopicValidatorclient生产者和消费者的实现DefaultMQProducer, DefaultMQPushConsumerstore消息存储引擎DefaultMessageStore, CommitLog, ConsumeQueuenamesrv路由中心负责Topic路由注册与发现NamesrvController, RouteInfoManagerbrokerBroker服务端处理消息读写BrokerController, SendMessageProcessorproxy5.x新增兼容gRPC协议的代理层ProxyControllerexample示例代码快速上手QuickStarttest各类测试用例按模块分包distribution部署相关的脚本和配置启动脚本、conf目录如果只看核心通信链路remoting和store是两座大山namesrv最简单但也最容易读懂建议第一个啃它——花半天时间把 NamesrvController 和 RouteInfoManager 看完你会对整个RocketMQ的“注册与发现”机制有一个非常具象的认识。2.2 启动入口与Controller体系RocketMQ的代码里到处能看到Controller这种东西。它不是一个Spring的Controller而是一个“总控组件”负责把该模块的所有子组件装配起来、启动起来、关闭掉。以Broker为例BrokerController是整个Broker进程的“总开关”。在它的initialize()方法里你会看到它依次创建了消息存储模块MessageStore、Broker对外通信服务NettyRemotingServer、各种消息处理器SendMessageProcessor, PullMessageProcessor等、定时任务定时向NameServer注册、定时打印指标等。start()方法则按照严格的顺序启动这些组件。这种Controller模式的优点在于你不需要满世界找组件之间的依赖关系直接看Controller的initialize方法就能知道一个进程里都有哪些“器官”。这也是我建议你入门的时候先看启动入口的原因比追着一条消息链路去反查组件要高效得多。3. 路由中心NameServer的底层实现3.1 NameServer 启动时做了什么NameServer是RocketMQ里最“轻量”的模块但它的地位不低——所有生产者和消费者都要从它这里获取Topic的路由信息。源码中入口是NamesrvStartup它会创建NamesrvController然后调用initialize()和start()。真正干活的是RouteInfoManager。这个类的成员变量直接暴露了NameServer的核心数据结构private final HashMapString/* topic */, ListQueueData topicQueueTable; private final HashMapString/* brokerName */, BrokerData brokerAddrTable; private final HashMapString/* clusterName */, SetString brokerNameSet; private final HashMapString/* brokerAddr */, BrokerLiveInfo brokerLiveTable; private final HashMapString/* brokerAddr */, ListString/* filterServer */ filterServerTable;这几张表就是NameServer内存里维护的全部“地图数据”。看完这个类你对RocketMQ路由机制的理解会比背面试题深刻得多——原来NameServer就是几个HashMap在支撑。3.2 Broker 注册与心跳保活Broker启动之后会启动一个定时任务每隔30秒向NameServer发送一次心跳包。这个逻辑在BrokerController里可以看到最终通过RemotingClient发送一个HEART_BEAT请求到NameServer。而NameServer这边有一个ScanBrokerHousekeepingService的定时任务每10秒扫描一次brokerLiveTable如果发现某个Broker的lastUpdateTimestamp已经超过120秒没有更新就将其从所有路由表中移除。这个设计其实非常聪明——NameServer不主动探测Broker全靠心跳的被动更新配合一个宽松的过期时间既简单又不会误删。我在实际排查问题的时候经常用这个原理去判断“Broker是不是已经和NameServer失联了”比看进程死没死更准确。4. 存储层消息真正落盘的地方4.1 从 CommitLog 到 ConsumeQueue存储层是RocketMQ最核心、也最值得花时间读的部分。先说几个关键文件CommitLog所有消息的“总账本”消息真正落盘的地方按顺序写。ConsumeQueue每个Queue一个文件保存消息在CommitLog中的物理偏移量offset、消息大小和Message Tag的哈希值。IndexFile按照消息Key建立的索引用于按Key查询消息。MappedFileQueue对一组MappedFile的管理封装。MappedFile基于内存映射的文件封装是存储层读写的基本单位。一条消息的存储路径大概是这样的Broker的SendMessageProcessor收到消息后交给DefaultMessageStore.asyncPutMessage()这个方法会先将消息追加到CommitLog然后通过一个后台线程ReputMessageService将消息的摘要信息offset、size、tag hash分发到对应的ConsumeQueue中。也就是说CommitLog是唯一真正存消息内容的地方ConsumeQueue只是索引。这个设计的好处非常明显顺序写CommitLog的性能远高于随机写多个文件的性能而且由于ConsumeQueue很小可以从容地做异步构建。如果你在面试里被问到“RocketMQ为什么快”一定要把“CommitLog顺序写异步构建ConsumeQueue”这个点讲清楚。4.2 刷盘机制同步刷盘与异步刷盘在MessageStoreConfig里有两个关键配置flushDiskType和flushIntervalCommitLog。前者决定刷盘方式——同步刷盘SYNC_FLUSH还是异步刷盘ASYNC_FLUSH后者是异步刷盘的周期。同步刷盘并不是“每条消息都强制调用fsync”而是通过GroupCommitService把一批消息的写请求攒一下然后统一刷盘刷完返回给生产者确认。这样既保证了消息不丢又尽量提升了性能。异步刷盘则是写入PageCache就返回成功由后台定时任务把脏页刷到磁盘上。默认是每500ms刷一次当然如果你对消息可靠性要求极高就改成同步刷盘。我有个做金融支付的朋友他们的生产环境是同步刷盘据他说写入TPS在单Broker上仍然能到上万说明同步刷盘的实际代价并没有想象中那么可怕。RocketMQ读写CommitLog时都用了内存映射MappedByteBuffer这个在源码中对应MappedFile的map()方法。内存映射的巧妙之处在于它让操作文件像操作内存一样高效省去了一次用户态到内核态的复制。5. 消息发送与接收的主链路拆解5.1 生产者发送消息时客户端做了什么生产者的入口是DefaultMQProducer.send()它内部会经过DefaultMQProducerImpl这个核心实现类。整个发送链路大致可以分成几步根据Topic从本地缓存的路由信息中获取TopicPublishInfo如果本地没有缓存则向NameServer请求。根据消息的MessageQueueSelector挑选一个队列默认是轮询。通过MQClientAPIImpl.sendMessage()将消息封装成RemotingCommand交给Netty发往Broker。等待Broker返回结果如果是SEND_OK则发送完成。如果你用Debug模式跟踪一次消息发送你会发现MQClientInstance真是一个“大管家”它既管理生产者和Consumer的客户端实例也维护着与NameServer、Broker的所有连接还跑着各种定时任务比如更新路由信息、清理超时请求等。阅读了这个类你就明白了为什么RocketMQ的客户端虽然API很简洁但能支撑这么复杂的场景。5.2 Broker 端收到消息后的处理流程Broker端的入口是NettyRemotingServerNetty收到请求后根据请求的code找到对应的processor。消息发送请求的code是SEND_MESSAGE对应SendMessageProcessor。SendMessageProcessor.processRequest()里做了几件关键事情校验消息的Topic是否存在如果不存在则尝试自动创建Topic。调用MessageStore.putMessage()将消息写入CommitLog。根据写入结果构造SendMessageResponse返回给生产者。这里有一个面试官很爱问的点消息写入CommitLog时如何保证顺序。答案是在CommitLog.asyncPutMessage()中通过putMessageLock对写入操作加锁同一时刻只有一个线程在追加消息。这是RocketMQ提升性能的关键手段之一。还要注意一个细节在SendMessageProcessor里有一个判断msg.isWaitStoreMsgOK()的逻辑这是同步发送和异步发送的一个重要分水岭。同步发送时Broker会等待刷盘完成或至少写入PageCache后再返回而异步发送则只返回一个提交成功的状态具体是否落盘由后台线程决定。6. 消息消费与重平衡机制6.1 推模式还是拉模式本质都是拉大家在用RocketMQ的DefaultMQPushConsumer时感觉像是Broker在“推送”消息但打开源码就会发现消费者本质上是主动去Broker拉消息的。Push和Pull的区别只在于——Push模式在本地封装了一个长轮询机制拉不到消息时会阻塞在Broker端挂起的请求上等有新消息时再立即返回。消费端的核心类是DefaultMQPushConsumerImpl它启动后会创建一个后台线程PullMessageService不断从ProcessQueue拿到待拉取的MessageQueue然后发送拉取请求到Broker。Broker端的PullMessageProcessor是处理这些拉取请求的入口。它从ConsumeQueue中获取消息的物理偏移量再到CommitLog中读取真正的消息内容。如果当前没有新消息它不会立刻返回空结果而是把请求挂起来等到有新消息时再唤醒——这就是长轮询的核心机制对应源码中PullRequestHoldService。从这里你应该能体会到RocketMQ的“实时性”不是靠Broker主动推而是靠大量客户端同时挂着长轮询请求来“等”消息。理解了这一点你对RocketMQ消费延迟的来源和调优思路就会更清晰。6.2 重平衡队列的分配与冲突处理重平衡Rebalance是消费端最容易出问题、也最值得研究的部分。它发生在消费者实例变化或Topic队列数量变化的时候目标是让队列在所有消费者之间重新分配。核心实现在RebalanceImpl.rebalanceByTopic()。它分两步从Broker获取当前Topic的队列信息。从Broker获取当前ConsumerGroup下所有在线消费者的ClientID列表。按照分配策略把队列分配给自己和其他消费者。RocketMQ内置了几种分配策略比如平均分配算法AllocateMessageQueueAveragely一致性哈希算法AllocateMessageQueueConsistentHash。默认用的是平均分配算法你可以通过消费者参数AllocateMessageQueueStrategy来改。重平衡过程中的一个常见问题就是消息重复消费。因为重平衡会让某个队列从消费者A移到消费者B如果A还没处理完队列里的消息B也会开始拉取两边就可能有重复。所以RocketMQ的消费语义是“至少一次”而不是“恰好一次”。要想去重只能在业务端做幂等处理。6.3 消费进度的保存与消息堆积RocketMQ每个Consumer Group都会保存当前消费到哪条消息的进度也就是Consumer Offset。这个offset不是在客户端本地保存而是作为一个特殊Topic__consumer_offset存储在Broker上这样即使消费者挂了换个节点也能恢复进度。这个机制的入口在ConsumerManageProcessor和ConsumerOffsetManager中。每次消费成功之后客户端会定时上报最新的消费位点到BrokerBroker则负责持久化这些位点。这里有个常用的调优点如果线上出现了消息堆积你可以用mqadmin consumerProgress命令查看消费位点和最新消息位点之间的差距。在实际排查堆积的时候很多人一上来就看日志效率很低。我更习惯先看consumerProgress确认消费位点卡住不动再去看线程状态或下游依赖是否超时。有了源码的底子你在排查问题时就能预判这个命令背后的实现逻辑而不是瞎跑命令。7. 高可用与主从同步机制7.1 主从复制的方式与源码实现RocketMQ的高可用依赖主从同步。Broker的主从关系由brokerId决定0表示主节点非0表示从节点。主从同步有两种模式同步复制SYNC_MASTER和异步复制ASYNC_MASTER由配置brokerRole指定。同步复制的关键源码在CommitLog和HAService中。主节点写入消息后会通过HAService中的HAClient将数据推送给从节点同时等待从节点返回确认确认之后才向生产者返回SEND_OK。异步复制则不等从节点确认写入主节点就直接返回。从节点在启动时会自动从主节点同步数据这个过程在HAClient.run()中实现。它会先向主节点发送自己的最大偏移量主节点从该偏移量开始逐步推送数据。如果主从之间的网络发生了长时间分区从节点会反复重连数据进度差值会不断拉大直到网络恢复。7.2 主从切换与读写分离的“真相”RocketMQ 4.x 的机制里如果主节点挂了从节点不会自动升级为主节点而是需要你手动通过mqadmin命令将某个从节点设置为brokerId0来触发切换。5.x 中虽然有自动容灾的演进但生产环境里依然需要依赖监控和运维工具来配合。这里还有一个很有意思的点消费者是可以从从节点拉消息的前提是主节点的slaveReadEnable配置开启。这样当主节点压力大的时候可以把一部分读流量分流到从节点。但生产者写入只能走主节点因为只有主节点才能提供消息写入的服务。明白主从的数据布局之后你就知道为什么部署RocketMQ至少是一主一从的架构而不是像某些简单系统那样只部署一台。在主从数据不同步期间如果主节点宕机消息可能丢——这是同步复制和异步复制的本质差别你得根据业务对数据可靠性的要求去做取舍。8. 基于源码的常见问题排查与心得8.1 从源码层面看最容易踩的坑读了一段时间源码后我发现网上很多RocketMQ的“玄学问题”其实在源码里都有明确答案。这里挑几个我踩过的坑说下。Producer发送超时但Broker实际已写入。这个在同步发送模式下很容易遇到。原因在于发送超时时间默认是3000ms如果Broker端刷盘慢一点、或者GC停顿一下客户端就可能提前超时。但Broker端可能已经把消息写进去了这时候如果业务方直接做失败重试就可能造成重复消息。所以我的建议是重要场景务必在消费端做幂等。消费堆积时不要盲目加消费者。增加消费者确实能提升消费能力但前提是你的下游处理能力没到瓶颈。如果消费慢是因为下游RPC调用慢或数据库慢增加消费者只能让下游压力更大甚至引起雪崩。有了源码基础这类问题的排查可以先看ConsumeRequest的线程池队列积压情况再决定是加机器还是优化下流程。重平衡期间消费抖动。由于重平衡是“全量分配”的某次网络抖动或者GC导致心跳超时就可能触发一次全量rebalance让多个消费者同时暂停消费引起消费延迟瞬时升高。这在源码里都有体现。我的处理经验是合理调大heartbeatInterval和rebalanceLockWaitTime并对消费者实例做“优雅退出”的设计尽量降低重平衡的频次。8.2 一份实用的源码阅读顺序清单最后分享一下我个人推荐的阅读顺序。不需要按模块顺序硬啃按依赖关系逐层展开更容易理解remoting模块理解Netty封装和RemotingCommand协议。namesrv模块理解路由表结构和心跳保活。store模块的CommitLog与ConsumeQueue理解消息存储的核心。broker的启动与消息处理把存储和通信串起来。client的消息发送从生产者视角看一次完整发送。client的消息消费与rebalance从消费端视角补齐最后一块拼图。有余力再看主从同步HA机制、过滤、事务消息、延迟消息。每一层理解透之后再去下一层整体的效率要比“一把梭”高很多。我通常在读一个模块时会开两个窗口一个看源码一个开一个B站或者博客的讲解视频。先看讲解建立全局观再回源码验证细节这样会省不少力气。8.3 最后一个体会有人问我读RocketMQ源码到底有没有用。我的看法是如果不读源码你能熟练使用RocketMQ也能排查绝大多数问题但读了源码之后你在面对诡异问题时会有一种“开天眼”的感觉——因为你能猜到它内部大概发生了什么用的是哪条链路的哪个类。这种感觉在面试里也很有用当你随口说出“可以在RebalanceImpl里加个日志看看这次重平衡是由哪个消费者触发的”面试官基本会给你加不少印象分。读源码不需要追求一次全懂更不需要背注释和类名。把它当成一个持续迭代的过程今天看懂一条发送链路明天看懂一个存储机制积累一两个月之后整个RocketMQ在你眼里就会从“黑盒”变成“白盒”。那时候你再去看任何消息队列的面试题都会觉得——不过如此。
返回列表