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

Java中RabbitMQ延遲隊列實現(xiàn)詳解

 更新時間:2023年09月20日 10:13:11   作者:CD4356  
這篇文章主要介紹了Java中RabbitMQ延遲隊列實現(xiàn)詳解,消息過期后,根據(jù)routing-key的不同,又會被死信交換機路由到不同的死信隊列中,消費者只需要監(jiān)聽對應的死信隊列進行消費即可,需要的朋友可以參考下

一、RabbitMQ延遲隊列實現(xiàn)

1.1、RabbitMQ延遲隊列實現(xiàn)流程

cd

  1. 生產(chǎn)者生產(chǎn)一條延遲消息,根據(jù)延遲時間的不同,利用不同的routing-key將消息路由到不同的延遲隊列,每個隊列都設置了不同的 TTL 屬性 ( TTL ( Time To Live ) 生存時間 ),并綁定到同一個死信交換機中。
  2. 消息過期后,根據(jù)routing-key的不同,又會被死信交換機路由到不同的死信隊列中,消費者只需要監(jiān)聽對應的死信隊列進行消費即可。

1.2、配置RabbitMQ連接

#[ RabbitMQ相關配置 ]
#rabbitmq服務器IP
spring.rabbitmq.host=安裝RabbitMQ的服務器IP
#rabbitmq服務器端口(默認為5672)
spring.rabbitmq.port=5672
#用戶名
spring.rabbitmq.username=guest
#用戶密碼
spring.rabbitmq.password=guest
#虛擬主機(一個RabbitMQ服務可以配置多個虛擬主機,每一個虛擬機主機之間是相互隔離,相互獨立的,授權用戶到指定的virtual-host就可以發(fā)送消息到指定隊列)
#vhost虛擬主機地址( 默認為/ )
spring.rabbitmq.virtual-host=/

1.3、創(chuàng)建配置類

配置兩個交換機、四個隊列、以及根據(jù)路由鍵配置交換機和隊列的綁定關系

