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

springcloud中RabbitMQ死信隊(duì)列與延遲交換機(jī)實(shí)現(xiàn)方法

 更新時(shí)間:2022年05月31日 09:58:08   作者:wu55555  
死信隊(duì)列是消息隊(duì)列中非常重要的概念,同時(shí)我們需要業(yè)務(wù)場(chǎng)景中都需要延遲發(fā)送的概念,比如12306中的30分鐘后未支付訂單取消,那么本期,我們就來(lái)講解死信隊(duì)列,以及如何通過(guò)延遲交換機(jī)來(lái)實(shí)現(xiàn)延遲發(fā)送的需求,感興趣的朋友一起看看吧

0.引言

死信隊(duì)列是消息隊(duì)列中非常重要的概念,同時(shí)我們需要業(yè)務(wù)場(chǎng)景中都需要延遲發(fā)送的概念,比如12306中的30分鐘后未支付訂單取消。那么本期,我們就來(lái)講解死信隊(duì)列,以及如何通過(guò)延遲交換機(jī)來(lái)實(shí)現(xiàn)延遲發(fā)送的需求。

1. 死信隊(duì)列

1.2 什么是死信?

理解死信隊(duì)列前,我們先講解什么是死信,所謂死信就是沒(méi)有被成功消費(fèi)的消息,但并不是所有未成功消費(fèi)的消息都是死信消息,死信消息的產(chǎn)生來(lái)源于以下三種途徑: (1)消息被消費(fèi)者拒絕,參數(shù)requeue設(shè)置為false的消息 (2)過(guò)期的消息,過(guò)期消息分為兩種: a. 發(fā)送消息時(shí),設(shè)置了某一條消息的生存時(shí)間(message TTL),如果生存時(shí)間到了,消息還沒(méi)有被消費(fèi),就會(huì)被標(biāo)注為死信消息 b. 設(shè)置了隊(duì)列的消息生存時(shí)間,針對(duì)隊(duì)列中所有的消息,如果生存時(shí)間到了,消息還沒(méi)有被消費(fèi),就會(huì)被標(biāo)注為死信消息 (3)當(dāng)隊(duì)列達(dá)到了最大長(zhǎng)度后,再發(fā)送過(guò)來(lái)的消息就會(huì)直接變成死信消息

1.3 什么是死信隊(duì)列?

直接來(lái)講,用來(lái)盛裝死信的隊(duì)列就是死信隊(duì)列,好像是一句廢話,所以其重點(diǎn)在于理解死信的概念。

死信隊(duì)列的作用: (1)隊(duì)列在已滿的情況下,會(huì)將消息發(fā)送到死信隊(duì)列中,這樣消息就不會(huì)丟失了,回頭再?gòu)乃佬抨?duì)列里將消息取出來(lái)進(jìn)行消費(fèi)即可 (2)可以基于死信隊(duì)列實(shí)現(xiàn)延遲消費(fèi)的效果。具體的實(shí)現(xiàn)我們后續(xù)講解

1.4 創(chuàng)建死信交換機(jī)、死信隊(duì)列

死信交換機(jī)、死信隊(duì)列其實(shí)都是普通的交換機(jī)、隊(duì)列,只是專門聲明出來(lái)用于存儲(chǔ)死信消息的。我們只需要通過(guò)deadLetterExchange方法來(lái)聲明死信交換機(jī),然后用deadLetterRoutingKey方法來(lái)聲明死信隊(duì)列

如下代碼所示,我們創(chuàng)建了test.queue、test.exchangedead.queuedead.exchange,并且在test.queue中將死信交換機(jī)和死信路由指定到了測(cè)試隊(duì)列中

注意:涉及到修改隊(duì)列、交換機(jī)屬性的,如果該隊(duì)列、交換機(jī)已經(jīng)存在需要將其刪除后才能生效,否則可能還會(huì)報(bào)錯(cuò)。

@Configuration
public class RabbitMqConfig {

