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

RocketMQ?Push?消費(fèi)模型示例詳解

 更新時(shí)間:2022年09月20日 15:23:42   作者:磊叔的技術(shù)博客  
這篇文章主要為大家介紹了RocketMQ?Push?消費(fèi)模型示例詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪

Push 模式是指由 Server 端來(lái)控制消息的推送,即當(dāng)有消息到 Server 之后,會(huì)將消息主動(dòng)投遞給 client(Consumer 端)。

使用 DefaultMQPushConsumer 消費(fèi)消息

下面是使用 DefaultMQPushConsumer 消費(fèi)消息的官方示例代碼:

// 初始化consumer,并設(shè)置consumer group name
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("MyGroup");
// 設(shè)置NameServer地址
consumer.setNamesrvAddr("localhost:9876");
//訂閱一個(gè)或多個(gè)topic,并指定tag過濾條件,這里指定*表示接收所有tag的消息
consumer.subscribe("TopicTest", "*");
//注冊(cè)回調(diào)接口來(lái)處理從Broker中收到的消息
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);
        // 返回消息消費(fèi)狀態(tài),ConsumeConcurrentlyStatus.CONSUME_SUCCESS 為消費(fèi)成功
        return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
    }
});
// 啟動(dòng)Consumer
consumer.start();

這里看到主要是通過 consumer 注冊(cè)回調(diào)接口來(lái)處理從 Broker 中收到的消息。這種監(jiān)聽回調(diào)的機(jī)制很容易想到是一種觀察者模式或者事件機(jī)制;對(duì)于這種 C-S 模型的架構(gòu)來(lái)說,如果要做到 Server 在有新消息時(shí)立即推送給 Client,那么 Client 和 Server 之間應(yīng)該是有連接存在的,Client 端開放端口來(lái) watch Server 的推送。這里好論證,即可以查看當(dāng)前 Client 端所在進(jìn)程開啟了什么端口即可,通過如下指令查看:

  • 1、先通過 jps 查看 Consumer Client 的進(jìn)程號(hào)
?  rocketmq-4.9.4 git:(06f2208a3) jps
10722 Jps
4676 rocketmq-dashboard-1.0.1-SNAPSHOT.jar
1766
4121 BrokerStartup
4009 NamesrvStartup
9419 PushConsumer
9692 RemoteMavenServer36

可以看到 PushConsumer 的進(jìn)程號(hào)是 9419

  • 2、通過 lsof 命令查看進(jìn)程端口占用
?  rocketmq-4.9.4 git:(06f2208a3) lsof -nP -p 9419| grep LISTEN
?  

這里沒有看到 PushConsumer 有開啟端口。同樣,這里可以看看 Broker 的進(jìn)程端口占用

?  rocketmq-4.9.4 git:(06f2208a3) lsof -nP -p 4121| grep LISTEN
java    4121 glmapper  137u    IPv6 0xca1142b0f200067d        0t0                 TCP *:10912 (LISTEN)
java    4121 glmapper  141u    IPv6 0xca1142b0f1fc8cfd        0t0                 TCP *:10911 (LISTEN)
java    4121 glmapper  142u    IPv6 0xca1142b0f1fc935d        0t0                 TCP *:10909 (LISTEN)

所以得到一個(gè)初步的結(jié)論是,在 Push 模式下,Consumer Client 并沒有啟動(dòng)端口來(lái)接收 Server 的消息推送。 那么 RocketMQ 是怎么實(shí)現(xiàn)的?

基于長(zhǎng)輪詢機(jī)制的偽 push 實(shí)現(xiàn)

真正的 Push 方式,是 Server 端接收到消息后,主動(dòng)把消息推送給 Client 端,這種情況一般需要 Client 和 Server 之間建立長(zhǎng)連接。通過前面的分析,Client 既然沒有開啟端口用于接收 Server 的信息推送,那么只有一種可能就是 Client 自己去拉了消息,但是這種主動(dòng)拉消息的方式是對(duì)于用戶無(wú)感的,從使用上體驗(yàn)上來(lái)看,做到了和 push 一樣的效果;這種機(jī)制就是“長(zhǎng)輪詢”。

