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

圖文并茂講解RocketMQ消息類別

 更新時間:2022年12月27日 15:35:42   作者:一個雙子座的Java攻城獅  
這篇文章主要介紹了圖文并茂講解RocketMQ消息類別,RocketMQ對于消息提供了很多用法,包括:同步消息、異步消息、單向發(fā)送、順序消息、延時消息、批量消息、過濾消息、事務(wù)消息等

1、同步消息

即時性較強(qiáng),重要的消息,且必須有回執(zhí)的消息,例如短信,通知(轉(zhuǎn)賬成功)

生產(chǎn)者:

public class Producer {
    public static void main(String[] args) throws Exception{
        DefaultMQProducer producer=new DefaultMQProducer("group1");
        producer.setNamesrvAddr("192.168.23.127:9876");
        producer.start();
        for (int i = 1; i <= 5; i++) {
            Message msg = new Message("topic2",("同步消息:hello rocketmq "+i).getBytes("UTF-8"));
            //同步消息發(fā)送
            SendResult result = producer.send(msg);
            System.out.println("返回結(jié)果:"+result);
        }
        producer.shutdown();
    }
}

消費(fèi)者:

public class Consumer {
    public static void main(String[] args) throws Exception{
        DefaultMQPushConsumer consumer=new DefaultMQPushConsumer("group1");
        consumer.setNamesrvAddr("192.168.23.127:9876");
        consumer.subscribe("topic2","*");
        consumer.registerMessageListener(new MessageListenerConcurrently() {
            @Override
            public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> list, ConsumeConcurrentlyContext consumeConcurrentlyContext) {
                for (MessageExt msg : list) {
                    //System.out.println("收到消息:"+msg);
                    System.out.println("消息:"+new String(msg.getBody()));
                }
                return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;// 成功處理, mq 收到這個 標(biāo)記后相同的消息講不會再次發(fā)給消費(fèi)者
            }
        });
        consumer.start();// 開啟多線程 監(jiān)控消息,持續(xù)運(yùn)行
        System.out.println("接收消息服務(wù)已運(yùn)行");
    }
}

測試:

2、異步消息

即時性較弱,但需要有回執(zhí)的消息,例如訂單中的某些信息

生產(chǎn)者:

public class Producer {
    public static void main(String[] args) throws Exception{
        DefaultMQProducer producer=new DefaultMQProducer("group1");
        producer.setNamesrvAddr("192.168.23.127:9876");
        producer.start();
        for (int i = 1; i <= 5; i++) {
            //異步消息發(fā)送
            Message msg = new Message("topic2",("異步消息:hello rocketmq "+i).getBytes("UTF-8"));
            producer.send(msg, new SendCallback() {
                //表示成功返回結(jié)果
                @Override
                public void onSuccess(SendResult sendResult) {
                    System.out.println(sendResult);
                }
                //表示發(fā)送消息失敗
                @Override
                public void onException(Throwable throwable) {
                    System.out.println(throwable);
                }
            });
        }
        //添加一個休眠操作,確保異步消息返回后能夠輸出
        TimeUnit.SECONDS.sleep(10);
        producer.shutdown();
    }
}

消費(fèi)者:

public class Consumer {
    public static void main(String[] args) throws Exception{
        DefaultMQPushConsumer consumer=new DefaultMQPushConsumer("group1");
        consumer.setNamesrvAddr("192.168.23.127:9876");
        consumer.subscribe("topic2","*");
        consumer.registerMessageListener(new MessageListenerConcurrently() {
            @Override
            public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> list, ConsumeConcurrentlyContext consumeConcurrentlyContext) {
                for (MessageExt msg : list) {
                    //System.out.println("收到消息:"+msg);
                    System.out.println("消息:"+new String(msg.getBody()));
                }
                return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;// 成功處理, mq 收到這個 標(biāo)記后相同的消息講不會再次發(fā)給消費(fèi)者
            }
        });
        consumer.start();// 開啟多線程 監(jiān)控消息,持續(xù)運(yùn)行
        System.out.println("接收消息服務(wù)已運(yùn)行");
    }
}

測試:

3、單向消息

不需要有回執(zhí)的消息,例如日志類消息

生產(chǎn)者:

public class Producer {
    public static void main(String[] args) throws Exception{
        DefaultMQProducer producer=new DefaultMQProducer("group1");
        producer.setNamesrvAddr("192.168.23.127:9876");
        producer.start();
        for (int i = 1; i <= 5; i++) {
            //單向消息
            Message msg = new Message("topic2",("單向消息:hello rocketmq "+i).getBytes("UTF-8"));
            producer.sendOneway(msg);
        }
        //添加一個休眠操作,確保異步消息返回后能夠輸出
        TimeUnit.SECONDS.sleep(10);
        producer.shutdown();
    }
}

消費(fèi)者代碼同上

測試:

總結(jié) 同步消息

SendResult result = producer.send(msg);

異步消息(回調(diào)處理結(jié)果必須在生產(chǎn)者進(jìn)程結(jié)束前執(zhí)行,否則回調(diào)無法正確執(zhí)行)

		producer.send(msg, new SendCallback() {
                //表示成功返回結(jié)果
                @Override
                public void onSuccess(SendResult sendResult) {
                    System.out.println(sendResult);
                }
                //表示發(fā)送消息失敗
                @Override
                public void onException(Throwable throwable) {
                    System.out.println(throwable);
                }
            });

