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

RocketMQ中的消費者啟動流程解讀

 更新時間:2023年10月11日 10:19:42   作者:CBeann  
這篇文章主要介紹了RocketMQ中的消費者啟動流程解讀,RocketMQ是一款高性能、高可靠性的分布式消息中間件,消費者是RocketMQ中的重要組成部分,消費者負(fù)責(zé)從消息隊列中獲取消息并進(jìn)行處理,需要的朋友可以參考下

RocketMQ消費者啟動

問題

消費者啟動的時候,去哪拿的消息呢?

問題答案

請?zhí)砑訄D片描述

(1)當(dāng)broker啟動的時候,會把broker的地址端口、broker上的主題信息、主題隊列信息發(fā)送到nameserver(如圖中1)

(2)消費者Client啟動的時候會去nameserver拿toipc、topic隊列以及對應(yīng)的broker信息,拿到以后把信息存儲到本地(如圖中2)

(3)消費者會給所有的broker發(fā)送心跳,并且附帶自己的消費者組信息和ClientID信息,此時broker中就有消費者組對應(yīng)的ClientID集合(如圖中3)

(4)消費者啟動后會reblance,有訂閱的主題隊列列表,并且通過broker可以拿到消費者組的ClientID集合,兩個集合做rebalance,就可以拿到當(dāng)前消費者對應(yīng)消費的主題隊列

(5) 消費者知道自己消費的主題隊列,就可以根據(jù)隊列信息通過Netty發(fā)送消息

跟源碼

注意

本文是消費者啟動流程,所以不去關(guān)注broker和nameserver的啟動流程,這樣關(guān)注點比較集中,因此步驟(1)本文不做描述。

消費者啟動時怎么拿到toipc的信息

消費者啟動的時候會調(diào)用 MQClientInstance###start()方法,start()方法里有會調(diào)用 MQClientInstance###startScheduledTask()方法,里面的一段代碼如下,會每隔一段時間更新一下topic路由信息

//MQClientInstance###startScheduledTask()
 this.scheduledExecutorService.scheduleAtFixedRate(new Runnable() {
            @Override
            public void run() {
                try {
                    MQClientInstance.this.updateTopicRouteInfoFromNameServer();
                } catch (Exception e) {
                    log.error("ScheduledTask updateTopicRouteInfoFromNameServer exception", e);
                }
            }
        }, 10, this.clientConfig.getPollNameServerInterval(), TimeUnit.MILLISECONDS);

會把路由信息保存到本地的一個HashMap里,這樣消費者就拿到了topic的信息并且會把broker的信息保存下來

//MQClientInstance###updateTopicRouteInfoFromNameServer(final String topic, boolean isDefault,DefaultMQProducer defaultMQProducer)
//根據(jù)主題從nameserver獲取topic信息
topicRouteData = this.mQClientAPIImpl.getTopicRouteInfoFromNameServer(topic, 1000 * 3);
//MQClientInstance###updateTopicRouteInfoFromNameServer(final String topic, boolean isDefault,DefaultMQProducer defaultMQProducer)
//把主題和主題隊列相關(guān)的broker保存下來
TopicRouteData cloneTopicRouteData = topicRouteData.cloneTopicRouteData();
                            for (BrokerData bd : topicRouteData.getBrokerDatas()) {
                                this.brokerAddrTable.put(bd.getBrokerName(), bd.getBrokerAddrs());
                            }

總結(jié):消費者拿到主題的隊列列表和broker信息

消費者給broker發(fā)現(xiàn)心跳的作用

MQClientInstance###startScheduledTask()方法,里面的一段代碼如下,會每隔一段時間給所有的broker發(fā)送心跳消息

//MQClientInstance###startScheduledTask()
this.scheduledExecutorService.scheduleAtFixedRate(new Runnable() {
            @Override
            public void run() {
                try {
                    MQClientInstance.this.cleanOfflineBroker();
                    MQClientInstance.this.sendHeartbeatToAllBrokerWithLock();
                } catch (Exception e) {
                    log.error("ScheduledTask sendHeartbeatToAllBroker exception", e);
                }
            }
        }, 1000, this.clientConfig.getHeartbeatBrokerInterval(), TimeUnit.MILLISECONDS);