為啥不用長(zhǎng)連接方式,讓 Server 主動(dòng) Push 呢?其實(shí)很好理解,對(duì)于一個(gè)提供隊(duì)列服務(wù)的 Server 來(lái)說,用 Push方式主動(dòng)推送有兩個(gè)問題:

  • 1、會(huì)增加 Server 端的工作量,進(jìn)而影響 Server 的性能
  • 2、Client 的處理能力存在差異,Client 的狀態(tài)不受 Server 控制,如果 Client 不能及時(shí)處理 Server 推送過來(lái)的消息,會(huì)造成各種潛在問題

客戶端側(cè)發(fā)起的長(zhǎng)輪詢請(qǐng)求

下圖是初始化相關(guān)資源的過程,DefaultMQPushConsumer 是面向用戶使用的 API client 類,內(nèi)部處理實(shí)際上是委托給 DefaultMQPushConsumerImpl 來(lái)處理的。DefaultMQPushConsumerImpl#start 時(shí),會(huì)初始化 MQClientInstance ,MQClientInstance 初始化過程中又會(huì)初始化一堆資源,比如請(qǐng)求-響應(yīng)的通道,開啟各種各樣的調(diào)度任務(wù)(定期拉去 NameServerAddress、定期更新 Topic 路由信息、定期清理 Offline狀態(tài)的 Broker、定期發(fā)送心跳給 Broker、定期持久化所有 Consumer Offset等等),開啟 pullMessageService,開啟 rebalance Service 等等。大致的調(diào)用鏈如下圖

下面這個(gè)代碼片段是 pullMessageService 的 run 方法(pullMessageService 是 Runnable 子類)

@Override
public void run() {
    log.info(this.getServiceName() + " service started");
    while (!this.isStopped()) {
        try {
            // 從 pullRequestQueue 中取 pullRequest
            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 會(huì)一直從 pullRequestQueue 中取 pullRequest,然后執(zhí)行 pullMessage 請(qǐng)求。實(shí)際上 MessageQueue 是和 pullRequest 一一對(duì)應(yīng)的 ,pullRequest 全部存儲(chǔ)到該 Consumer 的 pullRequestQueue 隊(duì)列里面;消費(fèi)者會(huì)不停的從 PullRequest 的隊(duì)列里取 request 然后向broker 請(qǐng)求消息。

這里還有一個(gè)問題是隊(duì)列取出之后什么時(shí)候放回去的?在 pullMessage 的回調(diào)方法中,如果正常得到了 broker 的響應(yīng),那么會(huì)把 PullRequest放回隊(duì)列,相關(guān)代碼可以從 org.apache.rocketmq.client.consumer.PullCallbackonSuccess 方法中得到答案。

服務(wù)端阻塞請(qǐng)求

服務(wù)端處理 pullRequest 請(qǐng)求的是 PullMessageProcessor,當(dāng)沒有消息時(shí),則通過 PullRequestHoldService 將當(dāng)前請(qǐng)求先 hold 住。

case ResponseCode.PULL_NOT_FOUND:
    if (brokerAllowSuspend && hasSuspendFlag) {
        long pollingTimeMills = suspendTimeoutMillisLong;
        // 如果是 LongPolling,則 hold 住
        if (!this.brokerController.getBrokerConfig().isLongPollingEnable()) {
            pollingTimeMills = this.brokerController.getBrokerConfig().getShortPollingTimeMills();
        }
        String topic = requestHeader.getTopic();
        long offset = requestHeader.getQueueOffset();
        int queueId = requestHeader.getQueueId();
        PullRequest pullRequest = new PullRequest(request, channel, pollingTimeMills,
            this.brokerController.getMessageStore().now(), offset, subscriptionData, messageFilter);
        this.brokerController.getPullRequestHoldService().suspendPullRequest(topic, queueId, pullRequest);
        response = null;
        break;
    }

PullRequestHoldService 中會(huì)將所有的 PullRequest 緩存到 pullRequestTable。PullRequestHoldService 也是一個(gè) task,默認(rèn)每次 hold 5s 然后再去檢查是否有新的消息過來(lái),如果有新的消息到來(lái),則喚醒對(duì)應(yīng)的線程來(lái)將消息返回給客戶端。

// 已省略無(wú)關(guān)代碼
public void run() {
    // loop
    while (!this.isStopped()) {
        // default hold 5s
        if (this.brokerController.getBrokerConfig().isLongPollingEnable()) {
            this.waitForRunning(5 * 1000);
        } else {
            this.waitForRunning(this.brokerController.getBrokerConfig().getShortPollingTimeMills());
        }
        long beginLockTimestamp = this.systemClock.now();
        // 檢查是否有新的消息到達(dá)
        this.checkHoldRequest();
        long costTime = this.systemClock.now() - beginLockTimestamp;
        if (costTime > 5 * 1000) {
            log.info("[NOTIFYME] check hold request cost {} ms.", costTime);
        }
	}
}

客戶端回調(diào)處理

我們?cè)诰帉?consumer 代碼時(shí),基于 push 模式是通過如下方式來(lái)監(jiān)聽消息的

//注冊(cè)回調(diào)接口來(lái)處理從Broker中收到的消息
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);
        // 返回消息消費(fèi)狀態(tài),ConsumeConcurrentlyStatus.CONSUME_SUCCESS 為消費(fèi)成功
        return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
    }
});

