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

Springboot RabbitMQ 消息隊列使用示例詳解

 更新時間:2024年06月05日 10:12:34   作者:bj_wasin  
本文通過示例代碼介紹了Springboot RabbitMQ 消息隊列使用,對大家的學(xué)習(xí)或工作具有一定的參考借鑒價值,感興趣的朋友跟隨小編一起看看吧

一、概念介紹:

RabbitMQ中幾個重要的概念介紹:

  • Channels:信道,多路復(fù)用連接中的一條獨立的雙向數(shù)據(jù)流通道。信道是建立在真實的 TCP 連接內(nèi)地虛擬連接,AMQP 命令都是通過信道發(fā)出去的,不管是發(fā)布消息、訂閱隊列還是接收消息,這些動作都是通過信道完成。因為對于操作系統(tǒng)來說建立和銷毀 TCP 都是非常昂貴的開銷,所以引入了信道的概念,以復(fù)用一條 TCP 連接。
  • Exchanges:交換器,用來接收生產(chǎn)者發(fā)送的消息并將這些消息路由給服務(wù)器中的隊列。
  • 交換機類型主要有以下幾種:
  • Direct Exchange(直連交換機):這種類型的交換機根據(jù)消息的Routing Key(路由鍵)進行精確匹配,只有綁定了相同路由鍵的隊列才會收到消息。適用于點對點的消息傳遞場景。
  • Fanout Exchange(扇形交換機):這種類型的交換機采用廣播模式,它會將消息發(fā)送給所有綁定到該交換機的隊列,不管消息的路由鍵是什么。適用于消息需要被多個消費者處理的場景。
  • Topic Exchange(主題交換機):這種類型的交換機支持基于模式匹配的路由鍵,可以使用通配符*(匹配一個單詞)和#(匹配零個或多個單詞)進行匹配。適用于實現(xiàn)更復(fù)雜的消息路由邏輯。
  • Headers Exchange(頭交換機):這種類型的交換機不處理路由鍵,而是根據(jù)發(fā)送的消息內(nèi)容中的headers屬性進行匹配。適用于需要在消息頭中攜帶額外信息的場景。
  • Queues:消息隊列,用來保存消息直到發(fā)送給消費者。它是消息的容器,也是消息的終點。一個消息可投入一個或多個隊列。消息一直在隊列里面,等待消費者連接到這個隊列將其取走。

二、引入依賴:

 <dependency>
      <groupId>org.springframework.boot</groupId>
      <artifactId>spring-boot-starter-amqp</artifactId>
 </dependency>

三、添加配置信息

spring:
  rabbitmq:
    host: 127.0.0.1
    port: 5672
    username: guest
    password: guest
    listener:
      simple:
        acknowledge-mode: manual  # 手動提交

四、Direct Exchange(直連交換機)模式

1、新建配置文件 RabbitDirectConfig類

package com.example.direct;
import org.springframework.amqp.core.Binding;
import org.springframework.amqp.core.BindingBuilder;
import org.springframework.amqp.core.DirectExchange;
import org.springframework.amqp.core.Queue;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
/**
 * @author wasin
 * @version 1.0
 * @date 2024/6/4
 * @description: 直連交換機--這種類型的交換機根據(jù)消息的Routing Key(路由鍵)進行精確匹配,
 * 只有綁定了相同路由鍵的隊列才會收到消息。適用于點對點的消息傳遞場景
 */
@Configuration
public class RabbitDirectConfig {
    /**
     * 隊列名稱
     */
    public static final String QUEUE_MESSAGE ="QUEUE_MESSAGE";
    public static final String QUEUE_USER ="QUEUE_USER";
    /**
     * 交換機
     */
    public static final String EXCHANGE="EXCHANGE_01";
    /**
     * 路由
     */
    public static final String ROUTING_KEY="ROUTING_KEY_01";
    @Bean
    public Queue queue01() {
        return new Queue(QUEUE_MESSAGE, //隊列名稱
                true, //是否持久化
                false, //是否排他
                false //是否自動刪除
        );
    }
    @Bean
    public Queue queue02() {
        return new Queue(QUEUE_USER, //隊列名稱
                true, //是否持久化
                false, //是否排他
                false //是否自動刪除
        );
    }
    @Bean
    public DirectExchange exchange01() {
        return new DirectExchange(EXCHANGE,
                true, //是否持久化
                false //是否排他
        );
    }
    @Bean
    public Binding demoBinding() {
        return BindingBuilder.bind(queue01()).to(exchange01()).with(ROUTING_KEY);
    }
    @Bean
    public Binding demoBinding2() {
        return BindingBuilder.bind(queue02()).to(exchange01()).with(ROUTING_KEY);
    }
}

