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

Java中RocketMq的消費方式詳解

 更新時間:2023年10月11日 09:11:19   作者:獵戶星座。  
這篇文章主要介紹了Java中RocketMq的消費方式詳解,RocketMQ的消費方式都是基于拉模式拉取消息的,而在這其中有一種長輪詢機制(對普通輪詢的一種優(yōu)化),來平衡上面Push/Pull模型的各自缺點,需要的朋友可以參考下

 一、如何選擇消息消費的方式—Pull or Push?

1.1 MQ中Pull和Push的兩種消費方式

對于任何一款消息中間件而言,消費者客戶端一般有兩種方式從消息中間件獲取消息并消費:

(1)Push方式:由消息中間件(MQ消息服務(wù)器代理)主動地將消息推送給消費者;采用Push方式,可以盡可能實時地將消息發(fā)送給消費者進行消費。但是,在消費者的處理消息的能力較弱的時候(比如,消費者端的業(yè)務(wù)系統(tǒng)處理一條消息的流程比較復(fù)雜,其中的調(diào)用鏈路比較多導(dǎo)致消費時間比較久。概括起來地說就是“慢消費問題”),而MQ不斷地向消費者Push消息,消費者端的緩沖區(qū)可能會溢出,導(dǎo)致異常;

(2)Pull方式:由消費者客戶端主動向消息中間件(MQ消息服務(wù)器代理)拉取消息;采用Pull方式,如何設(shè)置Pull消息的頻率需要重點去考慮,舉個例子來說,可能1分鐘內(nèi)連續(xù)來了1000條消息,然后2小時內(nèi)沒有新消息產(chǎn)生(概括起來說就是“消息延遲與忙等待”)。如果每次Pull的時間間隔比較久,會增加消息的延遲,即消息到達消費者的時間加長,MQ中消息的堆積量變大;若每次Pull的時間間隔較短,但是在一段時間內(nèi)MQ中并沒有任何消息可以消費,那么會產(chǎn)生很多無效的Pull請求的RPC開銷,影響MQ整體的網(wǎng)絡(luò)性能;

1.2 RocketMQ消息消費的長輪詢機制

思考題: 上面簡要說明了Push和Pull兩種消息消費方式的概念和各自特點。如果長時間沒有消息,而消費者端又不停的發(fā)送Pull請求不就會導(dǎo)致RocketMQ中Broker端負載很高嗎?那么在RocketMQ中如何解決以做到高效的消息消費呢?

通過研究源碼可知,RocketMQ的消費方式都是基于拉模式拉取消息的,而在這其中有一種長輪詢機制(對普通輪詢的一種優(yōu)化),來平衡上面Push/Pull模型的各自缺點。基本設(shè)計思路是:消費者如果第一次嘗試Pull消息失?。ū热纾築roker端沒有可以消費的消息),并不立即給消費者客戶端返回Response的響應(yīng),而是先hold住并且掛起請求(將請求保存至pullRequestTable本地緩存變量中),然后Broker端的后臺獨立線程—PullRequestHoldService會從pullRequestTable本地緩存變量中不斷地去取,具體的做法是查詢待拉取消息的偏移量是否小于消費隊列最大偏移量,如果條件成立則說明有新消息達到Broker端(這里,在RocketMQ的Broker端會有一個后臺獨立線程—ReputMessageService不停地構(gòu)建ConsumeQueue/IndexFile數(shù)據(jù),同時取出hold住的請求并進行二次處理),則通過重新調(diào)用一次業(yè)務(wù)處理器—PullMessageProcessor的處理請求方法—processRequest()來重新嘗試拉取消息(此處,每隔5S重試一次,默認長輪詢整體的時間設(shè)置為30s)。

RocketMQ消息Pull的長輪詢機制的關(guān)鍵在于Broker端的PullRequestHoldService和ReputMessageService兩個后臺線程。對于RocketMQ的長輪詢(LongPolling)消費模式后面會專門詳細介紹。

二、RocketMQ中兩種消費方式的demo代碼

