最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

RocketMQ NameServer保障數(shù)據(jù)一致性實現(xiàn)方法講解

 更新時間:2022年12月12日 08:55:39   作者:小王曾是少年  
這篇文章主要介紹了RocketMQ NameServer保障數(shù)據(jù)一致性實現(xiàn)方法,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪

路由注冊角度

對于ZooKeeper這樣的強(qiáng)一致性組件,使用主從分離的架構(gòu),數(shù)據(jù)只寫到主節(jié)點,主從之間的數(shù)據(jù)同步通過內(nèi)部機(jī)制來進(jìn)行數(shù)據(jù)復(fù)制。

對于RocketMQ來說,NameServer節(jié)點之間是互相不進(jìn)行通信的,這樣也就無法進(jìn)行數(shù)據(jù)復(fù)制。RocketMQ采用的機(jī)制是:在Broker節(jié)點啟動的時候,輪詢所有的NameServer節(jié)點,并與每個NameServer節(jié)點建立長連接,發(fā)送注冊請求。

相應(yīng)的,NameServer節(jié)點內(nèi)部也會維護(hù)一個Broker列表,用來動態(tài)存儲Broker的信息,做服務(wù)發(fā)現(xiàn)。

與此同時,Broker使用心跳機(jī)制來向所有NameServer節(jié)點證明自己是存活的,即定期發(fā)送心跳包;收到心跳包之后,NameServer節(jié)點會更新這個Broker的最新存活時間。

注意: NameServer節(jié)點在處理心跳包時,存在多個請求同時處理同一張表的情況,為了保證并發(fā)安全性,RocketMQ引入了讀寫鎖(ReadWriteLock),保證了多個Producer并發(fā)讀取路由信息不受影響,但同一時刻只能處理一個Broker發(fā)來的心跳包,這也符合讀多寫少的經(jīng)典場景。

路由剔除

正常情況下:

如果Broker下線,則會與NameServer斷開長連接,底層基于Netty的通道關(guān)閉監(jiān)聽器會監(jiān)聽到連接斷開事件,然后將這個Broker信息剔除。

異常情況下:

NameServer有一個周期為10s的定時任務(wù),定期掃描Broker表,如果超過120s沒有收到某個Broker的心跳包,則會判定其失效并移除。

對于日常運維的需求,RocketMQ提供了優(yōu)雅剔除路由信息的方式,即可以先禁止Broker的寫權(quán)限,這樣發(fā)送到這個Broker的請求都會收到一個NO_PERMISSION的響應(yīng),客戶端自動重試其他的Broker。

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

生產(chǎn)者視角:

一般是在發(fā)送第一條消息時,才會根據(jù)TopicNameServer獲取路由信息

消費者視角:

訂閱的Topic一般是固定的,所以在啟動時就會拉取

針對路由信息可能變化的場景,RocketMQ提供了定時拉取Topic最新路由信息的機(jī)制,以應(yīng)對Broker集群發(fā)生變化的場景。

DefaultMQProducerDefaultMQConsumer有一個pollNameServerInterval的配置項,用于指定從NameServer獲取路由信息的周期,其底層依賴MQClientInstance類,MQClientInstance類中的updateTopicRouteInfoFromNameServer方法,可以根據(jù)指定的時間間隔,周期性地從NameServer里拉取路由信息。在拉取時,會將當(dāng)前啟動的ProducerConsumer需要用到的Topic列表放到一個集合里,逐個進(jìn)行更新,源碼如下:

