RocketMQ 生產(chǎn)者 rebalence 原理

概述

生產(chǎn)者 producer 在發(fā)送消息的時(shí)候旗笔,每個(gè)消息發(fā)送到 broker 只存儲(chǔ)在某一個(gè) quene 上。那么 producer 是怎么選擇 queue 呢傲武?

下面主要通過(guò)以下5種方式進(jìn)行分析逸寓。
1、自定義 MessageQueueSelector 實(shí)現(xiàn)
2瘦黑、SelectMessageQueueByHash hash 選擇 queue。
3、 SelectMessageQueueByRandom 隨機(jī)選擇 queue幸斥。
4匹摇、 SelectMessageQueueByMachineRoom 機(jī)房選擇queue。
5甲葬、默認(rèn)發(fā)送隊(duì)列選擇實(shí)現(xiàn)

1廊勃、自定義 MessageQueueSelector 實(shí)現(xiàn)

下面這個(gè)示例是 rocketmq 官網(wǎng)上的一個(gè)示例。

public class Producer {
    public static void main(String[] args) throws UnsupportedEncodingException {
        try {
            MQProducer producer = new DefaultMQProducer("please_rename_unique_group_name");
            producer.start();

            String[] tags = new String[] {"TagA", "TagB", "TagC", "TagD", "TagE"};
            for (int i = 0; i < 100; i++) {
                int orderId = i % 10;
                Message msg =
                    new Message("TopicTest", tags[i % tags.length], "KEY" + i,
                        ("Hello RocketMQ " + i).getBytes(RemotingHelper.DEFAULT_CHARSET));
                SendResult sendResult = producer.send(msg, new MessageQueueSelector() {
                    @Override
                    public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
                        Integer id = (Integer) arg;
                        int index = id % mqs.size();
                        return mqs.get(index);
                    }
                }, orderId);

                System.out.printf("%s%n", sendResult);
            }

            producer.shutdown();
        } catch (MQClientException | RemotingException | MQBrokerException | InterruptedException e) {
            e.printStackTrace();
        }
    }
}

從示例中可以看到 producer.send(msg, new MessageQueueSelector(){}, orderId)
在發(fā)送的時(shí)候 自定義了一個(gè) MessageQueueSelector演顾。

MessageQueueSelector 的 selelct(List<MessageQueue> mqs, Message msg, Object arg) 方法中有三個(gè)參數(shù)供搀。

  • List<MessageQueue> mqs :topic 中的所有 queue 的集合。
  • Message msg:發(fā)送的消息
  • Object arg:上面示例中 send 方法的第三個(gè)參數(shù)钠至。

通過(guò)實(shí)現(xiàn) select 方法葛虐,通過(guò) arg 參數(shù)進(jìn)行取模 mqs.size() 進(jìn)行選擇隊(duì)列。

RocketMQ 已實(shí)現(xiàn)的 MessageQueueSelector

rocketmq 源碼中已經(jīng)提供了幾種 MessageQueueSelector 的實(shí)現(xiàn)棉钧。如下圖:


  • SelectMessageQueueByHash:通過(guò) hash 進(jìn)行選擇 queue屿脐。
  • SelectMessageQueueByRandom:隨機(jī)選擇 queue。
  • SelectMessageQueueByMachineRoom:機(jī)房選擇queue宪卿。

2的诵、SelectMessageQueueByHash

public class SelectMessageQueueByHash implements MessageQueueSelector {

    @Override
    public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
        int value = arg.hashCode();
        if (value < 0) {
            value = Math.abs(value);
        }

        value = value % mqs.size();
        return mqs.get(value);
    }
}

通過(guò) arg 的 hash,通過(guò) mqs.size() 進(jìn)行取模佑钾,來(lái)選擇要存儲(chǔ)的隊(duì)列西疤。

3、SelectMessageQueueByRandom

public class SelectMessageQueueByRandom implements MessageQueueSelector {
    private Random random = new Random(System.currentTimeMillis());

    @Override
    public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
        int value = random.nextInt(mqs.size());
        return mqs.get(value);
    }
}

隨機(jī)產(chǎn)生一個(gè)小于等于 mqs.size() 的隨機(jī)正整數(shù)休溶,來(lái)選擇要存儲(chǔ)的隊(duì)列代赁。

4、SelectMessageQueueByMachineRoom

public class SelectMessageQueueByMachineRoom implements MessageQueueSelector {
    private Set<String> consumeridcs;

    @Override
    public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
        return null;
    }

    public Set<String> getConsumeridcs() {
        return consumeridcs;
    }

    public void setConsumeridcs(Set<String> consumeridcs) {
        this.consumeridcs = consumeridcs;
    }
}

這個(gè)未實(shí)現(xiàn)兽掰,還是要通過(guò)自己的場(chǎng)景進(jìn)行實(shí)現(xiàn)芭碍。

