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

RocketMQ消息發(fā)送流程源碼剖析

 更新時間:2022年08月01日 11:31:50   作者:奔跑的毛球  
這篇文章主要為大家介紹了RocketMQ消息發(fā)送流程源碼剖析,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步,早日升職加薪

正文

就是說,我們打了個比方,把RocketMQ比作碼頭上的一個小房子,來送孩子登船的家長比作生產(chǎn)者,拉走孩子們的船夫比作消費者,所以,RocketMQ的故事就這么展開了。

這節(jié)我們研究研究,消息的發(fā)送流程。也就是說,消息孩子從進門到坐到message queue座位上都經(jīng)歷了啥。

父母把消息孩子送到碼頭之后,門口的門童defaultMQProducerImpl.send()接過孩子,進入到MQ房子內(nèi)部,然后引導(dǎo)孩子進入Broker候船大廳內(nèi)的message queue座位上就坐。這就是消息發(fā)送的流程了。

而且孩子在剛被門童接到之后,就被規(guī)定了能在候船大廳待多久,默認是3秒。也就是說,要是再小房子內(nèi)等了三秒沒走,就離開吧,你怕是沒想明白自己來干啥的。這就是消息的超時時間。

讀源碼

1 調(diào)用defaultMQProducerImpl.send()

public SendResult send(
    Message msg) throws MQClientException, RemotingException, MQBrokerException, InterruptedException {
    return send(msg, this.defaultMQProducer.getSendMsgTimeout());
}

2 設(shè)置過期時間

public SendResult send(Message msg,
    long timeout) throws MQClientException, RemotingException, MQBrokerException, InterruptedException {
    return this.sendDefaultImpl(msg, CommunicationMode.SYNC, null, timeout);
}

3 執(zhí)行defaultMQProducerImpl.sendDefaultImpl()方法

private SendResult sendDefaultImpl(
    Message msg,
    final CommunicationMode communicationMode,
    final SendCallback sendCallback,
    final long timeout
) throws MQClientException, RemotingException, MQBrokerException, InterruptedException {}

這里看看這幾個參數(shù),

  • communicationMode 是通信模式,同步異步還是單向
  • sendCallback 是針對異步模式的,異步模式需要設(shè)置發(fā)送完成后的回調(diào)。

sendDefaultImpl是發(fā)送消息的核心方法。

這里消息孩子進到第一個卡口,先要檢查送孩子來的家長是否還能聯(lián)系上,若是能聯(lián)系到,就繼續(xù)。要是聯(lián)系不到,這孩子豈不是被拋棄了,不敢接不敢接,送到孤兒院吧。

然后需要檢查消息孩子了,首先是檢查孩子還在不在,別扔個衣服跑了。
然后看看孩子指定的這個topic,不能說我想去內(nèi)個topic哈,必須是實實在在的名字。而且上頭也規(guī)定了,這個topic的名字也不能太長,也不能包含特殊字符。已有的一些領(lǐng)導(dǎo)定過的也不能用哈。
接下來就是檢查孩子的body了,之前說body就是孩子的技能,首先,技能為空,不行不行,啥都不會是不行的。再者太長也不行,你唱首歌兩年,這沒法玩。

檢查message不為null

檢查topic

  • topic不能為空
  • topic不能太長
  • 不能包含特殊字符

檢查話題的名字是否被系統(tǒng)已占用

檢查body

  • 檢查是否為空
  • 檢查長度是否過長,最大為4MB 這樣

下邊我們看看sendDefaultImpl這個方法。給他拆成一段一段的看。

1 兩個校驗

//校驗生產(chǎn)者服務(wù)是ok的,可以聯(lián)系到的
this.makeSureStateOK();
//校驗消息的參數(shù)
Validators.checkMessage(msg, this.defaultMQProducer);
  • 第一個檢查,檢查生產(chǎn)者服務(wù)是否是正常工作的,若是不正常工作,就拋出異常。
private void makeSureStateOK() throws MQClientException {
    if (this.serviceState != ServiceState.RUNNING) {
        throw new MQClientException("The producer service state not OK, "
            + this.serviceState
            + FAQUrl.suggestTodo(FAQUrl.CLIENT_SERVICE_NOT_OK),
            null);
    }
}
  • 第二個檢查,檢查消息本身是否為空,檢查topic,檢查消息的body