(1)Pull模式的Consumer端代碼如下:

        DefaultMQPullConsumer consumer = new DefaultMQPullConsumer("please_rename_unique_group_name_5");
        consumer.setNamesrvAddr("127.0.0.1:9876");
        consumer.setInstanceName("consumer");
        consumer.start();
        Set<MessageQueue> mqs = consumer.fetchSubscribeMessageQueues("TopicTest111");
        for (MessageQueue mq : mqs) {
            System.out.printf("Consume from the queue: %s%n", mq);
            SINGLE_MQ:
            while (true) {
                try {
                    PullResult pullResult =
                        consumer.pullBlockIfNotFound(mq, null, getMessageQueueOffset(mq), 32);
                    System.out.printf("%s%n", pullResult);
                    putMessageQueueOffset(mq, pullResult.getNextBeginOffset());
                    switch (pullResult.getPullStatus()) {
                        case FOUND:
                            System.out.println(pullResult.getMsgFoundList().get(0).toString());
                            break;
                        case NO_NEW_MSG:
                            break SINGLE_MQ;
                        case NO_MATCHED_MSG:
                        case OFFSET_ILLEGAL:
                            break;
                        default:
                            break;
                    }
                } catch (Exception e) {
                    //TODO
                }
            }
        }
        consumer.shutdown();

在示例代碼中,可以看到業(yè)務(wù)工程在Consumer啟動后,Consumer主動獲取MessageQueue的Set集合,遍歷該集合中的每一個隊列,發(fā)送Pull的請求(參數(shù)中帶有隊列中的消息偏移量),同時需要Consumer端自己保存消息消費的offset偏移量至本地變量中。

在Pull模式下,需要業(yè)務(wù)應(yīng)用代碼自身去完成比較多的事情,因此在實際應(yīng)用中用的較少。(2)Push模式的Consumer端代碼如下:

DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("CID_JODIE_1");
        consumer.subscribe("TopicTest111", "*");
        consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
        consumer.setInstanceName("consumer1");
        consumer.setNamesrvAddr("127.0.0.1:9876");
        consumer.registerMessageListener(new MessageListenerConcurrently() {
            @Override
            public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, ConsumeConcurrentlyContext context) {
                System.out.printf("%s Receive New Messages: %s %n", Thread.currentThread().getName(), msgs);
                return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
            }
        });
        consumer.start();

在示例代碼中,業(yè)務(wù)工程的應(yīng)用程序使用Push方式進行消費時,Consumer端注冊了一個監(jiān)聽器,Consumer在收到消息后主動調(diào)用這個監(jiān)聽器完成消費并進行對應(yīng)的業(yè)務(wù)邏輯處理。

由此可見,業(yè)務(wù)應(yīng)用代碼只需要完成消息消費即可,無需參與MQ本身的一些任務(wù)處理(ps:業(yè)務(wù)代碼顯得更為簡潔一些)。

三、RocketMQ中消費者Push方式的啟動流程

這一節(jié)主要先講下RocketMQ消費者的啟動流程,看下在啟動的時候究竟完成了什么樣的操作。由于RocketMQ的DefaultMQPushConsumer和DefaultMQPullConsumer啟動流程大部分類似,而DefaultMQPushConsumer更為復(fù)雜一些,因此這一節(jié)內(nèi)容主要講的是DefaultMQPushConsumer啟動流程。Push方式的Consumer啟動流程的時序圖如下圖所示:

 從上面的時序圖上可以看出,Push方式的Consumer啟動流程完成的任務(wù)比較多,主要任務(wù)如下:

(1)設(shè)置consumerGroup、NameServer服務(wù)地址、消費起始偏移地址并根據(jù)參數(shù)Topic構(gòu)建Consumer端的SubscriptionData(訂閱關(guān)系值);

(2)在Consumer端注冊消費者監(jiān)聽器,當消息到來時完成消費消息;

(3)啟動defaultMQPushConsumerImpl實例,主要完成前置校驗、復(fù)制訂閱關(guān)系(將defaultMQPushConsumer的訂閱關(guān)系復(fù)制至rebalanceImpl中,包括retryTopic(重試主題)對應(yīng)的訂閱關(guān)系)、創(chuàng)建MQClientInstance實例、設(shè)置rebalanceImpl的各個屬性值、pullAPIWrapper包裝類對象的初始化、初始化offsetStore實例并加載消費進度、啟動消息消費服務(wù)線程以及在MQClientInstance中注冊consumer等任務(wù);