import org.springframework.amqp.core.*;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import java.util.HashMap;
import java.util.Map;
@Configuration
public class RabbitMQConfiguration {
    //延遲交換機
    public static final String DELAY_EXCHANGE = "delay_exchange";
    //延遲隊列A
    public static final String DELAY_QUEUE_A = "delay_queue_a";
    //延遲隊列B
    public static final String DELAY_QUEUE_B = "delay_queue_b";
    //延遲路由鍵10S
    public static final String DELAY_QUEUE_10S_ROUTING_KEY = "delay_queue_10s_routing_key";
    //延遲路由鍵60S
    public static final String DELAY_QUEUE_60S_ROUTING_KEY = "delay_queue_60s_routing_key";
    //死信交換機
    public static final String DEAD_LETTER_EXCHANGE = "dead_letter_exchange";
    //死信隊列A
    public static final String DEAD_LETTER_QUEUE_A = "dead_letter_queue_a";
    //死信隊列B
    public static final String DEAD_LETTER_QUEUE_B = "dead_letter_queue_b";
    //死信路由鍵10S
    public static final String DEAD_LETTER_QUEUE_10S_ROUTING_KEY = "dead_letter_queue_10s_routing_key";
    //死信路由鍵60S
    public static final String DEAD_LETTER_QUEUE_60S_ROUTING_KEY = "dead_letter_queue_60s_routing_key";
    //延遲交換機
    @Bean("delayExchange")
    public DirectExchange delayExchange(){
        return new DirectExchange(DELAY_EXCHANGE, true, false);
    }
    //延遲隊列A
    @Bean("delayQueueA")
    public Queue delayQueueA(){
        Map<String, Object> args = new HashMap<>();
        //設置延遲隊列綁定的死信交換機
        args.put("x-dead-letter-exchange", DEAD_LETTER_EXCHANGE);
        //設置延遲隊列綁定的死信路由鍵
        args.put("x-dead-letter-routing-key", DEAD_LETTER_QUEUE_10S_ROUTING_KEY);
        //設置延遲隊列的 TTL 消息存活時間
        args.put("x-message-ttl", 10*1000);
        return new Queue(DELAY_QUEUE_A, true, false, false, args);
    }
    //延遲隊列B
    @Bean("delayQueueB")
    public Queue delayQueueB(){
        Map<String, Object> args = new HashMap<>();
        //設置延遲隊列綁定的死信交換機
        args.put("x-dead-letter-exchange", DEAD_LETTER_EXCHANGE);
        //設置延遲隊列綁定的死信路由鍵
        args.put("x-dead-letter-routing-key", DEAD_LETTER_QUEUE_60S_ROUTING_KEY);
        //設置延遲隊列的 TTL 消息存活時間
        args.put("x-message-ttl", 60*1000);
        return new Queue(DELAY_QUEUE_B, true, false, false, args);
    }
    //延遲隊列A的綁定關系
    @Bean("delayBindingA")
    public Binding delayBindingA(@Qualifier("delayQueueA")Queue queue,
                                 @Qualifier("delayExchange")DirectExchange exchange){
        return BindingBuilder.bind(queue).to(exchange).with(DELAY_QUEUE_10S_ROUTING_KEY);
    }
    //延遲隊列B的綁定關系
    @Bean("delayBindingB")
    public Binding delayBindingB(@Qualifier("delayQueueB")Queue queue,
                                 @Qualifier("delayExchange")DirectExchange exchange){
        return BindingBuilder.bind(queue).to(exchange).with(DELAY_QUEUE_60S_ROUTING_KEY);
    }
    //死信交換機
    @Bean("deadLetterExchange")
    public DirectExchange deadLetterExchange(){
        return new DirectExchange(DEAD_LETTER_EXCHANGE, true, false);
    }
    //死信隊列A
    @Bean("deadLetterQueueA")
    public Queue deadLetterQueueA(){
        return new Queue(DEAD_LETTER_QUEUE_A, true);
    }
    //死信隊列B
    @Bean("deadLetterQueueB")
    public Queue deadLetterQueueB(){
        return new Queue(DEAD_LETTER_QUEUE_B, true);
    }
    //死信隊列A的綁定關系
    @Bean("deadLetterBindingA")
    public Binding deadLetterBindingA(@Qualifier("deadLetterQueueA")Queue queue,
                                 @Qualifier("deadLetterExchange")DirectExchange exchange){
        return BindingBuilder.bind(queue).to(exchange).with(DEAD_LETTER_QUEUE_10S_ROUTING_KEY);
    }
    //死信隊列B的綁定關系
    @Bean("deadLetterBindingB")
    public Binding deadLetterBindingB(@Qualifier("deadLetterQueueB")Queue queue,
                                      @Qualifier("deadLetterExchange")DirectExchange exchange){
        return BindingBuilder.bind(queue).to(exchange).with(DEAD_LETTER_QUEUE_60S_ROUTING_KEY);
    }
}

1.4、創(chuàng)建一個枚舉類來配置延遲類型

@Getter
@AllArgsConstructor
public enum DelayTypeEnum {
    //10s
    DELAY_10s(1),
    //60s
    DELAY_60s(2);
    private Integer type;
    /**
     * 延遲類型
     * @param type
     * @return 延遲類型
     */
    public static DelayTypeEnum getDelayTypeEnum(Integer type){
        if(Objects.equals(type, DELAY_10s.type)){
            return DELAY_10s;
        }
        if(Objects.equals(type, DELAY_60s.type)){
            return DELAY_60s;
        }
        return null;
    }
}

1.5、創(chuàng)建生產(chǎn)者類發(fā)送消息

import com.cd.springbootrabbitmq.enums.DelayTypeEnum;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import static com.cd.springbootrabbitmq.config.RabbitMQConfiguration.DELAY_EXCHANGE;
import static com.cd.springbootrabbitmq.config.RabbitMQConfiguration.DELAY_QUEUE_10S_ROUTING_KEY;
import static com.cd.springbootrabbitmq.config.RabbitMQConfiguration.DELAY_QUEUE_60S_ROUTING_KEY;
/**
 * 延遲消息生產(chǎn)者
 */
