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

RabbitMQ中的延遲隊(duì)列機(jī)制詳解

 更新時間:2023年09月20日 11:10:39   作者:煎丶包  
這篇文章主要介紹了RabbitMQ中的延遲隊(duì)列機(jī)制詳解,延時隊(duì)列內(nèi)部是有序的,最重要的特性就體現(xiàn)在它的延時屬性上,延時隊(duì)列中的元素是希望,在指定時間到了以后或之前取出和處理,簡單來說,延時隊(duì)列就是用來存放需要在指定時間被處理的元素的隊(duì)列,需要的朋友可以參考下

一、延遲隊(duì)列

延時隊(duì)列內(nèi)部是有序的,最重要的特性就體現(xiàn)在它的延時屬性上,延時隊(duì)列中的元素是希望 在指定時間到了以后或之前取出和處理,簡單來說,延時隊(duì)列就是用來存放需要在指定時間被處理的元素的隊(duì)列。

二、隊(duì)列TTL

在這里插入圖片描述

創(chuàng)建一個配置類,聲明并配置交換機(jī)和隊(duì)列

@Configuration
public class TtlQueueConfig {
    //普通交換機(jī)名稱
    public static final String NORMAL_EXCHANGE = "X";
    //死信交換機(jī)名稱
    public static final String DEAD_EXCHANGE = "Y";
    //普通隊(duì)列名稱
    public static final String NORMAL_QUEUE_A = "QA";
    public static final String NORMAL_QUEUE_B = "QA";
    //死信隊(duì)列名稱
    public static final String DEAD_QUEUE = "QD";
    //聲明普通交換機(jī)
    @Bean("xExchange")
    public DirectExchange xExchange() {
        return new DirectExchange(NORMAL_EXCHANGE);
    }
    //聲明死信交換機(jī)
    @Bean("yExchange")
    public DirectExchange yExchange() {
        return new DirectExchange(DEAD_EXCHANGE);
    }
    //聲明普通隊(duì)列,TTL為10s
    @Bean("QA")
    public Queue qA() {
        Map<String, Object> arguments = new HashMap<>();
        //設(shè)置死信交換機(jī)
        arguments.put("x-dead-letter-exchange", DEAD_EXCHANGE);
        //設(shè)置死信RoutingKey
        arguments.put("x-dead-letter-routing-key", "YD");
        //設(shè)置TTL
        arguments.put("x-message-ttl", 10000);
        return QueueBuilder.durable(NORMAL_QUEUE_A).withArguments(arguments).build();
    }
    //聲明普通隊(duì)列,TTL為10s
    @Bean("QB")
    public Queue qB() {
        Map<String, Object> arguments = new HashMap<>();
        //設(shè)置死信交換機(jī)
        arguments.put("x-dead-letter-exchange", DEAD_EXCHANGE);
        //設(shè)置死信RoutingKey
        arguments.put("x-dead-letter-routing-key", "YD");
        //設(shè)置TTL
        arguments.put("x-message-ttl", 40000);
        return QueueBuilder.durable(NORMAL_QUEUE_B).withArguments(arguments).build();
    }
    //聲明死信隊(duì)列
    @Bean("QD")
    public Queue qD() {
        return QueueBuilder.durable(DEAD_QUEUE).build();
    }
    //綁定對應(yīng)的交換機(jī)和隊(duì)列
    @Bean
    public Binding queueABindingX(@Qualifier("QA") Queue QA,
                                  @Qualifier("xExchange") DirectExchange xExchange) {
        return BindingBuilder.bind(QA).to(xExchange).with("XA");
    }
    @Bean
    public Binding queueBBindingX(@Qualifier("QB") Queue QB,
                                  @Qualifier("xExchange") DirectExchange xExchange) {
        return BindingBuilder.bind(QB).to(xExchange).with("XB");
    }
    @Bean
    public Binding queueDBindingY(@Qualifier("QD") Queue QD,
                                  @Qualifier("yExchange") DirectExchange yExchange) {
        return BindingBuilder.bind(QD).to(yExchange).with("YD");
    }
}

創(chuàng)建一個生產(chǎn)者