(4)啟動MQClientInstance實例,其中包括完成客戶端網(wǎng)絡(luò)通信線程、拉取消息服務(wù)線程、負載均衡服務(wù)線程和若干個定時任務(wù)的啟動;

(5)向所有的Broker端發(fā)送心跳(采用加鎖方式);

(6)最后,喚醒負載均衡服務(wù)線程在Consumer端開始負載均衡;

四、RocketMQ中Pull和Push兩種消費模式流程簡析

RocketMQ提供了兩種消費模式,Push和Pull,大多數(shù)場景使用的是Push模式,在源碼中這兩種模式分別對應(yīng)的是DefaultMQPushConsumer類和DefaultMQPullConsumer類。Push模式實際上在內(nèi)部還是使用的Pull方式實現(xiàn)的,通過Pull不斷地輪詢Broker獲取消息,當不存在新消息時,Broker端會掛起Pull請求,直到有新消息產(chǎn)生才取消掛起,返回新消息。

(1)RocketMQ的Pull消費模式流程簡析 RocketMQ的Pull模式相對來得簡單,從上面的demo代碼中可以看出,業(yè)務(wù)應(yīng)用代碼通過由Topic獲取到的MessageQueue直接拉取消息(最后真正執(zhí)行的是PullAPIWrapper的pullKernelImpl()方法,通過發(fā)送拉取消息的RPC請求給Broker端)。其中,消息消費的偏移量需要Consumer端自己去維護。

(2)RocketMQ的Push消費模式流程簡析 在本文前面已經(jīng)提到過了,從嚴格意義上說,RocketMQ并沒有實現(xiàn)真正的消息消費的Push模式,而是對Pull模式進行了一定的優(yōu)化,一方面在Consumer端開啟后臺獨立的線程—PullMessageService不斷地從阻塞隊列—pullRequestQueue中獲取PullRequest請求并通過網(wǎng)絡(luò)通信模塊發(fā)送Pull消息的RPC請求給Broker端。另外一方面,后臺獨立線程—rebalanceService根據(jù)Topic中消息隊列個數(shù)和當前消費組內(nèi)消費者個數(shù)進行負載均衡,將產(chǎn)生的對應(yīng)PullRequest實例放入阻塞隊列—pullRequestQueue中。這里算是比較典型的生產(chǎn)者-消費者模型,實現(xiàn)了準實時的自動消息拉取。然后,再根據(jù)業(yè)務(wù)反饋是否成功消費來推動消費進度。 在Broker端,PullMessageProcessor業(yè)務(wù)處理器收到Pull消息的RPC請求后,通過MessageStore實例從commitLog獲取消息。如1.2節(jié)內(nèi)容所述,如果第一次嘗試Pull消息失?。ū热鏐roker端沒有可以消費的消息),則通過長輪詢機制先hold住并且掛起該請求,然后通過Broker端的后臺線程PullRequestHoldService重新嘗試和后臺線程ReputMessageService的二次處理。

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