通過前面的分析,對(duì)于如何通過“長(zhǎng)輪詢”實(shí)現(xiàn)偽“push” 有了大概得了解;客戶端通過一個(gè)定時(shí)任務(wù)不斷向 Broker 發(fā)請(qǐng)求,Broker 在沒有消息時(shí)先 hold 住一小段時(shí)間,當(dāng)有新的消息時(shí)會(huì)立即將消息返回給 consumer;本節(jié)就主要探討 consumer 在收到消息之后的處理邏輯,以及是怎么觸發(fā) MessageListener 回調(diào)執(zhí)行的。

客戶端發(fā)起請(qǐng)求的底層邏輯

以異步調(diào)用為例,代碼在

org.apache.rocketmq.client.impl.MQClientAPIImpl#pullMessageAsync中,截取部分代碼如下:

this.remotingClient.invokeAsync(addr, request, timeoutMillis, new InvokeCallback() {
    @Override
    public void operationComplete(ResponseFuture responseFuture) {
        RemotingCommand response = responseFuture.getResponseCommand();
        if (response != null) {
            try {
                PullResult pullResult = MQClientAPIImpl.this.processPullResponse(response, addr);
                assert pullResult != null;
                // 成功回調(diào)
                pullCallback.onSuccess(pullResult);
            } catch (Exception e) {
                // 異常回調(diào)
                pullCallback.onException(e);
            }
        } else {
            if (!responseFuture.isSendRequestOK()) {
                 // 異?;卣{(diào)
                pullCallback.onException(new MQClientException("send request failed to " + addr + ". Request: " + request, responseFuture.getCause()));
            } else if (responseFuture.isTimeout()) {
                 // 異?;卣{(diào)
                pullCallback.onException(new MQClientException("wait response from " + addr + " timeout :" + responseFuture.getTimeoutMillis() + "ms" + ". Request: " + request,
                    responseFuture.getCause()));
            } else {
                 // 異常回調(diào)
                pullCallback.onException(new MQClientException("unknown reason. addr: " + addr + ", timeoutMillis: " + timeoutMillis + ". Request: " + request, responseFuture.getCause()));
            }
        }
    }
});

PullCallback 回調(diào)

PullCallback 回調(diào)邏輯在 org.apache.rocketmq.client.impl.consumer.DefaultMQPushConsumerImpl#pullMessage方法中,以正常返回消息為例:

// 已省略無(wú)關(guān)代碼
public void onSuccess(PullResult pullResult) {
    // 將接收到的消息 交給 consumeMessageService 處理
    DefaultMQPushConsumerImpl.this.consumeMessageService.submitConsumeRequest(
        pullResult.getMsgFoundList(),
        processQueue,
        pullRequest.getMessageQueue(),
        dispatchToConsume);
    // 將 pullRequest 放回 pullRequestQueue
 DefaultMQPushConsumerImpl.this.executePullRequestImmediately(pullRequest);
}

ConsumeRequest 是一個(gè) Runnable,submitConsumeRequest 就是將返回結(jié)果丟在一個(gè)單獨(dú)的線程池中去處理返回結(jié)果的。ConsumeRequest 的 run 方法中,會(huì)拿到 messageListener,然后執(zhí)行 consumeMessage 方法。

總結(jié)

到此,關(guān)于 RocketMQ push 消費(fèi)模型基本就探討完了。從實(shí)現(xiàn)機(jī)制上來(lái)看,push 本質(zhì)上并不是在建立雙向通道的前提下,由 Server 主動(dòng)推送給 Client 的,而是由 Client 端觸發(fā) pullRequest 請(qǐng)求,以長(zhǎng)輪詢的方式“偽裝”的結(jié)果。從代碼上來(lái),RocketMQ 代碼中使用了非常多的異步機(jī)制,如 pullRequestQueue 來(lái)解耦發(fā)送請(qǐng)求和等待結(jié)果,各種定時(shí)任務(wù)等等。

整體看,PushConsumer 采用了 長(zhǎng)輪詢+超時(shí)時(shí)間+Pull的模式, 這種方式帶來(lái)的好處總結(jié)如下

  • 1、減少 Broker 的壓力,避免由于不同 Consumer 消費(fèi)能力導(dǎo)致 Broker 出現(xiàn)問題
  • 2、確保了 Consumer 不會(huì)負(fù)載過高,Consumer 在校驗(yàn)自己的緩存消息沒有超過閾值才會(huì)去從 Broker 拉取消息,Broker 不會(huì)主動(dòng)推過來(lái)
  • 3、兼顧了消息的即時(shí)性,Broker 在沒有消息的時(shí)候會(huì)先 hold 一小段時(shí)間,有消息會(huì)立即喚起線程將消息返回給 Consumer
  • 4、Broker 端無(wú)效請(qǐng)求的次數(shù)大大降低,Broker 在沒有消息時(shí)會(huì)掛起 PullRequest,而 Consumer 在未接收到Response 且未超時(shí)時(shí),也不會(huì)重新發(fā)起 PullRequest