@Slf4j
@RestController
public class SendMsgController {
    @Autowired
    RabbitTemplate rabbitTemplate;
    @GetMapping("/ttl/sendMsg/{message}")
    public void sendMsg(@PathVariable String message) {
        log.info("當(dāng)前時間:{}, 發(fā)送一條信息給兩個隊(duì)列:{}", new Date().toString(), message);
        rabbitTemplate.convertAndSend("X","XA","消息來自TTL為10s的隊(duì)列QA:" + message);
        rabbitTemplate.convertAndSend("X","XB","消息來自TTL為40s的隊(duì)列QB:" + message);
    }
}

創(chuàng)建一個消費(fèi)者

@Slf4j
@Component
public class DeadLetterQueueConsumer {
    //接收消息
    @RabbitListener(queues = "QD")
    public void receivedQD(Message message, Channel channel) {
        String msg = new String(message.getBody());
        log.info("當(dāng)前時間:{}, 收到死信隊(duì)列的消息:{}", new Date().toString(), message);
    }
}

瀏覽器發(fā)送消息

在這里插入圖片描述

消費(fèi)者分別過了10s和40s接收到了消息

在這里插入圖片描述

三、延遲隊(duì)列的優(yōu)化

不同的延遲時間需要設(shè)置不同的 TTL ,可以優(yōu)化聲明一個通用的 QC 隊(duì)列,具體的延遲時間有生產(chǎn)者決定

在這里插入圖片描述

在配置類 TtlQueueConfig 中配置通用隊(duì)列 QC

    //通用隊(duì)列名稱
    public static final String Generic_QUEUE_C = "QC";
    //聲明通用隊(duì)列
    @Bean("QC")
    public Queue qC() {
        Map<String, Object> arguments = new HashMap<>();
        //設(shè)置死信交換機(jī)
        arguments.put("x-dead-letter-exchange", DEAD_EXCHANGE);
        //設(shè)置死信RoutingKey
        arguments.put("x-dead-letter-routing-key", "YD");
        //因?yàn)槭峭ㄓ藐?duì)列,所以不設(shè)置TTL,由生產(chǎn)者指定消息的TTL
        return QueueBuilder.durable(Generic_QUEUE_C).withArguments(arguments).build();
    }
    //綁定通用隊(duì)列和普通交換機(jī)
    @Bean
    public Binding queueCBindingX(@Qualifier("QC") Queue QC,
                                  @Qualifier("xExchange") DirectExchange xExchange) {
        return BindingBuilder.bind(QC).to(xExchange).with("XC");
    }
    //綁定通用隊(duì)列和死信交換機(jī)
    @Bean
    public Binding queueCBindingY(@Qualifier("QC") Queue QC,
                                  @Qualifier("yExchange") DirectExchange yExchange) {
        return BindingBuilder.bind(QC).to(yExchange).with("YD");
    }

生產(chǎn)者發(fā)送消息,并指定 TTL 時長

    //發(fā)送消息,并指定消息的TTL
    @GetMapping("/ttl/sendExpirationMsg/{message}/{ttlTime}")
    public void sendExpirationMsg(@PathVariable("message") String message, @PathVariable("ttlTime") String ttlTime) {
        log.info("當(dāng)前時間:{}, 發(fā)送一條TTL為{}ms的消息給隊(duì)列QC:{}", new Date().toString(), ttlTime, message);
        rabbitTemplate.convertAndSend("X", "XC", message, msg -> {
            //設(shè)置消息的TTL時長
            msg.getMessageProperties().setExpiration(ttlTime);
            return msg;
        });
    }

發(fā)送兩條消息

在這里插入圖片描述

在這里插入圖片描述

消費(fèi)者接收消息

在這里插入圖片描述

但是,如果連續(xù)發(fā)送兩條消息,如果使用在消息屬性上設(shè)置 TTL 的方式,消息可能并不會按時“死亡“,因?yàn)?RabbitMQ 只會檢查第一個消息是否過期,如果過期則丟到死信隊(duì)列,如果第一個消息的延時時長很長,而第二個消息的延時時長很短,第二個消息并不會優(yōu)先得到執(zhí)行。結(jié)果會導(dǎo)致第二條消息消費(fèi)者收到時間有誤。

在這里插入圖片描述

四、基于 RabbitMQ 插件實(shí)現(xiàn)延遲隊(duì)列

如果不能實(shí)現(xiàn)在消息粒度上的 TTL ,并使其在設(shè)置的 TTL 時間及時死亡,就無法設(shè)計(jì)成一個通用的延時隊(duì)列??梢允褂没?RabbitMQ 插件來實(shí)現(xiàn)延遲隊(duì)列,從而解決這個問題。

