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

使用Java代碼實現(xiàn)RocketMQ的生產(chǎn)與消費消息

 更新時間:2024年07月31日 08:43:23   作者:小威要向諸佬學習呀  
這篇文章介紹一下其他的小組件以及使用Java代碼實現(xiàn)生產(chǎn)者對消息的生成,消費者消費消息等知識點,并通過代碼示例介紹的非常詳細,對大家的學習或工作有一定的幫助,需要的朋友可以參考下

RocketMQ其他組件

在RocketMQ中,除了生產(chǎn)者,消費者,還有一些其他的小組件,接下來逐一介紹一下他們。

監(jiān)聽器(Listener)

定義:監(jiān)聽器是消費者用于處理消息的組件。在PushConsumer(推)模式下,消費者客戶端必須設置消費監(jiān)聽器,以便在接收到消息時執(zhí)行相應的處理邏輯。(比如一會兒下面的代碼)

偏移量(Offset)

定義:偏移量是指在消費消息時,記錄消費者已經(jīng)消費到的消息位置的值。每個消息都有一個唯一的偏移量值,它代表了消息在消息隊列中的位置。

偏移量具有很大的作用:它能夠保證消費者在重啟或宕機后能夠從上次消費的位置繼續(xù)消費消息,避免重復消費或漏消費。

簡而言之就是它能告訴消費者已經(jīng)消費到哪一條消息??!

  • 集群模式:在集群消費模式下,消息隊列的消費進度保存在Broker端。消費者每次消費完消息后,會將最新的消費進度同步到Broker,以便在消費者重啟或者故障轉(zhuǎn)移的時候能夠從上一次消費的位置繼續(xù)消費。
  • 廣播模式:在廣播消費模式下,消息隊列的消費進度保存在消費者本地。因為廣播模式下每條消息都會被所有消費者消費,所以不需要在Broker端保存消費進度。

所以,偏移量的的實現(xiàn)方式有兩種:包括存儲在本地文件(OffsetStore)和存儲在Broker中這兩種方式。(這樣一看,清晰了吧)

命名服務器(NameServer)

定義:命名服務器是RocketMQ中的輕量級路由服務,存儲生產(chǎn)者和消費者與Broker之間的路由信息。

它的作用:提供Broker的動態(tài)注冊與發(fā)現(xiàn)服務,生產(chǎn)者和消費者通過NameServer查詢Broker的路由信息,從而進行消息的投遞和消費。

消息組成

  • Topic:消息主題,對不同的業(yè)務消息進行分類。
  • Tag:消息標簽,進一步區(qū)分某個Topic下的消息分類。使用Tag可以實現(xiàn)對Topic中的消息進行過濾。消費者可以根據(jù)Tag來訂閱自己感興趣的消息,而不是接收Topic下的所有消息。
  • Message Body:消息體,消息的實際內(nèi)容。
  • Keys:消息的鍵值,標識消息的唯一性。在RocketMQ中,每個消息都可以設置Keys字段,以便在需要的時候根據(jù)Keys來查詢或者定位消息。
  • 屬性:除了上面滴,RocketMQ的消息還可以包含一系列的屬性信息,比如消息的發(fā)送時間、生產(chǎn)者信息等等。這些屬性信息以鍵值對的形式存在,隨著消息一起被存儲和傳輸。

實現(xiàn)生產(chǎn)與消費消息

按之前的步驟搭建完成RocketMQ集群后...

首先我們創(chuàng)建一個空的maven工程,在pom.xml文件中添加RocketMQ的依賴(RocketMQ的依賴版本需要與虛擬機中的保持一致,這里選擇和之前一樣的4.7.1版本):

        <dependency>
            <groupId>org.apache.rocketmq</groupId>
            <artifactId>rocketmq-client</artifactId>
            <version>4.7.1</version>
        </dependency>

生產(chǎn)消息

然后編寫生產(chǎn)者生產(chǎn)消息的代碼:

// 1.創(chuàng)建一個DefaultMQProducer實例,指定生產(chǎn)者組名 
DefaultMQProducer producer = new DefaultMQProducer("my-producer-group1");  
// 2.設置NameServer的地址
producer.setNamesrvAddr("192.168.220.135:9876");    
// 3.啟動生產(chǎn)者實例  
producer.start();   
// 4.使用for循環(huán)發(fā)送10條消息  
for (int i = 0; i < 10; i++) {  
    // 創(chuàng)建一條消息,指定Topic為"MyTopic1",Tag標簽為"TagA",消息體為"hello rocketmq"加上循環(huán)變量的值,同時把字符串轉(zhuǎn)換為字節(jié)數(shù)組  
    Message message = new Message("MyTopic1","TagA",("hello rocketmq"+i).getBytes(StandardCharsets.UTF_8));       
    // 5.發(fā)送消息并接收發(fā)送結(jié)果  
    SendResult sendResult = producer.send(message);  
    // 打印發(fā)送結(jié)果,包括消息ID、發(fā)送狀態(tài)等信息  
    System.out.println(sendResult);  
}  
// 6.發(fā)送完所有消息后,關(guān)閉生產(chǎn)者實例,釋放資源  
producer.shutdown();

生產(chǎn)者生產(chǎn)消息和消費者消費消息這塊的代碼都相對較為簡單,已經(jīng)在代碼塊中加了注釋,這里就不再贅述了。

這個時候就可以訪問虛擬機+端口號來搜索到發(fā)送的消息詳情了!

消費消息

// 1.和生產(chǎn)者一樣,創(chuàng)建一個DefaultMQPushConsumer實例
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("my-consumer-group1");  
// 2.設置NameServer的地址,消費者通過這個地址與NameServer進行通信,來獲取Broker的地址信息  
consumer.setNamesrvAddr("192.168.220.135:9876");  
// 訂閱一個或多個Topic,以及Tag來過濾需要消費的消息。我們訂閱了"MyTopic1",使用"*"來匹配此Topic下的所有Tag  
consumer.subscribe("MyTopic1", "*");  
// 3.注冊消息監(jiān)聽器,用于處理從Broker接收到的消息。使用MessageListenerConcurrently接口的實現(xiàn),表示并行消費  
consumer.registerMessageListener(new MessageListenerConcurrently() {  
    @Override  
    // 4.當收到消息時,方法會被調(diào)用。
    public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, ConsumeConcurrentlyContext consumeConcurrentlyContext) {  
        // 5.遍歷消息列表,并打印每條消息的內(nèi)容(注意:這里直接打印msg對象不會得到預期的消息內(nèi)容字符串)  
        for (MessageExt msg : msgs) {  
            // 所以我們打印msg.getBody()的內(nèi)容,為了保留消息原樣  
            System.out.println("已收到消息" + msg);  
        }  
        // 6.返回消費狀態(tài),這里表示消息已成功消費  
        return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;  
    }  
});  
// 7.啟動消費者實例  
consumer.start();  
// 8.打印日志消費者已啟動  
System.out.println("消費者已啟動");

這里需要注意,msg包含了消息的詳細信息,包括消息體、標簽、屬性等等。如果想打印消息內(nèi)容,應該使用msg.getBody()方法獲取消息體的字節(jié)數(shù)組,并且把它轉(zhuǎn)換為字符串(如果消息體是文本的話)。

本篇文章到這里就結(jié)束了,后續(xù)會繼續(xù)分享RocketMQ相關(guān)的知識,感謝各位小伙伴們的支持!