@Component
public class DelayMessageProducer {
    @Autowired
    private RabbitTemplate rabbitTemplate;
    /**
     * 發(fā)送延遲消息
     * @param message  要發(fā)送的消息
     * @param type  延遲類型(延時10s的延遲隊列 或 延時60s的延遲隊列)
     */
    public void sendDelayMessage(String message, DelayTypeEnum type){
        switch (type){
            case DELAY_10s:
                rabbitTemplate.convertAndSend(DELAY_EXCHANGE, DELAY_QUEUE_10S_ROUTING_KEY, message);
                break;
            case DELAY_60s:
                rabbitTemplate.convertAndSend(DELAY_EXCHANGE, DELAY_QUEUE_60S_ROUTING_KEY, message);
                break;
            default:
                break;
        }
    }
}

1.6、創(chuàng)建消費者類消費消息

import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import java.time.LocalDateTime;
import static com.cd.springbootrabbitmq.config.RabbitMQConfiguration.DEAD_LETTER_QUEUE_A;
import static com.cd.springbootrabbitmq.config.RabbitMQConfiguration.DEAD_LETTER_QUEUE_B;
@Slf4j
@Component
public class DeadLetterQueueConsumer {
    @Autowired
    private RabbitTemplate rabbitTemplate;
    /**
     * 監(jiān)聽死信隊列A
     * @param message  接收的信息
     */
    //@RabbitListener(queues = "dead_letter_queue_a")
    @RabbitListener(queues = DEAD_LETTER_QUEUE_A)
    public void receiveA(Message message) {
        String msg = new String(message.getBody());
        // 記錄日志
        log.info("當前時間:{},死信隊列A收到的消息:{}", LocalDateTime.now(), msg);
    }
    /**
     * 監(jiān)聽死信隊列B
     * @param message  接收的信息
     */
    //@RabbitListener(queues = "dead_letter_queue_b")
    @RabbitListener(queues = DEAD_LETTER_QUEUE_B)
    public void receiveB(Message message){
        String msg = new String(message.getBody());
        // 記錄日志
        log.info("當前時間:{},死信隊列B收到的消息:{}", LocalDateTime.now(), msg);
    }
}

1.7、創(chuàng)建控制類

import com.cd.springbootrabbitmq.enums.DelayTypeEnum;
import com.cd.springbootrabbitmq.producer.DelayMessageProducer;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import java.time.LocalDateTime;
import java.util.Objects;
@Slf4j
@RestController
@RequestMapping("/rabbitmq")
public class RabbitMQController {
    @Autowired
    private DelayMessageProducer producer;
    @RequestMapping("/send")
    public void send(String message, Integer delayType){
        // 記錄日志
        log.info("當前時間:{},消息:{},延遲類型:{}", LocalDateTime.now(), message, delayType);
        // 發(fā)送延遲消息
        producer.sendDelayMessage(message, Objects.requireNonNull(DelayTypeEnum.getDelayTypeEnum(delayType)));
    }
}

1.8、測試

在瀏覽器中先后提交下面兩個請求:

1)localhost:8080/rabbitmq/send?message=測試自定義延遲處理60s&delayType=2

2)localhost:8080/rabbitmq/send?message=測試自定義延遲處理10s&delayType=1

查看idea控制臺:

cd

到此這篇關于Java中RabbitMQ延遲隊列實現(xiàn)詳解的文章就介紹到這了,更多相關RabbitMQ延遲隊列內容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關文章希望大家以后多多支持腳本之家!

