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

Spring Boot 集成 RocketMQ 全流程指南(從依賴引入到消息收發(fā))

 更新時(shí)間:2025年04月15日 16:48:14   作者:qq_17153885  
本文將通過(guò) 手動(dòng)連接 和 配置連接 兩種方式,詳細(xì)講解如何在 Spring Boot 中集成 RocketMQ,實(shí)現(xiàn)消息的同步與異步發(fā)送,并提供完整示例代碼,感興趣的朋友一起看看吧

前言

在分布式系統(tǒng)中,消息中間件是解耦服務(wù)、實(shí)現(xiàn)異步通信的核心組件。RocketMQ 作為阿里巴巴開源的高性能分布式消息中間件,憑借其高吞吐、低延遲、高可靠等特性,成為企業(yè)級(jí)應(yīng)用的首選。而 Spring Boot 通過(guò)其“約定優(yōu)于配置”的設(shè)計(jì)理念,極大簡(jiǎn)化了項(xiàng)目開發(fā)的復(fù)雜度。本文將通過(guò) 手動(dòng)連接配置連接 兩種方式,詳細(xì)講解如何在 Spring Boot 中集成 RocketMQ,實(shí)現(xiàn)消息的同步與異步發(fā)送,并提供完整示例代碼。

一、環(huán)境準(zhǔn)備

在開始前,請(qǐng)確保:

  • JDK 17、Maven 3.6+、Spring Boot 2.7+。
  • 安裝RocketMQ服務(wù)(本地或遠(yuǎn)程),推薦使用RocketMQ Docker鏡像快速搭建(可參考之前文章)。

二、示例—Springboot集成mq(手動(dòng)連接)

通過(guò)編碼方式初始化生產(chǎn)者,適用于需要?jiǎng)討B(tài)控制資源的場(chǎng)景。

2.1 新建項(xiàng)目

‍

2.2 引入依賴

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

image

2.3 生產(chǎn)者發(fā)送消息

  • 構(gòu)建一個(gè)消息生產(chǎn)者DefaultMQProducer實(shí)例,然后指定生產(chǎn)者組為jihaiProducer;
  • 指定NameServer的地址:服務(wù)器的ip:9876,因?yàn)樾枰獜腘ameServer拉取Broker的信息
  • producer.start() 啟動(dòng)生產(chǎn)者
  • 構(gòu)建一個(gè)內(nèi)容為:技海拾貝的消息1,然后指定這個(gè)消息往jihaishibei這個(gè)topic發(fā)送
  • producer.send(msg):發(fā)送消息,打印結(jié)果
  • 關(guān)閉生產(chǎn)者
public class Producer {
    public static void main(String[] args) throws Exception {
        //創(chuàng)建一個(gè)生產(chǎn)者,指定生產(chǎn)者組為jihaiProducer
        DefaultMQProducer producer = new DefaultMQProducer("jihaiProducer");
        // 指定NameServer的地址
        producer.setNamesrvAddr("localhost:9876");
        // 第一次發(fā)送可能會(huì)超時(shí),設(shè)置的比較大
        producer.setSendMsgTimeout(60000);
        // 啟動(dòng)生產(chǎn)者
        producer.start();
        // 創(chuàng)建一條消息
        // topic為 jihaishibei
        // 消息內(nèi)容為 技海拾貝的消息1
        // tags 為 TagA
        Message msg = new Message("jihaishibei", "TagA", "技海拾貝的消息1 ".getBytes(RemotingHelper.DEFAULT_CHARSET));
        // 發(fā)送消息并得到消息的發(fā)送結(jié)果,然后打印
        SendResult sendResult = producer.send(msg);
        System.out.printf("%s%n", sendResult);
        // 關(guān)閉生產(chǎn)者
        producer.shutdown();
    }
}

image

啟動(dòng),發(fā)送消息

image

在控制臺(tái)可以看到這條消息

image

image

這里就能看到發(fā)送消息的詳細(xì)信息。

左下角消息的消費(fèi)的消費(fèi),因?yàn)槲覀冞€沒(méi)有消費(fèi)者訂閱這個(gè)topic,所以左下角沒(méi)數(shù)據(jù)。

2.4 消費(fèi)者消費(fèi)消息

  • 創(chuàng)建一個(gè)消費(fèi)者實(shí)例對(duì)象,指定消費(fèi)者組為jihaiConsumer
  • 指定NameServer的地址:服務(wù)器的ip:9876
  • 訂閱 jihaishibei這個(gè)topic的所有信息
  • consumer.registerMessageListener ,這個(gè)很重要,是注冊(cè)一個(gè)監(jiān)聽器,這個(gè)監(jiān)聽器是當(dāng)有消息的時(shí)候就會(huì)回調(diào)這個(gè)監(jiān)聽器,處理消息,所以需要用戶實(shí)現(xiàn)這個(gè)接口,然后處理消息。
  • 啟動(dòng)消費(fèi)者