    private static final String TEST_EXCHANGE = "test.exchange";
    private static final String TEST_QUEUE = "test.queue";
    private static final String TEST_ROUTING_KEY = "test.routing.key";
    private static final String DEAD_EXCHANGE = "dead.exchange";
    private static final String DEAD_QUEUE = "dead.queue";
    private static final String DEAD_ROUTING_KEY = "dead.routing.key";
    
    @Bean
    public Queue deadQueue(){
        return new Queue(DEAD_QUEUE);
    }
    public DirectExchange deadExchange(){
    // 設(shè)置演示,使用了直接交換機(jī)Direct,大家可以根據(jù)自己的業(yè)務(wù)情況聲明為其他類型的交換機(jī)
        return new DirectExchange(DEAD_EXCHANGE);
    public Binding deadBinding(Queue deadQueue,Exchange deadExchange){
        return BindingBuilder.bind(deadQueue).to(deadExchange).with(DEAD_ROUTING_KEY).noargs();
    public Queue testQueue(){
        return QueueBuilder.durable(TEST_QUEUE).deadLetterExchange(DEAD_EXCHANGE).deadLetterRoutingKey(DEAD_ROUTING_KEY).build();
    public DirectExchange testExchange(){
        return new DirectExchange(TEST_EXCHANGE);
    public Binding testQueueBing(Queue testQueue, DirectExchange testExchange){
        return BindingBuilder.bind(testQueue).to(testExchange).with(TEST_ROUTING_KEY);
}

1.5 實(shí)現(xiàn)死信消息

1.5.1 基于消費(fèi)者進(jìn)行reject或nack實(shí)現(xiàn)死信消息

@Component
public class QueueListener {
    @RabbitListener(queues = RabbitMqConfig.TEST_QUEUE)
    public void handler(MyMessage messageInfo, Message message, Channel channel) {
        try{
            System.out.println("接收的消息:"+messageInfo.toString());
            // requeue參數(shù)設(shè)置為false 設(shè)置死信消息
            channel.basicReject(message.getMessageProperties().getDeliveryTag(),false);
            // multiple和requeue設(shè)置為false 設(shè)置死信消息
            channel.basicNack(message.getMessageProperties().getDeliveryTag(),false,false);
            // 返回ack 確認(rèn)接收到消息
//            channel.basicAck(message.getMessageProperties().getDeliveryTag(),false);
        }catch (IOException e){
            try {
                channel.basicRecover();
            } catch (IOException ex) {
                ex.printStackTrace();
                log.error("消息處理失?。簕}",e.getMessage());
            }
        }
    }
}

1.5.2 基于生存時(shí)間實(shí)現(xiàn)

(1)發(fā)送消息時(shí)設(shè)置生存時(shí)間

    @GetMapping("sendTestQueueWithExpiration")
    public String sendTestQueueWithExpiration(){
        MyMessage message = new MyMessage(1L,"物流提醒","到達(dá)裝貨區(qū)域,注意上傳憑證",new Date());
        rabbitTemplate.convertAndSend(RabbitMqConfig.TEST_EXCHANGE,RabbitMqConfig.TEST_ROUTING_KEY, message,msg -> {
            msg.getMessageProperties().setExpiration("5000");
            return msg;
        });
        return "發(fā)送成功";
    }

(2)隊(duì)列設(shè)置生存時(shí)間

    @Bean
    public Queue testQueue(){
        return QueueBuilder.durable(TEST_QUEUE)
                .deadLetterExchange(DEAD_EXCHANGE)
                .deadLetterRoutingKey(DEAD_ROUTING_KEY)
                // 10s 過(guò)期
                .ttl(10000)
                .build();
    }

1.5.3 基于隊(duì)列max_length實(shí)現(xiàn)

    @Bean
    public Queue testQueue(){
        return QueueBuilder.durable(TEST_QUEUE)
                .deadLetterExchange(DEAD_EXCHANGE)
                .deadLetterRoutingKey(DEAD_ROUTING_KEY)
                // 容量最大100條
                .maxLength(100)
                .build();
    }

1.6 基于死信隊(duì)列實(shí)現(xiàn)消息延遲發(fā)送

上述我們說(shuō)過(guò)死信隊(duì)列還可以消息延遲發(fā)送,其思路就是: (1)消息發(fā)送時(shí)設(shè)置消息的生存時(shí)間,其生存時(shí)間就是我們想要延遲的時(shí)間 (2)消息者監(jiān)控死信隊(duì)列進(jìn)行消費(fèi)

正常隊(duì)列的消息因?yàn)闆](méi)有消費(fèi)者消費(fèi),同時(shí)又指定了生存時(shí)間,到達(dá)時(shí)間后消息轉(zhuǎn)發(fā)到死信隊(duì)列中,消費(fèi)者監(jiān)聽(tīng)了死信隊(duì)列從而將其消費(fèi)掉。

基于死信隊(duì)列實(shí)現(xiàn)消息延遲發(fā)送的問(wèn)題

如果有兩個(gè)消息,一個(gè)是5s生存時(shí)間,一個(gè)是10s生存時(shí)間,當(dāng)我們先發(fā)送了10s生存時(shí)間的消息到queue中時(shí),因?yàn)閞abbitmq只會(huì)監(jiān)控隊(duì)列最外側(cè)的消息的生存時(shí)間,也就是監(jiān)控10s生存時(shí)間的消息,而5s生存時(shí)間的消息只會(huì)在最外側(cè)的10s消息到期后才會(huì)監(jiān)控,也就導(dǎo)致我實(shí)際需要5s生存的消息,實(shí)際需要10s才監(jiān)聽(tīng)到了。

所以呢,基于死信隊(duì)列實(shí)現(xiàn)的延遲消息,只使用于延遲時(shí)間一致的消息。

為了適配更多的延遲場(chǎng)景,已經(jīng)更加簡(jiǎn)單的實(shí)現(xiàn)延遲消息,我們引入了延遲交換機(jī)

2. 延遲交換機(jī)

延遲交換機(jī)并不是rabbitmq自帶的功能,而是要通過(guò)安裝延遲交換機(jī)插件delayed_message_exchange來(lái)實(shí)現(xiàn)

其插件的安裝我們之間已經(jīng)講解過(guò),不再累敘,可以參考如下博文 springcloud:安裝rabbitmq并配置延遲隊(duì)列插件

通過(guò)延遲交換機(jī)實(shí)現(xiàn)的延遲消息,其重點(diǎn)主要在交換機(jī)上,隊(duì)列就是普通隊(duì)列,消息發(fā)送到交換機(jī)上后,會(huì)記錄消息的延遲時(shí)間,到達(dá)時(shí)間后才會(huì)發(fā)送到隊(duì)列中,這樣消費(fèi)者通過(guò)監(jiān)控隊(duì)列,就能在指定時(shí)間獲取到消息

因此延遲交換機(jī)與普通交換機(jī)的實(shí)現(xiàn),只在創(chuàng)建交換機(jī)時(shí),其他的操作與普通交換機(jī)無(wú)異,因此使用起來(lái)也很方便

創(chuàng)建延遲交換機(jī),通過(guò)x-delayed-type屬性聲明交換機(jī)類型,可以是direct也可以是topic,具體支持4中交換機(jī)類型,如果不清楚的可以參考之前的博文

@Configuration
public class RabbitMqDelayConfig {
    public static final String DELAY_EXCHANGE = "delay.exchange";
    public static final String DELAY_QUEUE = "delay.queue";
    public static final String DELAY_ROUTING_KEY = "delay.routing.key";
    @Bean
    public Exchange delayExchange(){
        Map<String, Object> arguments = new HashMap<>(1);
        arguments.put("x-delayed-type","direct");
        return new CustomExchange(DELAY_EXCHANGE,"x-delayed-message",true,false,arguments);
    }
    @Bean
    public Queue delayQueue(){
        return new Queue(DELAY_QUEUE);
    }
    @Bean
    public Binding delayBinding(Queue delayQueue, Exchange delayExchange){
        return BindingBuilder.bind(delayQueue).to(delayExchange).with(DELAY_ROUTING_KEY).noargs();
    }

}

發(fā)送消息時(shí)指定延遲時(shí)間,單位毫秒

rabbitTemplate.convertAndSend(DelayedConfig.DELAYED_EXCHANGE, "delayed.abc", "xxxx", new MessagePostProcessor() {
            @Override
            public Message postProcessMessage(Message message) throws AmqpException {
                message.getMessageProperties().setDelay(30000);
                return message;
            }
        });

我們還可以將該方法封裝為工具類方法,方便之后調(diào)用

/**
	 * 發(fā)送 延遲隊(duì)列
	 * @param exchange 交換機(jī)
	 * @param routeKey 路由
	 * @param message 消息
	 * @param delaySecond 延遲秒數(shù)
	 */
	public void send(String exchange, String routeKey, Object message, int delaySecond){  
		rabbitTemplate.convertAndSend(exchange,routeKey,message,msg -> {
			// 消息持久化 
			msg.getMessageProperties().setDelay(delaySecond * 1000);
			return msg;
		});
	}

3. 應(yīng)用場(chǎng)景

延遲消息的應(yīng)用場(chǎng)景豐富,除了我們開(kāi)篇所說(shuō)的30分鐘未支付自動(dòng)取消訂單,還比如到貨后72小時(shí)未簽收自動(dòng)簽收

基本上所有需要延遲觸發(fā)的業(yè)務(wù)場(chǎng)景都可以用rabbitmq延遲隊(duì)列來(lái)實(shí)現(xiàn)。

4. 練習(xí)題

對(duì)于剛接觸rabbitmq的同學(xué),這里我提供一個(gè)練習(xí)題給大家,也讓大家在實(shí)操中加強(qiáng)對(duì)于rabbitmq的理解:

需求:訂單到貨后72小時(shí)未簽收,自動(dòng)簽收 講解:我們這里要實(shí)現(xiàn)訂單到貨后的自動(dòng)簽收功能,訂單到貨后會(huì)觸發(fā)發(fā)送自動(dòng)簽收消息的方法,訂單已簽收的狀態(tài)status為2,到貨狀態(tài)為1,如果72小時(shí)前已經(jīng)簽收了即status被更新為2了,那么需要取消自動(dòng)簽收(不執(zhí)行自動(dòng)簽收,即忽略自動(dòng)簽收消息)

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

相關(guān)文章

  • SpringCloud超詳細(xì)講解微服務(wù)網(wǎng)關(guān)Zuul基礎(chǔ)

    SpringCloud超詳細(xì)講解微服務(wù)網(wǎng)關(guān)Zuul基礎(chǔ)

    這篇文章主要介紹了SpringCloud?Zuul微服務(wù)網(wǎng)關(guān),負(fù)載均衡,熔斷和限流,本文給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2022-10-10
  • Redisson分布式鎖的源碼解讀分享

    Redisson分布式鎖的源碼解讀分享

    Redisson是一個(gè)在Redis的基礎(chǔ)上實(shí)現(xiàn)的Java駐內(nèi)存數(shù)據(jù)網(wǎng)格(In-Memory?Data?Grid)。Redisson有一樣功能是可重入的分布式鎖。本文來(lái)討論一下這個(gè)功能的特點(diǎn)以及源碼分析
    2022-11-11
  • IDEA 中創(chuàng)建并部署 JavaWeb 程序的方法步驟(圖文)

    IDEA 中創(chuàng)建并部署 JavaWeb 程序的方法步驟(圖文)

    本文主要介紹了IDEA 中創(chuàng)建并部署 JavaWeb 程序的方法步驟,文中通過(guò)示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2022-02-02
  • Javafx利用fxml變換場(chǎng)景的實(shí)現(xiàn)示例

    Javafx利用fxml變換場(chǎng)景的實(shí)現(xiàn)示例

    本文主要介紹了Javafx利用fxml變換場(chǎng)景的實(shí)現(xiàn)示例,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧
    2024-07-07
  • SpringBoot中的Bean的初始化與銷毀順序解析

    SpringBoot中的Bean的初始化與銷毀順序解析

    這篇文章主要介紹了SpringBoot中的Bean的初始化與銷毀順序,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2021-08-08
  • 將SpringBoot項(xiàng)目無(wú)縫部署到Tomcat服務(wù)器的操作流程

    將SpringBoot項(xiàng)目無(wú)縫部署到Tomcat服務(wù)器的操作流程

    SpringBoot 是一個(gè)用來(lái)簡(jiǎn)化 Spring 應(yīng)用初始搭建以及開(kāi)發(fā)過(guò)程的框架,我們可以通過(guò)內(nèi)置的 Tomcat 容器來(lái)輕松地運(yùn)行我們的應(yīng)用,本文給大家介紹 SpringBoot 項(xiàng)目部署到獨(dú)立 Tomcat 服務(wù)器的操作流程,需要的朋友可以參考下
    2024-05-05
  • 如何解決java.lang.IllegalStateException: Target host is null的問(wèn)題

    如何解決java.lang.IllegalStateException: Target host&n

    文章描述了通過(guò)MocoRunner模擬接口,并使用properties文件和ResourceBundle讀取配置文件進(jìn)行g(shù)et請(qǐng)求的過(guò)程,在執(zhí)行過(guò)程中遇到了目標(biāo)主機(jī)為空的錯(cuò)誤,通過(guò)檢查和修正url拼接問(wèn)題解決了該錯(cuò)誤
    2024-12-12
  • IDEA Error:java: 無(wú)效的源發(fā)行版: 17錯(cuò)誤

    IDEA Error:java: 無(wú)效的源發(fā)行版: 17錯(cuò)誤

    本文主要介紹了IDEA Error:java: 無(wú)效的源發(fā)行版: 17錯(cuò)誤,這個(gè)錯(cuò)誤是因?yàn)槟腎DEA編譯器不支持Java 17版本,您需要更新您的IDEA編譯器或者將您的Java版本降級(jí)到IDEA支持的版本,本文就來(lái)詳細(xì)的介紹一下
    2023-08-08
  • servlet 解決亂碼問(wèn)題

