<kbd id="afajh"><form id="afajh"></form></kbd>
<strong id="afajh"><dl id="afajh"></dl></strong>
    <del id="afajh"><form id="afajh"></form></del>
        1. <th id="afajh"><progress id="afajh"></progress></th>
          <b id="afajh"><abbr id="afajh"></abbr></b>
          <th id="afajh"><progress id="afajh"></progress></th>

          一文讀懂RocketMQ的高可用機(jī)制——集群管理高可用

          共 29477字,需瀏覽 59分鐘

           ·

          2021-03-26 20:51

          點(diǎn)擊上方老周聊架構(gòu)關(guān)注我



          一、前言

          在前三篇我們介紹了

          一文讀懂RocketMQ的高可用機(jī)制——消息存儲(chǔ)高可用

          一文讀懂RocketMQ的高可用機(jī)制——消息發(fā)送高可用

          一文讀懂RocketMQ的高可用機(jī)制——消息消費(fèi)高可用


          這一篇我們來說一下集群管理是如何保證高可用的。

          集群管理的高可用主要體現(xiàn)在 NameServer 的設(shè)計(jì)上,NameServer 承擔(dān)著路由注冊(cè)中心的作用。當(dāng)部分 NameServer 節(jié)點(diǎn)宕機(jī)時(shí)不會(huì)有什么糟糕的影響,只剩一個(gè) NameServer 節(jié)點(diǎn) RocketMQ 集群也能正常運(yùn)行,即使 NameServer 全部宕機(jī),也不影響已經(jīng)運(yùn)行的 Broker、Producer 和 Consumer。

          在說 NameServer 之前,我們是否有以下幾點(diǎn)思考。既然作為路由注冊(cè)中心,那有哪些路由信息注冊(cè)到了 NameServer?生產(chǎn)者如何知道消息要發(fā)送到哪臺(tái)消息服務(wù)器呢?當(dāng) Broker 不可用后,NameServer 并不會(huì)立即將變更后的注冊(cè)信息推送至 Client(Producer/Consumer),此時(shí) RocketMQ 如何保證 Client 正常發(fā)送/消費(fèi)消息?

          帶著這幾個(gè)疑問來開啟我們的 RocketMQ 集群管理高可用之旅。

          二、架構(gòu)設(shè)計(jì)

          1、NameServer 互相獨(dú)立,彼此沒有通信關(guān)系,單臺(tái) NameServer 掛掉,不影響其他 NameServer。

          2、NameServer 不去連接別的機(jī)器,不主動(dòng)推消息。

          3、單個(gè) Broker(Master、Slave) 與所有 NameServer 進(jìn)行定時(shí)注冊(cè),以便告知 NameServer 自己還活著。

          • Broker 每隔 30 秒向所有 NameServer 發(fā)送心跳,心跳包含了自身的 topic 配置信息。

          • NameServer 每隔 10 秒,掃描所有還存活的 broker 連接,如果某個(gè)連接的最后更新時(shí)間與當(dāng)前時(shí)間差值超過 2 分鐘,則斷開此連接,NameServer 也會(huì)斷開此 broker 下所有與 slave 的連接。同時(shí)更新 topic 與隊(duì)列的對(duì)應(yīng)關(guān)系,但不通知生產(chǎn)者和消費(fèi)者。

          • Broker slave 同步或者異步從 Broker master 上拷貝數(shù)據(jù)。


          4、Consumer 隨機(jī)與一個(gè) NameServer 建立長(zhǎng)連接,如果該 NameServer 斷開,則從 NameServer 列表中查找下一個(gè)進(jìn)行連接。

          • Consumer 主要從 NameServer 中根據(jù) Topic 查詢 Broker 的地址,查到就會(huì)緩存到客戶端,并向提供 Topic 服務(wù)的 Master、Slave 建立長(zhǎng)連接,且定時(shí)向 Master、Slave 發(fā)送心跳。

          • 如果 Broker 宕機(jī),則 NameServer 會(huì)將其剔除,而 Consumer 端的定時(shí)任務(wù) MQClientInstance.this.updateTopicRouteInfoFromNameServer 每 30 秒執(zhí)行一次,將 Topic 對(duì)應(yīng)的 Broker 地址拉取下來,此地址只有 Slave 地址了,此時(shí) Consumer 從 Slave 上消費(fèi)。

          • 消費(fèi)者與 Master 和 Slave 都建有連接,在不同場(chǎng)景有不同的消費(fèi)規(guī)則。


          5、Producer 隨機(jī)與一個(gè) NameServer 建立長(zhǎng)連接,每隔 30 秒(此處時(shí)間可配置)從 NameServer 獲取 Topic 的最新隊(duì)列情況,如果某個(gè) Broker Master 宕機(jī),Producer 最多 30 秒才能感知,在這個(gè)期間,發(fā)往該 broker master 的消息失敗。Producer 向提供 Topic 服務(wù)的 Master 建立長(zhǎng)連接,且定時(shí)向 Master 發(fā)送心跳。

          • 生產(chǎn)者與所有的 master 連接,但不能向 slave 寫入。

          • 客戶端是先從 NameServer 尋址的,得到可用 Broker 的 IP 和端口信息,然后據(jù)此信息連接 broker。


          綜上所述,NameServer 在 RocketMQ 中的作用:

          • NameServer 用來保存活躍的 broker 列表,包括 Master 和 Slave 。

          • NameServer 用來保存所有 topic 和該 topic 所有隊(duì)列的列表。

          • NameServer 用來保存所有 broker 的 Filter 列表。

          • 命名服務(wù)器為客戶端,包括生產(chǎn)者,消費(fèi)者和命令行客戶端提供最新的路由信息。

          三、啟動(dòng)流程

          NameServer 啟動(dòng)流程重點(diǎn)關(guān)注兩部分:路由信息維護(hù)和網(wǎng)絡(luò)通信(包括心跳),涉及到的核心類如下圖所示。NameServerStartup 是 NameServer 的啟動(dòng)類,負(fù)責(zé)解析配置文件、加載運(yùn)行時(shí)參數(shù)信息和初始化并啟動(dòng) NameServerController。NameServerController 是 NameServer 的核心控制器,其通過 RouteInfoManager 管理路由信息,通過 RemotingServer 與 RocketMQ 其他組件(Broker、Producer 和 Consumer)通信。


          NameServer 的配置參數(shù)包括 NamesrvConfig 和 NettyServerConfig:

          NamesrvConfig 為 NameServer 業(yè)務(wù)參數(shù),如 RocketMQ 主目錄路徑、KV 配置屬性持久化路徑、是否支持順序消息等;

          NettyServerConfig 為 NameServer 網(wǎng)絡(luò)參數(shù),如監(jiān)聽端口、IO 線程池線程個(gè)數(shù)、業(yè)務(wù)線程池線程個(gè)數(shù)、是否開啟 epoll 等。

          在分析源碼之前我們先來看張時(shí)序圖:

          啟動(dòng)類:

          org.apache.rocketmq.namesrv.NamesrvStartup

          1、解析配置文件,填充 NameServerConfig、NettyServerConfig 屬性值,并創(chuàng)建 NamesrvController。

          代碼:

          NamesrvStartup#createNamesrvController

          public static NamesrvController createNamesrvController(String[] args) throws IOException, JoranException {    System.setProperty(RemotingCommand.REMOTING_VERSION_KEY, Integer.toString(MQVersion.CURRENT_VERSION));
          Options options = ServerUtil.buildCommandlineOptions(new Options()); commandLine = ServerUtil.parseCmdLine("mqnamesrv", args, buildCommandlineOptions(options), new PosixParser()); if (null == commandLine) { System.exit(-1); return null; } // 創(chuàng)建NamesrvConfig final NamesrvConfig namesrvConfig = new NamesrvConfig(); // 創(chuàng)建NettyServerConfig final NettyServerConfig nettyServerConfig = new NettyServerConfig(); // 設(shè)置啟動(dòng)端口號(hào) nettyServerConfig.setListenPort(9876); // 解析啟動(dòng)-c參數(shù) if (commandLine.hasOption('c')) { String file = commandLine.getOptionValue('c'); if (file != null) { InputStream in = new BufferedInputStream(new FileInputStream(file)); properties = new Properties(); properties.load(in); MixAll.properties2Object(properties, namesrvConfig); MixAll.properties2Object(properties, nettyServerConfig);
          namesrvConfig.setConfigStorePath(file);
          System.out.printf("load config properties file OK, %s%n", file); in.close(); } } // 解析啟動(dòng)-p參數(shù) if (commandLine.hasOption('p')) { InternalLogger console = InternalLoggerFactory.getLogger(LoggerName.NAMESRV_CONSOLE_NAME); MixAll.printObjectProperties(console, namesrvConfig); MixAll.printObjectProperties(console, nettyServerConfig); System.exit(0); } // 將啟動(dòng)參數(shù)填充到namesrvConfig,nettyServerConfig MixAll.properties2Object(ServerUtil.commandLine2Properties(commandLine), namesrvConfig);
          if (null == namesrvConfig.getRocketmqHome()) { System.out.printf("Please set the %s variable in your environment to match the location of the RocketMQ installation%n", MixAll.ROCKETMQ_HOME_ENV); System.exit(-2); }
          LoggerContext lc = (LoggerContext) LoggerFactory.getILoggerFactory(); JoranConfigurator configurator = new JoranConfigurator(); configurator.setContext(lc); lc.reset(); configurator.doConfigure(namesrvConfig.getRocketmqHome() + "/conf/logback_namesrv.xml");
          log = InternalLoggerFactory.getLogger(LoggerName.NAMESRV_LOGGER_NAME);
          MixAll.printObjectProperties(log, namesrvConfig); MixAll.printObjectProperties(log, nettyServerConfig); // 創(chuàng)建NameServerController final NamesrvController controller = new NamesrvController(namesrvConfig, nettyServerConfig);
          // remember all configs to prevent discard controller.getConfiguration().registerConfig(properties);
          return controller;}


          2、NamesrvConfig 屬性

          • rocketmqHome:rocketmq 主目錄

          • kvConfig:NameServer 存儲(chǔ) KV 配置屬性的持久化路徑

          • configStorePath:nameServer 默認(rèn)配置文件路徑

          • orderMessageEnable:是否支持順序消息


          3、NettyServerConfig 屬性

          • listenPort:NameServer 監(jiān)聽端口,該值默認(rèn)會(huì)被初始化為 9876;

          • serverWorkerThreads:Netty 業(yè)務(wù)線程池線程個(gè)數(shù);

          • serverCallbackExecutorThreads:Netty public 任務(wù)線程池線程個(gè)數(shù), Netty 網(wǎng)絡(luò)設(shè)計(jì),根據(jù)業(yè)務(wù)類型會(huì)創(chuàng)建不同的線程池,比如處理消息發(fā)送、消息消費(fèi)、心跳檢測(cè)等。如果該業(yè)務(wù)類型未注冊(cè)線程池,則由 public 線程池執(zhí)行;

          • serverSelectorThreads:IO 線程池個(gè)數(shù),主要是 NameServer、Broker 端解析請(qǐng)求、返回相應(yīng)的線程個(gè)數(shù),這類線程主要是處理網(wǎng)路請(qǐng)求的,解析請(qǐng)求包,然后轉(zhuǎn)發(fā)到各個(gè)業(yè)務(wù)線程池完成具體的操作,然后將結(jié)果返回給調(diào)用方;

          • serverOnewaySemaphoreValue:send oneway 消息請(qǐng)求并發(fā)讀(Broker端參數(shù));

          • serverAsyncSemaphoreValue:異步消息發(fā)送最大并發(fā)度;

          • serverChannelMaxIdleTimeSeconds:網(wǎng)絡(luò)連接最大的空閑時(shí)間,默認(rèn) 120s

          • serverSocketSndBufSize:網(wǎng)絡(luò)socket發(fā)送緩沖區(qū)大小;

          • serverSocketRcvBufSize:網(wǎng)絡(luò)接收端緩存區(qū)大小;

          • serverPooledByteBufAllocatorEnable:ByteBuffer 是否開啟緩存;

          • useEpollNativeSelector:是否啟用 Epoll IO 模型


          4、根據(jù)啟動(dòng)屬性創(chuàng)建 NamesrvController 實(shí)例,并初始化該實(shí)例。NameServerController 實(shí)例為 NameServer 核心控制器。

          代碼:

          NamesrvController#initialize

          public boolean initialize() {    // 加載KV配置    this.kvConfigManager.load();    // 創(chuàng)建NettyServer網(wǎng)絡(luò)處理對(duì)象    this.remotingServer = new NettyRemotingServer(this.nettyServerConfig, this.brokerHousekeepingService);    // 開啟定時(shí)任務(wù):每隔10s掃描一次Broker,移除不活躍的Broker    this.remotingExecutor =        Executors.newFixedThreadPool(nettyServerConfig.getServerWorkerThreads(), new ThreadFactoryImpl("RemotingExecutorThread_"));
          this.registerProcessor();
          this.scheduledExecutorService.scheduleAtFixedRate(new Runnable() {
          @Override public void run() { NamesrvController.this.routeInfoManager.scanNotActiveBroker(); } }, 5, 10, TimeUnit.SECONDS); // 開啟定時(shí)任務(wù):每隔10min打印一次KV配置 this.scheduledExecutorService.scheduleAtFixedRate(new Runnable() {
          @Override public void run() { NamesrvController.this.kvConfigManager.printAllPeriodically(); } }, 1, 10, TimeUnit.MINUTES);
          if (TlsSystemConfig.tlsMode != TlsMode.DISABLED) {        // Register a listener to reload SslContext try { fileWatchService = new FileWatchService( new String[] { TlsSystemConfig.tlsServerCertPath, TlsSystemConfig.tlsServerKeyPath, TlsSystemConfig.tlsServerTrustCertPath }, new FileWatchService.Listener() { boolean certChanged, keyChanged = false; @Override public void onChanged(String path) { if (path.equals(TlsSystemConfig.tlsServerTrustCertPath)) { log.info("The trust certificate changed, reload the ssl context"); reloadServerSslContext(); } if (path.equals(TlsSystemConfig.tlsServerCertPath)) { certChanged = true; } if (path.equals(TlsSystemConfig.tlsServerKeyPath)) { keyChanged = true; } if (certChanged && keyChanged) { log.info("The certificate and private key changed, reload the ssl context"); certChanged = keyChanged = false; reloadServerSslContext(); } } private void reloadServerSslContext() { ((NettyRemotingServer) remotingServer).loadSslContext(); } }); } catch (Exception e) { log.warn("FileWatchService created error, can't load the certificate dynamically"); } }
          return true;}


          5、在 JVM 進(jìn)程關(guān)閉之前,先將線程池關(guān)閉,及時(shí)釋放資源。

          代碼:

          NamesrvStartup#start

          public static NamesrvController start(final NamesrvController controller) throws Exception {    if (null == controller) {        throw new IllegalArgumentException("NamesrvController is null");    }
          boolean initResult = controller.initialize(); if (!initResult) { controller.shutdown(); System.exit(-3); } // 注冊(cè)JVM鉤子函數(shù)代碼 Runtime.getRuntime().addShutdownHook(new ShutdownHookThread(log, new Callable<Void>() { @Override public Void call() throws Exception { // 釋放資源 controller.shutdown(); return null; } }));
          controller.start();
          return controller;}


          四、路由管理

          NameServer 作為路由注冊(cè)中心,其核心作用是為 Client 提供消息發(fā)送/消費(fèi)的路由信息。Master Broker 節(jié)點(diǎn)和所有 NameServer 建立長(zhǎng)連接,每個(gè) NameServer 節(jié)點(diǎn)擁有所有 Topic 對(duì)應(yīng) Queue 以及 Broker 的映射關(guān)系,包括路由注冊(cè)、路由刪除等,具體由 RouteInfoManager 負(fù)責(zé)管理。

          1、路由元信息

          代碼:

          org.apache.rocketmq.namesrv.routeinfo.RouteInfoManager

          • topicQueueTable

            維護(hù)了 Topic 和其對(duì)應(yīng)消息隊(duì)列的映射關(guān)系,QueueData 記錄了一條隊(duì)列的元信息:所在 Broker、讀隊(duì)列數(shù)量、寫隊(duì)列數(shù)量等。

          • brokerAddrTable

            維護(hù)了 Broker Name 和 Broker 元信息的映射關(guān)系,Broker 通常以 Master-Slave 架構(gòu)部署,BrokerData 記錄了同一個(gè) Broker Name 下所有節(jié)點(diǎn)的地址信息。

          • clusterAddrTable

            維護(hù)了 Broker 的集群信息。

          • brokerLiveTable

            維護(hù)了 Broker 的存活信息。NameServer 在收到來自 Broker 的心跳消息后,更新 BrokerLiveInfo 中的 lastUpdateTimestamp,如果 NameServer 長(zhǎng)時(shí)間未收到 Broker 的心跳信息,NameServer 就會(huì)將其移除。

          • filterServerTable

            Broker 上的 FilterServer 列表,用于類模式消息過濾。


          還是有點(diǎn)抽象是不是,沒關(guān)系,看下面這張圖與上面的存儲(chǔ)關(guān)系對(duì)比下你就很清晰了。

          2、NameServer 請(qǐng)求處理

          NettyRequestProcessor 定義了 RocketMQ 處理網(wǎng)絡(luò)請(qǐng)求的接口,DefaultRequestProcessor 實(shí)現(xiàn)了該接口并負(fù)責(zé)處理 NameServer 收到的網(wǎng)絡(luò)請(qǐng)求。NettyRequestProcessor 定義如下。

          // org.apache.rocketmq.remoting.netty.NettyRequestProcessorpublic interface NettyRequestProcessor {    RemotingCommand processRequest(ChannelHandlerContext ctx, RemotingCommand request)        throws Exception;
          boolean rejectRequest();}


          RemotingCommand 包含請(qǐng)求編碼(code),針對(duì)不同的請(qǐng)求編碼執(zhí)行對(duì)應(yīng)的請(qǐng)求處理邏輯。請(qǐng)求編碼定義在 RequestCode 常量類中,本節(jié)以 REGISTER_BROKER(Broker 注冊(cè) & 心跳)為例,結(jié)合 Broker 的啟動(dòng)流程闡述請(qǐng)求從被 Broker 發(fā)出到被 NameServer 處理的流程。

          Broker 啟動(dòng)類 BrokerStartup 在初始化核心控制器 BrokerController 階段注冊(cè)定時(shí)任務(wù),定時(shí)發(fā)送 HTTP 請(qǐng)求獲取 NameServer 地址列表并存儲(chǔ)于 remotingClient。此處考慮一個(gè)問題:在棄用 ZooKeeper 后,RocketMQ 不存在注冊(cè)中心供 NameServer 注冊(cè),那么 NameServer 地址列表是如何維護(hù)的?在獲取過時(shí)的地址列表后,RocketMQ 如何持續(xù)保證可用性?BrokerController 定時(shí)獲取地址列表核心邏輯精簡(jiǎn)如下。

          org.apache.rocketmq.broker.BrokerController#initialize

          org.apache.rocketmq.broker.out.BrokerOuterAPI#fetchNameServerAddr


          org.apache.rocketmq.common.namesrv.TopAddressing#fetchNSAddr(boolean, long)

          Broker 在獲取到 NameServer 地址列表后,針對(duì)每個(gè)地址開啟一個(gè)線程,將自身信息同步(默認(rèn))注冊(cè)至 NameServer。Broker 利用 CountDownLatch 等待所有線程注冊(cè)工作完成后,繼續(xù)執(zhí)行后續(xù)的工作。下面我們就來看下 Broker 是如何路由注冊(cè)到 NameServer 上的。

          3、路由注冊(cè)

          3.1 發(fā)送心跳包

          RocketMQ 路由注冊(cè)是通過 Broker 與 NameServer 的心跳功能實(shí)現(xiàn)的。Broker 啟動(dòng)時(shí)向集群中所有的 NameServer 發(fā)送心跳信息,每隔 30s 向集群中所有 NameServer 發(fā)送心跳包,NameServer 收到心跳包時(shí)會(huì)更新 brokerLiveTable 緩存中 BrokerLiveInfo 的 lastUpdataTimeStamp 信息,然后 NameServer 每隔 10s 掃描 brokerLiveTable,如果連續(xù) 120s 沒有收到心跳包,NameServer 將移除 Broker 的路由信息同時(shí)關(guān)閉 Socket 連接。

          org.apache.rocketmq.broker.BrokerController#start


          org.apache.rocketmq.broker.out.BrokerOuterAPI#registerBrokerAll

          public List<RegisterBrokerResult> registerBrokerAll(    final String clusterName,    final String brokerAddr,    final String brokerName,    final long brokerId,    final String haServerAddr,    final TopicConfigSerializeWrapper topicConfigWrapper,    final List<String> filterServerList,    final boolean oneway,    final int timeoutMills,    final boolean compressed) {
          final List<RegisterBrokerResult> registerBrokerResultList = Lists.newArrayList(); // 獲得nameServer地址信息 List<String> nameServerAddressList = this.remotingClient.getNameServerAddressList(); // 遍歷所有nameserver列表 if (nameServerAddressList != null && nameServerAddressList.size() > 0) { // 封裝請(qǐng)求頭 final RegisterBrokerRequestHeader requestHeader = new RegisterBrokerRequestHeader(); requestHeader.setBrokerAddr(brokerAddr); requestHeader.setBrokerId(brokerId); requestHeader.setBrokerName(brokerName); requestHeader.setClusterName(clusterName); requestHeader.setHaServerAddr(haServerAddr); requestHeader.setCompressed(compressed); // 封裝請(qǐng)求體 RegisterBrokerBody requestBody = new RegisterBrokerBody(); requestBody.setTopicConfigSerializeWrapper(topicConfigWrapper); requestBody.setFilterServerList(filterServerList); final byte[] body = requestBody.encode(compressed); final int bodyCrc32 = UtilAll.crc32(body); requestHeader.setBodyCrc32(bodyCrc32); final CountDownLatch countDownLatch = new CountDownLatch(nameServerAddressList.size()); for (final String namesrvAddr : nameServerAddressList) { brokerOuterExecutor.execute(new Runnable() { @Override public void run() { try { // 分別向NameServer注冊(cè) RegisterBrokerResult result = registerBroker(namesrvAddr,oneway, timeoutMills,requestHeader,body); if (result != null) { registerBrokerResultList.add(result); }
          log.info("register broker[{}]to name server {} OK", brokerId, namesrvAddr); } catch (Exception e) { log.warn("registerBroker Exception, {}", namesrvAddr, e); } finally { countDownLatch.countDown(); } } }); }
          try { countDownLatch.await(timeoutMills, TimeUnit.MILLISECONDS); } catch (InterruptedException e) { } }
          return registerBrokerResultList;}


          org.apache.rocketmq.broker.out.BrokerOuterAPI#registerBroker

          3.2 處理心跳包

          DefaultRequestProcessor 網(wǎng)路處理類解析請(qǐng)求 類型,如果請(qǐng)求類型是為 REGISTER_BROKER,則將請(qǐng)求轉(zhuǎn)發(fā)到 RouteInfoManager#regiesterBroker。

          org.apache.rocketmq.namesrv.processor.DefaultRequestProcessor#processRequest


          org.apache.rocketmq.namesrv.routeinfo.RouteInfoManager#registerBroker

          維護(hù)路由信息

          public RegisterBrokerResult registerBroker(    final String clusterName,    final String brokerAddr,    final String brokerName,    final long brokerId,    final String haServerAddr,    final TopicConfigSerializeWrapper topicConfigWrapper,    final List<String> filterServerList,    final Channel channel) {    RegisterBrokerResult result = new RegisterBrokerResult();    try {        try {            // 加鎖            this.lock.writeLock().lockInterruptibly();            // 維護(hù)clusterAddrTable            Set<String> brokerNames = this.clusterAddrTable.get(clusterName);            if (null == brokerNames) {                brokerNames = new HashSet<String>();                this.clusterAddrTable.put(clusterName, brokerNames);            }            brokerNames.add(brokerName);
          boolean registerFirst = false; // 維護(hù)brokerAddrTable BrokerData brokerData = this.brokerAddrTable.get(brokerName); // 第一次注冊(cè),則創(chuàng)建brokerData if (null == brokerData) { registerFirst = true; brokerData = new BrokerData(clusterName, brokerName, new HashMap<Long, String>()); this.brokerAddrTable.put(brokerName, brokerData); } // 非第一次注冊(cè),更新Broker Map<Long, String> brokerAddrsMap = brokerData.getBrokerAddrs(); //Switch slave to master: first remove <1, IP:PORT> in namesrv, then add <0, IP:PORT> //The same IP:PORT must only have one record in brokerAddrTable Iterator<Entry<Long, String>> it = brokerAddrsMap.entrySet().iterator(); while (it.hasNext()) { Entry<Long, String> item = it.next(); if (null != brokerAddr && brokerAddr.equals(item.getValue()) && brokerId != item.getKey()) { it.remove(); } }
          String oldAddr = brokerData.getBrokerAddrs().put(brokerId, brokerAddr); registerFirst = registerFirst || (null == oldAddr);            // 維護(hù)topicQueueTable if (null != topicConfigWrapper && MixAll.MASTER_ID == brokerId) { if (this.isBrokerTopicConfigChanged(brokerAddr, topicConfigWrapper.getDataVersion()) || registerFirst) { ConcurrentMap<String, TopicConfig> tcTable = topicConfigWrapper.getTopicConfigTable(); if (tcTable != null) { for (Map.Entry<String, TopicConfig> entry : tcTable.entrySet()) { this.createAndUpdateQueueData(brokerName, entry.getValue()); } } } } // 維護(hù)brokerLiveTable BrokerLiveInfo prevBrokerLiveInfo = this.brokerLiveTable.put(brokerAddr, new BrokerLiveInfo( System.currentTimeMillis(), topicConfigWrapper.getDataVersion(), channel, haServerAddr)); if (null == prevBrokerLiveInfo) { log.info("new broker registered, {} HAServer: {}", brokerAddr, haServerAddr); } // 維護(hù)filterServerList if (filterServerList != null) { if (filterServerList.isEmpty()) { this.filterServerTable.remove(brokerAddr); } else { this.filterServerTable.put(brokerAddr, filterServerList); } }
          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); } } } } finally { this.lock.writeLock().unlock(); } } catch (Exception e) { log.error("registerBroker Exception", e); }
          return result;}


          4、路由刪除

          Broker 每隔 30s 向 NameServer 發(fā)送一個(gè)心跳包,心跳包包含 BrokerId、Broker 地址、Broker 名稱, Broker 所屬集群名稱、 Broker 關(guān)聯(lián)的 FilterServer 列表。但是如果 Broker 宕機(jī),NameServer 無法收到心跳包,此時(shí) NameServer 如何來剔除這些失效的 Broker 呢? NameServer 會(huì)每隔 10s 掃描 brokerLiveTable 狀態(tài)表,如果 BrokerLive 的 lastUpdateTimestamp 的時(shí)間戳距當(dāng)前時(shí)間超過 120s,則認(rèn)為 Broker 失效,移除該 Broker ,關(guān)閉與 Broker 連接,同時(shí)更新 topicQueueTable、brokerAddrTable 、brokerLiveTable、filterServerTable 。

          RocketMQ 有兩個(gè)觸發(fā)點(diǎn)來刪除路由信息:

          • NameServer 定期掃描 brokerLiveTable 檢測(cè)上次心跳包與當(dāng)前系統(tǒng)的時(shí)間差,如果時(shí)間超過 120s,則需要移除 broker。

          • Broker 在正常關(guān)閉的情況下,會(huì)執(zhí)行 unregisterBroker 指令。

          這兩種方式路由刪除的方法都是一樣的,就是從相關(guān)路由表中刪除與該 broker 相關(guān)的信息。


          org.apache.rocketmq.namesrv.NamesrvController#initialize


          public void scanNotActiveBroker() {  // 獲得brokerLiveTable  Iterator<Entry<String, BrokerLiveInfo>> it = this.brokerLiveTable.entrySet().iterator();  // 遍歷brokerLiveTable  while (it.hasNext()) {      Entry<String, BrokerLiveInfo> next = it.next();      long last = next.getValue().getLastUpdateTimestamp();      // 如果收到心跳包的時(shí)間距當(dāng)時(shí)時(shí)間是否超過120s      if ((last + BROKER_CHANNEL_EXPIRED_TIME) < System.currentTimeMillis()) {          // 關(guān)閉連接          RemotingUtil.closeChannel(next.getValue().getChannel());          // 移除broker          it.remove();          log.warn("The broker channel expired, {} {}ms", next.getKey(), BROKER_CHANNEL_EXPIRED_TIME);          // 維護(hù)路由表          this.onChannelDestroy(next.getKey(), next.getValue().getChannel());      }  }}


          public void onChannelDestroy(String remoteAddr, Channel channel) {    String brokerAddrFound = null;    if (channel != null) {        try {            try {                this.lock.readLock().lockInterruptibly();                Iterator<Entry<String, BrokerLiveInfo>> itBrokerLiveTable =                    this.brokerLiveTable.entrySet().iterator();                while (itBrokerLiveTable.hasNext()) {                    Entry<String, BrokerLiveInfo> entry = itBrokerLiveTable.next();                    if (entry.getValue().getChannel() == channel) {                        brokerAddrFound = entry.getKey();                        break;                    }                }            } finally {                this.lock.readLock().unlock();            }        } catch (Exception e) {            log.error("onChannelDestroy Exception", e);        }    }
          if (null == brokerAddrFound) { brokerAddrFound = remoteAddr; } else { log.info("the broker's channel destroyed, {}, clean it's data structure at once", brokerAddrFound); }
          if (brokerAddrFound != null && brokerAddrFound.length() > 0) {
          try { try { // 申請(qǐng)寫鎖,根據(jù)brokerAddress從brokerLiveTable和filterServerTable移除 this.lock.writeLock().lockInterruptibly(); this.brokerLiveTable.remove(brokerAddrFound); this.filterServerTable.remove(brokerAddrFound); // 維護(hù)brokerAddrTable String brokerNameFound = null; boolean removeBrokerName = false; Iterator<Entry<String, BrokerData>> itBrokerAddrTable = this.brokerAddrTable.entrySet().iterator(); // 遍歷brokerAddrTable while (itBrokerAddrTable.hasNext() && (null == brokerNameFound)) { BrokerData brokerData = itBrokerAddrTable.next().getValue(); // 遍歷broker地址 Iterator<Entry<Long, String>> it = brokerData.getBrokerAddrs().entrySet().iterator(); while (it.hasNext()) { Entry<Long, String> entry = it.next(); Long brokerId = entry.getKey(); String brokerAddr = entry.getValue(); // 根據(jù)broker地址移除brokerAddr if (brokerAddr.equals(brokerAddrFound)) { brokerNameFound = brokerData.getBrokerName(); it.remove(); log.info("remove brokerAddr[{}, {}] from brokerAddrTable, because channel destroyed", brokerId, brokerAddr); break; } } // 如果當(dāng)前主題只包含待移除的broker,則移除該topic if (brokerData.getBrokerAddrs().isEmpty()) { removeBrokerName = true; itBrokerAddrTable.remove(); log.info("remove brokerName[{}] from brokerAddrTable, because channel destroyed", brokerData.getBrokerName()); } } // 維護(hù)clusterAddrTable if (brokerNameFound != null && removeBrokerName) { Iterator<Entry<String, Set<String>>> it = this.clusterAddrTable.entrySet().iterator(); // 遍歷clusterAddrTable while (it.hasNext()) { Entry<String, Set<String>> entry = it.next(); // 獲得集群名稱 String clusterName = entry.getKey(); // 獲得集群中brokerName集合 Set<String> brokerNames = entry.getValue(); // 從brokerNames中移除brokerNameFound boolean removed = brokerNames.remove(brokerNameFound); if (removed) { log.info("remove brokerName[{}], clusterName[{}] from clusterAddrTable, because channel destroyed", brokerNameFound, clusterName);
          if (brokerNames.isEmpty()) { log.info("remove the clusterName[{}] from clusterAddrTable, because channel destroyed and no broker in this cluster", clusterName); // 如果集群中不包含任何broker,則移除該集群 it.remove(); }
          break; } } } // 維護(hù)topicQueueTable隊(duì)列 if (removeBrokerName) { // 遍歷topicQueueTable Iterator<Entry<String, List<QueueData>>> itTopicQueueTable = this.topicQueueTable.entrySet().iterator(); while (itTopicQueueTable.hasNext()) { Entry<String, List<QueueData>> entry = itTopicQueueTable.next(); // 主題名稱 String topic = entry.getKey(); // 隊(duì)列集合 List<QueueData> queueDataList = entry.getValue(); // 遍歷該主題隊(duì)列 Iterator<QueueData> itQueueData = queueDataList.iterator(); while (itQueueData.hasNext()) { // 從隊(duì)列中移除為活躍broker信息 QueueData queueData = itQueueData.next(); if (queueData.getBrokerName().equals(brokerNameFound)) { itQueueData.remove(); log.info("remove topic[{} {}], from topicQueueTable, because channel destroyed", topic, queueData); } } // 如果該topic的隊(duì)列為空,則移除該topic if (queueDataList.isEmpty()) { itTopicQueueTable.remove(); log.info("remove topic[{}] all queue, from topicQueueTable, because channel destroyed", topic); } } } } finally { // 釋放寫鎖 this.lock.writeLock().unlock(); } } catch (Exception e) { log.error("onChannelDestroy Exception", e); } }}


          5、路由發(fā)現(xiàn)

          RocketMQ 路由發(fā)現(xiàn)是非實(shí)時(shí)的,當(dāng) Topic 路由出現(xiàn)變化后,NameServer 不會(huì)主動(dòng)推送給客戶端,而是由客戶端定時(shí)拉取主題最新的路由。

          org.apache.rocketmq.namesrv.processor.DefaultRequestProcessor#getRouteInfoByTopic


          public RemotingCommand getRouteInfoByTopic(ChannelHandlerContext ctx,    RemotingCommand request) throws RemotingCommandException {    final RemotingCommand response = RemotingCommand.createResponseCommand(null);    final GetRouteInfoRequestHeader requestHeader =        (GetRouteInfoRequestHeader) request.decodeCommandCustomHeader(GetRouteInfoRequestHeader.class);    // 調(diào)用RouteInfoManager的方法,從路由表topicQueueTable、brokerAddrTable、 filterServerTable中    // 分別填充TopicRouteData的List<QueueData>、List<BrokerData>、 filterServer    TopicRouteData topicRouteData = this.namesrvController.getRouteInfoManager().pickupTopicRouteData(requestHeader.getTopic());    // 如果找到主題對(duì)應(yīng)你的路由信息并且該主題為順序消息,則從NameServer KVConfig中獲取 關(guān)于順序消息相關(guān)的配置填充路由信息    if (topicRouteData != null) {        if (this.namesrvController.getNamesrvConfig().isOrderMessageEnable()) {            String orderTopicConf =                this.namesrvController.getKvConfigManager().getKVConfig(NamesrvUtil.NAMESPACE_ORDER_TOPIC_CONFIG,                    requestHeader.getTopic());            topicRouteData.setOrderTopicConf(orderTopicConf);        }
          byte[] content = topicRouteData.encode(); response.setBody(content); response.setCode(ResponseCode.SUCCESS); response.setRemark(null); return response; }
          response.setCode(ResponseCode.TOPIC_NOT_EXIST); response.setRemark("No topic route info in name server for the topic: " + requestHeader.getTopic() + FAQUrl.suggestTodo(FAQUrl.APPLY_TOPIC_URL)); return response;}


          五、總結(jié)

          最后我們?cè)賮砜聪挛覀兦把蕴岬哪侨齻€(gè)問題。既然作為路由注冊(cè)中心,那有哪些路由信息注冊(cè)到了 NameServer?生產(chǎn)者如何知道消息要發(fā)送到哪臺(tái)消息服務(wù)器呢?當(dāng) Broker 不可用后,NameServer 并不會(huì)立即將變更后的注冊(cè)信息推送至 Client(Producer/Consumer),此時(shí) RocketMQ 如何保證 Client 正常發(fā)送/消費(fèi)消息?

          第一個(gè)問題在給定 Topic 的情況下,Client 根據(jù)負(fù)載均衡策略選擇合適的消息隊(duì)列,進(jìn)一步獲取到對(duì)應(yīng)的 Broker 地址信息,具體有哪些路由信息注冊(cè)到了 NameServer,可以看上面的路由元信息 RouteInfoManager 類。

          第二個(gè)問題可以歸納為由于路由注冊(cè)、路由刪除以及路由發(fā)現(xiàn),使生產(chǎn)者如何知道消息要發(fā)送到哪臺(tái)消息服務(wù)器。

          第三個(gè)問題對(duì)于 Broker 無效的場(chǎng)景,RocketMQ 犧牲了 C,選擇了 AP,即通過 Client 重試保證可用性,由此產(chǎn)生的重復(fù)消息問題,通過 Client 冪等處理邏輯來規(guī)避。

          最后老周再上一張圖,涵蓋了上述整個(gè)核心流程。到此為止,RocketMQ 的高可用機(jī)制——消息存儲(chǔ)高可用、消息發(fā)送高可用、消息消費(fèi)高可用以及集群管理高可用都分析完畢,希望對(duì)你有幫助。



          歡迎大家關(guān)注我的公眾號(hào)【老周聊架構(gòu)】,Java后端主流技術(shù)棧的原理、源碼分析、架構(gòu)以及各種互聯(lián)網(wǎng)高并發(fā)、高性能、高可用的解決方案。

          喜歡的話,點(diǎn)贊、再看、分享三連。

          點(diǎn)個(gè)在看你最好看



          瀏覽 153
          點(diǎn)贊
          評(píng)論
          收藏
          分享

          手機(jī)掃一掃分享

          分享
          舉報(bào)
          評(píng)論
          圖片
          表情
          推薦
          點(diǎn)贊
          評(píng)論
          收藏
          分享

          手機(jī)掃一掃分享

          分享
          舉報(bào)
          <kbd id="afajh"><form id="afajh"></form></kbd>
          <strong id="afajh"><dl id="afajh"></dl></strong>
            <del id="afajh"><form id="afajh"></form></del>
                1. <th id="afajh"><progress id="afajh"></progress></th>
                  <b id="afajh"><abbr id="afajh"></abbr></b>
                  <th id="afajh"><progress id="afajh"></progress></th>
                  精品人妻中文字幕 | 韩国一级毛 | 仙儿媛护士 | 干少妇电影无码 | 天天操天天干天天 |