源码深入篇: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() 将映射区的数据写入磁盘。ConsumeQueue 构建流程源码剖析ReputMessageServiceConsumeQueue 是 CommitLog 的索引文件由 ReputMessageService 这个后台线程负责构建。ReputMessageService 的核心机制Broker 启动时会开启一个线程每毫秒执行一次 doReput() 方法它有一个属性 reputFromOffset记录消息重放的偏移量每次执行时从 reputFromOffset 开始读取 CommitLog 中的消息逐条解析消息构建 ConsumeQueue 和 IndexFileConsumeQueue 条目格式每个 ConsumeQueue 条目固定 20 个字节字段 长度 说明CommitLog 偏移量 8 字节 消息在 CommitLog 中的物理位置消息长度 4 字节 消息的字节数Tag 哈希码 8 字节 Tag 的哈希值用于消息过滤工作流程图有新数据无新数据Broker 启动启动 ReputMessageService 线程每 1ms 执行一次 doReput从 reputFromOffset检查 CommitLog 是否有新数据逐条解析消息构建 ConsumeQueue 条目offset size tagHash构建 IndexFileKey → Offset 映射更新 reputFromOffset短暂休眠Rebalance 机制源码剖析Rebalance重平衡是消费者组内 Queue 分配的动态调整机制。触发入口RebalanceService 是一个线程任务在消费者客户端启动时被调用。广播模式集群模式RebalanceService.rundoRebalance遍历 consumerTable获取每个 MQConsumerInnerrebalanceByTopic执行重平衡消费模式简化版分配所有 Queue 分配给所有 Consumer核心分配逻辑根据分配策略计算AllocateMessageQueueStrategy平均分配 / 环形分配更新 ProcessQueue创建/释放 PullRequest完成核心实现负载均衡的最终执行逻辑在 rebalanceByTopic() 方法中。广播模式和集群模式的实现不同集群模式是核心实现。分配策略策略 实现类 特点平均分配 AllocateMessageQueueAveragely 尽可能均匀分配默认策略环形分配 AllocateMessageQueueAveragelyByCircle 轮流分配类似发牌事务消息源码剖析半消息与回查事务消息是 RocketMQ 最复杂的特性之一其核心是半消息Half Message 和事务回查Transaction Check 机制。半消息的存储半消息发送成功后会进入 RocketMQ 内部的 RMQ_SYS_TRANS_HALF_TOPIC 的 ConsumeQueue 中。这个消息只存储在 CommitLog 中但在 ConsumeQueue 中不可见因此消费者无法消费到它。事务消息的完整流程业务应用BrokerProducer业务应用BrokerProducer进入回查流程loop[每 60 秒回查]alt[COMMIT][ROLLBACK][UNKNOWN]发送半消息写入 RMQ_SYS_TRANS_HALF_TOPICCommitLog 可见ConsumeQueue 不可见半消息发送成功执行本地事务返回事务状态6a. 提交事务7a. 将消息转发到目标 Topic6b. 回滚事务7b. 删除半消息8c. 发起回查请求9c. 检查本地事务状态10c. 返回状态11c. 提交最终事务状态事务回查机制如果 Producer 在发送半消息后未能及时通知 Broker 提交或回滚Broker 会定期默认 60 秒回查 Producer 的事务状态。check 方法会对半消息进行过滤如超过 72 小时的事务消息算作过期只保留符合条件的半消息进行回查。消息重试与死信源码剖析消息重试机制RocketMQ 的消费者在消息处理失败时会根据返回状态码自动触发重试默认情况下消息最多重试 16 次并发消费每次重试的时间间隔呈指数增长如 10s → 30s → 1min → 2min…重试与死信的流转是否是否消息消费消费成功提交 Offset进入重试队列%RETRY%ConsumerGroup延迟重试时间逐次递增重试次数≤ 16 次进入死信队列%DLQ%ConsumerGroup需要人工介入处理并发消费与顺序消费的重试差异并发消费失败的消息进入重试队列不影响其他消息的消费顺序消费一条消息失败会阻塞该 Queue 的后续消息直到重试成功或进入 DLQ死信队列当消息重试次数超过最大重试次数默认 16 次消息会被放入死信主题%DLQ%ConsumerGroup。死信队列中的消息需要人工介入处理。Netty 通信层源码剖析RocketMQ 的底层通信完全基于 Netty 实现。整体架构Broker 端Netty 服务器负责与客户端的连接请求处理Producer/Consumer 端Netty 客户端负责与 Broker 的通信及请求响应处理Netty 多线程模型RocketMQ 在 Netty 基础之上采用了多线程分离设计将 I/O 线程和业务处理线程分开。核心类类 职责NettyRemotingServer 服务端实现底层基于 ServerBootstrapNettyRemotingClient 客户端实现NettyServerConfig / NettyClientConfig 通信配置连接感知Broker 通过 Netty 的 ChannelInboundHandlerAdapter#channelInactive() 可以实时感知到 Consumer/Producer 的下线。这为 Rebalance 和故障剔除提供了基础。消息过滤源码剖析RocketMQ 支持两种消息过滤方式Tag 过滤和 SQL92 过滤。Tag 过滤根据消息的 Tag 进行过滤性能极高在 ConsumeQueue 中存储了 Tag 的哈希码8 字节过滤时只需比对哈希值一条消息只能有一个 Tag这是它的主要限制SQL92 过滤使用 SQL92 语法作为过滤规则表达式可以过滤消息的属性和 Tag在 SQL 语法中Tag 的属性名称为 TAGS比 Tag 过滤更灵活但性能开销更大需要设置 Broker 配置项 enablePropertyFiltertrue默认为 false两种过滤方式的对比对比维度 Tag 过滤 SQL92 过滤过滤依据 Tag 字符串 用户自定义属性 Tag性能 极高哈希比对 较低解析 SQL 遍历属性灵活性 低只能一个 Tag 高复杂条件组合Broker 配置 默认开启 需 enablePropertyFiltertrue过滤表达式类型在源码中定义为 ExpressionType.TAG 和 ExpressionType.SQL92。SQL92 表达式需要先编译检查合法性再使用编译后的表达式进行计算。源码阅读实战建议读完上面这些模块的源码剖析你可能跃跃欲试了。这里给你几个实战建议搭建源码调试环境从 GitHub 克隆 RocketMQ 源码用 IDEA 导入 Maven 项目先启动 NamesrvStartup再启动 BrokerStartup运行 example 模块中的示例代码进行调试2. 阅读顺序建议阶段 模块 目的第一阶段 remoting 理解网络通信基础第二阶段 namesrv 理解路由注册与发现第三阶段 store 理解存储核心CommitLog ConsumeQueue第四阶段 broker 理解服务端业务逻辑第五阶段 client 理解生产者和消费者3. 调试断点建议Producer 发送DefaultMQProducerImpl#sendDefaultImplConsumer 拉取PullMessageService#runBroker 写入CommitLog#putMessageBroker 拉取PullMessageProcessor#processRequestRebalanceRebalanceService#doRebalance4. 善用日志RocketMQ 的日志非常详细在 ~/logs/rocketmqlogs/ 目录下broker.logBroker 运行日志namesrv.logNameServer 日志store.log存储相关日志rocketmq_client.log客户端日志小结这篇文章我们完整走了一遍 RocketMQ 源码的“地图”通过 8 张流程图 代码片段搞清楚了源码工程结构各模块的职责划分从哪里入手NameServer路由管理的 5 个核心数据结构、心跳注册流程Broker分层架构、启动流程、存储设计理念消息发送4 个核心步骤、Producer 启动流程、定时任务机制消息拉取与消费Push 消费者启动、长轮询的本质CommitLog 写入MappedFile 机制、锁策略、刷盘实现ConsumeQueue 构建ReputMessageService 的“消息重放”机制Rebalance触发入口、分配策略、广播与集群模式的区别事务消息半消息存储、回查机制的完整流程消息重试与死信16 次重试、指数退避、DLQ 处理Netty 通信层多线程模型、连接感知消息过滤Tag 与 SQL92 的原理与对比