public class Consumer {
    public static void main(String[] args) throws InterruptedException, MQClientException {
        // 通過(guò)push模式消費(fèi)消息,指定消費(fèi)者組
        DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("jihaiConsumer");
        // 指定NameServer的地址
        consumer.setNamesrvAddr("localhost:9876");
        // 訂閱這個(gè)topic下的所有的消息
        consumer.subscribe("jihaishibei", "*");
        // 注冊(cè)一個(gè)消費(fèi)的監(jiān)聽器,當(dāng)有消息的時(shí)候,會(huì)回調(diào)這個(gè)監(jiān)聽器來(lái)消費(fèi)消息
        consumer.registerMessageListener(new MessageListenerConcurrently() {
            @Override
            public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs,
                                                            ConsumeConcurrentlyContext context) {
                for (MessageExt msg : msgs) {
                    System.out.printf("消費(fèi)消息:%s", new String(msg.getBody()) + "\n");
                }
                return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
            }
        });
        // 啟動(dòng)消費(fèi)者
        consumer.start();
        System.out.printf("Consumer Started.%n");
    }
}

image

啟動(dòng)服務(wù),進(jìn)行消費(fèi)

image

在控制臺(tái),發(fā)現(xiàn)被jihaiConsumer這個(gè)消費(fèi)者組給消費(fèi)了。

image

三、示例2—Springboot集成mq(配置連接)

在 Spring Boot 中,可以通過(guò)配置文件簡(jiǎn)化 RocketMQ 的連接配置。以下是在 application.yml? 文件中進(jìn)行的配置:

3.1 配置文件修改

image

image

image

spring:
  application:
    name: rocket-mq-demo
rocketmq:
  name-server: 127.0.0.1:9876
  producer:
    group: rocket-mq-demo-producer
    send-message-timeout: 10000
  comsumer:
    group: rocket-mq-demo-comsumer
    send-message-timeout: 10000

image

3.2 添加依賴

<dependency>
    <groupId>org.apache.rocketmq</groupId>
    <artifactId>rocketmq-spring-boot-starter</artifactId>
    <version>2.3.0</version>  
</dependency>

根據(jù)需要選擇最新版本,從中央倉(cāng)庫(kù)可以查看?https://central.sonatype.com/artifact/org.apache.rocketmq/rocketmq-spring-boot-starter?

image

image

備注:如果添加rocketmq-client依賴,先注釋這個(gè)依賴

3.3 消費(fèi)者service類

import org.apache.rocketmq.spring.annotation.MessageModel;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Service;
/**
 * messageModel=MessageModel.CLUSTERING
 * 監(jiān)聽模式,有消息就會(huì)消費(fèi)
 */
@Service
@RocketMQMessageListener(topic = "jihaishibei-topic", consumerGroup = "rocket-mq-demo-comsumer", messageModel = MessageModel.CLUSTERING)
public class RocketMQConsumer implements RocketMQListener<String> {
    @Override
    public void onMessage(String s) {
        System.out.printf("收到消息: %s\n", s);
    }
}

3.4 生產(chǎn)者service類

import org.apache.rocketmq.client.producer.SendCallback;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Service;
@Service
public class RocketMQProducer {
    @Autowired
    private RocketMQTemplate rocketMQTemplate;
    private final String topic = "jihaishibei-topic";
    // 1.同步發(fā)送消息
    // 同步發(fā)送是指發(fā)送方發(fā)送一條消息后,會(huì)等待服務(wù)器返回確認(rèn)信息后再進(jìn)行后續(xù)操作。這種方式適用于需要可靠性保證的場(chǎng)景。
    public void createAndSend(String message){
        rocketMQTemplate.convertAndSend(topic, message);
        System.out.printf("同步發(fā)送結(jié)果: %s\n", message);
    }
    // 1.同步發(fā)送消息
    // 同步發(fā)送是指發(fā)送方發(fā)送一條消息后,會(huì)等待服務(wù)器返回確認(rèn)信息后再進(jìn)行后續(xù)操作。這種方式適用于需要可靠性保證的場(chǎng)景。
    public void sendSyncMessage(String message){
        SendResult sendResult = rocketMQTemplate.syncSend(topic, MessageBuilder.withPayload(message).build());
        System.out.println(sendResult.getMsgId());
        System.out.printf("同步發(fā)送結(jié)果: %s\n", message);
    }
    // 2.異步發(fā)送消息
    // 異步發(fā)送是指發(fā)送方發(fā)送消息后,不等待服務(wù)器返回確認(rèn)信息,而是通過(guò)回調(diào)接口處理返回結(jié)果。這種方式適用于對(duì)響應(yīng)時(shí)間要求較高的場(chǎng)景。
    public void sendAsyncMessage(String message){
        rocketMQTemplate.asyncSend(topic, MessageBuilder.withPayload(message).build(), new SendCallback() {
            @Override
            public void onSuccess(SendResult sendResult) {
                System.out.printf("異步發(fā)送成功: %s\n", sendResult);
            }
            @Override
            public void onException(Throwable throwable) {
                System.out.printf("異步發(fā)送失敗: %s\n", throwable.getMessage());
            }
        });
    }
    // 3.單向發(fā)送消息
    // 單向發(fā)送是指發(fā)送方只負(fù)責(zé)發(fā)送消息,不關(guān)心服務(wù)器的響應(yīng)。該方式適用于對(duì)可靠性要求不高的場(chǎng)景,如日志收集。
    public void sendOneWayMessage(String message){
        rocketMQTemplate.sendOneWay(topic, MessageBuilder.withPayload(message).build());
        System.out.println("單向消息發(fā)送成功");
    }
}

