MQ系列4:NameServer 原理解析

简介: MQ系列4:NameServer 原理解析

1 关于NameServer

上一节的 MQ系列3:RocketMQ 架构分析,我们大致介绍了 RocketMQ的基本组件构成,包括 NameServer、Broker、Producer以及Consumer四部分。

NameServer,指的是服务可以根据给定的名字来进行资源或对象的地址定位,并获取有关的属性信息。在Rocket中也一样,NameServer是 RocketMQ 的服务注册中心(类似于 Kafka 集群 后面的 Zookeeper 集群一样, 对集群元数据进行管理),根据元数据(ip、port和router信息)来唯一定位服务。RocketMQ 需要先启动 NameServer ,再启动 Rocket 中的 Broker。

2 NameServer运行流程

2.1 注册

注册发生在Broker启动之后,启动后快速与NameServer建立长连接,并每30s对NameService发送一次心跳包,Broker会将自己的IP Address、Port、Router 等信息随着心跳一并注册到 NameServer中。

这里的RouterInfo 主要指Broker下包含哪些Topic信息,这种映射关系方便后面消息的生产和消费的时候进行寻址。

image.png

注册使用到的核心数据结构如下:

HashMap brokerAddrTable

  • HashMap 的 Key 是 Broker 的名称,存储了一个Broker服务所对应的属性信息。
  • Value 是个对象,数据结构如下:

字段

类型

说明

cluster

String

所属的集群名称

broker

String

broker的名称

brokerAddress

HashMap

Broker的IP地址列表,包含一个Master IP地址列表 和 多个Slave IP地址列表