基于 RabbitMQ 插件實(shí)現(xiàn)延遲,是交換機(jī)實(shí)現(xiàn)延遲,而不再是隊(duì)列實(shí)現(xiàn)延遲

在這里插入圖片描述

在這里插入圖片描述

創(chuàng)建一個基于插件的延遲隊(duì)列配置類 DelayedQueueConfig

@Configuration
public class DelayedQueueConfig {
    //交換機(jī)名稱
    public static final String DELAYED_EXCHANGE_NAME = "delayed_exchange";
    //隊(duì)列名稱
    public static final String DELAYED_QUEUE_NAME = "delayed_queue";
    //routingKey
    public static final String DELAYED_ROUTING_KEY = "delayed_routingKey";
    //聲明交換機(jī)
    @Bean
    public CustomExchange delayedExchange() {
        Map<String, Object> arguments = new HashMap<>();
        arguments.put("x-delayed-type", "direct");  //設(shè)置延遲類型
        return new CustomExchange(DELAYED_EXCHANGE_NAME, "x-delayed-message", true, false, arguments);
    }
    //聲明隊(duì)列
    @Bean
    public Queue delayedQueue() {
        return new Queue(DELAYED_QUEUE_NAME);
    }
    //綁定隊(duì)列和交換機(jī)
    @Bean
    public Binding delayedQueueBindingDelayedExchange(@Qualifier("delayedQueue") Queue delayedQueue,
                                                      @Qualifier("delayedExchange") CustomExchange delayedExchange) {
        return BindingBuilder.bind(delayedQueue).to(delayedExchange).with(DELAYED_ROUTING_KEY).noargs();
    }
}

創(chuàng)建生產(chǎn)者發(fā)送延遲消息

    //基于插件發(fā)送消息
    @GetMapping("/ttl/sendDelayedMsg/{message}/{delayedTime}")
    public void sendDelayedMsg(@PathVariable("message") String message, @PathVariable("delayedTime") Integer delayedTime) {
        log.info("當(dāng)前時間:{}, 發(fā)送一條時長為{}ms的消息給延遲隊(duì)列delayed_queue:{}", new Date().toString(), delayedTime, message);
        rabbitTemplate.convertAndSend(DelayedQueueConfig.DELAYED_EXCHANGE_NAME, DelayedQueueConfig.DELAYED_ROUTING_KEY, message, msg -> {
            //設(shè)置消息的延遲時長
            msg.getMessageProperties().setDelay(delayedTime);
            return msg;
        });
    }

創(chuàng)建消費(fèi)者

@Slf4j
@Component
public class DelayedQueueConsumer {
    //監(jiān)聽消息
    @RabbitListener(queues = {DelayedQueueConfig.DELAYED_QUEUE_NAME})
    public void receiveDelayQueue(Message message){
        String msg = new String(message.getBody());
        log.info("當(dāng)前時間:{}, 收到延遲隊(duì)列的消息:{}", new Date().toString(), msg);
    }
}

當(dāng)連續(xù)發(fā)送兩條不同延遲時長的消息時,消費(fèi)者會先接收到延遲時長短的那條消息,再接收延遲時長長的那條消息。

在這里插入圖片描述

實(shí)現(xiàn)延遲隊(duì)列,一種是基于死信隊(duì)列的方式,一種是基于RabbitMQ插件的方式。

延時隊(duì)列在需要延時處理的場景下非常有用,使用 RabbitMQ 來實(shí)現(xiàn)延時隊(duì)列可以很好的利用RabbitMQ 的特性,如消息可靠發(fā)送、消息可靠投遞、死信隊(duì)列來保障消息至少被消費(fèi)一次以及未被正確處理的消息不會被丟棄。另外,通過 RabbitMQ 集群的特性,可以很好的解決單點(diǎn)故障問題,不會因?yàn)閱蝹€節(jié)點(diǎn)掛掉導(dǎo)致延時隊(duì)列不可用或者消息丟失。

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