/**
* 更新單個Topic路由信息
*/
public boolean updateTopicRouteInfoFromNameServer(final String topic) {
    return updateTopicRouteInfoFromNameServer(topic, false, null);
}
/**
* 更新單個Topic路由信息
*/
public boolean updateTopicRouteInfoFromNameServer(final String topic, boolean isDefault,
    DefaultMQProducer defaultMQProducer) {
    try {
        if (this.lockNamesrv.tryLock(LOCK_TIMEOUT_MILLIS, TimeUnit.MILLISECONDS)) {
            try {
                TopicRouteData topicRouteData;
                if (isDefault && defaultMQProducer != null) {
                	// 使用默認(rèn)TopicKey獲取TopicRouteData
                    topicRouteData = this.mQClientAPIImpl.getDefaultTopicRouteInfoFromNameServer(defaultMQProducer.getCreateTopicKey(), 1000 * 3);
                    if (topicRouteData != null) {
                        for (QueueData data : topicRouteData.getQueueDatas()) {
                            int queueNums = Math.min(defaultMQProducer.getDefaultTopicQueueNums(), data.getReadQueueNums());
                            data.setReadQueueNums(queueNums);
                            data.setWriteQueueNums(queueNums);
                        }
                    }
                } else {
                    topicRouteData = this.mQClientAPIImpl.getTopicRouteInfoFromNameServer(topic, 1000 * 3);
                }
                if (topicRouteData != null) {
                    TopicRouteData old = this.topicRouteTable.get(topic);
                    boolean changed = topicRouteDataIsChange(old, topicRouteData);
                    if (!changed) {
                        changed = this.isNeedUpdateTopicRouteInfo(topic);
                    } else {
                        log.info("the topic[{}] route info changed, old[{}] ,new[{}]", topic, old, topicRouteData);
                    }
                    if (changed) {
                    	// 克隆出一個實例cloneTopicRouteData : topicRouteData會被設(shè)置到下面的publishInfo/subscribeInfo 
                        TopicRouteData cloneTopicRouteData = topicRouteData.cloneTopicRouteData();
						// 更新Broker地址相關(guān)信息,當(dāng)某個Broker心跳超時后,會被從brokerAddrTable中移除
                        for (BrokerData bd : topicRouteData.getBrokerDatas()) {
                            this.brokerAddrTable.put(bd.getBrokerName(), bd.getBrokerAddrs());
                        }
                        // Update Pub info
                        {
                            TopicPublishInfo publishInfo = topicRouteData2TopicPublishInfo(topic, topicRouteData);
                            publishInfo.setHaveTopicRouterInfo(true);
                            Iterator<Entry<String, MQProducerInner>> it = this.producerTable.entrySet().iterator();
                            while (it.hasNext()) {
                                Entry<String, MQProducerInner> entry = it.next();
                                MQProducerInner impl = entry.getValue();
                                if (impl != null) {
                                    impl.updateTopicPublishInfo(topic, publishInfo);
                                }
                            }
                        }
                        // Update sub info
                        {
                            Set<MessageQueue> subscribeInfo = topicRouteData2TopicSubscribeInfo(topic, topicRouteData);
                            Iterator<Entry<String, MQConsumerInner>> it = this.consumerTable.entrySet().iterator();
                            while (it.hasNext()) {
                                Entry<String, MQConsumerInner> entry = it.next();
                                MQConsumerInner impl = entry.getValue();
                                if (impl != null) {
                                    impl.updateTopicSubscribeInfo(topic, subscribeInfo);
                                }
                            }
                        }
                        log.info("topicRouteTable.put. Topic = {}, TopicRouteData[{}]", topic, cloneTopicRouteData);
                        this.topicRouteTable.put(topic, cloneTopicRouteData);
                        return true;
                    }
                } else {
                    log.warn("updateTopicRouteInfoFromNameServer, getTopicRouteInfoFromNameServer return null, Topic: {}. [{}]", topic, this.clientId);
                }
            } catch (MQClientException e) {
                if (!topic.startsWith(MixAll.RETRY_GROUP_TOPIC_PREFIX)) {
                    log.warn("updateTopicRouteInfoFromNameServer Exception", e);
                }
            } catch (RemotingException e) {
                log.error("updateTopicRouteInfoFromNameServer Exception", e);
                throw new IllegalStateException(e);
            } finally {
                this.lockNamesrv.unlock();
            }
        } else {
            log.warn("updateTopicRouteInfoFromNameServer tryLock timeout {}ms. [{}]", LOCK_TIMEOUT_MILLIS, this.clientId);
        }
    } catch (InterruptedException e) {
        log.warn("updateTopicRouteInfoFromNameServer Exception", e);
    }
    return false;
}

當(dāng)Broker宕機(jī)時,還可以通過客戶端的重試機(jī)制來解決,避免因為定時更新路由信息不及時導(dǎo)致的服務(wù)宕機(jī)~~