以上就是RocketMQ Push 消費(fèi)模型示例詳解的詳細(xì)內(nèi)容,更多關(guān)于RocketMQ Push 消費(fèi)模型的資料請(qǐng)關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

  • Java讀取Excel文件內(nèi)容的簡(jiǎn)單實(shí)例

    Java讀取Excel文件內(nèi)容的簡(jiǎn)單實(shí)例

    這篇文章主要介紹了Java讀取Excel文件內(nèi)容的簡(jiǎn)單實(shí)例,有需要的朋友可以參考一下
    2013-11-11
  • 詳解Spring中bean實(shí)例化的三種方式

    詳解Spring中bean實(shí)例化的三種方式

    本篇文章主要介紹了詳解Spring中bean實(shí)例化的三種方式,小編覺得挺不錯(cuò)的,現(xiàn)在分享給大家,也給大家做個(gè)參考。一起跟隨小編過來(lái)看看吧
    2017-04-04
  • SpringBoot集成Druid配置(yaml版本配置文件)詳解

    SpringBoot集成Druid配置(yaml版本配置文件)詳解

    這篇文章主要介紹了SpringBoot集成Druid配置(yaml版本配置文件),本文通過實(shí)例代碼給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2020-12-12
  • Java超詳細(xì)講解類的繼承

    Java超詳細(xì)講解類的繼承

    繼承就是可以直接使用前輩的屬性和方法。自然界如果沒有繼承,那一切都是處于混沌狀態(tài)。多態(tài)是同一個(gè)行為具有多個(gè)不同表現(xiàn)形式或形態(tài)的能力。多態(tài)就是同一個(gè)接口,使用不同的實(shí)例而執(zhí)行不同操作
    2022-04-04
  • Java?JSON處理庫(kù)之Gson的用法詳解

    Java?JSON處理庫(kù)之Gson的用法詳解

    Gson是Google開發(fā)的一款Java?JSON處理庫(kù),旨在簡(jiǎn)化Java開發(fā)人員操作JSON數(shù)據(jù)的過程,本文就來(lái)和大家簡(jiǎn)單聊聊Gson的原理與具體使用吧
    2023-05-05
  • IDEA中設(shè)置Tab健為4個(gè)空格的方法

    IDEA中設(shè)置Tab健為4個(gè)空格的方法

    這篇文章給大家介紹了代碼縮進(jìn)用空格還是Tab?(IDEA中設(shè)置Tab健為4個(gè)空格)的相關(guān)知識(shí),本文給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友參考下吧
    2021-03-03
  • SpringMVC Idea 搭建 部署war的詳細(xì)過程

    SpringMVC Idea 搭建 部署war的詳細(xì)過程

    本文介紹了如何在IntelliJ IDEA中使用Maven模板創(chuàng)建一個(gè)Web項(xiàng)目,并詳細(xì)說明了如何配置web.xml、創(chuàng)建springmvc-servlet.xml和application.properties文件,以及如何使用Maven打包生成WAR文件并部署到Tomcat服務(wù)器,感興趣的朋友跟隨小編一起看看吧
    2025-01-01
  • Java中內(nèi)存問題之OOM詳解

    Java中內(nèi)存問題之OOM詳解

    這篇文章主要介紹了Java中內(nèi)存管理的OOM詳解,OOM,全稱“Out?Of?Memory”,翻譯成中文就是“內(nèi)存用完了”,來(lái)源于java.lang.OutOfMemoryError,當(dāng)JVM因?yàn)闆]有足夠的內(nèi)存來(lái)為對(duì)象分配空間并且垃圾回收器也已經(jīng)沒有空間可回收時(shí),就會(huì)拋出這個(gè)error,需要的朋友可以參考下
    2023-08-08
  • 詳細(xì)分析Java中String、StringBuffer、StringBuilder類的性能

    詳細(xì)分析Java中String、StringBuffer、StringBuilder類的性能

    在Java中,String類和StringBuffer類以及StringBuilder類都能用于創(chuàng)建字符串對(duì)象,而在分別操作這些對(duì)象時(shí)我們會(huì)發(fā)現(xiàn)JVM執(zhí)行它們的性能并不相同,下面我們就來(lái)詳細(xì)分析Java中String、StringBuffer、StringBuilder類的性能
    2016-05-05
  • 解決JSONObject.toJSONString()輸出null的問題

    解決JSONObject.toJSONString()輸出null的問題

    這篇文章主要介紹了解決JSONObject.toJSONString()輸出null的問題,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2022-02-02

最新評(píng)論

瓦房店市| 长武县| 宣威市| 黑山县| 博兴县| 邮箱| 永城市| 东城区| 台中市| 皮山县| 漳州市| 永新县| 远安县| 滦南县| 和田市| 江门市| 休宁县| 阳原县| 屯昌县| 濮阳县| 江口县| 鄢陵县| 新巴尔虎左旗| 普兰县| 囊谦县| 南涧| 综艺| 武城县| 南汇区| 通州市| 措美县| 合肥市| 达州市| 古交市| 北流市| 吉首市| 阜阳市| 亚东县| 孝义市| 云霄县| 蒲江县|