相關(guān)文章

  • springboot創(chuàng)建文件夾失敗的解決

    springboot創(chuàng)建文件夾失敗的解決

    這篇文章主要介紹了springboot創(chuàng)建文件夾失敗的解決方案,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2022-01-01
  • 使用@Value 注入 List 類型的配置屬性需要注意的 BUG

    使用@Value 注入 List 類型的配置屬性需要注意的 BUG

    這篇文章主要介紹了使用@Value 注入 List 類型的配置屬性需要注意的 BUG,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-08-08
  • Spring+SpringMVC+MyBatis深入學(xué)習及搭建(一)之MyBatis的基礎(chǔ)知識

    Spring+SpringMVC+MyBatis深入學(xué)習及搭建(一)之MyBatis的基礎(chǔ)知識

    這篇文章主要介紹了Spring+SpringMVC+MyBatis深入學(xué)習及搭建(一)之MyBatis的基礎(chǔ)知識,需要的朋友可以參考下
    2017-05-05
  • IDEA自帶Maven插件找不到settings.xml配置文件

    IDEA自帶Maven插件找不到settings.xml配置文件

    IDEA自帶了Maven插件,最近發(fā)現(xiàn)了一個問題,IDEA自帶Maven插件找不到settings.xml配置文件,本文就來詳細的介紹一下解決方法,感興趣的可以了解一下
    2023-11-11
  • Java中的main函數(shù)的詳細介紹

    Java中的main函數(shù)的詳細介紹

    這篇文章主要介紹了Java中的main函數(shù)的詳細介紹的相關(guān)資料,main()函數(shù)在java程序中必出現(xiàn)的函數(shù),這里就講解下使用方法,需要的朋友可以參考下
    2017-09-09
  • java 數(shù)據(jù)庫連接與增刪改查操作實例詳解

    java 數(shù)據(jù)庫連接與增刪改查操作實例詳解

    這篇文章主要介紹了java 數(shù)據(jù)庫連接與增刪改查操作,結(jié)合實例形式詳細分析了java使用jdbc進行數(shù)據(jù)庫連接及增刪改查等相關(guān)操作實現(xiàn)技巧與注意事項,需要的朋友可以參考下
    2019-11-11
  • Tomcat?8.5?+mysql?5.7+jdk1.8開發(fā)JavaSE的金牌榜小項目

    Tomcat?8.5?+mysql?5.7+jdk1.8開發(fā)JavaSE的金牌榜小項目

    這篇文章主要介紹了Tomcat?8.5?+mysql?5.7+jdk1.8開發(fā)JavaSE的金牌榜小項目,本文通過圖文實例相結(jié)合給大家介紹的非常詳細,對大家的學(xué)習或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2022-05-05
  • Spring Cloud OpenFeign實現(xiàn)動態(tài)服務(wù)名調(diào)用的示例代碼

    Spring Cloud OpenFeign實現(xiàn)動態(tài)服務(wù)名調(diào)用的示例代碼

    在微服務(wù)架構(gòu)中,我們經(jīng)常需要根據(jù)動態(tài)傳入的服務(wù)名來遠程調(diào)用其他服務(wù),例如,你的業(yè)務(wù)中可能有多個子服務(wù):service-1、service-2……需要動態(tài)決定調(diào)用哪個,所以本文給大家介紹了Spring Cloud OpenFeign 實現(xiàn)動態(tài)服務(wù)名調(diào)用指南,需要的朋友可以參考下
    2025-06-06
  • Java LongAdder原理解析與實戰(zhàn)應(yīng)用小結(jié)

    Java LongAdder原理解析與實戰(zhàn)應(yīng)用小結(jié)

    LongAdder是Java 8中java.util.concurrent.atomic包引入的高性能計數(shù)器類,專為高并發(fā)場景下的數(shù)值累加操作優(yōu)化設(shè)計,本文給大家介紹Java LongAdder原理解析與實戰(zhàn)應(yīng)用小結(jié),感興趣的朋友一起看看吧
    2025-06-06
  • idea2023.3安裝及配置詳細圖文教程

    idea2023.3安裝及配置詳細圖文教程

    IDEA全稱IntelliJ?IDEA,是Java語言對的集成開發(fā)環(huán)境,IDEA在業(yè)界被認為是公認最好的Java開發(fā)工具,這篇文章主要給大家介紹了關(guān)于idea2023.3安裝及配置的相關(guān)資料,需要的朋友可以參考下
    2023-11-11

最新評論

马山县| 杂多县| 青海省| 尼勒克县| 焦作市| 宁陕县| 长沙县| 扎兰屯市| 江都市| 城口县| 景泰县| 景德镇市| 稷山县| 昭通市| 泗水县| 娱乐| 恩施市| 新巴尔虎右旗| 古交市| 益阳市| 准格尔旗| 囊谦县| 体育| 通化县| 章丘市| 东至县| 娄烦县| 凤城市| 呼伦贝尔市| 怀集县| 府谷县| 高密市| 忻州市| 保亭| 进贤县| 合江县| 婺源县| 阳泉市| 栾城县| 安阳市| 慈利县|