2、添加消息生產(chǎn)者 Producer類

package com.example.direct;
import com.example.entity.User;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
/**
 * @author wasin
 * @version 1.0
 * @date 2024/6/4
 * @description:
 */
@Component
public class Producer {
    @Resource
    RabbitTemplate rabbitTemplate;
    public void sendMessageByExchangeANdRoute(String message){
        rabbitTemplate.convertAndSend(RabbitDirectConfig.EXCHANGE, RabbitDirectConfig.ROUTING_KEY,message);
    }
    /**
     * 默認交換器,隱式地綁定到每個隊列,路由鍵等于隊列名稱。
     * @param message
     */
    public void sendMessageByQueue(String message){
        rabbitTemplate.convertAndSend(RabbitDirectConfig.QUEUE_MESSAGE,message);
    }
    public void sendMessage(User user){
        rabbitTemplate.convertAndSend(RabbitDirectConfig.QUEUE_USER,user);
    }
}

3、添加消息消費者

package com.example.direct;
import com.example.entity.User;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
/**
 * @author wasin
 * @version 1.0
 * @date 2024/6/4
 * @description:
 */
@Component
public class Consumer {
    @RabbitListener(queues = RabbitDirectConfig.QUEUE_USER)
    public void onMessage(User user){
        System.out.println("收到的實體bean消息:"+user);
    }
    @RabbitListener(queues = RabbitDirectConfig.QUEUE_MESSAGE)
    public void onMessage2(String message){
        System.out.println("收到的字符串消息:"+message);
    }
}

4、 測試

package com.example;
import com.example.entity.User;
import com.example.direct.Producer;
import com.example.fanout.FanoutProducer;
import com.example.topic.TopicProducer;
import org.junit.jupiter.api.Test;
import org.springframework.boot.test.context.SpringBootTest;
import javax.annotation.Resource;
@SpringBootTest
class SpringbootRabbitMqApplicationTests {
    @Resource
    Producer producer;
    @Test
    public void sendMessage() throws InterruptedException {
        producer.sendMessageByQueue("哈哈");
        producer.sendMessage(new User().setAge(10).setName("wasin"));
    }
}

五、Topic Exchange(主題交換機)模式

1、新建RabbitTopicConfig類

package com.example.topic;
import org.springframework.amqp.core.*;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
/**
 * @author wasin
 * @version 1.0
 * @date 2024/6/4
 * @description: 主題交換機--這種類型的交換機支持基于模式匹配的路由鍵,
 * 可以使用通配符*(匹配一個單詞)和#(匹配零個或多個單詞)進行匹配。適用于實現(xiàn)更復(fù)雜的消息路由邏輯。
 */
@Configuration
public class RabbitTopicConfig {
    /**
     * 交換機
     */
    public static final String EXCHANGE = "EXCHANGE_TOPIC1";
    /**
     * 隊列名稱
     */
    public static final String QUEUE_TOPIC1 = "QUEUE_TOPIC";
    /**
     * 路由
     * "*" 與 "#",用于做模糊匹配。其中 "*" 用于匹配一個單詞,"#" 用于匹配多個單詞(可以是零個)
     * 可以匹配 aa.wasin.aa.bb  wasin.aa.bb  wasin.aa ....
     * aa.bb.wasin.cc 無法匹配
     */
    public static final String ROUTING_KEY1 = "*.wasin.#";
    @Bean
    public Queue queue() {
        return new Queue(QUEUE_TOPIC1, //隊列名稱
                true, //是否持久化
                false, //是否排他
                false //是否自動刪除
        );
    }
    @Bean
    public TopicExchange exchange() {
        return new TopicExchange(EXCHANGE,
                true, //是否持久化
                false //是否排他
        );
    }
    @Bean
    public Binding binding() {
        return BindingBuilder.bind(queue()).to(exchange()).with(ROUTING_KEY1);
    }
}

2、新建 消息生產(chǎn)者和發(fā)送者

TopicProducer類

package com.example.topic;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
/**
 * @author wasin
 * @version 1.0
 * @date 2024/6/4
 * @description:
 */