5、默認(rèn)是輪詢(xún)進(jìn)行發(fā)送消息

如果直接調(diào)用 SendResult send(final Message msg) 方法孽尽,RocketMQ 是如何選擇隊(duì)列的呢窖壕?

public MessageQueue selectOneMessageQueue(final TopicPublishInfo tpInfo, final String lastBrokerName) {
        if (this.sendLatencyFaultEnable) {
            try {
                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);
                    if (latencyFaultTolerance.isAvailable(mq.getBrokerName())) {
                        if (null == lastBrokerName || mq.getBrokerName().equals(lastBrokerName))
                            return mq;
                    }
                }

                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 {
                    latencyFaultTolerance.remove(notBestBroker);
                }
            } catch (Exception e) {
                log.error("Error occurred when selecting message queue", e);
            }

            return tpInfo.selectOneMessageQueue();
        }

        return tpInfo.selectOneMessageQueue(lastBrokerName);
    }

1、int index = tpInfo.getSendWhichQueue().getAndIncrement();獲取 一個(gè)自增的index杉女。
2瞻讽、然后 int pos = Math.abs(index++) % tpInfo.getMessageQueueList().size(); 進(jìn)行選擇一個(gè) queue。

通過(guò)上面的代碼可以看出熏挎,默認(rèn)是通過(guò)輪詢(xún)的方式進(jìn)行選擇發(fā)送隊(duì)列的速勇。

ThreadLocalIndex 實(shí)現(xiàn)

public class ThreadLocalIndex {
    private final ThreadLocal<Integer> threadLocalIndex = new ThreadLocal<Integer>();
    private final Random random = new Random();

    public int getAndIncrement() {
        Integer index = this.threadLocalIndex.get();
        if (null == index) {
            index = Math.abs(random.nextInt());
            if (index < 0)
                index = 0;
            this.threadLocalIndex.set(index);
        }

        index = Math.abs(index + 1);
        if (index < 0)
            index = 0;

        this.threadLocalIndex.set(index);
        return index;
    }

    @Override
    public String toString() {
        return "ThreadLocalIndex{" +
            "threadLocalIndex=" + threadLocalIndex.get() +
            '}';
    }
}

從 getAndIncrement() 方法中,可以看出婆瓜。
為每個(gè)線程分配一個(gè)隨機(jī)數(shù)快集,然后每次調(diào)用都自增 1。

