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

RocketMQ普通消息實(shí)戰(zhàn)演練詳解

 更新時(shí)間:2022年08月22日 14:59:31   作者:奔跑的毛球  
這篇文章主要為大家介紹了RocketMQ普通消息實(shí)戰(zhàn)演練詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪

引言

之前研究了RocketMQ的源碼,在這里將各種消息發(fā)送與消費(fèi)的demo進(jìn)行舉例,方便以后使用的時(shí)候CV。

相關(guān)的配置,安裝和啟動(dòng)在這篇文章有相關(guān)講解  http://www.fzitv.net/article/260237.htm

普通消息同步發(fā)送

同步消息是指發(fā)送出消息后,同步等待,直到接收到Broker發(fā)送成功的響應(yīng)才會(huì)繼續(xù)發(fā)送下一個(gè)消息。這個(gè)方式可以確保消息發(fā)送到Broker成功,一些重要的消息可以使用此方式,比如重要的通知。

public static void main(String[] args) throws Exception {
    //實(shí)例化消息生產(chǎn)者對(duì)象
    DefaultMQProducer producer = new DefaultMQProducer("group_luke");
    //設(shè)置NameSever地址
    producer.setNamesrvAddr("127.0.0.1:9876");
    //啟動(dòng)Producer實(shí)例
    producer.start();
    for (int i = 0; i < 10; i++) {
        Message msg = new Message("topic_luke", "tag", ("這是第"+i+"條消息。").getBytes(StandardCharsets.UTF_8));
        //同步發(fā)送方式
        SendResult send = producer.send(msg);
        //確認(rèn)返回
        System.out.println(send);
    }
    //關(guān)閉producer
    producer.shutdown();
}

普通消息異步發(fā)送

異步消息發(fā)送方在發(fā)送了一條消息后,不等接收方發(fā)回響應(yīng),接著進(jìn)行第二條消息發(fā)送。發(fā)送方通過(guò)回調(diào)接口的方式接收服務(wù)器響應(yīng),并對(duì)響應(yīng)結(jié)果進(jìn)行處理。

public static void main(String[] args) throws Exception {
    //實(shí)例化消息生產(chǎn)者對(duì)象
    DefaultMQProducer producer = new DefaultMQProducer("group_luke");
    //設(shè)置NameSever地址
    producer.setNamesrvAddr("127.0.0.1:9876");
    //啟動(dòng)Producer實(shí)例
    producer.start();
    for (int i = 0; i < 10; i++) {
        Message msg = new Message("topic_luke", "tag", ("這是第"+i+"條消息。").getBytes(StandardCharsets.UTF_8));
        //SendCallback會(huì)接收異步返回結(jié)果的回調(diào)
        producer.send(msg, new SendCallback() {
            @Override
            public void onSuccess(SendResult sendResult) {
                System.out.println(sendResult);
            }
            @Override
            public void onException(Throwable throwable) {
                throwable.printStackTrace();
            }
        });
    }
    //若是過(guò)早關(guān)閉producer,會(huì)拋出The producer service state not OK, SHUTDOWN_ALREADY的錯(cuò)
    Thread.sleep(10000);
    //關(guān)閉producer
    producer.shutdown();
}

普通消息單向發(fā)送

單項(xiàng)發(fā)送不關(guān)心發(fā)送的結(jié)果,只發(fā)送請(qǐng)求不等待應(yīng)答。發(fā)送消息耗時(shí)極短。

public static void main(String[] args) throws Exception {
    //實(shí)例化消息生產(chǎn)者對(duì)象
    DefaultMQProducer producer = new DefaultMQProducer("group_luke");
    //設(shè)置NameSever地址
    producer.setNamesrvAddr("127.0.0.1:9876");
    //啟動(dòng)Producer實(shí)例
    producer.start();
    for (int i = 0; i < 10; i++) {
        Message msg = new Message("topic_luke", "tag", ("這是第"+i+"條消息。").getBytes(StandardCharsets.UTF_8));
        //同步發(fā)送方式
        producer.sendOneway(msg);
    }
    //關(guān)閉producer
    producer.shutdown();
}