3.5 測(cè)試controller類

@RequestMapping("api")
@RestController
public class RocketController {
    @Autowired
    private RocketMQProducer rocketMQProducer;
    @GetMapping("/createAndSend")
    public String createAndSend(@RequestParam String message) {
        rocketMQProducer.createAndSend(message);
        return "同步消息發(fā)送成功";
    }
    @GetMapping("/sendSync")
    public String sendSync(@RequestParam String message) {
        rocketMQProducer.sendSyncMessage(message);
        return "同步消息發(fā)送成功";
    }
    @GetMapping("/sendAsync")
    public String sendAsync(@RequestParam String message) {
        rocketMQProducer.sendAsyncMessage(message);
        return "異步消息發(fā)送中";
    }
    @GetMapping("/sendOneWay")
    public String sendOneWay(@RequestParam String message) {
        rocketMQProducer.sendOneWayMessage(message);
        return "單向消息發(fā)送成功";
    }
}

3.6 啟動(dòng)服務(wù)

image

3.7 測(cè)試 同步消息1

image

image

image

同步消息2

image

image

image

異步消息

image

image

image

單向發(fā)送消息

image

image

image

四、結(jié)束語(yǔ)

本文通過(guò)手動(dòng)連接與配置連接兩種方式,展示了Spring Boot與RocketMQ的集成實(shí)踐。手動(dòng)連接幫助開發(fā)者理解底層API邏輯,而Spring Boot的配置化集成則極大簡(jiǎn)化了開發(fā)流程。無(wú)論是同步消息的可靠性保障,還是異步消息的性能優(yōu)化,RocketMQ均能與Spring Boot無(wú)縫協(xié)作,為分布式系統(tǒng)提供高效的消息通信能力。

未來(lái)可進(jìn)一步探索集群部署、消息重試機(jī)制及監(jiān)控告警,以實(shí)現(xiàn)更健壯的消息服務(wù)。希望本文能為開發(fā)者快速構(gòu)建高可用的消息系統(tǒng)提供參考!?