@Component
public class TopicProducer {
    @Resource
    RabbitTemplate rabbitTemplate;
    /**
     * @param routeKey 路由
     * @param message 消息
     */
    public void sendMessageByQueue(String routeKey, String message){
        rabbitTemplate.convertAndSend(RabbitTopicConfig.EXCHANGE,routeKey,message);
    }
}

TopicConsumer類

package com.example.topic;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
/**
 * @author wasin
 * @version 1.0
 * @date 2024/6/4
 * @description:
 */
@Slf4j
@Component
public class TopicConsumer {
    @RabbitListener(queues = RabbitTopicConfig.QUEUE_TOPIC1)
    public void onMessage2(String message){
        log.info("topic收到的字符串消息:{}",message);
    }
}

六、Fanout Exchange(扇形交換機)模式

1、 新建 RabbitFanoutConfig類

package com.example.fanout;
import org.springframework.amqp.core.*;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
/**
 * @author wasin
 * @version 1.0
 * @date 2024/6/4
 * @description: 扇形交換機--這種類型的交換機采用廣播模式,它會將消息發(fā)送給所有綁定到該交換機的隊列,
 * 不管消息的路由鍵是什么。適用于消息需要被多個消費者處理的場景。
 */
@Configuration
public class RabbitFanoutConfig {
    /**
     * 交換機
     */
    public static final String EXCHANGE = "EXCHANGE_FANOUT";
    /**
     * 隊列名稱
     */
    public static final String QUEUE_FANOUT1 = "QUEUE_FANOUT";
    /**
     * 隊列名稱
     */
    public static final String QUEUE_FANOUT2 = "QUEUE_FANOUT2";
    @Bean
    public Queue queueFanout1() {
        return new Queue(QUEUE_FANOUT1, //隊列名稱
                true, //是否持久化
                false, //是否排他
                false //是否自動刪除
        );
    }
    @Bean
    public Queue queueFanout2() {
        return new Queue(QUEUE_FANOUT2, //隊列名稱
                true, //是否持久化
                false, //是否排他
                false //是否自動刪除
        );
    }
    @Bean
    public FanoutExchange exchangeFanout() {
        return new FanoutExchange(EXCHANGE,
                true, //是否持久化
                false //是否排他
        );
    }
    @Bean
    public Binding bindingFanout() {
        return BindingBuilder.bind(queueFanout1()).to(exchangeFanout());
    }
    @Bean
    public Binding bindingFanout2() {
        return BindingBuilder.bind(queueFanout2()).to(exchangeFanout());
    }
}

2、新建 消息生產(chǎn)者和發(fā)送者

FanoutProducer類:

package com.example.fanout;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
/**
 * @author wasin
 * @version 1.0
 * @date 2024/6/4
 * @description:
 */
@Component
public class FanoutProducer {
    @Resource
    RabbitTemplate rabbitTemplate;
    /**
     * @param message 消息
     */
    public void sendMessageByQueue(String message) {
        rabbitTemplate.convertAndSend(RabbitFanoutConfig.EXCHANGE, "", message);
    }
}

FanoutConsumer類

package com.example.fanout;
import com.rabbitmq.client.Channel;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.amqp.support.AmqpHeaders;
import org.springframework.messaging.handler.annotation.Header;
import org.springframework.stereotype.Component;
import java.io.IOException;
/**
 * @author wasin
 * @version 1.0
 * @date 2024/6/4
 * @description:
 */
@Slf4j
@Component
public class FanoutConsumer {
    /**
     * 手動提交
     * @param message
     * @param channel
     * @param tag
     * @throws IOException
     */
    @RabbitListener(queues = RabbitFanoutConfig.QUEUE_FANOUT1)
    public void onMessage(String message, Channel channel, @Header(AmqpHeaders.DELIVERY_TAG) long tag) throws IOException {
        log.info("fanout1收到的字符串消息:{}",message);
        channel.basicAck(tag,false);
    }
    @RabbitListener(queues = RabbitFanoutConfig.QUEUE_FANOUT2)
    public void onMessage2(String message){
        log.info("fanout2到的字符串消息:{}",message);
    }
}

