五伊群、RocketMQ-Consumer啟動流程

一、概述

一個最簡單的Consumer的啟動代碼如下:

public static void main(String[] args) throws Exception{
            // Instantiate with specified consumer group name.
            DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("c1");
            // Specify name server addresses.
            consumer.setNamesrvAddr("10.1.11.155:9876");
            // Subscribe one more more topics to consume.
            consumer.subscribe("qqq", "TagA||TagB");
            // Register callback to execute on arrival of messages fetched from brokers.
            consumer.registerMessageListener(new MessageListenerConcurrently() {
                @Override
                public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs,
                                                                ConsumeConcurrentlyContext context) {
                    msgs.forEach(m-> System.out.println(new String(m.getBody())));
                    return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
                }
            });
            //Launch the consumer instance.
            consumer.start();
            System.out.printf("Consumer Started.%n");
        }

最重要的有幾個步驟:

  • new DefaultMQPushConsumer("c1");
  • consumer.subscribe("qqq", "TagA||TagB");
  • registerMessageListener()
  • start()
    下面就這幾個方法進(jìn)行深入

二问拘、實例化一個DefaultMQPushConsumer

通過下面的代碼就能夠?qū)嵗粋€DefaultMQPushConsumer

DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("c1");

這個默認(rèn)構(gòu)造方法其實調(diào)用了下面的構(gòu)造方法

public DefaultMQPushConsumer(final String consumerGroup, RPCHook rpcHook,
        AllocateMessageQueueStrategy allocateMessageQueueStrategy) {
        this.consumerGroup = consumerGroup;
        this.allocateMessageQueueStrategy = allocateMessageQueueStrategy;
        defaultMQPushConsumerImpl = new DefaultMQPushConsumerImpl(this, rpcHook);
    }

一共三個參數(shù):

  • consumerGroup
    消費(fèi)者組
  • rpcHook
    遠(yuǎn)程調(diào)用的hook赦抖,我看Rocket的德行一般就在RPC調(diào)用之前,針對訪問控制(ACL)做了一個前置
  • allocateMessageQueueStrategy
    隊列負(fù)載均衡策略宦赠,默認(rèn)會用 AllocateMessageQueueAveragely(平均算法)
    系統(tǒng)提供了6種策略:后續(xù)研究下
    • AllocateMachineRoomNearby
    • AllocateMessageQueueAveragely
    • AllocateMessageQueueAveragelyByCircle
    • AllocateMessageQueueByConfig
    • AllocateMessageQueueByMachineRoom
    • AllocateMessageQueueConsistentHash

三陪毡、定義監(jiān)聽的主題及subExpression

代碼如下:

consumer.subscribe("qqq", "TagA||TagB");

解析一下topic及subExpression米母,生成一個SubscriptionData實體(包含了topic、subExpression的value和hash值)毡琉,并放到map中

四铁瞒、注冊監(jiān)聽處理器 registerMessageListener

consumer.registerMessageListener(new MessageListenerConcurrently() {
                @Override
                public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs,
                                                                ConsumeConcurrentlyContext context) {
                    // msgs.forEach(m-> System.out.println(new String(m.getBody())));
                    return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
                }
            });

新建了MessageListenerConcurrently接口的匿名類,并實現(xiàn)了consumeMessage()方法桅滋,注意慧耍,這里MessageListenerConcurrently是一個接口,繼承自MessageListener丐谋,同樣繼承自MessageListener的還有MessageListenerOrderly(順序消費(fèi)監(jiān)聽器)

五芍碧、啟動 start()

這個是核心的核心了

1、this.checkConfig();

檢查group等參數(shù)是否合法

2号俐、this.copySubscription();

復(fù)制監(jiān)聽規(guī)則泌豆,其實就是把之前生成的SubscriptionData復(fù)制一份,交給RebalancePushImpl(這個應(yīng)該是負(fù)載均衡消費(fèi)的核心類)吏饿,除了復(fù)制監(jiān)聽規(guī)則踪危,還會把以下信息復(fù)制過去:

  • setConsumerGroup 復(fù)制組名稱
  • setMessageModel 復(fù)制消費(fèi)模式:廣播/集群
  • setAllocateMessageQueueStrategy 復(fù)制負(fù)載均衡算法
  • setmQClientFactory(在下面3中初始化) 復(fù)制mqClient

3、實例化一個MQClientInstance

this.mQClientFactory = MQClientManager.getInstance().getAndCreateMQClientInstance(this.defaultMQPushConsumer, this.rpcHook);

4猪落、指定offsetStore

如果是集群模式贞远,則會實例化一個LocalFileOffsetStore
如果是廣播模式,則會實例化一個RemoteBrokerOffsetStore
都會調(diào)用load的方法笨忌,但是RemoteBrokerOffsetStore是一個空實現(xiàn)

 switch (this.defaultMQPushConsumer.getMessageModel()) {
                        case BROADCASTING:
                            this.offsetStore = new LocalFileOffsetStore(this.mQClientFactory, this.defaultMQPushConsumer.getConsumerGroup());
                            break;
                        case CLUSTERING:
                            this.offsetStore = new RemoteBrokerOffsetStore(this.mQClientFactory, this.defaultMQPushConsumer.getConsumerGroup());
                            break;
                        default:
                            break;
                    }
      this.defaultMQPushConsumer.setOffsetStore(this.offsetStore);
}
 this.offsetStore.load();

