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

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

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

引言

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

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

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

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

public static void main(String[] args) throws Exception {
    //實(shí)例化消息生產(chǎn)者對象
    DefaultMQProducer producer = new DefaultMQProducer("group_luke");
    //設(shè)置NameSever地址
    producer.setNamesrvAddr("127.0.0.1:9876");
    //啟動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ā)送方通過回調(diào)接口的方式接收服務(wù)器響應(yīng),并對響應(yīng)結(jié)果進(jìn)行處理。

public static void main(String[] args) throws Exception {
    //實(shí)例化消息生產(chǎn)者對象
    DefaultMQProducer producer = new DefaultMQProducer("group_luke");
    //設(shè)置NameSever地址
    producer.setNamesrvAddr("127.0.0.1:9876");
    //啟動Producer實(shí)例
    producer.start();
    for (int i = 0; i < 10; i++) {
        Message msg = new Message("topic_luke", "tag", ("這是第"+i+"條消息。").getBytes(StandardCharsets.UTF_8));
        //SendCallback會接收異步返回結(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ān)閉producer,會拋出The producer service state not OK, SHUTDOWN_ALREADY的錯
    Thread.sleep(10000);
    //關(guān)閉producer
    producer.shutdown();
}

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

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

public static void main(String[] args) throws Exception {
    //實(shí)例化消息生產(chǎn)者對象
    DefaultMQProducer producer = new DefaultMQProducer("group_luke");
    //設(shè)置NameSever地址
    producer.setNamesrvAddr("127.0.0.1:9876");
    //啟動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)消息,同一個Group下的多個Consumer共同消費(fèi)Queue里的Message,每個Consumer處理的消息不同。

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

例如某個Topic有3個隊(duì)列,其中一個Consumer Group 有 3 個實(shí)例,那么每個實(shí)例只消費(fèi)其中的1個隊(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);
    // 注冊回調(diào)實(shí)現(xiàn)類來處理從broker拉取回來的消息
    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;
        }
    });
    // 啟動消費(fèi)者實(shí)例
    consumer.start();
    System.out.printf("Consumer Started.%n");
}

廣播消費(fèi)模式

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

實(shí)際上,是一個消費(fèi)組下的每個消費(fèi)者實(shí)例都獲取到了topic下面的每個Message Queue去拉取消費(fèi)。所以消息會投遞到每個消費(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);
    // 注冊回調(diào)實(shí)現(xiàn)類來處理從broker拉取回來的消息
    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;
        }
    });
    // 啟動消費(fèi)者實(shí)例
    consumer.start();
    System.out.printf("Consumer Started.%n");
}

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

相關(guān)文章

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

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

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

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

    這篇文章主要介紹了關(guān)于springBoot yml文件的list讀取問題總結(jié),具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    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搭建全局異常攔截,本文通過詳細(xì)的介紹與代碼的展示,詳細(xì)的說明了如何搭建該項(xiàng)目,包括創(chuàng)建,啟動和測試步驟,需要的朋友可以參考下
    2021-06-06
  • Java代理模式的深入了解

    Java代理模式的深入了解

    這篇文章主要為大家介紹了Java代理模式,具有一定的參考價值,感興趣的小伙伴們可以參考一下,希望能夠給你帶來幫助
    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來方便用戶創(chuàng)建Project,為了避免在創(chuàng)建project時重復(fù)的拷貝和修改,我們通過自定義Archetype來規(guī)范顯得還蠻有必要,下面就來介紹一下,感興趣的可以了解一下
    2025-01-01
  • 淺談Java枚舉的作用與好處

    淺談Java枚舉的作用與好處

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

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

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

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

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

最新評論

新营市| 竹北市| 四平市| 宣化县| 墨脱县| 嘉荫县| 乐昌市| 托克逊县| 兴义市| 本溪市| 虹口区| 饶河县| 交口县| 高阳县| 环江| 平塘县| 马尔康县| 巴青县| 武隆县| 文成县| 龙江县| 云安县| 渑池县| 库伦旗| 浪卡子县| 屏南县| 平陆县| 东源县| 新野县| 花莲县| 砀山县| 桃江县| 休宁县| 修文县| 岫岩| 行唐县| 漯河市| 阳西县| 玉林市| 织金县| 女性|