以上就是使用Java代碼實現(xiàn)RocketMQ的生產(chǎn)與消費消息的詳細內(nèi)容,更多關(guān)于Java實現(xiàn)RocketMQ生產(chǎn)與消費的資料請關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

  • spring rocketmq集成方案

    spring rocketmq集成方案

    本文詳細介紹了如何在Spring Boot項目中集成RocketMQ 5.x,包括前置條件、核心依賴、配置、生產(chǎn)者和消費者實現(xiàn)、測試驗證以及關(guān)鍵注意事項,感興趣的朋友跟隨小編一起看看吧
    2026-03-03
  • 徹底理解Java 中的ThreadLocal

    徹底理解Java 中的ThreadLocal

    這篇文章主要介紹了徹底理解Java 中的ThreadLocal的相關(guān)資料,需要的朋友可以參考下
    2017-07-07
  • spring?boot只需兩步優(yōu)雅整合activiti示例解析

    spring?boot只需兩步優(yōu)雅整合activiti示例解析

    這篇文章主要主要來教大家spring?boot優(yōu)雅整合activiti只需兩步就可完成測操作示例詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助祝大家多多進步
    2022-03-03
  • SpringBoot輕松實現(xiàn)ip解析(含源碼)

    SpringBoot輕松實現(xiàn)ip解析(含源碼)

    IP地址一般以數(shù)字形式表示,如192.168.0.1,IP解析是將這個數(shù)字IP轉(zhuǎn)換為包含地區(qū)、城市、運營商等信息的字符串形式,如“廣東省深圳市 電信”,這樣更方便人理解和使用,本文給大家介紹了SpringBoot如何輕松實現(xiàn)ip解析,需要的朋友可以參考下
    2023-10-10
  • 使用Java實現(xiàn)將ppt轉(zhuǎn)換為文本

    使用Java實現(xiàn)將ppt轉(zhuǎn)換為文本

    這篇文章主要為大家詳細介紹了如何使用Java實現(xiàn)將ppt轉(zhuǎn)換為文本,文中的示例代碼簡潔易懂,具有一定的借鑒價值,有需要的小伙伴可以參考下
    2024-01-01
  • SpringBoot 開發(fā)提速神器 Lombok+MybatisPlus+SwaggerUI

    SpringBoot 開發(fā)提速神器 Lombok+MybatisPlus+SwaggerUI

    這篇文章主要介紹了SpringBoot 開發(fā)提速神器 Lombok+MybatisPlus+SwaggerUI,本文給大家介紹的非常詳細,對大家的學習或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2021-03-03
  • Java中的CyclicBarrier同步屏障詳解

    Java中的CyclicBarrier同步屏障詳解

    這篇文章主要介紹了Java中的CyclicBarrier同步屏障詳解,CyclicBarrier也叫同步屏障,在JDK1.5被引入,可以讓一組線程達到一個屏障時被阻塞,直到最后一個線程達到屏障時,屏障才會開門,所有被阻塞的線程才會繼續(xù)執(zhí)行,需要的朋友可以參考下
    2023-09-09
  • Mybatis 中的<![CDATA[ ]]>淺析

    Mybatis 中的<![CDATA[ ]]>淺析

    本文給大家解析使用<![CDATA[ ]]>解決xml文件不被轉(zhuǎn)義的問題, 對mybatis 中的<![CDATA[ ]]>相關(guān)知識感興趣的朋友一起看看吧
    2017-09-09
  • Java實用小技能之快速創(chuàng)建List常用幾種方式

    Java實用小技能之快速創(chuàng)建List常用幾種方式

    java集合可以說無論是面試、刷題還是工作中都是非常常用的,下面這篇文章主要給大家介紹了關(guān)于Java實用小技能之快速創(chuàng)建List常用的幾種方式,文中通過實例代碼介紹的非常詳細,需要的朋友可以參考下
    2022-12-12
  • Java程序員的10道常見的XML面試問答題(XML術(shù)語詳解)

    Java程序員的10道常見的XML面試問答題(XML術(shù)語詳解)

    包括web開發(fā)人員的Java面試在內(nèi)的各種面試中,XML面試題在各種編程工作的面試中很常見。XML是一種成熟的技術(shù),經(jīng)常作為從一個平臺到其他平臺傳輸數(shù)據(jù)的標準
    2014-04-04

最新評論

华池县| 金塔县| 谢通门县| 兴海县| 柏乡县| 墨脱县| 武冈市| 金湖县| 龙门县| 固阳县| 永州市| 金沙县| 嘉荫县| 砚山县| 镇江市| 鹿邑县| 大竹县| 苏尼特右旗| 崇仁县| 焉耆| 昌吉市| 秦皇岛市| 淮阳县| 云霄县| 青州市| 华安县| 航空| 巫溪县| 双流县| 丽水市| 乌苏市| 长岛县| 合阳县| 博乐市| 光泽县| 新巴尔虎右旗| 榕江县| 阳泉市| 德清县| 阿坝县| 灵璧县|