5蓝仲、實例化ConsumeMessageService

主要有兩種:并發(fā)消費(fèi)service及順序消費(fèi)service

  • 并發(fā)消費(fèi) ConsumeMessageOrderlyService.start();
    start() 會啟動一個定時任務(wù),延遲15分鐘蜜唾,每15分鐘執(zhí)行一次杂曲,清理失效的消息處理線程
  • 順序消費(fèi) ConsumeMessageConcurrentlyService
    start() 會啟動一個定時任務(wù),延遲1秒袁余,每20秒執(zhí)行一次 lockMQPeriodically()方法擎勘,這個方法應(yīng)該是順序消費(fèi)的精髓,可以深入研讀一下

6颖榜、啟動客戶端 mQClientFactory.start();

producer啟動的時候也會調(diào)用這個方法
上面1 ~ 5 的步驟棚饵,大部分是為了這個start()做準(zhǔn)備,這個start做了很多事情掩完,如下:

  • this.mQClientAPIImpl.fetchNameServerAddr();
    如果沒有指定namesrv的地址噪漾,則會通過RPC獲取,這個最終是通過調(diào)用一個外部的URL連接來獲取的且蓬,這個鏈接可以自己定義欣硼,官方也推薦這種方法,最底層也最靈活恶阴。
  • this.mQClientAPIImpl.start();
    啟動netty客戶端诈胜,可以和brokernamesrv通訊的通道
  • this.startScheduledTask();
    啟動一批自動任務(wù)
    • MQClientInstance.this.mQClientAPIImpl.fetchNameServerAddr();
      定期獲取namesrv的地址豹障,前面在啟動的時候已經(jīng)獲取一次了
    • MQClientInstance.this.updateTopicRouteInfoFromNameServer();
      定期從namesrv拉取meta信息及topic等路由信息
    • cleanOfflineBroker()、sendHeartbeatToAllBrokerWithLock()
      定期檢查下線的broker及發(fā)送心跳
    • MQClientInstance.this.persistAllConsumerOffset();
      定期固化消費(fèi)者的offset
    • MQClientInstance.this.adjustThreadPool();
      定期調(diào)整處理線程池的大小焦匈,這個4.4版本好像是個空實現(xiàn)血公,沒用
  • this.pullMessageService.start();
    啟動客戶端主動pull的服務(wù)
  • this.rebalanceService.start();
    啟動負(fù)載均衡服務(wù)

最后還會調(diào)用一次getDefaultMQProducerImpl().start(false),沒錯缓熟,就是producerstart累魔,但是參數(shù)是false,不會最終調(diào)用producerstart够滑,就把配置刷過去了垦写。

7、其他方法

  • this.updateTopicSubscribeInfoWhenSubscriptionChanged();
  • this.mQClientFactory.checkClientInBroker();
  • this.mQClientFactory.sendHeartbeatToAllBrokerWithLock();
  • this.mQClientFactory.rebalanceImmediately();
    這幾個方法后續(xù)研究下

到此為止彰触,consumer的啟動流程完畢