public static void checkMessage(Message msg, DefaultMQProducer defaultMQProducer) throws MQClientException {
    if (null == msg) {
        throw new MQClientException(ResponseCode.MESSAGE_ILLEGAL, "the message is null");
    }
    // 這里校驗Topic的時候,校驗了不能為空,長度和特殊字符
    Validators.checkTopic(msg.getTopic());
    //這里則校驗了一些不允許使用的topic名字
    Validators.isNotAllowedSendTopic(msg.getTopic());
    // body不為空
    if (null == msg.getBody()) {
        throw new MQClientException(ResponseCode.MESSAGE_ILLEGAL, "the message body is null");
    }
    // body長度不為0
    if (0 == msg.getBody().length) {
        throw new MQClientException(ResponseCode.MESSAGE_ILLEGAL, "the message body length is zero");
    }
    // body 長度不能過長
    if (msg.getBody().length > defaultMQProducer.getMaxMessageSize()) {
        throw new MQClientException(ResponseCode.MESSAGE_ILLEGAL,
            "the message body size over max value, MAX: " + defaultMQProducer.getMaxMessageSize());
    }
}

2 獲取topic路由信息

嗯,這里孩子終于通過了檢查,服務(wù)人員開始帶著他去找自己指定的topic區(qū)域,指定是自己指定,劃分還是工作人員劃分的。咱總得知道這個topic區(qū)域在哪吧。

先去緩存筆記里找,有沒有這個區(qū)域的信息,若是沒有這個topic,就新建一個,然后更新到緩存筆記里邊。若有topic但是不知道在哪,就找name server大腦去申請這個topic在哪的信息。

執(zhí)行tryToFindTopicPublishInfo方法去獲取Topic的路由信息,若是不存在就新建,若是有topic但是緩存中沒有路由信息,則通過name server獲取路由信息。

TopicPublishInfo topicPublishInfo = this.tryToFindTopicPublishInfo(msg.getTopic());
private TopicPublishInfo tryToFindTopicPublishInfo(final String topic) {
    //獲取topic信息
    TopicPublishInfo topicPublishInfo = this.topicPublishInfoTable.get(topic);
    //不存在
    if (null == topicPublishInfo || !topicPublishInfo.ok()) {
        //新建
        this.topicPublishInfoTable.putIfAbsent(topic, new TopicPublishInfo());
        //修改topic的路由信息并更新到本地
        this.mQClientFactory.updateTopicRouteInfoFromNameServer(topic);
        topicPublishInfo = this.topicPublishInfoTable.get(topic);
    }
    //包含路由信息就直接返回
    if (topicPublishInfo.isHaveTopicRouterInfo() || topicPublishInfo.ok()) {
        return topicPublishInfo;
    } else {
        //不包含路由信息則向name server申請,修改topic的路由信息并更新到本地
        this.mQClientFactory.updateTopicRouteInfoFromNameServer(topic, true, this.defaultMQProducer);
        topicPublishInfo = this.topicPublishInfoTable.get(topic);
        return topicPublishInfo;
    }
}

3 計算重試次數(shù)

這就是計算消息孩子可以嘗試去找地方坐幾次,沒坐上,欸,我又來了,沒坐上,欸,我又來了。

這行代碼就是計算重試次數(shù)的,根據(jù)communicationMode傳入的值,同步異步還是單向的來決定重試次數(shù)是幾次。 很明顯,若是同步的,就會嘗試三次。若是異步的或者單向的就只發(fā)送一次。

int timesTotal = communicationMode == CommunicationMode.SYNC ? 1 + this.defaultMQProducer.getRetryTimesWhenSendFailed() : 1;

4 執(zhí)行隊列選擇方法

我們之前說了,Broker類似于候船大廳,為了均分壓力,每次都要進與上次不同的候船大廳。

執(zhí)行selectOneMessageQueue方法通過Queue將消息發(fā)送到與上次不同的一個Broker。也可以通過 sendLatencyFaultEnable判斷是否啟用延遲容錯開關(guān)

MessageQueue mqSelected = this.selectOneMessageQueue(topicPublishInfo, lastBrokerName);

5 發(fā)送消息

這就是走過巷道坐到屬于自己的座位上了

然后就通過sendKernelImpl發(fā)送消息了,這是發(fā)送消息的核心方法。會準備通信層的入?yún)?,并將請求發(fā)送給通信層,內(nèi)部實現(xiàn)是基于Netty的。

sendResult = this.sendKernelImpl(msg, mq, communicationMode, sendCallback, topicPublishInfo, timeout - costTime);