" Broker-A":{ "cluster":"Broker-Cluster", "brokerName":"Broker-A", "cluster":{  // 1主2从    "0":"192.168.0.1:1234",    "1":"192.168.0.2:1234",    "2":"192.168.0.3:1234" } }

2.2 注册信息更新

当你对你的Broker中的Topic信息进行更新了(增、删、改)怎么办,你才需要重新将信息注册到NameServer中。

  • 如果你创建了新的 Topic,Broker会向 NameServer 发送注册信息,接收到信息后会对每个Master 角色的Broker ,创建一个新的QueueData对象。
  • 如果你修改了Topic,则NameServer 会先把旧的 QueueData 删除,在加一个新的 QueueData。
  • 如果你删除了Topic,则NameServer 会将对应的 QueueData 删除。

image.png

使用到的核心数据结构如下:

HashMap> topicQueueTable

  • HashMap 的 Key 是 Topic 的名称,里面存储了Topic的所有属性信息。
  • Value 是个列表,列表的数据类型是 QueueData,列表的length就是Topic中的 Master角色的 Broker 个数。
  • QueueData的结构如下

字段

类型

说明

brokerName

String

broker名称

readQueueNums

Long

读Queue的数量

writeQueueNums

Long

写Queue的数量

perm

Integer

权限 PRIORITY = 3, READ = 2, WRITE = 1 , INHERIT = 0

topicSyncFlag

Long

同步的位置标识

{   "topic-test":[ // topic名称,注意下面会用到    {      "brokerName":"Broker-A", "readQueueNums":37, "writeQueueNums":37, "perm":6,  // 读写权限 "topicSynFlag":12    },    {      "brokerName":"Broker-B", "readQueueNums":37, "writeQueueNums":37, "perm":6,  // 读写权限 "topicSynFlag":12    }   ] }

参考RocketMQ源码如下,这边加了注释,方便理解:

   /**      * 创建或者更新 MessageQueue 的数据      * @param brokerName      * @param topicConfig      */     private void createAndUpdateQueueData(final String brokerName, final TopicConfig topicConfig) {         QueueData queueData = new QueueData();         queueData.setBrokerName(brokerName); // broker 名称         queueData.setWriteQueueNums(topicConfig.getWriteQueueNums());   // 读Queue的数量         queueData.setReadQueueNums(topicConfig.getReadQueueNums());  // 写Queue的数量         queueData.setPerm(topicConfig.getPerm());  // 权限: PRIORITY = 3, READ = 2, WRITE = 1 , INHERIT = 0         queueData.setTopicSynFlag(topicConfig.getTopicSysFlag());         List<QueueData> queueDataList = this.topicQueueTable.get(topicConfig.getTopicName());         if (null == queueDataList) {  // 新增             queueDataList = new LinkedList<QueueData>();             queueDataList.add(queueData);             this.topicQueueTable.put(topicConfig.getTopicName(), queueDataList);             log.info("new topic registerd, {} {}", topicConfig.getTopicName(), queueData);         } else {   // 更新             boolean addNewOne = true;                Iterator<QueueData> it = queueDataList.iterator();             while (it.hasNext()) {                 QueueData qd = it.next();                 if (qd.getBrokerName().equals(brokerName)) {                     if (qd.equals(queueData)) {                         addNewOne = false;                     } else {                         log.info("topic changed, {} OLD: {} NEW: {}", topicConfig.getTopicName(), qd,                                 queueData);                         it.remove();   // 先删除                     }                 }             }             if (addNewOne) {                 queueDataList.add(queueData);   // 再添加             }         }     }

2.3 异常清理

如果Broker挂掉,那么再被消息的生产者和消费者使用就会有问题了。这时候需要对已经宕掉的Broker进行清理,确保NamServer中注册的Broker服务信息都是Alive的。它的做法是这样的:

  • 前面我们说了,Broker每30s对NameService发送一次心跳包给NameServer
  • NameServer接收到心跳包的时候,会将当前时间戳更新到

brokerLiveTable 表的

lastUpdateTimestamp 字段中。

  • NameServer中会启动一个定时任务
  • 每10s(记住这边扫描是10s间隔,与上面心跳包区分开)扫描 一下

brokerLiveTable 表

  • 检查

lastUpdateTimestamp字段,如果时间戳与当前时间相隔超过 120s(即两分钟),则认为 Broker 已经宕了,并会将broker清除出NameServer的注册表。

使用到的核心数据结构如下:

HashMap brokerLiveTable

  • HashMap 的 Key 是 Broker服务器的地址信息(IP+Port),里面存储了该Broker服务器的基本信息。
  • Value 是个对象,结构如下:

字段

类型

说明

lastUpdateTimestamp

Long

最后一次收到心跳包的时间戳

dataVersion

DataVersion

数据版本号对象

channel

Channel

netty的Channel,IO数据交互媒介

haServerAddr

String

master地址,初次请求的时候值为空,slave向NameServer注册之后返回

2.4 消息生产和消费

上面的步骤都完成之后,NameServer这个 "中央大脑" 正式开始投入使用。这时候 ,消息的生产和消费具体是怎么做的呢?

  • Producer 或者 Consumer 启动之后会和 NameServer 建立长连接
  • 定时(默认为每30s)从 NameServer 获取Routers信息,并将路由信息保存至Producer或者Consumer的本地。
  • Producer发送一条消息

hello-brand 到 topic (

topic-test) 中

  • 因为名称为

topic-test 的 topic 存在于多个 broker中,所以需要如下几个步骤,才能找到具体的地址:

  • 先 根据 topic 名称

topic-test 查询

topicQueueTable , 选择一个并获取它的broker信息(包含brokerName)

  • 再根据已经获取到的brokerName 查询

brokerAddressTable 获取具体的Broker IP地址(一般包含1个Master和n个Slave的IP地址)

  • 拿到IP地址之后,生产者与broker建立连接,并发送消息
  • 消费者同理

3 总结

上述的流程图比较清晰的描述如下运转流程:

image.png

  • NameServer 作为整个 RocketMQ 的“中央大脑” ,负责对集群元数据进行管理,所以 RocketMQ 需要先启动 NameServer 再启动 Rocket 中的 Broker。
  • Broker 启动后,与 NameServer 保持长连接,每 30s 发送一次发送心跳包,来确保Broker是否存活。并将 Broker 信息 ( IP+、端口等信息)以及Broker中存储的Topic信息上报。注册成功后,NameServer 集群中就有 Topic 跟 Broker 的映射关系。
  • NameServer有个定时任务,每10s扫描下

brokerLiveTable表 , 如果检测到某个Broker 宕机(因为使用心跳机制, 如果检测超120s(两分钟)无上报心跳),则从路由注册表中将其移除。

  • 生产者在发送某个主题的消息之前先从 NamerServer 获取 Broker 服务器地址列表(通过topic名称查询

topicQueueTable获得broker名称,通过broker名称查询

brokerAddressTable获取具体的Broker IP地址),然后根据负载均衡算法从列表中选择1台Broker ,建立连接通道,进行消息发送。

  • 消费者在订阅某个topic的消息之前从 NamerServer 获取 Broker 服务器地址列表(同上),包括关联的全部Topic队列信息。进而获取当前订阅 Topic 存在哪些 Broker 上,然后直接跟 Broker 建立连接通道,开始消费数据。
  • 生产者和消费者默认每30s 从 NamerServer 获取 Broker 服务器地址列表,以及关联的所有Topic队列信息,更新到Client本地。

参考:

https://zhuanlan.zhihu.com/p/388807516

相关实践学习
快速体验阿里云云消息队列RocketMQ版
本实验将带您快速体验使用云消息队列RocketMQ版Serverless系列实例进行获取接入点、创建Topic、创建订阅组、收发消息、查看消息轨迹和仪表盘。
消息队列 MNS 入门课程
1、消息队列MNS简介 本节课介绍消息队列的MNS的基础概念 2、消息队列MNS特性 本节课介绍消息队列的MNS的主要特性 3、MNS的最佳实践及场景应用 本节课介绍消息队列的MNS的最佳实践及场景应用案例 4、手把手系列:消息队列MNS实操讲 本节课介绍消息队列的MNS的实际操作演示 5、动手实验:基于MNS,0基础轻松构建 Web Client 本节课带您一起基于MNS,0基础轻松构建 Web Client
相关文章
|
运维 持续交付 云计算
深入解析云计算中的微服务架构:原理、优势与实践
深入解析云计算中的微服务架构:原理、优势与实践
932 86
|
消息中间件 存储 缓存
RocketMQ原理—4.消息读写的性能优化
本文详细解析了RocketMQ消息队列的核心原理与性能优化机制,涵盖Producer消息分发、Broker高并发写入、Consumer拉取消息流程等内容。重点探讨了基于队列的消息分发、Hash有序分发、CommitLog内存写入优化、ConsumeQueue物理存储设计等关键技术点。同时分析了数据丢失场景及解决方案,如同步刷盘与JVM OffHeap缓存分离策略,并总结了写入与读取流程的性能优化方法,为理解和优化分布式消息系统提供了全面指导。
RocketMQ原理—4.消息读写的性能优化
|
安全 算法 网络协议
解析:HTTPS通过SSL/TLS证书加密的原理与逻辑
HTTPS通过SSL/TLS证书加密,结合对称与非对称加密及数字证书验证实现安全通信。首先,服务器发送含公钥的数字证书,客户端验证其合法性后生成随机数并用公钥加密发送给服务器,双方据此生成相同的对称密钥。后续通信使用对称加密确保高效性和安全性。同时,数字证书验证服务器身份,防止中间人攻击;哈希算法和数字签名确保数据完整性,防止篡改。整个流程保障了身份认证、数据加密和完整性保护。
|
存储 缓存 算法
HashMap深度解析:从原理到实战
HashMap,作为Java集合框架中的一个核心组件,以其高效的键值对存储和检索机制,在软件开发中扮演着举足轻重的角色。作为一名资深的AI工程师,深入理解HashMap的原理、历史、业务场景以及实战应用,对于提升数据处理和算法实现的效率至关重要。本文将通过手绘结构图、流程图,结合Java代码示例,全方位解析HashMap,帮助读者从理论到实践全面掌握这一关键技术。
580 14
|
消息中间件 存储 设计模式
RocketMQ原理—5.高可用+高并发+高性能架构
本文主要从高可用架构、高并发架构、高性能架构三个方面来介绍RocketMQ的原理。
3723 21
RocketMQ原理—5.高可用+高并发+高性能架构
|
存储 消息中间件 缓存
RocketMQ原理—3.源码设计简单分析下
本文介绍了Producer作为生产者是如何创建出来的、启动时是如何准备好相关资源的、如何从拉取Topic元数据的、如何选择MessageQueue的、与Broker是如何进行网络通信的,Broker收到一条消息后是如何存储的、如何实时更新索引文件的、如何实现同步刷盘以及异步刷盘的、如何清理存储较久的磁盘数据的,Consumer作为消费者是如何创建和启动的、消费者组的多个Consumer会如何分配消息、Consumer会如何从Broker拉取一批消息。
677 11
RocketMQ原理—3.源码设计简单分析下
|
存储 消息中间件 网络协议
RocketMQ原理—1.RocketMQ整体运行原理
本文详细解析了RocketMQ的整体运行原理,涵盖从生产者到消费者的全流程。首先介绍生产者发送消息的机制,包括Topic与MessageQueue的关系及写入策略;接着分析Broker如何通过CommitLog和ConsumeQueue实现消息持久化,并探讨同步与异步刷盘的优缺点。同时,讲解基于DLedger技术的主从同步原理,确保高可用性。消费者部分则重点讨论消费模式(集群 vs 广播)、拉取消息策略及负载均衡机制。网络通信层面,基于Netty的高性能架构通过多线程池分工协作提升并发能力。最后,揭示mmap与PageCache技术优化文件读写的细节,总结了RocketMQ的核心运行机制。
RocketMQ原理—1.RocketMQ整体运行原理
|
机器学习/深度学习 算法 数据挖掘
解析静态代理IP改善游戏体验的原理
静态代理IP通过提高网络稳定性和降低延迟,优化游戏体验。具体表现在加快游戏网络速度、实时玩家数据分析、优化游戏设计、简化更新流程、维护网络稳定性、提高连接可靠性、支持地区特性及提升访问速度等方面,确保更流畅、高效的游戏体验。
425 22
解析静态代理IP改善游戏体验的原理
|
编解码 缓存 Prometheus
「ximagine」业余爱好者的非专业显示器测试流程规范,同时也是本账号输出内容的数据来源!如何测试显示器?荒岛整理总结出多种测试方法和注意事项,以及粗浅的原理解析!
本期内容为「ximagine」频道《显示器测试流程》的规范及标准,我们主要使用Calman、DisplayCAL、i1Profiler等软件及CA410、Spyder X、i1Pro 2等设备,是我们目前制作内容数据的重要来源,我们深知所做的仍是比较表面的活儿,和工程师、科研人员相比有着不小的差距,测试并不复杂,但是相当繁琐,收集整理测试无不花费大量时间精力,内容不完善或者有错误的地方,希望大佬指出我们好改进!
1504 16
「ximagine」业余爱好者的非专业显示器测试流程规范,同时也是本账号输出内容的数据来源!如何测试显示器?荒岛整理总结出多种测试方法和注意事项,以及粗浅的原理解析!
|
机器学习/深度学习 数据可视化 PyTorch
深入解析图神经网络注意力机制:数学原理与可视化实现
本文深入解析了图神经网络(GNNs)中自注意力机制的内部运作原理,通过可视化和数学推导揭示其工作机制。文章采用“位置-转移图”概念框架,并使用NumPy实现代码示例,逐步拆解自注意力层的计算过程。文中详细展示了从节点特征矩阵、邻接矩阵到生成注意力权重的具体步骤,并通过四个类(GAL1至GAL4)模拟了整个计算流程。最终,结合实际PyTorch Geometric库中的代码,对比分析了核心逻辑,为理解GNN自注意力机制提供了清晰的学习路径。
1024 7
深入解析图神经网络注意力机制:数学原理与可视化实现

热门文章

最新文章

推荐镜像

更多
  • DNS