    servlet 解決亂碼問(wèn)題

    這篇文章主要介紹了servlet 解決亂碼問(wèn)題 ,需要的朋友可以參考下
    2015-04-04
  • 如何用java對(duì)接微信小程序下單后的發(fā)貨接口

    如何用java對(duì)接微信小程序下單后的發(fā)貨接口

    這篇文章主要介紹了在微信小程序后臺(tái)實(shí)現(xiàn)發(fā)貨通知的步驟,包括獲取Access_token、使用RestTemplate調(diào)用發(fā)貨接口、處理AccessToken緩存以及發(fā)貨成功后的提醒,需要的朋友可以參考下
    2025-03-03

最新評(píng)論

夹江县| 泸西县| 喀喇沁旗| 肇州县| 万源市| 临高县| 德安县| 宁乡县| 西乌珠穆沁旗| 融水| 保山市| 荃湾区| 江山市| 武隆县| 仁怀市| 宝山区| 平乐县| 淳安县| 黄骅市| 孟津县| 西宁市| 西吉县| 通辽市| 屏边| 辽阳市| 新干县| 信丰县| 南丰县| 樟树市| 西贡区| 江达县| 从江县| 金山区| 涟源市| 衡东县| 隆尧县| 察哈| 崇信县| 增城市| 宜君县| 义马市|