到此這篇關(guān)于Spring Boot 集成 RocketMQ 全流程指南(從依賴引入到消息收發(fā))的文章就介紹到這了,更多相關(guān)Spring Boot 集成 RocketMQ內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • 一文帶你掌握J(rèn)ava8中Lambda表達(dá)式 函數(shù)式接口及方法構(gòu)造器數(shù)組的引用

    一文帶你掌握J(rèn)ava8中Lambda表達(dá)式 函數(shù)式接口及方法構(gòu)造器數(shù)組的引用

    Java 8 (又稱為 jdk 1.8) 是 Java 語(yǔ)言開發(fā)的一個(gè)主要版本。 Oracle 公司于 2014 年 3 月 18 日發(fā)布 Java 8 ,它支持函數(shù)式編程,新的 JavaScript 引擎,新的日期 API,新的Stream API 等
    2021-10-10
  • java反射機(jī)制給實(shí)體類相同字段自動(dòng)賦值實(shí)例

    java反射機(jī)制給實(shí)體類相同字段自動(dòng)賦值實(shí)例

    這篇文章主要介紹了java反射機(jī)制給實(shí)體類相同字段自動(dòng)賦值實(shí)例,具有
    2020-08-08
  • SpringBoot基于MyBatisPlus實(shí)現(xiàn)公共字段自動(dòng)填充

    SpringBoot基于MyBatisPlus實(shí)現(xiàn)公共字段自動(dòng)填充

    本文介紹了一種在MyBatisPlus中實(shí)現(xiàn)公共字段自動(dòng)填充的方法,包括創(chuàng)建MetaObjectHandler實(shí)現(xiàn)類和在實(shí)體類中添加注解,具有一定的參考價(jià)值,感興趣的可以了解一下
    2025-08-08
  • Java如何獲取一個(gè)IP段內(nèi)的所有IP地址

    Java如何獲取一個(gè)IP段內(nèi)的所有IP地址

    這篇文章主要為大家詳細(xì)介紹了Java如何根據(jù)起始和結(jié)束的IP地址獲取IP段內(nèi)所有IP地址,文中的示例代碼講解詳細(xì),感興趣的小伙伴可以跟隨小編一起學(xué)習(xí)一下
    2024-11-11
  • Java復(fù)雜鏈表的復(fù)制詳解

    Java復(fù)雜鏈表的復(fù)制詳解

    復(fù)雜鏈表指的是一個(gè)鏈表有若干個(gè)結(jié)點(diǎn),每個(gè)結(jié)點(diǎn)有一個(gè)數(shù)據(jù)域用于存放數(shù)據(jù),還有兩個(gè)指針域,其中一個(gè)指向下一個(gè)節(jié)點(diǎn),還有一個(gè)隨機(jī)指向當(dāng)前復(fù)雜鏈表中的任意一個(gè)節(jié)點(diǎn)或者是一個(gè)空結(jié)點(diǎn),我們來(lái)探究一下在Java中復(fù)雜鏈表的復(fù)制
    2022-01-01
  • Java調(diào)用C++動(dòng)態(tài)庫(kù)DLL進(jìn)行無(wú)圖像無(wú)實(shí)體處理

    Java調(diào)用C++動(dòng)態(tài)庫(kù)DLL進(jìn)行無(wú)圖像無(wú)實(shí)體處理

    這篇文章主要為大家詳細(xì)介紹了如何通過(guò) JNI(Java Native Interface)在 Java 中調(diào)用一個(gè)用 C++ 編寫的分割算法庫(kù),實(shí)現(xiàn)無(wú)圖像無(wú)實(shí)體處理,感興趣的小伙伴可以了解下
    2025-09-09
  • SpringBoot中的@Configuration注解詳解

    SpringBoot中的@Configuration注解詳解

    這篇文章主要介紹了SpringBoot中的@Configuration注解詳解,Spring Boot推薦使用JAVA配置來(lái)完全代替XML 配置,JAVA配置就是通過(guò) @Configuration和 @Bean兩個(gè)注解實(shí)現(xiàn)的,需要的朋友可以參考下
    2023-08-08
  • SpringBoot項(xiàng)目如何添加2FA雙因素身份認(rèn)證

    SpringBoot項(xiàng)目如何添加2FA雙因素身份認(rèn)證

    雙因素身份驗(yàn)證2FA是一種安全系統(tǒng),要求用戶提供兩種不同的身份驗(yàn)證方式才能訪問(wèn)某個(gè)系統(tǒng)或服務(wù),國(guó)內(nèi)普遍做短信驗(yàn)證碼這種的用的比較少,不過(guò)在國(guó)外的網(wǎng)站中使用雙因素身份驗(yàn)證的還是很多的,這篇文章主要介紹了SpringBoot項(xiàng)目如何添加2FA雙因素身份認(rèn)證,需要的朋友參考下
    2024-04-04
  • Java中List與數(shù)組之間的相互轉(zhuǎn)換

    Java中List與數(shù)組之間的相互轉(zhuǎn)換

    在日常Java學(xué)習(xí)或項(xiàng)目開發(fā)中,經(jīng)常會(huì)遇到需要int[]數(shù)組和List列表相互轉(zhuǎn)換的場(chǎng)景,然而往往一時(shí)難以想到有哪些方法,最后可能會(huì)使用暴力逐個(gè)轉(zhuǎn)換法,往往不是我們所滿意的,下面這篇文章主要給大家介紹了關(guān)于Java中List與數(shù)組之間的相互轉(zhuǎn)換,需要的朋友可以參考下
    2023-05-05
  • 查詢Java程序日志的實(shí)現(xiàn)方式

    查詢Java程序日志的實(shí)現(xiàn)方式

    在Linux中查詢Java日志需定位路徑(如/var/log/、/opt/、~/, 使用find)、用tail/grep查看內(nèi)容、journalctl查系統(tǒng)服務(wù)日志、調(diào)整日志級(jí)別、監(jiān)控變化及管理日志輪轉(zhuǎn)
    2025-09-09

最新評(píng)論

枣阳市| 太白县| 青田县| 湖口县| 顺昌县| 淅川县| 平度市| 松阳县| 宝山区| 新宁县| 积石山| 定远县| 慈利县| 梧州市| 吐鲁番市| 富源县| 洪泽县| 九江市| 武定县| 镇沅| 门源| 前郭尔| 河间市| 缙云县| 北安市| 吉木乃县| 延吉市| 蚌埠市| 乌鲁木齐县| 磐安县| 肃南| 开阳县| 扎兰屯市| 洪泽县| 芮城县| 宜黄县| 娄烦县| 海原县| 昆山市| 太湖县| 南川市|