集群消費(fèi)模式

消費(fèi)者采用負(fù)載均衡的方式消費(fèi)消息,同一個(gè)Group下的多個(gè)Consumer共同消費(fèi)Queue里的Message,每個(gè)Consumer處理的消息不同。

一個(gè)Consumer Group中的各個(gè)Consumer實(shí)例分共同消費(fèi)消息,即一條消息只會(huì)投遞到一個(gè)Group下面的一個(gè)實(shí)例,并且只消費(fèi)一遍。

例如某個(gè)Topic有3個(gè)隊(duì)列,其中一個(gè)Consumer Group 有 3 個(gè)實(shí)例,那么每個(gè)實(shí)例只消費(fèi)其中的1個(gè)隊(duì)列。集群消費(fèi)模式是消費(fèi)者默認(rèn)的消費(fèi)方式。

public static void main(String[] args) throws Exception {
    //實(shí)例化消息消費(fèi)者
    DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("group_luke");
    //指定nameserver地址
    consumer.setNamesrvAddr("127.0.0.1:9876");
    //訂閱topic,"*"表示所有tag
    consumer.subscribe("topic_luke","*");
    consumer.setMessageModel(MessageModel.CLUSTERING);
    // 注冊(cè)回調(diào)實(shí)現(xiàn)類來(lái)處理從broker拉取回來(lái)的消息
    consumer.registerMessageListener(new MessageListenerConcurrently() {
        @SneakyThrows
        @Override
        public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, ConsumeConcurrentlyContext context) {
            for (MessageExt msg : msgs) {
                System.out.println(new String(msg.getBody()));
            }
            // 標(biāo)記該消息已經(jīng)被成功消費(fèi)
            return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
        }
    });
    // 啟動(dòng)消費(fèi)者實(shí)例
    consumer.start();
    System.out.printf("Consumer Started.%n");
}

廣播消費(fèi)模式

廣播消費(fèi)模式中把消息對(duì)一個(gè)Group下的各個(gè)Consumer實(shí)例都投遞一遍。也就是說(shuō)消息也會(huì)被 Group 中的每個(gè)Consumer都消費(fèi)一次。

實(shí)際上,是一個(gè)消費(fèi)組下的每個(gè)消費(fèi)者實(shí)例都獲取到了topic下面的每個(gè)Message Queue去拉取消費(fèi)。所以消息會(huì)投遞到每個(gè)消費(fèi)者實(shí)例。

public static void main(String[] args) throws Exception {
    //實(shí)例化消息消費(fèi)者
    DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("group_luke");
    //指定nameserver地址
    consumer.setNamesrvAddr("127.0.0.1:9876");
    //訂閱topic,"*"表示所有tag
    consumer.subscribe("topic_luke","*");
    consumer.setMessageModel(MessageModel.BROADCASTING);
    // 注冊(cè)回調(diào)實(shí)現(xiàn)類來(lái)處理從broker拉取回來(lái)的消息
    consumer.registerMessageListener(new MessageListenerConcurrently() {
        @SneakyThrows
        @Override
        public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, ConsumeConcurrentlyContext context) {
            for (MessageExt msg : msgs) {
                System.out.println(new String(msg.getBody()));
            }
            // 標(biāo)記該消息已經(jīng)被成功消費(fèi)
            return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
        }
    });
    // 啟動(dòng)消費(fèi)者實(shí)例
    consumer.start();
    System.out.printf("Consumer Started.%n");
}