單向消息

producer.sendOneway(msg);

到此這篇關(guān)于圖文并茂講解RocketMQ消息類別的文章就介紹到這了,更多相關(guān)RocketMQ消息類別內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • Spring?Boot中@Validated注解不生效問題匯總大全

    Spring?Boot中@Validated注解不生效問題匯總大全

    這篇文章主要給大家介紹了關(guān)于Spring?Boot中@Validated注解不生效問題匯總的相關(guān)資料,@Validated注解是Spring框架中的一個注解,用于在方法參數(shù)上添加參數(shù)校驗(yàn)規(guī)則,需要的朋友可以參考下
    2023-07-07
  • 如何對?Excel?表格中提取的數(shù)據(jù)進(jìn)行批量更新

    如何對?Excel?表格中提取的數(shù)據(jù)進(jìn)行批量更新

    這篇文章主要介紹了如何對Excel表格中提取的數(shù)據(jù)進(jìn)行批量更新操作,本文通過示例代碼介紹的非常詳細(xì),感興趣的朋友跟隨小編一起看看吧
    2024-06-06
  • 詳解springboot如何更新json串里面的內(nèi)容

    詳解springboot如何更新json串里面的內(nèi)容

    這篇文章主要為大家介紹了springboot 如何更新json串里面的內(nèi)容,文中有詳細(xì)的解決方案供大家參考,對大家的學(xué)習(xí)或工作有一定的幫助,需要的朋友可以參考下
    2023-10-10
  • Java中@JSONField注解用法、場景與實(shí)踐詳解

    Java中@JSONField注解用法、場景與實(shí)踐詳解

    這篇文章主要給大家介紹了關(guān)于Java中@JSONField注解用法、場景與實(shí)踐的相關(guān)資料,并結(jié)合實(shí)際應(yīng)用場景,幫助開發(fā)者在項(xiàng)目中更高效地處理JSON數(shù)據(jù),文中通過代碼介紹的非常詳細(xì),需要的朋友可以參考下
    2024-12-12
  • SpringBoot整合Redis實(shí)現(xiàn)刷票過濾功能

    SpringBoot整合Redis實(shí)現(xiàn)刷票過濾功能

    隨著互聯(lián)網(wǎng)的不斷發(fā)展,網(wǎng)站或APP的用戶流量增加,也衍生出了一些惡意刷量等問題,給數(shù)據(jù)分析及運(yùn)營帶來極大的困難,所以本文使用SpringBoot和Redis實(shí)現(xiàn)一個刷票過濾功能,需要的可以參考一下
    2023-06-06
  • springboot2.0.0配置多數(shù)據(jù)源出現(xiàn)jdbcUrl is required with driverClassName的錯誤

    springboot2.0.0配置多數(shù)據(jù)源出現(xiàn)jdbcUrl is required with driverClassN

    這篇文章主要介紹了springboot2.0.0配置多數(shù)據(jù)源出現(xiàn)jdbcUrl is required with driverClassName的錯誤,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2020-11-11
  • Spring中為bean指定InitMethod和DestroyMethod的執(zhí)行方法

    Spring中為bean指定InitMethod和DestroyMethod的執(zhí)行方法

    在Spring中,那些組成應(yīng)用程序的主體及由Spring IoC容器所管理的對象,被稱之為bean,接下來通過本文給大家介紹Spring中為bean指定InitMethod和DestroyMethod的執(zhí)行方法,感興趣的朋友一起看看吧
    2021-11-11
  • java原生動態(tài)生成驗(yàn)證碼

    java原生動態(tài)生成驗(yàn)證碼

    這篇文章主要為大家詳細(xì)介紹了java原生動態(tài)生成驗(yàn)證碼,文中示例代碼介紹的非常詳細(xì),具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2020-10-10
  • Java實(shí)現(xiàn)提取Word文檔表格數(shù)據(jù)

    Java實(shí)現(xiàn)提取Word文檔表格數(shù)據(jù)

    使用Java實(shí)現(xiàn)Word文檔表格數(shù)據(jù)的提取,可以確保數(shù)據(jù)處理的一致性和準(zhǔn)確性,同時大大減少所需的時間和成本,下面我們來看看具體實(shí)現(xiàn)方法吧
    2025-01-01
  • Java使用FileReader讀取文件詳解

    Java使用FileReader讀取文件詳解

    本文將為大家介紹FileReader類的基本用法,包括如何創(chuàng)建FileReader對象,如何讀取文件,以及如何關(guān)閉流,感興趣的小伙伴可以跟隨小編一起了解一下
    2023-09-09

最新評論

东阳市| 昌黎县| 鄱阳县| 张家港市| 元谋县| 苍南县| 开封市| 阳朔县| 多伦县| 江源县| 漠河县| 株洲市| 张北县| 苏尼特左旗| 焦作市| 永寿县| 中宁县| 城固县| 垫江县| 赤水市| 商丘市| 余庆县| 新野县| 九江市| 榆树市| 锡林郭勒盟| 久治县| 景洪市| 翁源县| 明水县| 云梦县| 玉田县| 威海市| 蚌埠市| 黔西| 志丹县| 治县。| 屯门区| 黄平县| 翼城县| 常宁市|