那么發(fā)送的心跳包中攜帶什么信息呢?如代碼中所示,攜帶clientID和組名稱

//MQClientInstance###prepareHeartbeatData
private HeartbeatData prepareHeartbeatData() {
        HeartbeatData heartbeatData = new HeartbeatData();
        // clientID
        //放入了當(dāng)前消費者的clientID
        //放入了當(dāng)前消費者的clientID
        //放入了當(dāng)前消費者的clientID
        heartbeatData.setClientID(this.clientId);
        // Consumer
        for (Map.Entry<String, MQConsumerInner> entry : this.consumerTable.entrySet()) {
            MQConsumerInner impl = entry.getValue();
            if (impl != null) {
                ConsumerData consumerData = new ConsumerData(); 
                //放入了當(dāng)前消費者的組名稱
                //放入了當(dāng)前消費者的組名稱
                //放入了當(dāng)前消費者的組名稱
                //放入了當(dāng)前消費者的組名稱
                consumerData.setGroupName(impl.groupName());
                consumerData.setConsumeType(impl.consumeType());
                consumerData.setMessageModel(impl.messageModel());
                consumerData.setConsumeFromWhere(impl.consumeFromWhere());
                consumerData.getSubscriptionDataSet().addAll(impl.subscriptions());
                consumerData.setUnitMode(impl.isUnitMode());
                heartbeatData.getConsumerDataSet().add(consumerData);
            }
        }
        // Producer
        for (Map.Entry<String/* group */, MQProducerInner> entry : this.producerTable.entrySet()) {
            MQProducerInner impl = entry.getValue();
            if (impl != null) {
                ProducerData producerData = new ProducerData();
                producerData.setGroupName(entry.getKey());
                heartbeatData.getProducerDataSet().add(producerData);
            }
        }
        return heartbeatData;
    }

此時broker拿到心跳消息怎么處理的呢?有一部分邏輯如下面代碼所示,記錄一下消費者信息