以上就是RocketMQ普通消息實(shí)戰(zhàn)演練詳解的詳細(xì)內(nèi)容,更多關(guān)于RocketMQ普通消息的資料請(qǐng)關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

  • 詳解JVM的內(nèi)存對(duì)象介紹[創(chuàng)建和訪問(wèn)]

    詳解JVM的內(nèi)存對(duì)象介紹[創(chuàng)建和訪問(wèn)]

    這篇文章主要介紹了JVM的內(nèi)存對(duì)象介紹[創(chuàng)建和訪問(wèn)],文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧
    2019-03-03
  • 關(guān)于springBoot yml文件的list讀取問(wèn)題總結(jié)(親測(cè))

    關(guān)于springBoot yml文件的list讀取問(wèn)題總結(jié)(親測(cè))

    這篇文章主要介紹了關(guān)于springBoot yml文件的list讀取問(wèn)題總結(jié),具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2021-12-12
  • 史上最全最強(qiáng)SpringMVC詳細(xì)示例實(shí)戰(zhàn)教程(圖文)

    史上最全最強(qiáng)SpringMVC詳細(xì)示例實(shí)戰(zhàn)教程(圖文)

    這篇文章主要介紹了史上最全最強(qiáng)SpringMVC詳細(xì)示例實(shí)戰(zhàn)教程(圖文),需要的朋友可以參考下
    2016-12-12
  • SpringBoot搭建全局異常攔截

    SpringBoot搭建全局異常攔截

    這篇文章主要介紹了SpringBoot搭建全局異常攔截,本文通過(guò)詳細(xì)的介紹與代碼的展示,詳細(xì)的說(shuō)明了如何搭建該項(xiàng)目,包括創(chuàng)建,啟動(dòng)和測(cè)試步驟,需要的朋友可以參考下
    2021-06-06
  • Java代理模式的深入了解

    Java代理模式的深入了解

    這篇文章主要為大家介紹了Java代理模式,具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下,希望能夠給你帶來(lái)幫助
    2022-01-01
  • Java判斷本機(jī)IP地址類型的方法

    Java判斷本機(jī)IP地址類型的方法

    Java判斷本機(jī)IP地址類型的方法,需要的朋友可以參考一下
    2013-03-03
  • maven中自定義MavenArchetype的實(shí)現(xiàn)

    maven中自定義MavenArchetype的實(shí)現(xiàn)

    Maven自身提供了許多Archetype來(lái)方便用戶創(chuàng)建Project,為了避免在創(chuàng)建project時(shí)重復(fù)的拷貝和修改,我們通過(guò)自定義Archetype來(lái)規(guī)范顯得還蠻有必要,下面就來(lái)介紹一下,感興趣的可以了解一下
    2025-01-01
  • 淺談Java枚舉的作用與好處

    淺談Java枚舉的作用與好處

    下面小編就為大家?guī)?lái)一篇淺談Java枚舉的作用與好處。小編覺(jué)得挺不錯(cuò)的,現(xiàn)在就分享給大家,也給大家做個(gè)參考。一起跟隨小編過(guò)來(lái)看看吧
    2016-07-07
  • 解析Java內(nèi)存分配和回收策略以及MinorGC、MajorGC、FullGC

    解析Java內(nèi)存分配和回收策略以及MinorGC、MajorGC、FullGC

    本節(jié)將會(huì)介紹一下:對(duì)象的內(nèi)存分配與回收策略;對(duì)象何時(shí)進(jìn)入新生代、老年代;MinorGC、MajorGC、FullGC的定義區(qū)別和觸發(fā)條件;還有通過(guò)圖示展示了GC的過(guò)程。
    2021-09-09
  • Spring 父類變量注入失敗的解決

    Spring 父類變量注入失敗的解決

    這篇文章主要介紹了Spring 父類變量注入失敗的解決方案,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2021-09-09

最新評(píng)論

大理市| 聂荣县| 尉犁县| 醴陵市| 乌拉特前旗| 甘肃省| 久治县| 大足县| 特克斯县| 卓尼县| 兴业县| 长治县| 虎林市| 阜新| 永丰县| 灌阳县| 浮山县| 安泽县| 克山县| 黎城县| 襄垣县| 陕西省| 峨眉山市| 皋兰县| 毕节市| 泾源县| 永吉县| 长治市| 壤塘县| 曲松县| 武平县| 大田县| 米易县| 灵宝市| 奉节县| 新龙县| 湘西| 广安市| 鸡西市| 忻城县| 古田县|