RocketMQ 核心源码精读指南
这是 RocketMQ 系列的最后一篇也是最硬核的一篇——源码阅读与分析。经过前五个阶段的学习你已经掌握了 RocketMQ 的架构原理、存储机制、发送消费流程、进阶特性和部署运维。可以说你已经是一名合格的 RocketMQ 开发者了。但“合格”和“精通”之间还隔着源码这道门槛。为什么建议你读源码三个原因遇到诡异问题时源码是你最可靠的“字典”——它能告诉你框架到底是怎么运作的性能调优时只有理解了底层实现才知道参数该怎么调面试时能讲清楚源码的实现细节是区分“会用”和“真懂”的分水岭今天这篇文章我会带你走一遍 RocketMQ 源码的“地图”——从工程结构到各个核心模块告诉你从哪里入手、看什么、怎么看。老规矩配合流程图和代码片段一步一图。十六、源码阅读与分析源码工程结构与模块划分在开始阅读源码之前我们先要搞清楚 RocketMQ 源码工程的整体布局。以 RocketMQ 5.x 主干分支为例源码目录结构如下rocketmq/├── broker/ # Broker 服务端核心模块├── client/ # 客户端实现Producer、Consumer├── common/ # 公共工具包常量、配置、工具类├── distribution/ # 发行包与安装配置脚本├── example/ # 示例代码├── filtersrv/ # 消息过滤服务器已废弃├── logging/ # 日志组件├── namesrv/ # NameServer 路由中心├── openmessaging/ # OpenMessaging 标准兼容├── proxy/ # 5.x 新增代理服务gRPC/HTTP├── remoting/ # 远程通信模块基于 Netty├── store/ # 消息存储底层实现├── test/ # 单元测试与集成测试└── tools/ # 运维管理命令行工具各模块职责速览模块 职责namesrv/ NameServer 路由中心实现 Broker 注册、路由管理、心跳检测broker/ Broker 服务端核心实现消息的接收、存储、转发、投递、消费进度管理store/ 消息存储底层CommitLog、ConsumeQueue、IndexFile 等remoting/ 远程通信基于 Netty 实现客户端与服务端的网络通信client/ 客户端 APIProducer 和 Consumer 的核心逻辑proxy/ 5.x 新增代理服务支持 gRPC 协议客户端收发消息controller/ 5.x 新增控制器帮助 Broker 做主从切换 阅读建议如果你第一次读 RocketMQ 源码建议按这个顺序入手remoting通信基础→ namesrv路由→ store存储→ broker服务端→ client客户端。由浅入深逐步推进。NameServer 核心源码剖析路由管理NameServer 是 RocketMQ 的“轻量级注册中心”。它的源码非常精简——只有八个类不到 1000 行代码。核心类RouteInfoManagerNameServer 的路由管理核心在 org.apache.rocketmq.namesrv.routeinfo.RouteInfoManager 中实现。它通过 5 个核心数据结构来维护路由元信息// 路由元信息的 5 个核心数据结构private final HashMapString/* topic/, List topicQueueTable;private final HashMapString/brokerName/, BrokerData brokerAddrTable;private final HashMapString/clusterName/, SetString/brokerName/ clusterAddrTable;private final HashMapString/brokerAddr/, BrokerLiveInfo brokerLiveTable;private final HashMapString/brokerAddr/, List/Filter Server */ filterServerTable;数据结构 作用topicQueueTable Topic → Queue 列表的映射brokerAddrTable Broker 名称 → Broker 地址信息的映射clusterAddrTable 集群名称 → Broker 名称集合的映射brokerLiveTable Broker 地址 → 存活信息的映射含心跳时间filterServerTable Broker 地址 → 过滤服务器列表的映射Broker 心跳注册流程Broker 启动后每隔 30 秒向所有 NameServer 发送心跳命令。源码中使用 CountDownLatch 实现多线程同步并发地向所有 NameServer 发送注册请求// Broker 向所有 NameServer 发送心跳源码简化for (final String namesrvAddr : nameServerAddressList) {brokerOuterExecutor.execute(() - {RegisterBrokerResult result registerBroker(namesrvAddr, …);// 处理注册结果});}countDownLatch.await(timeoutMills, TimeUnit.MILLISECONDS);NameServer 之间无状态、不通信NameServer 集群节点之间没有任何数据同步和通信。每个节点独立维护路由信息即使某个时刻各节点的数据不完全一致也不会影响消息的发送。这种设计极大简化了 NameServer 的实现也让它变得极为轻量和稳定。Broker 核心源码剖析消息存储、转发Broker 是 RocketMQ 最复杂的模块涉及消息的接收、存储、转发、投递和消费进度管理。Broker 的分层设计Broker 分层架构请求处理层SendMessageProcessor / PullMessageProcessor解析 RemotingCommand 的 RequestCode业务逻辑层DefaultMessageStoreputMessage / getMessage文件映射层MappedFile基于 MappedByteBuffer存储层CommitLog / ConsumeQueue / IndexFileBroker 启动流程加载持久化的配置信息消费进度、订阅信息等加载 DefaultMessageStore消息存储组件创建 MappedFileQueue 映射 CommitLog、ConsumeQueue、IndexFile 等文件创建并启动 BrokerController 控制器处理消息的发送和接收核心存储设计理念RocketMQ 将所有主题的消息不分主题一律顺序写入 CommitLog 文件。这与 Kafka 按分区存储的设计不同——Kafka 在 Topic 和分区数量增长时写入性能会下降而 RocketMQ 的表现则稳定得多。因此Kafka 适合 Topic 和分区较少的场景RocketMQ 更适合多 Topic、多消费端的业务场景。消息发送流程源码剖析消息发送的入口是 DefaultMQProducer其核心实现在 DefaultMQProducerImpl 中。发送流程的四个核心步骤imageProducer 启动流程检测配置判断生产者组是否合法创建客户端实例MQClientInstance 通过 MQClientManager 单例创建是非常核心的类每个实例有唯一的 clientId注册本地生产者将 Producer 注册到 MQClientInstance 的 producerTable 中启动客户端实例启动 Netty 通信模块、定时任务、负载均衡服务定时任务是 Producer 的核心机制之一发送心跳每隔 30 秒将客户端信息发送到 Broker更新路由定时从 NameServer 拉取最新的 Topic 路由信息发送消息的核心方法sendDefaultImpl获取主题发布信息topicPublishInfo根据路由算法选择一个消息队列selectOneMessageQueue调用 sendKernelImpl 发送消息封装成 SendResult消息拉取与消费流程源码剖析RocketMQ 的消费者有 DefaultMQPushConsumer 和 DefaultMQPullConsumer 两种但底层都是基于长轮询实现的。Push 消费者启动流程DefaultMQPushConsumerImpl#startconsumer.start加载偏移量广播模式存本地 / 集群模式存 Broker启动 Netty 客户端与 Broker 建立通信连接启动定时任务发送心跳、更新路由启动 PullMessageService异步拉取消息启动 RebalanceService负载均衡Consumer 就绪关键点消费者启动时需要加载各个 Topic 的偏移量。广播模式下偏移量存储在消费者本地集群模式下存储在 Broker 端。拉取消息的核心流程PullMessageService 线程不断从 pullRequestQueue 中取出 PullRequest向 Broker 发起拉取请求包含 ConsumerGroup、Topic、Queue、queueOffset 等信息Broker 通过长轮询机制响应有消息立即返回无消息挂起等待拉取到的消息提交到消费线程池处理Push 与 Pull 的本质Push 模式只是在客户端将消息拉取到本地后自动回调业务方的监听器执行消费逻辑。内核依然是 Pull 长轮询。CommitLog 写入流程源码剖析CommitLog 是 RocketMQ 存储的核心所有消息都顺序写入 CommitLog。CommitLog 的核心数据结构组件 说明MappedFile 单个文件的内存映射基于 MappedByteBufferMappedFileQueue 一组 MappedFile 的队列管理文件的滚动CommitLog 消息写入的入口封装了写入逻辑每个 CommitLog 文件默认 1GB文件名以起始偏移量命名如 00000000000000000000。消息写入流程同步刷盘异步刷盘Producer 发送消息Broker 接收请求SendMessageProcessor.processRequestDefaultMessageStore.putMessage消息存储入口CommitLog.putMessage执行写入获取写入锁可重入锁或自旋锁通过 MappedFile将消息追加到 PageCache更新写入指针释放锁刷盘策略MappedByteBuffer.force等待刷盘完成唤醒刷盘线程立即返回返回写入结果锁机制putMessage 会有多个线程并行处理需要加锁。可以通过配置选择使用可重入锁还是自旋锁useReentrantLockWhenPutMessage。刷盘的最终实现都是使用 NIO 中的 MappedByteBuffer.force() 将映射区的数据写入磁盘。