到此這篇關(guān)于Springboot RabbitMQ 消息隊列使用的文章就介紹到這了,更多相關(guān)Springboot RabbitMQ 消息隊列內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • javaweb開發(fā)提高效率利器JRebel詳解

    javaweb開發(fā)提高效率利器JRebel詳解

    這篇文章主要介紹了javaweb開發(fā)提高效率利器JRebel詳解,本文通過圖文并茂的形式給大家介紹的非常詳細,對大家的學(xué)習(xí)或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2021-04-04
  • SpringBoot封裝MinIO工具的實現(xiàn)步驟

    SpringBoot封裝MinIO工具的實現(xiàn)步驟

    MinIO是一款高性能、開源、兼容Amazon?S3?API的分布式對象存儲系統(tǒng),Spring?Boot通過封裝配置與工具類,標準化設(shè)計降低開發(fā)復(fù)雜度,提升開發(fā)效率與系統(tǒng)穩(wěn)定性,感興趣的可以了解一下
    2025-07-07
  • Scala之文件讀取、寫入、控制臺操作的方法示例

    Scala之文件讀取、寫入、控制臺操作的方法示例

    這篇文章主要介紹了Scala之文件讀取、寫入、控制臺操作的方法示例,文中通過示例代碼介紹的非常詳細,對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2019-06-06
  • Java計算代碼段執(zhí)行時間的詳細過程

    Java計算代碼段執(zhí)行時間的詳細過程

    java里計算代碼段執(zhí)行時間可以有兩種方法,一種是毫秒級別的計算,另一種是更精確的納秒級別的計算,這篇文章主要介紹了java計算代碼段執(zhí)行時間,需要的朋友可以參考下
    2023-02-02
  • Java多線程實現(xiàn)聊天客戶端和服務(wù)器

    Java多線程實現(xiàn)聊天客戶端和服務(wù)器

    這篇文章主要為大家詳細介紹了Java多線程聊天客戶端和服務(wù)器實現(xiàn)代碼,文中示例代碼介紹的非常詳細,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2016-10-10
  • AgentScope Java 核心架構(gòu)深度解析(推薦)

    AgentScope Java 核心架構(gòu)深度解析(推薦)

    AgentScopeJava是一個面向生產(chǎn)環(huán)境的智能體編程框架,它將大語言模型的推理能力、工具調(diào)用、記憶管理和多智能體協(xié)作整合在一起,本文從ReAct推理循環(huán)、工具系統(tǒng)、記憶管理、多智能體協(xié)作和生產(chǎn)就緒特性五個維度,深入剖析了框架的核心機制,感興趣的朋友一起看看吧
    2025-12-12
  • SpringSecurity實現(xiàn)前后端分離的示例詳解

    SpringSecurity實現(xiàn)前后端分離的示例詳解

    Spring Security默認提供賬號密碼認證方式,具體實現(xiàn)是在UsernamePasswordAuthenticationFilter 中,這篇文章主要介紹了SpringSecurity實現(xiàn)前后端分離的示例詳解,需要的朋友可以參考下
    2023-03-03
  • Java Scala之模式匹配與隱式轉(zhuǎn)換

    Java Scala之模式匹配與隱式轉(zhuǎn)換

    在Java中我們有switch case default這三個組成的基礎(chǔ)語法,在Scala中我們是有match和case組成 default的作用由case代替,本文詳細介紹了Scala的模式匹配與隱式轉(zhuǎn)換,感興趣的可以參考本文
    2023-04-04
  • java并發(fā)編程之同步器代碼示例

    java并發(fā)編程之同步器代碼示例

    這篇文章主要介紹了java并發(fā)編程之同步器代碼示例,分享了相關(guān)代碼,具有一定參考價值,需要的朋友可以了解下。
    2017-11-11
  • java 實現(xiàn)websocket的兩種方式實例詳解

    java 實現(xiàn)websocket的兩種方式實例詳解

    這篇文章主要介紹了java 實現(xiàn)websocket的兩種方式實例詳解,一種使用tomcat的websocket實現(xiàn),一種使用spring的websocket,本文通過代碼給大家介紹的非常詳細,需要的朋友可以參考下
    2018-07-07

最新評論

道孚县| 文安县| 成都市| 方正县| 封丘县| 中卫市| 海口市| 陕西省| 巴彦淖尔市| 临泉县| 吉首市| 仁化县| 昭觉县| 小金县| 加查县| 大同县| 曲麻莱县| 顺昌县| 桐乡市| 廊坊市| 普安县| 陕西省| 胶南市| 宜丰县| 广宗县| 昌宁县| 平果县| 武宣县| 郴州市| 万载县| 油尖旺区| 拉萨市| 盖州市| 潼关县| 弥勒县| 鄂尔多斯市| 成武县| 宣城市| 东丽区| 铜鼓县| 中宁县|