最后編輯于
?著作權(quán)歸作者所有,轉(zhuǎn)載或內(nèi)容合作請(qǐng)聯(lián)系作者
  • 序言:七十年代末个初,一起剝皮案震驚了整個(gè)濱河市,隨后出現(xiàn)的幾起案子猴蹂,更是在濱河造成了極大的恐慌院溺,老刑警劉巖,帶你破解...
    沈念sama閱讀 212,185評(píng)論 6 493
  • 序言:濱河連續(xù)發(fā)生了三起死亡事件磅轻,死亡現(xiàn)場(chǎng)離奇詭異珍逸,居然都是意外死亡,警方通過(guò)查閱死者的電腦和手機(jī)聋溜,發(fā)現(xiàn)死者居然都...
    沈念sama閱讀 90,445評(píng)論 3 385
  • 文/潘曉璐 我一進(jìn)店門(mén)谆膳,熙熙樓的掌柜王于貴愁眉苦臉地迎上來(lái),“玉大人撮躁,你說(shuō)我怎么就攤上這事漱病。” “怎么了把曼?”我有些...
    開(kāi)封第一講書(shū)人閱讀 157,684評(píng)論 0 348
  • 文/不壞的土叔 我叫張陵杨帽,是天一觀的道長(zhǎng)。 經(jīng)常有香客問(wèn)我嗤军,道長(zhǎng)注盈,這世上最難降的妖魔是什么? 我笑而不...
    開(kāi)封第一講書(shū)人閱讀 56,564評(píng)論 1 284
  • 正文 為了忘掉前任叙赚,我火速辦了婚禮老客,結(jié)果婚禮上,老公的妹妹穿的比我還像新娘纠俭。我一直安慰自己沿量,他們只是感情好,可當(dāng)我...
    茶點(diǎn)故事閱讀 65,681評(píng)論 6 386
  • 文/花漫 我一把揭開(kāi)白布冤荆。 她就那樣靜靜地躺著朴则,像睡著了一般。 火紅的嫁衣襯著肌膚如雪钓简。 梳的紋絲不亂的頭發(fā)上乌妒,一...
    開(kāi)封第一講書(shū)人閱讀 49,874評(píng)論 1 290
  • 那天,我揣著相機(jī)與錄音外邓,去河邊找鬼撤蚊。 笑死,一個(gè)胖子當(dāng)著我的面吹牛损话,可吹牛的內(nèi)容都是我干的侦啸。 我是一名探鬼主播槽唾,決...
    沈念sama閱讀 39,025評(píng)論 3 408
  • 文/蒼蘭香墨 我猛地睜開(kāi)眼,長(zhǎng)吁一口氣:“原來(lái)是場(chǎng)噩夢(mèng)啊……” “哼光涂!你這毒婦竟也來(lái)了庞萍?” 一聲冷哼從身側(cè)響起,我...
    開(kāi)封第一講書(shū)人閱讀 37,761評(píng)論 0 268
  • 序言:老撾萬(wàn)榮一對(duì)情侶失蹤忘闻,失蹤者是張志新(化名)和其女友劉穎钝计,沒(méi)想到半個(gè)月后,有當(dāng)?shù)厝嗽跇?shù)林里發(fā)現(xiàn)了一具尸體齐佳,經(jīng)...
    沈念sama閱讀 44,217評(píng)論 1 303
  • 正文 獨(dú)居荒郊野嶺守林人離奇死亡私恬,尸身上長(zhǎng)有42處帶血的膿包…… 初始之章·張勛 以下內(nèi)容為張勛視角 年9月15日...
    茶點(diǎn)故事閱讀 36,545評(píng)論 2 327
  • 正文 我和宋清朗相戀三年,在試婚紗的時(shí)候發(fā)現(xiàn)自己被綠了炼吴。 大學(xué)時(shí)的朋友給我發(fā)了我未婚夫和他白月光在一起吃飯的照片本鸣。...
    茶點(diǎn)故事閱讀 38,694評(píng)論 1 341
  • 序言:一個(gè)原本活蹦亂跳的男人離奇死亡,死狀恐怖硅蹦,靈堂內(nèi)的尸體忽然破棺而出永高,到底是詐尸還是另有隱情,我是刑警寧澤提针,帶...
    沈念sama閱讀 34,351評(píng)論 4 332
  • 正文 年R本政府宣布命爬,位于F島的核電站,受9級(jí)特大地震影響辐脖,放射性物質(zhì)發(fā)生泄漏饲宛。R本人自食惡果不足惜,卻給世界環(huán)境...
    茶點(diǎn)故事閱讀 39,988評(píng)論 3 315
  • 文/蒙蒙 一嗜价、第九天 我趴在偏房一處隱蔽的房頂上張望艇抠。 院中可真熱鬧,春花似錦久锥、人聲如沸家淤。這莊子的主人今日做“春日...
    開(kāi)封第一講書(shū)人閱讀 30,778評(píng)論 0 21
  • 文/蒼蘭香墨 我抬頭看了看天上的太陽(yáng)絮重。三九已至,卻和暖如春歹苦,著一層夾襖步出監(jiān)牢的瞬間青伤,已是汗流浹背。 一陣腳步聲響...
    開(kāi)封第一講書(shū)人閱讀 32,007評(píng)論 1 266
  • 我被黑心中介騙來(lái)泰國(guó)打工殴瘦, 沒(méi)想到剛下飛機(jī)就差點(diǎn)兒被人妖公主榨干…… 1. 我叫王不留狠角,地道東北人。 一個(gè)月前我還...
    沈念sama閱讀 46,427評(píng)論 2 360
  • 正文 我出身青樓蚪腋,卻偏偏與公主長(zhǎng)得像丰歌,于是被迫代替她去往敵國(guó)和親姨蟋。 傳聞我的和親對(duì)象是個(gè)殘疾皇子,可洞房花燭夜當(dāng)晚...
    茶點(diǎn)故事閱讀 43,580評(píng)論 2 349

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

  • 分布式開(kāi)放消息系統(tǒng)(RocketMQ)的原理與實(shí)踐 來(lái)源:http://www.reibang.com/p/453...
    meng_philip123閱讀 12,906評(píng)論 6 104
  • Spring Cloud為開(kāi)發(fā)人員提供了快速構(gòu)建分布式系統(tǒng)中一些常見(jiàn)模式的工具(例如配置管理立帖,服務(wù)發(fā)現(xiàn)芬探,斷路器,智...
    卡卡羅2017閱讀 134,633評(píng)論 18 139
  • RocketMQ是一款分布式厘惦、隊(duì)列模型的消息中間件,具有以下特點(diǎn): 能夠保證嚴(yán)格的消息順序 提供豐富的消息拉取模式...
    AI喬治閱讀 2,064評(píng)論 2 5
  • metaq是阿里團(tuán)隊(duì)的消息中間件哩簿,之前也有用過(guò)和了解過(guò)kafka宵蕉,據(jù)說(shuō)metaq是基于kafka的源碼改過(guò)來(lái)的,他...
    菜鳥(niǎo)小玄閱讀 32,865評(píng)論 0 14
  • 新試手
    伊苼藝飾閱讀 185評(píng)論 0 0