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

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

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

一、RabbitMQ延遲隊(duì)列實(shí)現(xiàn)

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

cd

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

1.2、配置RabbitMQ連接

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

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

配置兩個(gè)交換機(jī)、四個(gè)隊(duì)列、以及根據(jù)路由鍵配置交換機(jī)和隊(duì)列的綁定關(guān)系

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 {
    //延遲交換機(jī)
    public static final String DELAY_EXCHANGE = "delay_exchange";
    //延遲隊(duì)列A
    public static final String DELAY_QUEUE_A = "delay_queue_a";
    //延遲隊(duì)列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";
    //死信交換機(jī)
    public static final String DEAD_LETTER_EXCHANGE = "dead_letter_exchange";
    //死信隊(duì)列A
    public static final String DEAD_LETTER_QUEUE_A = "dead_letter_queue_a";
    //死信隊(duì)列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";
    //延遲交換機(jī)
    @Bean("delayExchange")
    public DirectExchange delayExchange(){
        return new DirectExchange(DELAY_EXCHANGE, true, false);
    }
    //延遲隊(duì)列A
    @Bean("delayQueueA")
    public Queue delayQueueA(){
        Map<String, Object> args = new HashMap<>();
        //設(shè)置延遲隊(duì)列綁定的死信交換機(jī)
        args.put("x-dead-letter-exchange", DEAD_LETTER_EXCHANGE);
        //設(shè)置延遲隊(duì)列綁定的死信路由鍵
        args.put("x-dead-letter-routing-key", DEAD_LETTER_QUEUE_10S_ROUTING_KEY);
        //設(shè)置延遲隊(duì)列的 TTL 消息存活時(shí)間
        args.put("x-message-ttl", 10*1000);
        return new Queue(DELAY_QUEUE_A, true, false, false, args);
    }
    //延遲隊(duì)列B
    @Bean("delayQueueB")
    public Queue delayQueueB(){
        Map<String, Object> args = new HashMap<>();
        //設(shè)置延遲隊(duì)列綁定的死信交換機(jī)
        args.put("x-dead-letter-exchange", DEAD_LETTER_EXCHANGE);
        //設(shè)置延遲隊(duì)列綁定的死信路由鍵
        args.put("x-dead-letter-routing-key", DEAD_LETTER_QUEUE_60S_ROUTING_KEY);
        //設(shè)置延遲隊(duì)列的 TTL 消息存活時(shí)間
        args.put("x-message-ttl", 60*1000);
        return new Queue(DELAY_QUEUE_B, true, false, false, args);
    }
    //延遲隊(duì)列A的綁定關(guān)系
    @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);
    }
    //延遲隊(duì)列B的綁定關(guān)系
    @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);
    }
    //死信交換機(jī)
    @Bean("deadLetterExchange")
    public DirectExchange deadLetterExchange(){
        return new DirectExchange(DEAD_LETTER_EXCHANGE, true, false);
    }
    //死信隊(duì)列A
    @Bean("deadLetterQueueA")
    public Queue deadLetterQueueA(){
        return new Queue(DEAD_LETTER_QUEUE_A, true);
    }
    //死信隊(duì)列B
    @Bean("deadLetterQueueB")
    public Queue deadLetterQueueB(){
        return new Queue(DEAD_LETTER_QUEUE_B, true);
    }
    //死信隊(duì)列A的綁定關(guān)系
    @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);
    }
    //死信隊(duì)列B的綁定關(guān)系
    @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)建一個(gè)枚舉類來配置延遲類型

@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  延遲類型(延時(shí)10s的延遲隊(duì)列 或 延時(shí)60s的延遲隊(duì)列)
     */
    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)建消費(fèi)者類消費(fèi)消息

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)聽死信隊(duì)列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("當(dāng)前時(shí)間:{},死信隊(duì)列A收到的消息:{}", LocalDateTime.now(), msg);
    }
    /**
     * 監(jiān)聽死信隊(duì)列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("當(dāng)前時(shí)間:{},死信隊(duì)列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("當(dāng)前時(shí)間:{},消息:{},延遲類型:{}", LocalDateTime.now(), message, delayType);
        // 發(fā)送延遲消息
        producer.sendDelayMessage(message, Objects.requireNonNull(DelayTypeEnum.getDelayTypeEnum(delayType)));
    }
}

1.8、測(cè)試

在瀏覽器中先后提交下面兩個(gè)請(qǐng)求:

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

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

查看idea控制臺(tái):

cd

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

相關(guān)文章

  • RocketMQ延遲消息超詳細(xì)講解

    RocketMQ延遲消息超詳細(xì)講解

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

    Java線程池如何實(shí)現(xiàn)精準(zhǔn)控制每秒API請(qǐng)求

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

    深入解析Java編程中final關(guān)鍵字的作用

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

最新評(píng)論

灵川县| 阿拉善左旗| 荣成市| 平罗县| 两当县| 克拉玛依市| 呼伦贝尔市| 崇信县| 建湖县| 内江市| 青岛市| 蒲城县| 康保县| 九江县| 高要市| 资源县| 罗甸县| 金沙县| 民勤县| 卢龙县| 广河县| 巴中市| 万源市| 尉氏县| 休宁县| 娱乐| 嘉义市| 巧家县| 称多县| 自治县| 鹤庆县| 乐昌市| 平陆县| 屏山县| 邢台县| 博罗县| 汤原县| 桑日县| 孟津县| 榆中县| 凌海市|