相關(guān)文章

  • java組件SmartUpload和FileUpload實(shí)現(xiàn)文件上傳功能

    java組件SmartUpload和FileUpload實(shí)現(xiàn)文件上傳功能

    這篇文章主要為大家詳細(xì)介紹了java組件SmartUpload和FileUpload實(shí)現(xiàn)文件上傳功能,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2017-11-11
  • 使用Spring Boot創(chuàng)建Web應(yīng)用程序的示例代碼

    使用Spring Boot創(chuàng)建Web應(yīng)用程序的示例代碼

    本篇文章主要介紹了使用Spring Boot創(chuàng)建Web應(yīng)用程序的示例代碼,我們將使用Spring Boot構(gòu)建一個簡單的Web應(yīng)用程序,并為其添加一些有用的服務(wù),小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2018-05-05
  • 解析ConcurrentHashMap: get、remove方法分析

    解析ConcurrentHashMap: get、remove方法分析

    ConcurrentHashMap是由Segment數(shù)組結(jié)構(gòu)和HashEntry數(shù)組結(jié)構(gòu)組成。Segment的結(jié)構(gòu)和HashMap類似,是一種數(shù)組和鏈表結(jié)構(gòu),今天給大家普及java面試常見問題---ConcurrentHashMap知識,一起看看吧
    2021-06-06
  • 詳談Spring是否支持對靜態(tài)方法進(jìn)行Aop增強(qiáng)

    詳談Spring是否支持對靜態(tài)方法進(jìn)行Aop增強(qiáng)

    這篇文章主要介紹了Spring是否支持對靜態(tài)方法進(jìn)行Aop增強(qiáng),具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-12-12
  • Java中對List去重 Stream去重的解決方法

    Java中對List去重 Stream去重的解決方法

    這篇文章主要介紹了Java中對List去重, Stream去重的問題解答,文中給大家介紹了Java中List集合去除重復(fù)數(shù)據(jù)的方法,需要的朋友可以參考下
    2018-04-04
  • Springcloud+Mybatis使用多數(shù)據(jù)源的四種方式(小結(jié))

    Springcloud+Mybatis使用多數(shù)據(jù)源的四種方式(小結(jié))

    這篇文章主要介紹了Springcloud+Mybatis使用多數(shù)據(jù)源的四種方式,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2020-09-09
  • 簡單講解奇偶排序算法及在Java數(shù)組中的實(shí)現(xiàn)

    簡單講解奇偶排序算法及在Java數(shù)組中的實(shí)現(xiàn)

    這篇文章主要介紹了奇偶排序算法及Java數(shù)組的實(shí)現(xiàn),奇偶排序的時間復(fù)雜度為O(N^2),需要的朋友可以參考下
    2016-04-04
  • 詳解Java中的泛型

    詳解Java中的泛型

    這篇文章主要介紹了Java中的泛型,當(dāng)我們不確定數(shù)據(jù)類型時,我們可以暫時使用一個字母 T代替數(shù)據(jù)類型,例如寫一個方法,但是我們不知道它是傳遞的是什么數(shù)據(jù)類型,我們就可以使用泛型,到時候只要指明T是什么數(shù)據(jù)類型,就可以使用了,需要的朋友可以參考下
    2023-05-05
  • java實(shí)現(xiàn)讀取txt文件中的內(nèi)容

    java實(shí)現(xiàn)讀取txt文件中的內(nèi)容

    本文通過一個具體的例子向大家展示了如何使用java實(shí)現(xiàn)讀取TXT文件里的內(nèi)容的方法以及思路,有需要的小伙伴可以參考下
    2016-03-03
  • Java判斷字符串是否含有亂碼實(shí)例代碼

    Java判斷字符串是否含有亂碼實(shí)例代碼

    本文通過實(shí)例代碼給大家介紹了Java判斷字符串是否含有亂碼的方法,代碼簡單易懂,非常不錯,具有一定的參考借鑒價值,需要的朋友參考下吧
    2018-11-11

最新評論

南乐县| 抚远县| 浮山县| 临江市| 阿城市| 左贡县| 台北县| 伊金霍洛旗| 新疆| 垣曲县| 武鸣县| 嘉峪关市| 杂多县| 栾城县| 中超| 澎湖县| 吉木萨尔县| 长泰县| 江阴市| 沧州市| 静宁县| 富民县| 鄢陵县| 山丹县| 奉贤区| 静宁县| 乡宁县| 会昌县| 二手房| 安达市| 朝阳县| 凤山市| 井陉县| 新疆| 巴塘县| 县级市| 长泰县| 沂水县| 海晏县| 措美县| 成都市|