以上就是RocketMQ消息發(fā)送流程源碼剖析的詳細內(nèi)容,更多關(guān)于RocketMQ消息發(fā)送流程的資料請關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

  • MyBatis-plus數(shù)據(jù)庫字段排序不準確的解決

    MyBatis-plus數(shù)據(jù)庫字段排序不準確的解決

    這篇文章主要介紹了MyBatis-plus數(shù)據(jù)庫字段排序不準確的解決方案,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2022-02-02
  • java中jvm逃逸問題分析

    java中jvm逃逸問題分析

    本篇文章給大家詳細總結(jié)了java中jvm逃逸問題的相關(guān)內(nèi)容,有興趣的朋友可以根據(jù)小編一起學(xué)習(xí)下。
    2018-02-02
  • Spring深入了解常用配置應(yīng)用

    Spring深入了解常用配置應(yīng)用

    這篇文章主要給大家介紹了關(guān)于Spring的常用配置,文中通過示例代碼介紹的非常詳細,對大家學(xué)習(xí)或者使用springboot具有一定的參考學(xué)習(xí)價值,需要的朋友可以參考下
    2022-07-07
  • SpringBoot 統(tǒng)一公共返回類的實現(xiàn)

    SpringBoot 統(tǒng)一公共返回類的實現(xiàn)

    本文主要介紹了SpringBoot 統(tǒng)一公共返回類的實現(xiàn),配置后臺的統(tǒng)一公共返回類,這樣做目的是為了統(tǒng)一返回信息,文中示例代碼介紹的很詳細,感興趣的可以了解一下
    2022-01-01
  • hadoop client與datanode的通信協(xié)議分析

    hadoop client與datanode的通信協(xié)議分析

    本文主要分析了hadoop客戶端read和write block的流程. 以及client和datanode通信的協(xié)議, 數(shù)據(jù)流格式等
    2012-11-11
  • SpringBoot 使用jwt進行身份驗證的方法示例

    SpringBoot 使用jwt進行身份驗證的方法示例

    這篇文章主要介紹了SpringBoot 使用jwt進行身份驗證的方法示例,小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2018-12-12
  • 模擬打印機排隊打印效果

    模擬打印機排隊打印效果

    本節(jié)主要介紹了模擬打印機排隊打印效果的具體實現(xiàn),感興趣的朋友可以參考下
    2014-07-07
  • IDEA之配置JDK、Git、Maven詳解

    IDEA之配置JDK、Git、Maven詳解

    文章總結(jié):本文介紹了如何在IDEA中配置JDK、Git和Maven,包括設(shè)置Java編譯器路徑、配置Git版本控制、修改Maven根目錄以加快jar包下載速度,并提供了一個解決方案以確保配置在新項目中生效
    2025-01-01
  • 基于Spring Security實現(xiàn)對密碼進行加密和校驗

    基于Spring Security實現(xiàn)對密碼進行加密和校驗

    我們在入門案例中,其實已經(jīng)是一個非常簡單的認證,但是用戶名是寫死的,密碼也需要從控制臺查看,很顯然實際中并不能這么做,下面的學(xué)習(xí)中,我們來實現(xiàn)基于內(nèi)存模型的認證以及用戶的自定義認證,密碼加密等內(nèi)容,需要的朋友可以參考下
    2024-07-07
  • 微服務(wù)Redis-Session共享登錄狀態(tài)的過程詳解

    微服務(wù)Redis-Session共享登錄狀態(tài)的過程詳解

    這篇文章主要介紹了微服務(wù)Redis-Session共享登錄狀態(tài),本文采取Spring security做登錄校驗,用redis做session共享,實現(xiàn)單服務(wù)登錄可靠性,微服務(wù)之間調(diào)用的可靠性與通用性,需要的朋友可以參考下
    2023-12-12

最新評論

南平市| 常山县| 北宁市| 洛川县| 洛南县| 建始县| 两当县| 容城县| 西吉县| 搜索| 封丘县| 芜湖县| 吉林市| 来安县| 泗洪县| 巴林左旗| 克山县| 邢台市| 绥阳县| 乐平市| 英吉沙县| 东乡族自治县| 沽源县| 兰考县| 牟定县| 六枝特区| 建阳市| 吉首市| 怀来县| 新龙县| 武川县| 准格尔旗| 洪洞县| 桐庐县| 鄂伦春自治旗| 邢台市| 柳林县| 龙胜| 彭阳县| 筠连县| 开封市|