?著作權(quán)歸作者所有,轉(zhuǎn)載或內(nèi)容合作請聯(lián)系作者
  • 序言:七十年代末梯澜,一起剝皮案震驚了整個濱河市,隨后出現(xiàn)的幾起案子渴析,更是在濱河造成了極大的恐慌,老刑警劉巖吮龄,帶你破解...
    沈念sama閱讀 211,265評論 6 490
  • 序言:濱河連續(xù)發(fā)生了三起死亡事件俭茧,死亡現(xiàn)場離奇詭異,居然都是意外死亡漓帚,警方通過查閱死者的電腦和手機(jī)母债,發(fā)現(xiàn)死者居然都...
    沈念sama閱讀 90,078評論 2 385
  • 文/潘曉璐 我一進(jìn)店門,熙熙樓的掌柜王于貴愁眉苦臉地迎上來尝抖,“玉大人毡们,你說我怎么就攤上這事∶亮桑” “怎么了衙熔?”我有些...
    開封第一講書人閱讀 156,852評論 0 347
  • 文/不壞的土叔 我叫張陵,是天一觀的道長搅荞。 經(jīng)常有香客問我红氯,道長,這世上最難降的妖魔是什么咕痛? 我笑而不...
    開封第一講書人閱讀 56,408評論 1 283
  • 正文 為了忘掉前任痢甘,我火速辦了婚禮,結(jié)果婚禮上茉贡,老公的妹妹穿的比我還像新娘塞栅。我一直安慰自己,他們只是感情好腔丧,可當(dāng)我...
    茶點(diǎn)故事閱讀 65,445評論 5 384
  • 文/花漫 我一把揭開白布放椰。 她就那樣靜靜地躺著作烟,像睡著了一般。 火紅的嫁衣襯著肌膚如雪庄敛。 梳的紋絲不亂的頭發(fā)上俗壹,一...
    開封第一講書人閱讀 49,772評論 1 290
  • 那天,我揣著相機(jī)與錄音藻烤,去河邊找鬼绷雏。 笑死,一個胖子當(dāng)著我的面吹牛怖亭,可吹牛的內(nèi)容都是我干的涎显。 我是一名探鬼主播,決...
    沈念sama閱讀 38,921評論 3 406
  • 文/蒼蘭香墨 我猛地睜開眼兴猩,長吁一口氣:“原來是場噩夢啊……” “哼期吓!你這毒婦竟也來了?” 一聲冷哼從身側(cè)響起倾芝,我...
    開封第一講書人閱讀 37,688評論 0 266
  • 序言:老撾萬榮一對情侶失蹤讨勤,失蹤者是張志新(化名)和其女友劉穎,沒想到半個月后晨另,有當(dāng)?shù)厝嗽跇淞掷锇l(fā)現(xiàn)了一具尸體潭千,經(jīng)...
    沈念sama閱讀 44,130評論 1 303
  • 正文 獨(dú)居荒郊野嶺守林人離奇死亡,尸身上長有42處帶血的膿包…… 初始之章·張勛 以下內(nèi)容為張勛視角 年9月15日...
    茶點(diǎn)故事閱讀 36,467評論 2 325
  • 正文 我和宋清朗相戀三年借尿,在試婚紗的時候發(fā)現(xiàn)自己被綠了刨晴。 大學(xué)時的朋友給我發(fā)了我未婚夫和他白月光在一起吃飯的照片。...
    茶點(diǎn)故事閱讀 38,617評論 1 340
  • 序言:一個原本活蹦亂跳的男人離奇死亡路翻,死狀恐怖狈癞,靈堂內(nèi)的尸體忽然破棺而出,到底是詐尸還是另有隱情茂契,我是刑警寧澤蝶桶,帶...
    沈念sama閱讀 34,276評論 4 329
  • 正文 年R本政府宣布,位于F島的核電站掉冶,受9級特大地震影響莫瞬,放射性物質(zhì)發(fā)生泄漏。R本人自食惡果不足惜郭蕉,卻給世界環(huán)境...
    茶點(diǎn)故事閱讀 39,882評論 3 312
  • 文/蒙蒙 一疼邀、第九天 我趴在偏房一處隱蔽的房頂上張望。 院中可真熱鬧召锈,春花似錦旁振、人聲如沸。這莊子的主人今日做“春日...
    開封第一講書人閱讀 30,740評論 0 21
  • 文/蒼蘭香墨 我抬頭看了看天上的太陽吉嚣。三九已至,卻和暖如春蹬铺,著一層夾襖步出監(jiān)牢的瞬間尝哆,已是汗流浹背。 一陣腳步聲響...
    開封第一講書人閱讀 31,967評論 1 265
  • 我被黑心中介騙來泰國打工甜攀, 沒想到剛下飛機(jī)就差點(diǎn)兒被人妖公主榨干…… 1. 我叫王不留秋泄,地道東北人。 一個月前我還...
    沈念sama閱讀 46,315評論 2 360
  • 正文 我出身青樓规阀,卻偏偏與公主長得像恒序,于是被迫代替她去往敵國和親。 傳聞我的和親對象是個殘疾皇子谁撼,可洞房花燭夜當(dāng)晚...
    茶點(diǎn)故事閱讀 43,486評論 2 348

推薦閱讀更多精彩內(nèi)容

  • 大致可以通過上述情況進(jìn)行排除 1.kafka服務(wù)器問題 查看日志是否有報錯歧胁,網(wǎng)絡(luò)訪問問題等。 2. kafka p...
    生活的探路者閱讀 7,578評論 0 10
  • MQ在我們?nèi)粘i_發(fā)過程中有著不可替代的作用厉碟,不僅可以幫助我們做到信息在系統(tǒng)間的傳遞喊巍,還能進(jìn)行系統(tǒng)間的解耦合,也就是...
    數(shù)齊閱讀 3,433評論 2 7
  • 每個人的想法不同 箍鼓, RocketMQ 介紹的時候就說 是阿里從他們使用的上 解耦出來 近一步簡化 便捷的 目...
    樓亭樵客閱讀 403評論 0 0
  • consumer 1.啟動 有別于其他消息中間件由broker做負(fù)載均衡并主動向consumer投遞消息玄糟,Rock...
    veShi文閱讀 4,929評論 0 2
  • Swift1> Swift和OC的區(qū)別1.1> Swift沒有地址/指針的概念1.2> 泛型1.3> 類型嚴(yán)謹(jǐn) 對...
    cosWriter閱讀 11,090評論 1 32