//ClientManageProcessor###heartBeat(ChannelHandlerContext ctx, RemotingCommand request)
```java
public RemotingCommand heartBeat(ChannelHandlerContext ctx, RemotingCommand request) {
        //省略
        for (ConsumerData data : heartbeatData.getConsumerDataSet()) {
            //省略
            boolean changed = this.brokerController.getConsumerManager().registerConsumer(
                data.getGroupName(),
                clientChannelInfo,
                data.getConsumeType(),
                data.getMessageModel(),
                data.getConsumeFromWhere(),
                data.getSubscriptionDataSet(),
                isNotifyConsumerIdsChangedEnable
            );
            //省略
        }
       //省略
    }

消費者怎么做reblance

MQClientInstance的start的方法里會開啟一個rebalance的線程,如下面代碼所示

//MQClientInstance###start()
public void start() throws MQClientException {
 //省略
 // Start rebalance service
 this.rebalanceService.start();
 //省略
}

跟RebalanceService的run()方法一直跟下去最后跟到RebalanceImpl的rebalanceByTopic方法,如下面代碼所示。根據(jù)主題隊列列表和消費者組集合去做一個Rebalance,最后的返回結(jié)果是當(dāng)前消費者需要消費的主題隊列。

//RebalanceImpl##rebalanceByTopic
private void rebalanceByTopic(final String topic, final boolean isOrder) {
                //獲取訂閱的主題的隊列
                //獲取訂閱的主題的隊列
                //獲取訂閱的主題的隊列
                Set<MessageQueue> mqSet = this.topicSubscribeInfoTable.get(topic);
                //獲取同消費者組的ClientID集合
                //獲取同消費者組的ClientID集合
                //獲取同消費者組的ClientID集合
                List<String> cidAll = this.mQClientFactory.findConsumerIdList(topic, consumerGroup);
                if (mqSet != null && cidAll != null) {
                    List<MessageQueue> mqAll = new ArrayList<MessageQueue>();
                    mqAll.addAll(mqSet);
                    //排序
                    //排序
                    //排序
                    Collections.sort(mqAll);
                    Collections.sort(cidAll);
                    AllocateMessageQueueStrategy strategy = this.allocateMessageQueueStrategy;
                    List<MessageQueue> allocateResult = null;
                    try {
                        //rebalance算法核心實現(xiàn),最后的結(jié)果是返回應(yīng)該消費的隊列
                        //rebalance算法核心實現(xiàn),最后的結(jié)果是返回應(yīng)該消費的隊列
                        //rebalance算法核心實現(xiàn),最后的結(jié)果是返回應(yīng)該消費的隊列
                        allocateResult = strategy.allocate(
                            this.consumerGroup,
                            this.mQClientFactory.getClientId(),
                            mqAll,
                            cidAll);
                    } catch (Throwable e) {
                    }
                    Set<MessageQueue> allocateResultSet = new HashSet<MessageQueue>();
                    if (allocateResult != null) {
                        //rebalance算法核心實現(xiàn),最后的結(jié)果是返回應(yīng)該消費的隊列
                        //rebalance算法核心實現(xiàn),最后的結(jié)果是返回應(yīng)該消費的隊列
                        //rebalance算法核心實現(xiàn),最后的結(jié)果是返回應(yīng)該消費的隊列
                        allocateResultSet.addAll(allocateResult);
                    }
                    //此處看下面的消費者怎么去拉消息
                    //此處看下面的消費者怎么去拉消息
                    //此處看下面的消費者怎么去拉消息
                    boolean changed = this.updateProcessQueueTableInRebalance(topic, allocateResultSet, isOrder);
        }
    }

總結(jié):消費者拿到主題的隊列列表和消費者組中ClientID集合,通過在消費者這變做rebalance,從而確定被分配的主題隊列集合

消費者怎么拉取消息

此處還是繼續(xù)跟上面的代碼,,然后執(zhí)行到下面的代碼,當(dāng)消費者確定自己被分配的主題隊列后,會把主題隊列封裝成PullRequest 并進(jìn)行dispatch

//RebalanceImpl###updateProcessQueueTableInRebalance
private boolean updateProcessQueueTableInRebalance(final String topic, final Set<MessageQueue> mqSet,
        final boolean isOrder) {
       //省列
        List<PullRequest> pullRequestList = new ArrayList<PullRequest>();
        for (MessageQueue mq : mqSet) {
                        //省略
                        PullRequest pullRequest = new PullRequest();
                        pullRequest.setConsumerGroup(consumerGroup);
                        pullRequest.setNextOffset(nextOffset);
                        pullRequest.setMessageQueue(mq);
                        pullRequest.setProcessQueue(pq);
                        pullRequestList.add(pullRequest);
                        changed = true;
            }
        }
         //省略
        //派發(fā)請求任務(wù)
        this.dispatchPullRequest(pullRequestList);
        return changed;
    }

下面跟RebalanceImpl###dispatchPullRequest方法,最后跟到下面的代碼,就是把PullRequest放入到一個阻塞隊列里。

//PullMessageService###executePullRequestImmediately
public void executePullRequestImmediately(final PullRequest pullRequest) {
        try {
            this.pullRequestQueue.put(pullRequest);
        } catch (InterruptedException e) {
            log.error("executePullRequestImmediately pullRequestQueue.put", e);
        }
    }

那么誰取阻塞隊列里的數(shù)據(jù)誰就是消費消息了? PullMessageService是一個線程,他的run方法里會取上面阻塞隊列里的PullRequest,如下面代碼所示

//PullMessageService###run()
public void run() {
        log.info(this.getServiceName() + " service started");
        while (!this.isStopped()) {
            try {
                PullRequest pullRequest = this.pullRequestQueue.take();
                this.pullMessage(pullRequest);
            } catch (InterruptedException ignored) {
            } catch (Exception e) {
                log.error("Pull Message Service Run Method exception", e);
            }
        }
        log.info(this.getServiceName() + " service end");
    }

從PullMessageService###pullMessage方法一直往下跟,就跟到下面的代碼

//DefaultMQPushConsumerImpl###pullMessage(final PullRequest pullRequest) 
public void pullMessage(final PullRequest pullRequest) {
        //省略
        final long beginTimestamp = System.currentTimeMillis();
        PullCallback pullCallback = new PullCallback() {
            //省略,但是重要,后面會說
            //省略,但是重要,后面會說
            //省略,但是重要,后面會說
            //省略,但是重要,后面會說
        };
           //省略
        try {
            //發(fā)送數(shù)據(jù)并且執(zhí)行回調(diào)方法,下面我們看一下回調(diào)方法的內(nèi)容就好好了
            //發(fā)送數(shù)據(jù)并且執(zhí)行回調(diào)方法,下面我們看一下回調(diào)方法的內(nèi)容就好好了
            //發(fā)送數(shù)據(jù)并且執(zhí)行回調(diào)方法,下面我們看一下回調(diào)方法的內(nèi)容就好好了
            this.pullAPIWrapper.pullKernelImpl(
                pullRequest.getMessageQueue(),
                subExpression,
                subscriptionData.getExpressionType(),
                subscriptionData.getSubVersion(),
                pullRequest.getNextOffset(),
                this.defaultMQPushConsumer.getPullBatchSize(),
                sysFlag,
                commitOffsetValue,
                BROKER_SUSPEND_MAX_TIME_MILLIS,
                CONSUMER_TIMEOUT_MILLIS_WHEN_SUSPEND,
                CommunicationMode.ASYNC,
                pullCallback
            );
        } catch (Exception e) {
            log.error("pullKernelImpl exception", e);
            this.executePullRequestLater(pullRequest, pullTimeDelayMillsWhenException);
        }
    }

那么回調(diào)方法是什么邏輯呢?代碼如下所示,發(fā)現(xiàn)數(shù)據(jù)并且submitConsumeRequest

PullCallback pullCallback = new PullCallback() {
            @Override
            public void onSuccess(PullResult pullResult) {
                if (pullResult != null) {
                    pullResult = DefaultMQPushConsumerImpl.this.pullAPIWrapper.processPullResult(pullRequest.getMessageQueue(), pullResult,
                        subscriptionData);
                    switch (pullResult.getPullStatus()) {
                        //發(fā)現(xiàn)數(shù)據(jù)
                        //發(fā)現(xiàn)數(shù)據(jù)
                        //發(fā)現(xiàn)數(shù)據(jù)
                        case FOUND:
                            long prevRequestOffset = pullRequest.getNextOffset();
                            pullRequest.setNextOffset(pullResult.getNextBeginOffset());
                            long pullRT = System.currentTimeMillis() - beginTimestamp;
                            DefaultMQPushConsumerImpl.this.getConsumerStatsManager().incPullRT(pullRequest.getConsumerGroup(),
                                pullRequest.getMessageQueue().getTopic(), pullRT);
                            long firstMsgOffset = Long.MAX_VALUE;
                            if (pullResult.getMsgFoundList() == null || pullResult.getMsgFoundList().isEmpty()) {
                                DefaultMQPushConsumerImpl.this.executePullRequestImmediately(pullRequest);
                            } else {
                                firstMsgOffset = pullResult.getMsgFoundList().get(0).getQueueOffset();
                                DefaultMQPushConsumerImpl.this.getConsumerStatsManager().incPullTPS(pullRequest.getConsumerGroup(),
                                    pullRequest.getMessageQueue().getTopic(), pullResult.getMsgFoundList().size());
                                boolean dispatchToConsume = processQueue.putMessage(pullResult.getMsgFoundList());
                                //跟進(jìn)去
                                //跟進(jìn)去
                                //跟進(jìn)去
                                DefaultMQPushConsumerImpl.this.consumeMessageService.submitConsumeRequest(
                                    pullResult.getMsgFoundList(),
                                    processQueue,
                                    pullRequest.getMessageQueue(),
                                    dispatchToConsume);
                                if (DefaultMQPushConsumerImpl.this.defaultMQPushConsumer.getPullInterval() > 0) {
                                    DefaultMQPushConsumerImpl.this.executePullRequestLater(pullRequest,
                                        DefaultMQPushConsumerImpl.this.defaultMQPushConsumer.getPullInterval());
                                } else {
                                    DefaultMQPushConsumerImpl.this.executePullRequestImmediately(pullRequest);
                                }
                            }
                            //省略
                            break;
                        case NO_NEW_MSG:
                        case NO_MATCHED_MSG:
                            //省略
                    }
                }
            }
        };

跟上面submitConsumeRequest方法的到下面的代碼,封裝成ConsumeRequest,其實ConsumerRequest是一個線程

//ConsumeMessageConcurrentlyService###submitConsumeRequest
public void submitConsumeRequest(
        final List<MessageExt> msgs,
        final ProcessQueue processQueue,
        final MessageQueue messageQueue,
        final boolean dispatchToConsume) {
        final int consumeBatchSize = this.defaultMQPushConsumer.getConsumeMessageBatchMaxSize();
        if (msgs.size() <= consumeBatchSize) {
            ConsumeRequest consumeRequest = new ConsumeRequest(msgs, processQueue, messageQueue);
            try {
                this.consumeExecutor.submit(consumeRequest);
            } catch (RejectedExecutionException e) {
                this.submitConsumeRequestLater(consumeRequest);
            }
        } else {
            //省略
            }
        }
    }
ConsumeRequest的run方法就會執(zhí)行我們注冊的listener方法,此時就消費到數(shù)據(jù)
```java
@Override
        public void run() {
            //省略
                status = listener.consumeMessage(Collections.unmodifiableList(msgs), context);
            //省略
            }
        }

