ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

RocketMQ NameServer 与 Broker 通信机制源码解析(心跳注册、路由表维护与失效剔除)

RocketMQ NameServer 与 Broker 通信机制源码解析(心跳注册、路由表维护与失效剔除) 文档教程知识库【免费下载链接】source-code-hunter 从源码层面剖析挖掘互联网行业主流技术的底层实现原理为广大开发者 “提升技术深度” 提供便利。目前开放 Spring 全家桶Mybatis、Netty、Dubbo 框架及 Redis、Tomcat 中间件等项目地址https://gitcode.com/doocs/source-code-hunter点击查看免费下载本文基于 RocketMQ 4.9.3 源码从NamesrvController、RouteInfoManager、BrokerController三个核心类入手完整剖析 NameServer 与 Broker 之间基于长连接的心跳注册、路由表读写锁设计以及失效 Broker 剔除机制的底层实现。读完本文你将掌握 RocketMQ 注册中心的通信时序、brokerLiveTable/topicQueueTable等四张路由表的具体维护逻辑以及生产环境中 Broker 上下线时路由元数据如何被正确增删。RocketMQ 中 NameServer 的角色与通信模型RocketMQ 采用无状态注册中心 轻量级 NameServer的设计NameServer 本身不参与消息的存储与转发它的职责是维护 Broker 的路由元数据为生产者和消费者提供 Topic 路由查询能力。生产者和消费者在启动与运行过程中都会周期性地向 NameServer 拉取路由信息相关客户端逻辑见 RocketMQ 消息发送流程 与 RocketMQ 生产者启动流程。NameServer 与 Broker 的通信建立在长连接之上整个协作可以概括为两条核心链路Broker 主动上报Broker 每隔 30 秒向集群中所有 NameServer 发送心跳包完成路由注册registerBroker。NameServer 主动巡检NameServer 每隔 10 秒扫描一次 Broker 状态表移除长时间未上报的失效 Broker。围绕这两条链路NameServer 内部维护了四张核心路由表全部封装在org.apache.rocketmq.namesrv.routeinfo.RouteInfoManager中路由表类型作用topicQueueTableHashMapString, MapString, QueueDataTopic 与消息队列的路由关系按 topicName 索引到该 Topic 下各 Broker 的 QueueDatabrokerAddrTableHashMapString, BrokerDataBroker 名称与 Broker 地址主/从节点地址列表的对应关系clusterAddrTableHashMapString, SetString集群名称与该集群下所有 Broker 名称的集合关系brokerLiveTableHashMapString, BrokerLiveInfo存活的 Broker 地址与其最近心跳时间等状态信息是执行路由删除的重要依据filterServerTableHashMapString, ListStringBroker 关联的 FilterServer 消息过滤服务器地址列表Broker 端心跳注册每 30 秒一次的定时上报Broker 启动后在org.apache.rocketmq.broker.BrokerController的初始化过程中会注册一个定时任务每隔一段时间向集群中所有NameServer 调用registerBrokerAll上报自身信息默认周期为 30 秒且周期会被限制在 10 秒到 60 秒之间this.scheduledExecutorService.scheduleAtFixedRate(new Runnable() { Override public void run() { try { BrokerController.this.registerBrokerAll(true, false, brokerConfig.isForceRegister()); } catch (Throwable e) { log.error(registerBrokerAll Exception, e); } } }, 1000 * 10, Math.max(10000, Math.min(brokerConfig.getRegisterNameServerPeriod(), 60000)), TimeUnit.MILLISECONDS);这段代码对应三个关键参数均可在brokerConfig中配置首次延迟启动后 10 秒首次执行给 Broker 留出初始化完成的时间执行周期Math.max(10000, Math.min(registerNameServerPeriod, 60000))即默认 30 秒、最小 10 秒、最大 60 秒forceRegister由brokerConfig.isForceRegister()控制是否强制注册忽略needRegister的比对结果。registerBrokerAll注册前的预处理registerBrokerAll是同步方法synchronized修饰保证同一时刻一个 Broker 只有一个线程在执行注册逻辑。它的处理流程如下public synchronized void registerBrokerAll(final boolean checkOrderConfig, boolean oneway, boolean forceRegister) { TopicConfigSerializeWrapper topicConfigWrapper this.getTopicConfigManager().buildTopicConfigSerializeWrapper(); if (!PermName.isWriteable(this.getBrokerConfig().getBrokerPermission()) || !PermName.isReadable(this.getBrokerConfig().getBrokerPermission())) { ConcurrentHashMapString, TopicConfig topicConfigTable new ConcurrentHashMapString, TopicConfig(); for (TopicConfig topicConfig : topicConfigWrapper.getTopicConfigTable().values()) { TopicConfig tmp new TopicConfig(topicConfig.getTopicName(), topicConfig.getReadQueueNums(), topicConfig.getWriteQueueNums(), this.brokerConfig.getBrokerPermission()); topicConfigTable.put(topicConfig.getTopicName(), tmp); } topicConfigWrapper.setTopicConfigTable(topicConfigTable); } if (forceRegister || needRegister(this.brokerConfig.getBrokerClusterName(), this.getBrokerAddr(), this.brokerConfig.getBrokerName(), this.brokerConfig.getBrokerId(), this.brokerConfig.getRegisterBrokerTimeoutMills())) { doRegisterBrokerAll(checkOrderConfig, oneway, topicConfigWrapper); } }关键逻辑有两处权限修正如果当前 Broker 被设置为只读或只写BrokerPermission非读写全开则上报的每个 TopicConfig 都会被强制替换为当前 Broker 的权限值防止 Topic 的读队列、写队列权限与 Broker 实际权限不一致注册判定forceRegister为 true 时无条件注册否则调用needRegister比对 Broker 集群名、地址、BrokerName、BrokerId 与上次注册数据是否一致只有发生变化时才真正发起注册doRegisterBrokerAll从而减少无效的心跳网络开销。NameServer 端请求处理从 Remoting 协议到路由表更新Broker 的心跳注册请求通过 Netty 长连接发送到 NameServer 后由org.apache.rocketmq.namesrv.processor.DefaultRequestProcessor负责解析请求类型如果请求类型为RequestCode.REGISTER_BROKER则请求最终转发到org.apache.rocketmq.namesrv.routeinfo.RouteInfoManager#registerBroker随后在RouteInfoManager中完成路由注册与状态更新。registerBroker是整个路由注册的核心方法总共分为五步。第一步加写锁保证集群与 Broker 注册表的一致性路由注册需要加写锁防止并发修改RouteInfoManager中的路由表。首先判断 Broker 所属集群是否存在如果不存在则创建集群然后将 broker 名加入集群this.lock.writeLock().lockInterruptibly(); // 判断 clusterAddrTable 中是否存在该集群不存在则创建 SetString brokerNames this.clusterAddrTable.get(clusterName); if (null brokerNames) { brokerNames new HashSet(); this.clusterAddrTable.put(clusterName, brokerNames); } brokerNames.add(brokerName);clusterAddrTable维护的是集群 → Broker 名称集合的映射这一步为后续集群维度的路由查询和 Broker 下线清理提供了基础索引。第二步维护 BrokerData区分首次注册与非首次注册从brokerAddrTable中根据 broker 名尝试获取 Broker 信息如果不存在则新建BrokerData放入brokerAddrTable此时registerFirst设置为true如果存在直接替换原先的 Broker 信息registerFirst设置为false表示非第一次注册。BrokerData brokerData this.brokerAddrTable.get(brokerName); if (null brokerData) { brokerData new BrokerData(clusterName, brokerName, new HashMap()); this.brokerAddrTable.put(brokerName, brokerData); registerFirst true; } // 维护 brokerData 中的主从地址列表 MapLong, String brokerAddrsMap brokerData.getBrokerAddrs(); brokerAddrsMap.put(brokerId, brokerAddr);BrokerData内部以brokerId为 key 保存地址列表0表示主节点MixAll.MASTER_ID非 0 表示从节点。同一个 BrokerName 下的主从节点共享一份BrokerData这也是后续从节点查找主节点地址的依据。第三步主节点注册时创建或更新 Topic 路由元数据如果 Broker 为主节点brokerId MixAll.MASTER_ID并且 Broker 的 topic 配置信息发生变化topicConfigWrapper.getDataVersion()版本号变化或者是初次注册registerFirst则需要创建或者更新 topic 的路由元数据填充topicQueueTableif (MixAll.MASTER_ID brokerId) { if (registerFirst || !topicConfigWrapper.getDataVersion().equals(this.dataVersion)) { this.dataVersion topicConfigWrapper.getDataVersion(); for (TopicConfig topicConfig : topicConfigWrapper.getTopicConfigTable().values()) { this.createAndUpdateQueueData(brokerName, topicConfig); } } }这里的DataVersion相当于 Topic 配置的版本号Broker 每次注册都会携带NameServer 通过版本号比对即可判断 Topic 配置是否发生变化避免在配置未变更时重复重建路由。createAndUpdateQueueData的实现根据 topicConfig 创建QueueData数据结构然后更新topicQueueTableprivate void createAndUpdateQueueData(final String brokerName, final TopicConfig topicConfig) { QueueData queueData new QueueData(); queueData.setBrokerName(brokerName); queueData.setWriteQueueNums(topicConfig.getWriteQueueNums()); queueData.setReadQueueNums(topicConfig.getReadQueueNums()); queueData.setPerm(topicConfig.getPerm()); queueData.setTopicSysFlag(topicConfig.getTopicSysFlag()); MapString, QueueData queueDataMap this.topicQueueTable.get(topicConfig.getTopicName()); if (null queueDataMap) { queueDataMap new HashMap(); queueDataMap.put(queueData.getBrokerName(), queueData); this.topicQueueTable.put(topicConfig.getTopicName(), queueDataMap); log.info(new topic registered, {} {}, topicConfig.getTopicName(), queueData); } else { QueueData old queueDataMap.put(queueData.getBrokerName(), queueData); if (old ! null !old.equals(queueData)) { log.info(topic changed, {} OLD: {} NEW: {}, topicConfig.getTopicName(), old, queueData); } } }QueueData封装了该 Broker 上某个 Topic 的读写队列数、权限Perm和系统标记是客户端从 NameServer 拉取路由后计算MessageQueue列表的直接依据对应 RocketMQ 消息发送流程 中的tryToFindTopicPublishInfo与topicRouteData2TopicPublishInfo。第四步更新 BrokerLiveInfo刷新存活状态无论是否首次注册每次心跳都需要更新brokerLiveTable记录 Broker 最近一次上报的时间戳这是判断 Broker 是否存活的唯一依据BrokerLiveInfo prevBrokerLiveInfo this.brokerLiveTable.put(brokerAddr, new BrokerLiveInfo( System.currentTimeMillis(), topicConfigWrapper.getDataVersion(), channel, haServerAddr));BrokerLiveInfo包含四个字段lastUpdateTimestamp最近一次心跳时间用于失效判定dataVersionTopic 配置版本号channel与该 Broker 的长连接 ChannelhaServerAddr主节点 HA高可用服务地址。同时如果该 Broker 之前注册过prevBrokerLiveInfo ! null则会将其haServerAddr赋值给本次注册结果保证主从切换后 HA 地址能够被正确传递。第五步注册 FilterServer 地址列表维护主从关系一个 Broker 上会关联多个FilterServer消息过滤服务器因此需要注册 Broker 的过滤器 Server 地址列表。如果此 Broker 为从节点brokerId ! MixAll.MASTER_ID则需要查找该 Broker 的主节点信息并更新对应的masterAddr属性if (MixAll.MASTER_ID ! brokerId) { String masterAddr brokerData.getBrokerAddrs().get(MixAll.MASTER_ID); if (masterAddr ! null) { BrokerLiveInfo brokerLiveInfo this.brokerLiveTable.get(masterAddr); if (brokerLiveInfo ! null) { result.setHaServerAddr(brokerLiveInfo.getHaServerAddr()); result.setMasterAddr(masterAddr); } } }从节点注册时会从brokerAddrTable中取出同 BrokerName 下brokerId 0的主节点地址再从brokerLiveTable中取得主节点的 HA 服务地址一并返回这样从节点就能知道自己所属的主节点为后续主从同步HA做好准备。读写锁保证消息发送高并发与路由注册串行化RouteInfoManager使用了一把ReentrantReadWriteLock来保护上述全部路由表路由注册、注销、失效剔除等写操作获取写锁生产者和消费者通过getTopicRouteInfo、getAllTopicList等查询接口并发读路由表时获取读锁。从源码结构看这一设计达成了两个目标允许多个消息发送者并发读操作保证消息发送时的高并发同一时刻 NameServer 只处理一个 Broker 心跳包多个心跳包请求串行执行避免并发注册导致路由表数据错乱。NameServer 如何剔除失效的 BrokerBroker 的状态信息存储在brokerLiveTable中NameServer 每收到一个心跳包将更新brokerLiveTable中关于 broker 的状态信息以及路由表topicQueueTable、brokerAddrTable、brokerLiveTable、filterServerTable。与之对应的NameServer 也必须具备把失联 Broker 从路由表中清除的能力否则路由表中将残留大量不可用的地址。剔除失效 Broker 有两条触发途径定时扫描被动剔除NameServer 每隔 10 秒扫描一次brokerLiveTable如果BrokerLiveInfo的lastUpdateTimestamp时间戳距当前时间超过 120 秒BROKER_CHANNEL_EXPIRED_TIME则认为 Broker 失效主动注销优雅下线如果 Broker 在正常关闭的情况下会发送unRegisterBroker指令主动注销。定时扫描scanNotActiveBrokerNameServer 在启动时通过NamesrvController注册了周期扫描任务注意这里是先延迟 5 秒然后每隔 10 秒执行一次this.scheduledExecutorService.scheduleAtFixedRate(NamesrvController.this.routeInfoManager::scanNotActiveBroker, 5, 10, TimeUnit.SECONDS);scanNotActiveBroker的实现如下遍历brokerLiveTable将超过 120 秒未上报心跳的 Broker 判定为失效关闭连接并从表中移除public int scanNotActiveBroker() { int removeCount 0; IteratorEntryString, BrokerLiveInfo it this.brokerLiveTable.entrySet().iterator(); while (it.hasNext()) { EntryString, BrokerLiveInfo next it.next(); long last next.getValue().getLastUpdateTimestamp(); if ((last BROKER_CHANNEL_EXPIRED_TIME) System.currentTimeMillis()) { RemotingUtil.closeChannel(next.getValue().getChannel()); it.remove(); log.warn(The broker channel expired, {} {}ms, next.getKey(), BROKER_CHANNEL_EXPIRED_TIME); this.onChannelDestroy(next.getKey(), next.getValue().getChannel()); removeCount; } } return removeCount; }其中BROKER_CHANNEL_EXPIRED_TIME的默认值为120 秒。这里存在一个隐含关系Broker 每 30 秒心跳一次NameServer 以 120 秒为容忍上限相当于允许连续漏掉约 4 次心跳才判定 Broker 失联从而对网络抖动等瞬时故障有一定的容错能力。调用onChannelDestroy的目的是让该 Channel 上相关的路由信息如与该 Broker 关联的消费进度通道等一并得到清理。主动注销unRegisterBroker 指令如果 broker 在正常关闭的情况下会发送unRegisterBroker指令请求类型对应RequestCode.UNREGISTER_BROKER最终同样由DefaultRequestProcessor转发到RouteInfoManager#unregisterBroker。不管是哪一种方式触发的路由删除处理逻辑是一样的都在unregisterBroker中完成共分四步外加锁的释放第一步申请写锁移除 brokerLiveTable 与 filterServerTable 中的信息this.lock.writeLock().lockInterruptibly(); BrokerLiveInfo brokerLiveInfo this.brokerLiveTable.remove(brokerAddr); log.info(unregisterBroker, remove from brokerLiveTable {}, {}, brokerLiveInfo ! null ? OK : Failed, brokerAddr); this.filterServerTable.remove(brokerAddr);先移除存活状态表和 FilterServer 列表brokerLiveInfo ! null与否直接决定了后续日志中该步骤是OK还是Failed。第二步维护 brokerAddrTable找到具体的 broker根据brokerName将其从BrokerData中移除如果移除之后该BrokerData不再包含任何 broker 地址主从全部下线则在brokerAddrTable中移除该 brokerName 对应的数据BrokerData brokerData this.brokerAddrTable.get(brokerName); if (null ! brokerData) { String addr brokerData.getBrokerAddrs().remove(brokerId); log.info(unregisterBroker, remove addr from brokerAddrTable {}, {}, addr ! null ? OK : Failed, brokerAddr); if (brokerData.getBrokerAddrs().isEmpty()) { this.brokerAddrTable.remove(brokerName); log.info(unregisterBroker, remove name from brokerAddrTable OK, {}, brokerName); removeBrokerName true; } }注意这里只移除当前brokerId对应的地址。例如主节点下线而从节点仍在则brokerAddrTable中该 BrokerName 仍会保留只是少了一个主节点地址removeBrokerName保持 false不会触发后续的集群和 Topic 清理。第三步根据 brokerName 维护 clusterAddrTable仅当上一步判定整个 BrokerName 已被移除removeBrokerName true时才将该 BrokerName 从所属集群中移除如果移除后集群中不包含任何 Broker则将该集群从clusterAddrTable中移除if (removeBrokerName) { SetString nameSet this.clusterAddrTable.get(clusterName); if (nameSet ! null) { boolean removed nameSet.remove(brokerName); log.info(unregisterBroker, remove name from clusterAddrTable {}, {}, removed ? OK : Failed, brokerName); if (nameSet.isEmpty()) { this.clusterAddrTable.remove(clusterName); log.info(unregisterBroker, remove cluster from clusterAddrTable {}, clusterName); } } }这样设计保证了集群维度信息的自洽集群 → Broker 集合、Broker → 地址集合两层关系保持级联一致。第四步根据 brokerName 清理 topicQueueTable遍历所有主题的队列如果队列中包含当前 broker 的队列则移除如果某个 topic 的所有队列都被移除queueDataMap.size() 0则从路由表中删除该 topicthis.topicQueueTable.forEach((topic, queueDataMap) - { QueueData old queueDataMap.remove(brokerName); if (old ! null) { log.info(removeTopicByBrokerName, remove one brokers topic {} {}, topic, old); } if (queueDataMap.size() 0) { noBrokerRegisterTopic.add(topic); log.info(removeTopicByBrokerName, remove the topic all queue {}, topic); } });被彻底清空的 topic 会加入noBrokerRegisterTopic集合这些 topic 将不再出现在 NameServer 的路由查询结果中客户端后续也自然无法再路由到这些主题。第五步释放写锁完成路由删除finally { this.lock.writeLock().unlock(); }整个删除过程在try/finally中包住写锁保证无论中间是否出现异常锁都能被正确释放避免死锁导致 NameServer 路由表长期不可写。与客户端侧的联动心跳与路由刷新的完整闭环NameServer 与 Broker 的通信并非孤立存在它和客户端侧的任务形成了完整的闭环。在客户端生产者/消费者共用的MQClientInstance启动过程中会启动一系列周期任务其中与本文主题直接相关的有详见 RocketMQ 生产者启动流程更新 NameServer 地址若未显式配置namesrvAddr每 2 分钟调用一次fetchNameServerAddr从 HTTP 服务获取 NameServer 地址列表更新 Topic 路由信息每 30 秒pollNameServerInterval默认 30000ms调用updateTopicRouteInfoFromNameServer从 NameServer 拉取最新路由并刷新本地缓存清理离线 Broker 并发送心跳每 30 秒heartbeatBrokerInterval默认 30000ms调用cleanOfflineBroker与sendHeartbeatToAllBrokerWithLock将本地不可达的 Broker 清理掉并向所有 Broker 上报客户端心跳。因此RocketMQ 的路由体系实际上是双层心跳 双层淘汰Broker → NameServer 层Broker 30 秒注册一次NameServer 120 秒判定失联并剔除客户端 → NameServer/Broker 层客户端 30 秒拉取一次路由并做本地清理保证消息发送与消费始终命中可用的 Broker 队列。当 Broker 失联被 NameServer 剔除后客户端在下一个 30 秒轮询周期内就会感知到路由变化从而把后续请求发送到其他可用 Broker配合 RocketMQ 消息发送流程 中的故障延迟机制还能进一步规避故障 Broker。这一设计让 RocketMQ 在不引入额外协调组件如 Zookeeper的前提下实现了注册中心的轻量化和路由数据的高可用容错。总结NameServer 与 Broker 的通信是 RocketMQ 路由体系的心脏其设计要点可归纳为长连接 定时心跳Broker 每 30 秒可配置范围 10~60 秒向所有 NameServer 上报NameServer 每 10 秒巡检一次超过 120 秒未上报即判定失效路由表的分层维护注册时按集群 → Broker 名称 → 主从地址 → Topic 队列 → 存活状态 → 过滤服务器逐层写入五张表删除时按相反次序逐层清理并通过DataVersion版本号避免无效重建读写锁的精细并发控制同一时刻只允许一个心跳包写入路由表但允许大量客户端并发读取兼顾了数据一致性与消息发送的高并发主动注销与被动剔除殊途同归无论 Broker 是优雅下线unRegisterBroker还是异常失联scanNotActiveBroker最终都汇聚到同一套路由删除逻辑保证路由表在任何情况下都能保持自洽。对于希望深入 RocketMQ 源码的读者可以从 RouteInfoManager 出发结合 RocketMQ 消息发送流程客户端如何消费路由表、RocketMQ Broker 处理拉取消息请求流程Broker 端如何校验与响应以及 RocketMQ 消息发送存储流程Broker 端落盘链路串联起注册中心 → Broker → 存储的完整数据流。赞分享文档教程知识库【免费下载链接】source-code-hunter 从源码层面剖析挖掘互联网行业主流技术的底层实现原理为广大开发者 “提升技术深度” 提供便利。目前开放 Spring 全家桶Mybatis、Netty、Dubbo 框架及 Redis、Tomcat 中间件等项目地址https://gitcode.com/doocs/source-code-hunter点击查看免费下载相关推荐RocketMQ NameServer 与 Broker 通信机制源码解析心跳注册、路由表维护与失效剔除RocketMQ NameServer 与 Broker 通信机制源码解析心跳注册、路由表维护与失效剔除 本篇文章基于 RocketMQ 4.9.3 源码对文档教程技术博客知识库Apache RocketMQ NameServer终极指南深入解析Topic与Broker注册机制Apache RocketMQ NameServer终极指南深入解析Topic与Broker注册机制 Apache RocketMQ NameServer是分消息队列后端微服务流处理Apache RocketMQ 架构设计详解从 NameServer 路由注册到 Broker 存储高可用Apache RocketMQ 架构设计详解从 NameServer 路由注册到 Broker 存储高可用 本篇文章以 Apache RocketMQ 官方架消息队列流处理后端创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表