rocketmq發(fā)送消息之選擇消息隊(duì)列

DefaultMQProducerImpl#selectOneMessageQueue選擇消息隊(duì)列

boolean callTimeout = false;
            MessageQueue mq = null;
            Exception exception = null;
            SendResult sendResult = null;
            #如果是同步發(fā)送消息早敬,有3次發(fā)送機(jī)會(huì)
            int timesTotal = communicationMode == CommunicationMode.SYNC ? 1 + this.defaultMQProducer.getRetryTimesWhenSendFailed() : 1;
            int times = 0;
            String[] brokersSent = new String[timesTotal];
            for (; times < timesTotal; times++) {
                String lastBrokerName = null == mq ? null : mq.getBrokerName();
                #選擇發(fā)送消息隊(duì)列入口
                MessageQueue mqSelected = this.selectOneMessageQueue(topicPublishInfo, lastBrokerName);
                if (mqSelected != null) {
                    mq = mqSelected;
                    brokersSent[times] = mq.getBrokerName();
                    try {
                        beginTimestampPrev = System.currentTimeMillis();
                        if (times > 0) {
                            //Reset topic with namespace during resend.
                            msg.setTopic(this.defaultMQProducer.withNamespace(msg.getTopic()));
                        }
                        long costTime = beginTimestampPrev - beginTimestampFirst;
                        if (timeout < costTime) {
                            callTimeout = true;
                            break;
                        }

                        sendResult = this.sendKernelImpl(msg, mq, communicationMode, sendCallback, topicPublishInfo, timeout - costTime);
                        endTimestamp = System.currentTimeMillis();
                        #將發(fā)送延遲信息更新到延遲機(jī)制的緩存中
                        this.updateFaultItem(mq.getBrokerName(), endTimestamp - beginTimestampPrev, false);
                        switch (communicationMode) {
                            case ASYNC:
                                return null;
                            case ONEWAY:

lastBrokerName 為上一次發(fā)送消息的broker信息,第一次發(fā)送時(shí)為null

public MessageQueue selectOneMessageQueue(final TopicPublishInfo tpInfo, final String lastBrokerName) {
        #如果啟用了故障延遲發(fā)送機(jī)制時(shí)使用該邏輯
        if (this.sendLatencyFaultEnable) {
            try {
                #隨機(jī)選擇一個(gè)隊(duì)列,并且判斷是否可用檬寂,如果可用則直接返回該隊(duì)列
                int index = tpInfo.getSendWhichQueue().getAndIncrement();
                for (int i = 0; i < tpInfo.getMessageQueueList().size(); i++) {
                    int pos = Math.abs(index++) % tpInfo.getMessageQueueList().size();
                    if (pos < 0)
                        pos = 0;
                    MessageQueue mq = tpInfo.getMessageQueueList().get(pos);
                    #當(dāng)前時(shí)間和延遲時(shí)間比較是否可用
                    if (latencyFaultTolerance.isAvailable(mq.getBrokerName())) {
                        if (null == lastBrokerName || mq.getBrokerName().equals(lastBrokerName))
                            return mq;
                    }
                }
                #如果沒(méi)找到可用的隊(duì)列把将,這隨機(jī)至少獲取一個(gè)隊(duì)列
                final String notBestBroker = latencyFaultTolerance.pickOneAtLeast();
                int writeQueueNums = tpInfo.getQueueIdByBroker(notBestBroker);
                if (writeQueueNums > 0) {
                    final MessageQueue mq = tpInfo.selectOneMessageQueue();
                    if (notBestBroker != null) {
                        mq.setBrokerName(notBestBroker);
                        mq.setQueueId(tpInfo.getSendWhichQueue().getAndIncrement() % writeQueueNums);
                    }
                    return mq;
                } else {
                    #此處不太明白移除機(jī)制
                    latencyFaultTolerance.remove(notBestBroker);
                }
            } catch (Exception e) {
                log.error("Error occurred when selecting message queue", e);
            }
            #輪詢選擇
            return tpInfo.selectOneMessageQueue();
        }
        #沒(méi)有使用故障延遲發(fā)送機(jī)制時(shí)瘤睹,使用此策略
        return tpInfo.selectOneMessageQueue(lastBrokerName);
    }
public MessageQueue selectOneMessageQueue(final String lastBrokerName) {
        if (lastBrokerName == null) {
            #如果是第一次發(fā)送隅要,則按照index+1的次序選擇
            return selectOneMessageQueue();
        } else {
            #非第一次選擇時(shí)弯菊,依次向后選擇当宴,并且不能選擇上次已經(jīng)發(fā)送失敗的broker
            int index = this.sendWhichQueue.getAndIncrement();
            for (int i = 0; i < this.messageQueueList.size(); i++) {
                int pos = Math.abs(index++) % this.messageQueueList.size();
                if (pos < 0)
                    pos = 0;
                MessageQueue mq = this.messageQueueList.get(pos);
                if (!mq.getBrokerName().equals(lastBrokerName)) {
                    return mq;
                }
            }
            return selectOneMessageQueue();
        }
    }
#更新延遲信息到緩存
public void updateFaultItem(final String brokerName, final long currentLatency, boolean isolation) {
        if (this.sendLatencyFaultEnable) {
            #isolation 如果為true畜吊,這采用30秒故障隔離時(shí)長(zhǎng),如果為false這采用實(shí)際延遲時(shí)間計(jì)算故障隔離時(shí)長(zhǎng)
            long duration = computeNotAvailableDuration(isolation ? 30000 : currentLatency);
            this.latencyFaultTolerance.updateFaultItem(brokerName, currentLatency, duration);
        }
    }

  #根據(jù)延遲時(shí)間計(jì)算具體故障隔離時(shí)長(zhǎng)
    private long computeNotAvailableDuration(final long currentLatency) {
        for (int i = latencyMax.length - 1; i >= 0; i--) {
            if (currentLatency >= latencyMax[i])
                return this.notAvailableDuration[i];
        }

        return 0;
    }
#更新延遲時(shí)間及故障隔離時(shí)長(zhǎng)户矢,故障隔離時(shí)長(zhǎng)會(huì)在isAvailable方法中使用
public void updateFaultItem(final String name, final long currentLatency, final long notAvailableDuration) {
        FaultItem old = this.faultItemTable.get(name);
        if (null == old) {
            final FaultItem faultItem = new FaultItem(name);
            faultItem.setCurrentLatency(currentLatency);
            faultItem.setStartTimestamp(System.currentTimeMillis() + notAvailableDuration);

            old = this.faultItemTable.putIfAbsent(name, faultItem);
            if (old != null) {
                old.setCurrentLatency(currentLatency);
                old.setStartTimestamp(System.currentTimeMillis() + notAvailableDuration);
            }
        } else {
            old.setCurrentLatency(currentLatency);
            old.setStartTimestamp(System.currentTimeMillis() + notAvailableDuration);
        }
    }


public boolean isAvailable(final String name) {
        final FaultItem faultItem = this.faultItemTable.get(name);
        if (faultItem != null) {
            return faultItem.isAvailable();
        }
        return true;
    }
最后編輯于
?著作權(quán)歸作者所有,轉(zhuǎn)載或內(nèi)容合作請(qǐng)聯(lián)系作者
  • 序言:七十年代末玲献,一起剝皮案震驚了整個(gè)濱河市,隨后出現(xiàn)的幾起案子梯浪,更是在濱河造成了極大的恐慌捌年,老刑警劉巖,帶你破解...
    沈念sama閱讀 210,978評(píng)論 6 490
  • 序言:濱河連續(xù)發(fā)生了三起死亡事件挂洛,死亡現(xiàn)場(chǎng)離奇詭異礼预,居然都是意外死亡,警方通過(guò)查閱死者的電腦和手機(jī)虏劲,發(fā)現(xiàn)死者居然都...
    沈念sama閱讀 89,954評(píng)論 2 384
  • 文/潘曉璐 我一進(jìn)店門托酸,熙熙樓的掌柜王于貴愁眉苦臉地迎上來(lái),“玉大人柒巫,你說(shuō)我怎么就攤上這事励堡。” “怎么了堡掏?”我有些...
    開(kāi)封第一講書人閱讀 156,623評(píng)論 0 345
  • 文/不壞的土叔 我叫張陵应结,是天一觀的道長(zhǎng)。 經(jīng)常有香客問(wèn)我泉唁,道長(zhǎng)摊趾,這世上最難降的妖魔是什么? 我笑而不...
    開(kāi)封第一講書人閱讀 56,324評(píng)論 1 282
  • 正文 為了忘掉前任游两,我火速辦了婚禮砾层,結(jié)果婚禮上,老公的妹妹穿的比我還像新娘贱案。我一直安慰自己肛炮,他們只是感情好止吐,可當(dāng)我...
    茶點(diǎn)故事閱讀 65,390評(píng)論 5 384
  • 文/花漫 我一把揭開(kāi)白布。 她就那樣靜靜地躺著侨糟,像睡著了一般碍扔。 火紅的嫁衣襯著肌膚如雪。 梳的紋絲不亂的頭發(fā)上秕重,一...
    開(kāi)封第一講書人閱讀 49,741評(píng)論 1 289
  • 那天不同,我揣著相機(jī)與錄音,去河邊找鬼溶耘。 笑死二拐,一個(gè)胖子當(dāng)著我的面吹牛,可吹牛的內(nèi)容都是我干的凳兵。 我是一名探鬼主播百新,決...
    沈念sama閱讀 38,892評(píng)論 3 405
  • 文/蒼蘭香墨 我猛地睜開(kāi)眼,長(zhǎng)吁一口氣:“原來(lái)是場(chǎng)噩夢(mèng)啊……” “哼庐扫!你這毒婦竟也來(lái)了饭望?” 一聲冷哼從身側(cè)響起,我...
    開(kāi)封第一講書人閱讀 37,655評(píng)論 0 266
  • 序言:老撾萬(wàn)榮一對(duì)情侶失蹤形庭,失蹤者是張志新(化名)和其女友劉穎铅辞,沒(méi)想到半個(gè)月后,有當(dāng)?shù)厝嗽跇?shù)林里發(fā)現(xiàn)了一具尸體萨醒,經(jīng)...
    沈念sama閱讀 44,104評(píng)論 1 303
  • 正文 獨(dú)居荒郊野嶺守林人離奇死亡巷挥,尸身上長(zhǎng)有42處帶血的膿包…… 初始之章·張勛 以下內(nèi)容為張勛視角 年9月15日...
    茶點(diǎn)故事閱讀 36,451評(píng)論 2 325
  • 正文 我和宋清朗相戀三年,在試婚紗的時(shí)候發(fā)現(xiàn)自己被綠了验靡。 大學(xué)時(shí)的朋友給我發(fā)了我未婚夫和他白月光在一起吃飯的照片倍宾。...
    茶點(diǎn)故事閱讀 38,569評(píng)論 1 340
  • 序言:一個(gè)原本活蹦亂跳的男人離奇死亡,死狀恐怖胜嗓,靈堂內(nèi)的尸體忽然破棺而出高职,到底是詐尸還是另有隱情,我是刑警寧澤辞州,帶...
    沈念sama閱讀 34,254評(píng)論 4 328
  • 正文 年R本政府宣布怔锌,位于F島的核電站,受9級(jí)特大地震影響变过,放射性物質(zhì)發(fā)生泄漏埃元。R本人自食惡果不足惜,卻給世界環(huán)境...
    茶點(diǎn)故事閱讀 39,834評(píng)論 3 312
  • 文/蒙蒙 一媚狰、第九天 我趴在偏房一處隱蔽的房頂上張望岛杀。 院中可真熱鬧,春花似錦崭孤、人聲如沸类嗤。這莊子的主人今日做“春日...
    開(kāi)封第一講書人閱讀 30,725評(píng)論 0 21
  • 文/蒼蘭香墨 我抬頭看了看天上的太陽(yáng)遗锣。三九已至货裹,卻和暖如春,著一層夾襖步出監(jiān)牢的瞬間精偿,已是汗流浹背弧圆。 一陣腳步聲響...
    開(kāi)封第一講書人閱讀 31,950評(píng)論 1 264
  • 我被黑心中介騙來(lái)泰國(guó)打工, 沒(méi)想到剛下飛機(jī)就差點(diǎn)兒被人妖公主榨干…… 1. 我叫王不留笔咽,地道東北人搔预。 一個(gè)月前我還...
    沈念sama閱讀 46,260評(píng)論 2 360
  • 正文 我出身青樓,卻偏偏與公主長(zhǎng)得像拓轻,于是被迫代替她去往敵國(guó)和親。 傳聞我的和親對(duì)象是個(gè)殘疾皇子经伙,可洞房花燭夜當(dāng)晚...
    茶點(diǎn)故事閱讀 43,446評(píng)論 2 348

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