在這里插入圖片描述

總結(jié): 如下圖所示,RebalanceService線程會根據(jù)情況把請求放在PullMessageService的pullRequestQueue阻塞隊列隊列里,隊列的每一個節(jié)點就是拉請求;PullMessageService線程就是不斷去pullRequestQueue里拿任務(wù)然后去看一下broker中有沒有數(shù)據(jù),如果有數(shù)據(jù)就消費。

在這里插入圖片描述

總結(jié)

(1)忽然發(fā)現(xiàn)nameserver在整個過程中的作用感覺不是很大,其實我感覺這種設(shè)計還挺好的,因為把所有的壓力都放在nameserver返回減少系統(tǒng)的健壯性。

(2)RocketMQ的rebalance是在消息消費者這邊實現(xiàn)的,這樣有一個很大的優(yōu)勢是減少nameserver和broker的壓力。那消費者是怎么實現(xiàn)rebalance的呢?通過一個參數(shù)為當(dāng)前消費者ID、主題隊列、消費者組ClientID列表的匹配算法,每次只要保證算法的冪等性就可以了。

(3)RocketMQ的rebalance的rebalance是根據(jù)單個主題去實現(xiàn)的,這樣的一個缺點是容易出現(xiàn)消費不平衡的問題。如下圖所示。