到此這篇關(guān)于RocketMQ NameServer保障數(shù)據(jù)一致性實現(xiàn)方法講解的文章就介紹到這了,更多相關(guān)RocketMQ NameServer內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • 一篇文章帶你初步認(rèn)識Maven

    一篇文章帶你初步認(rèn)識Maven

    這篇文章主要為大家初步認(rèn)識了Maven,具有一定的參考價值,感興趣的小伙伴們可以參考一下,希望能夠給你帶來幫助
    2022-01-01
  • Ubuntu下配置Tomcat服務(wù)器以及設(shè)置自動啟動的方法

    Ubuntu下配置Tomcat服務(wù)器以及設(shè)置自動啟動的方法

    這篇文章主要介紹了Ubuntu下配置Tomcat服務(wù)器以及設(shè)置自動啟動的方法,適用于Java的web程序開發(fā),需要的朋友可以參考下
    2015-10-10
  • Java8中Stream使用的一個注意事項

    Java8中Stream使用的一個注意事項

    最近在工作中發(fā)現(xiàn)了對于集合操作轉(zhuǎn)換的神器,java8新特性 stream,但在使用中遇到了一個非常重要的注意點,所以這篇文章主要給大家介紹了關(guān)于Java8中Stream使用過程中的一個注意事項,需要的朋友可以參考借鑒,下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧。
    2017-11-11
  • 使用Java通過OAuth協(xié)議驗證發(fā)送微博的教程

    使用Java通過OAuth協(xié)議驗證發(fā)送微博的教程

    這篇文章主要介紹了使用Java通過OAuth協(xié)議驗證發(fā)送微博的教程,使用到了新浪微博為Java開放的API weibo4j,需要的朋友可以參考下
    2016-02-02
  • @Scheduled 如何讀取動態(tài)配置文件

    @Scheduled 如何讀取動態(tài)配置文件

    這篇文章主要介紹了@Scheduled 如何讀取動態(tài)配置文件的操作,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-06-06
  • Spring框架IOC容器底層原理詳解

    Spring框架IOC容器底層原理詳解

    在java當(dāng)中一個類想要使用另一個類的方法,就必須在這個類當(dāng)中創(chuàng)建這個類的對象,Spring將創(chuàng)建對象的權(quán)利給了IOC,在IOC當(dāng)中創(chuàng)建了ABC三個對象,那么我們我們其他的類只需要調(diào)用集合,大大的解決了程序耦合性的問題
    2022-07-07
  • Spring Security 構(gòu)建rest服務(wù)實現(xiàn)rememberme 記住我功能

    Spring Security 構(gòu)建rest服務(wù)實現(xiàn)rememberme 記住我功能

    這篇文章主要介紹了Spring Security 構(gòu)建rest服務(wù)實現(xiàn)rememberme 記住我功能,需要的朋友可以參考下
    2018-03-03
  • java騰訊AI人臉對比對接代碼實例

    java騰訊AI人臉對比對接代碼實例

    這篇文章主要介紹了java騰訊AI人臉對比對接,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2019-03-03
  • 怎樣給Kafka新增分區(qū)

    怎樣給Kafka新增分區(qū)

    這篇文章主要介紹了怎樣給Kafka新增分區(qū)問題,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2022-12-12
  • Mybatis Generator最完美配置文件詳解(完整版)

    Mybatis Generator最完美配置文件詳解(完整版)

    今天小編給大家整理了一篇關(guān)于Mybatis Generator最完美配置文件詳解教程,非常不錯具有參考借鑒價值,感興趣的朋友一起學(xué)習(xí)吧
    2016-11-11

最新評論

留坝县| 安远县| 保亭| 安龙县| 静海县| 祁东县| 肇东市| 茂名市| 蓝山县| 北辰区| 建水县| 永兴县| 乌恰县| 天津市| 白沙| 新兴县| 方城县| 滦南县| 河东区| 景德镇市| 红安县| 三都| 东安县| 共和县| 广昌县| 和政县| 洞头县| 婺源县| 仪征市| 峨山| 房产| 盐山县| 重庆市| 盐亭县| 百色市| 丹寨县| 富宁县| 临高县| 武汉市| 梁河县| 比如县|