相關文章

  • RocketMQ延遲消息超詳細講解

    RocketMQ延遲消息超詳細講解

    延時消息是指發(fā)送到 RocketMQ 后不會馬上被消費者拉取到,而是等待固定的時間,才能被消費者拉取到。延時消息的使用場景很多,比如電商場景下關閉超時未支付的訂單,某些場景下需要在固定時間后發(fā)送提示消息
    2023-02-02
  • Java線程池如何實現(xiàn)精準控制每秒API請求

    Java線程池如何實現(xiàn)精準控制每秒API請求

    這篇文章主要介紹了Java線程池如何實現(xiàn)精準控制每秒API請求問題,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教
    2024-08-08
  • 深入解析Java編程中final關鍵字的作用

    深入解析Java編程中final關鍵字的作用

    final關鍵字正如其字面意思一樣,意味著最后,比如被final修飾后類不能集成、變量不能被再賦值等,以下我們就來深入解析Java編程中final關鍵字的作用:
    2016-06-06
  • Spring Cloud之服務監(jiān)控turbine的示例

    Spring Cloud之服務監(jiān)控turbine的示例

    這篇文章主要介紹了Spring Cloud之服務監(jiān)控turbine的示例,小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2018-05-05
  • SpringBoot集成easy-rules規(guī)則引擎流程詳解

    SpringBoot集成easy-rules規(guī)則引擎流程詳解

    這篇文章主要介紹了SpringBoot集成easy-rules規(guī)則引擎流程,合理的使用規(guī)則引擎可以極大的減少代碼復雜度,提升代碼可維護性。業(yè)界知名的開源規(guī)則引擎有Drools,功能豐富,但也比較龐大
    2023-03-03
  • Spring 靜態(tài)變量/構造函數(shù)注入失敗的解決方案

    Spring 靜態(tài)變量/構造函數(shù)注入失敗的解決方案

    我們經(jīng)常會遇到一下問題:Spring對靜態(tài)變量的注入為空、在構造函數(shù)中使用Spring容器中的Bean對象,得到的結果為空。不要擔心,本文將為大家介紹如何解決這些問題,跟隨小編來看看吧
    2021-11-11
  • java 遍歷Map及Map轉化為二維數(shù)組的實例

    java 遍歷Map及Map轉化為二維數(shù)組的實例

    這篇文章主要介紹了java 遍歷Map及Map轉化為二維數(shù)組的實例的相關資料,希望通過本文能幫助到大家,實現(xiàn)這樣的功能,需要的朋友可以參考下
    2017-08-08
  • JavaIO模型中的BIO,NIO和AIO詳解

    JavaIO模型中的BIO,NIO和AIO詳解

    這篇文章主要為大家詳細介紹了JavaIO模型中的BIO,NIO和AIO,文中示例代碼介紹的非常詳細,具有一定的參考價值,感興趣的小伙伴們可以參考一下,希望能夠給你帶來幫助
    2022-02-02
  • Echarts+SpringMvc顯示后臺實時數(shù)據(jù)

    Echarts+SpringMvc顯示后臺實時數(shù)據(jù)

    這篇文章主要為大家詳細介紹了Echarts+SpringMvc顯示后臺實時數(shù)據(jù),文中示例代碼介紹的非常詳細,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2019-12-12
  • SpringBoot中的條件注解使用示例詳解

    SpringBoot中的條件注解使用示例詳解

    SpringBoot條件注解用于動態(tài)控制Bean創(chuàng)建與配置加載,基于@Conditional機制,支持按類、Bean、屬性等條件判斷,廣泛應用于多數(shù)據(jù)源等場景,提升應用靈活性與智能化,接下來通過本文給大家講解SpringBoot中的條件注解使用,感興趣的朋友一起看看吧
    2025-08-08

最新評論

黄浦区| 万源市| SHOW| 曲阳县| 会宁县| 房山区| 阳信县| 颍上县| 临夏市| 曲阜市| 沙田区| 辽阳市| 丹棱县| 来安县| 疏勒县| 方山县| 南木林县| 沂南县| 贵阳市| 和硕县| 横山县| 武安市| 青神县| 乐平市| 喀喇| 墨玉县| 汕尾市| 文昌市| 沙洋县| 芜湖市| 叶城县| 龙口市| 繁峙县| 昭通市| 万年县| 固安县| 南木林县| 洪泽县| 日土县| 莲花县| 太康县|