在這里插入圖片描述

(4)RocketMQ是AP的,因為他的很操作都是都是通過線程池的定時任務(wù)去做的。

到此這篇關(guān)于RocketMQ中的消費者啟動流程解讀的文章就介紹到這了,更多相關(guān)RocketMQ消費者啟動內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • Spring整合Mybatis思路梳理總結(jié)

    Spring整合Mybatis思路梳理總結(jié)

    mybatis-plus是一個 Mybatis 的增強(qiáng)工具,在 Mybatis 的基礎(chǔ)上只做增強(qiáng)不做改變,為簡化開發(fā)、提高效率而生,下面這篇文章主要給大家介紹了關(guān)于SpringBoot整合Mybatis-plus案例及用法實例的相關(guān)資料,需要的朋友可以參考下
    2022-12-12
  • IDEA 通過腳本配置終端提示符樣式的方法

    IDEA 通過腳本配置終端提示符樣式的方法

    這篇文章給大家介紹IDEA通過腳本配置終端提示符樣式的方法,本文給大家介紹的非常詳細(xì),對大家的學(xué)習(xí)或工作具有一定的參考借鑒價值,需要的朋友參考下吧
    2025-08-08
  • Mybatis讀取和存儲json類型數(shù)據(jù)的實現(xiàn)

    Mybatis讀取和存儲json類型數(shù)據(jù)的實現(xiàn)

    本文主要介紹了Mybatis讀取和存儲json類型數(shù)據(jù)的實現(xiàn),文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2023-06-06
  • SpringBoot配置HTTPS及開發(fā)調(diào)試的操作方法

    SpringBoot配置HTTPS及開發(fā)調(diào)試的操作方法

    在實際開發(fā)過程中,如果后端需要啟用https訪問,通常項目啟動后配置nginx代理再配置https,前端調(diào)用時高版本的chrome還會因為證書未信任導(dǎo)致調(diào)用失敗,通過摸索整理一套開發(fā)調(diào)試下的https方案,下面給大家分享SpringBoot配置HTTPS及開發(fā)調(diào)試,感興趣的朋友跟隨小編一起看看吧
    2024-05-05
  • 如何把char數(shù)組轉(zhuǎn)換成String

    如何把char數(shù)組轉(zhuǎn)換成String

    這篇文章主要介紹了如何把char數(shù)組轉(zhuǎn)換成String問題,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2023-02-02
  • Redisson 分布式延時隊列 RedissonDelayedQueue 運行流程

    Redisson 分布式延時隊列 RedissonDelayedQueue 運行流程

    這篇文章主要介紹了Redisson分布式延時隊列 RedissonDelayedQueue運行流程,文章圍繞主題展開詳細(xì)的內(nèi)容介紹,具有一定的參考價值,需要的小伙伴可以參考一下
    2022-09-09
  • JVM(Java?Virtual?Machine,Java虛擬機(jī))的作用詳解

    JVM(Java?Virtual?Machine,Java虛擬機(jī))的作用詳解

    JVM是Java語言實現(xiàn)“一次編寫,到處運行”特性的基石,也是Java平臺的核心組成部分,其主要作用包括平臺無關(guān)性、內(nèi)存管理、運行Java程序、安全性以及性能優(yōu)化,通過這些功能,JVM確保了Java程序的可移植性、高效性和安全性
    2025-03-03
  • 利用Thumbnailator輕松實現(xiàn)圖片縮放、旋轉(zhuǎn)與加水印

    利用Thumbnailator輕松實現(xiàn)圖片縮放、旋轉(zhuǎn)與加水印

    java開發(fā)中經(jīng)常遇到對圖片的處理,JDK中也提供了對應(yīng)的工具類,不過處理起來很麻煩,Thumbnailator是一個優(yōu)秀的圖片處理的開源Java類庫,處理效果遠(yuǎn)比Java API的好,這篇文章主要介紹了利用Thumbnailator如何輕松的實現(xiàn)圖片縮放、旋轉(zhuǎn)與加水印,需要的朋友可以參考下
    2017-01-01
  • 新手了解java基礎(chǔ)知識(二)

    新手了解java基礎(chǔ)知識(二)

    這篇文章主要介紹了Java基礎(chǔ)知識,本文介紹了Java語言相關(guān)的基礎(chǔ)知識、歷史介紹、主要應(yīng)用方向等內(nèi)容,需要的朋友可以參考下,希望對你有所幫助
    2021-07-07
  • java 生成xml并轉(zhuǎn)為字符串的方法

    java 生成xml并轉(zhuǎn)為字符串的方法

    今天小編就為大家分享一篇java 生成xml并轉(zhuǎn)為字符串的方法,具有很好的參考價值,希望對大家有所幫助。一起跟隨小編過來看看吧
    2018-07-07

最新評論

宝鸡市| 台中县| 个旧市| 习水县| 兴城市| 离岛区| 遵义市| 墨脱县| 筠连县| 达日县| 兴化市| 丹阳市| 永春县| 夏邑县| 赤壁市| 阿拉尔市| 昔阳县| 苍山县| 岳池县| 西昌市| 泸溪县| 浦北县| 枝江市| 资阳市| 万载县| 商丘市| 井冈山市| 临湘市| 青阳县| 冀州市| 囊谦县| 炉霍县| 昌图县| 甘孜| 曲松县| 新乡县| 原平市| 定安县| 读书| 夹江县| 富川|