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

redis和rabbitmq實(shí)現(xiàn)延時隊(duì)列的示例代碼

 更新時間:2024年03月20日 11:03:52   作者:開心就好啦啦啦  
在高并發(fā)場景下,延遲隊(duì)列顯得尤為重要,本文主要介紹了兩種方式,redis和rabbitmq實(shí)現(xiàn)延時隊(duì)列,具有一定的參考價值,感興趣的可以了解一下

延遲隊(duì)列使用場景

1. 訂單超時處理:延遲隊(duì)列可以用于處理訂單超時問題。當(dāng)用戶下單后,將訂單信息放入延遲隊(duì)列,并設(shè)置一定的超時時間。如果在超時時間內(nèi)用戶未支付訂單,消費(fèi)者會從延遲隊(duì)列中獲取到該訂單,并執(zhí)行相應(yīng)的處理操作,如取消訂單、釋放庫存等。

2. 優(yōu)惠券過期提醒:延遲隊(duì)列可以用于優(yōu)惠券的過期提醒功能。將即將過期的優(yōu)惠券信息放入延遲隊(duì)列,并設(shè)置合適的延遲時間。當(dāng)延遲時間到達(dá)時,消費(fèi)者將提醒用戶優(yōu)惠券即將過期,引導(dǎo)用戶盡快使用。

3. 異步通知與提醒:延遲隊(duì)列可以用于異步通知和提醒功能。例如,當(dāng)用戶完成某個操作后,系統(tǒng)可以將相關(guān)通知消息放入延遲隊(duì)列,并設(shè)置一定的延遲時間,以便在合適的時機(jī)發(fā)送通知給用戶。

Redis中zset實(shí)現(xiàn)延時隊(duì)列

1. 創(chuàng)建延遲隊(duì)列服務(wù)類

  • 創(chuàng)建一個延遲隊(duì)列的服務(wù)類,例如DelayQueueService,用于操作Redis中的ZSet。這個服務(wù)類需要完成以下功能:
  • 將消息放入延遲隊(duì)列:將消息作為元素添加到ZSet中,設(shè)置對應(yīng)的延遲時間作為分?jǐn)?shù)。輪詢并處理已到期的消息:定時任務(wù)或者消息消費(fèi)者輪詢檢查ZSet中的元素,獲取到達(dá)指定時間的消息進(jìn)行處理。刪除已處理的消息:處理完消息后,從ZSet中將其刪除。
@Service
public class DelayQueueService {
    private static final String DELAY_QUEUE_KEY = "delay_queue";


    @Autowired
    private RedisTemplate<String, String> redisTemplate;

    public void addToDelayQueue(String message,long delayTime){
        redisTemplate.opsForZSet().add(DELAY_QUEUE_KEY,message,System.currentTimeMillis()+delayTime);
    }

    public void processDelayedMessage(){
        //reverseRangeByScore 從高到低
        //rangeByScore 從低到高
        Set<String> messages = redisTemplate.opsForZSet().rangeByScore(DELAY_QUEUE_KEY, 0, System.currentTimeMillis());
        for(String message:messages){
            //處理消息
            System.out.println(message);
            redisTemplate.opsForZSet().remove(DELAY_QUEUE_KEY,message);
        }

    }
}

2. 配置定時任務(wù)或消息消費(fèi)者

使用Spring Boot的定時任務(wù)或消息隊(duì)列框架,定時調(diào)用延遲隊(duì)列服務(wù)類的輪詢方法或監(jiān)聽指定的消息隊(duì)列,可以將輪訓(xùn)粒度放到1s一次。

@Component
public class DelayQueueSchedule {
    @Autowired
    private DelayQueueService delayQueueService;


    // 每隔一段時間進(jìn)行輪詢并處理延遲消息
    @Scheduled(fixedDelay = 1000)
    public void pollAndProcessDelayedMessages() {
        delayQueueService.pollAndProcessDelayedMessages();
    }
}

然后在啟動類上通過@EnableScheduling注解開啟任務(wù)調(diào)度能力。

缺點(diǎn):使用ZSET(有序集合,Sorted Set)來實(shí)現(xiàn)延遲任務(wù)調(diào)度(如訂單超時取消)是一種有效的方法,但它也有一些缺點(diǎn)和限制:

  • 內(nèi)存消耗:ZSET 在Redis中是一個有序集合,它需要占用一定的內(nèi)存來存儲成員和分?jǐn)?shù)。如果你需要存儲大量的延遲任務(wù),可能會導(dǎo)致內(nèi)存消耗較大。這可能會對Redis服務(wù)器的性能和成本產(chǎn)生影響,特別是在大規(guī)模應(yīng)用中。
  • 不適用于大規(guī)模延遲任務(wù):ZSET 可以處理相對較小數(shù)量的延遲任務(wù),但當(dāng)需要管理大規(guī)模延遲任務(wù)隊(duì)列時,可能會導(dǎo)致性能下降。在這種情況下,需要考慮更高效的延遲隊(duì)列解決方案,例如使用分布式消息隊(duì)列。
  • 無法動態(tài)修改延遲時間: 一旦將任務(wù)添加到ZSET中,你不能輕松地修改任務(wù)的延遲時間。如果需要在任務(wù)已經(jīng)添加后更改延遲時間,可能需要復(fù)雜的操作。
  • 沒有重試機(jī)制:ZSET 只能用于一次性延遲任務(wù),無法自動處理任務(wù)失敗后的重試。如果任務(wù)在執(zhí)行時失敗,你需要自己實(shí)現(xiàn)重試邏輯。
  • 沒有持久化: Redis是內(nèi)存數(shù)據(jù)庫,如果Redis服務(wù)器重啟或發(fā)生故障,已添加的延遲任務(wù)數(shù)據(jù)將丟失。雖然可以通過Redis持久化機(jī)制來部分解決這個問題,但仍然存在一定風(fēng)險。
  • 復(fù)雜性增加: 使用ZSET來管理延遲任務(wù)隊(duì)列需要編寫復(fù)雜的代碼來處理任務(wù)的添加、檢索和刪除。這可能增加應(yīng)用程序的復(fù)雜性。

Rabbitmq實(shí)現(xiàn)延遲隊(duì)列

死信,顧名思義就是無法被消費(fèi)的消息。一般來說,producer 將消息投遞到 broker 或者直接到queue 里了,consumer 從 queue 取出消息進(jìn)行消費(fèi),但某些時候由于特定的原因?qū)е聁ueu 中的某些消息無法被消費(fèi),這樣的消息如果沒有后續(xù)的處理,就變成了死信,有死信自然就有了死信隊(duì)列。

列出2種實(shí)現(xiàn)方式。
(1)使用Time To Live(TTL) + Dead Letter Exchanges(DLX)死信隊(duì)列組合實(shí)現(xiàn)延遲隊(duì)列的效果。
(2)使用RabbitMQ官方延遲插件rabbitmq_delayed_message_exchange,實(shí)現(xiàn)延時隊(duì)列效果。

由于TTL(生存時間)過期導(dǎo)致的死信,就是我們實(shí)現(xiàn)延遲隊(duì)列的的方式。
我們需要聲明如下形式的交互機(jī)和隊(duì)列,以及對應(yīng)的routing key,并進(jìn)行綁定:

請?zhí)砑訄D片描述

上圖綁定的代碼如下所示

@Configuration
public class DeadQueueConfig {
    //普通交換機(jī)及隊(duì)列
    public static final String X_EXCHANGE = "X";
    public static final String QUEUE_A = "QA";
    public static final String QUEUE_B = "QB";
    //死信交換機(jī)及隊(duì)列
    public static final String Y_DEAD_LETTER_EXCHANGE = "Y";
    public static final String DEAD_LETTER_QUEUE = "QD";
    //通用隊(duì)列
    public static final String QUEUE_C = "QC";


    // 聲明 xExchange
    @Bean("xExchange")
    public DirectExchange xExchange() {
        return new DirectExchange(X_EXCHANGE);
    }

    //聲明隊(duì)列 A ttl 為 10s 并綁定到對應(yīng)的死信交換機(jī)
    @Bean("queueA")
    public Queue queueA() {
        Map<String, Object> args = new HashMap<>(3);
        //聲明當(dāng)前隊(duì)列綁定的死信交換機(jī)
        args.put("x-dead-letter-exchange", Y_DEAD_LETTER_EXCHANGE);
        //聲明當(dāng)前隊(duì)列的死信路由 key
        args.put("x-dead-letter-routing-key", "YD");
        //聲明隊(duì)列的 TTL
        args.put("x-message-ttl", 10000);
        return QueueBuilder.durable(QUEUE_A).withArguments(args).build();
    }
    //聲明隊(duì)列A綁定X交換機(jī)  路由為XA
    @Bean
    public Binding queueABingX(@Qualifier("queueA") Queue queueA,
                               @Qualifier("xExchange") DirectExchange xExchange){
        return BindingBuilder.bind(queueA).to(xExchange).with("XA");
    }

    //聲明隊(duì)列 B ttl 為 40s 并綁定到對應(yīng)的死信交換機(jī)
    @Bean("queueB")
    public Queue queueB() {
        Map<String, Object> args = new HashMap<>(3);
        //聲明當(dāng)前隊(duì)列綁定的死信交換機(jī)
        args.put("x-dead-letter-exchange", Y_DEAD_LETTER_EXCHANGE);
        //聲明當(dāng)前隊(duì)列的死信路由 key
        args.put("x-dead-letter-routing-key", "YD");
        //聲明隊(duì)列的 TTL
        args.put("x-message-ttl", 40000);
        return QueueBuilder.durable(QUEUE_B).withArguments(args).build();
    }

    //聲明隊(duì)列 B 綁定 X 交換機(jī)
    @Bean
    public Binding queuebBindingX(@Qualifier("queueB") Queue queue1B,
                                  @Qualifier("xExchange") DirectExchange xExchange) {
        return BindingBuilder.bind(queue1B).to(xExchange).with("XB");
    }

    //聲明通用隊(duì)列C 不設(shè)ttl,由消費(fèi)者決定ttl
    @Bean("queueC")
    public Queue queueC() {
        Map<String, Object> args = new HashMap<>(3);
        //聲明當(dāng)前隊(duì)列綁定的死信交換機(jī)
        args.put("x-dead-letter-exchange", Y_DEAD_LETTER_EXCHANGE);
        //聲明當(dāng)前隊(duì)列的死信路由 key
        args.put("x-dead-letter-routing-key", "YD");
        return QueueBuilder.durable(QUEUE_C).withArguments(args).build();
    }
    // 聲明隊(duì)列 C 綁定 X 交換機(jī)
    @Bean
    public Binding queuecBindingX(@Qualifier("queueC") Queue queueC,
                                  @Qualifier("xExchange") DirectExchange xExchange) {
        return BindingBuilder.bind(queueC).to(xExchange).with("XC");
    }

    // 聲明 死信隊(duì)列交換機(jī)
    @Bean("yExchange")
    public DirectExchange yExchange() {
        return new DirectExchange(Y_DEAD_LETTER_EXCHANGE);
    }
    //聲明死信隊(duì)列 QD
    @Bean("queueD")
    public Queue queueD() {
        return new Queue(DEAD_LETTER_QUEUE,true);
    }
    //聲明死信隊(duì)列 QD 綁定關(guān)系
    @Bean
    public Binding deadLetterBindingQAD(@Qualifier("queueD") Queue queueD,
                                        @Qualifier("yExchange") DirectExchange yExchange) {
        return BindingBuilder.bind(queueD).to(yExchange).with("YD");
    }

}

其中,QD為死信隊(duì)列。當(dāng)QA和QB隊(duì)列中的消息,達(dá)到設(shè)定的TTL(10s和40s)后,將進(jìn)入指定的死信隊(duì)列QD。該方法缺點(diǎn)就是一個TTL對應(yīng)一個隊(duì)列

其中的QC作為通用的隊(duì)列,即在消費(fèi)者處指定消息對應(yīng)的TTL,TTL過期后轉(zhuǎn)入死信隊(duì)列。使用該通用隊(duì)列可以避免每增加一個新的時間需求,就要新增一個隊(duì)列的問題。但該方法由于隊(duì)列先進(jìn)先出的性質(zhì),會導(dǎo)致一定的問題:

即先發(fā)出一個TTL為10s的消息a,進(jìn)入隊(duì)列;再馬上發(fā)出一個TTL為2s的消息b,進(jìn)入隊(duì)列。由于隊(duì)列的性質(zhì),會在消息a的TTL結(jié)束后,a進(jìn)入死信隊(duì)列后,b才會進(jìn)入死信隊(duì)列。而不是根據(jù)TTL的時間,b比a先進(jìn)入死信隊(duì)列。

聲明交換機(jī)、隊(duì)列,并綁定成功后,編寫死信隊(duì)列消費(fèi)者代碼;

@Component
@Slf4j
public class DeadQueueConsumer {

    @RabbitListener(queues = "QD")
    public void receiveD(Message message, Channel channel) throws IOException {
        String msg = new String(message.getBody());
        log.info("當(dāng)前時間:{},收到死信隊(duì)列信息:{}", new Date().toString(), msg);
    }
}

在controller中編寫生產(chǎn)者代碼,進(jìn)行測試:

    @Autowired
    private RabbitTemplate rabbitTemplate;

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

結(jié)果如圖:

請?zhí)砑訄D片描述

測試通用隊(duì)列QC的效果:

@GetMapping("/send/{message}/{ttlTime}")
    public void sendMsg(@PathVariable String message, @PathVariable String ttlTime) {
        rabbitTemplate.convertAndSend("X", "XC", message, correlationData -> {
            correlationData.getMessageProperties().setExpiration(ttlTime);
            return correlationData;
        });
        log.info("當(dāng)前時間:{},發(fā)送一條時長{}毫秒 TTL 信息給隊(duì)列 C:{}", new Date(), ttlTime, message);
    }

結(jié)果如下圖

請?zhí)砑訄D片描述

可以看到, 兩條消息幾乎同時到達(dá)死信隊(duì)列,因?yàn)門TL為2s的消息由于被堵在TTL為10s的消息后導(dǎo)致。

到此這篇關(guān)于redis和rabbitmq實(shí)現(xiàn)延時隊(duì)列的示例代碼的文章就介紹到這了,更多相關(guān)redis rabbitmq 延時隊(duì)列內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • 將音頻文件轉(zhuǎn)二進(jìn)制分包存儲到Redis的實(shí)現(xiàn)方法(奇淫技巧操作)

    將音頻文件轉(zhuǎn)二進(jìn)制分包存儲到Redis的實(shí)現(xiàn)方法(奇淫技巧操作)

    這篇文章主要介紹了將音頻文件轉(zhuǎn)二進(jìn)制分包存儲到Redis的實(shí)現(xiàn)方法(奇淫技巧操作),本文通過實(shí)例代碼給大家介紹的非常詳細(xì),對大家的學(xué)習(xí)或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2020-07-07
  • redis中redisson實(shí)現(xiàn)鎖自動延時

    redis中redisson實(shí)現(xiàn)鎖自動延時

    redisson作為分布式鎖能夠解決分布式的加鎖解鎖問題,還能夠?qū)崿F(xiàn)鎖的設(shè)置存活時間以及自動續(xù)期,本文主要介紹了redis中redisson實(shí)現(xiàn)鎖自動延時,感興趣的可以了解一下
    2024-02-02
  • Redis基本數(shù)據(jù)類型String常用操作命令

    Redis基本數(shù)據(jù)類型String常用操作命令

    這篇文章主要為大家介紹了Redis基本數(shù)據(jù)類型String常用操作命令,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪
    2022-05-05
  • 基于Redis實(shí)現(xiàn)API接口訪問次數(shù)限制

    基于Redis實(shí)現(xiàn)API接口訪問次數(shù)限制

    日常開發(fā)中會有一個常見的需求,需要限制接口在單位時間內(nèi)的訪問次數(shù),比如說某個免費(fèi)的接口限制單個IP一分鐘內(nèi)只能訪問5次,該怎么實(shí)現(xiàn)呢,本文小編給大家介紹了如何基于Redis實(shí)現(xiàn)API接口訪問次數(shù)限制,需要的朋友可以參考下
    2024-11-11
  • Redis消息隊(duì)列、阻塞隊(duì)列、延時隊(duì)列的實(shí)現(xiàn)

    Redis消息隊(duì)列、阻塞隊(duì)列、延時隊(duì)列的實(shí)現(xiàn)

    Redis是一種常用的內(nèi)存數(shù)據(jù)庫,它提供了豐富的功能,通常用于數(shù)據(jù)緩存和分布式隊(duì)列,本文主要介紹了Redis消息隊(duì)列、阻塞隊(duì)列、延時隊(duì)列的實(shí)現(xiàn),感興趣的可以了解一下
    2023-11-11
  • Redis 實(shí)現(xiàn)分布式鎖時需要考慮的問題解決方案

    Redis 實(shí)現(xiàn)分布式鎖時需要考慮的問題解決方案

    本文詳細(xì)探討了使用Redis實(shí)現(xiàn)分布式鎖時需要考慮的問題,包括鎖的競爭、鎖的釋放、超時管理、網(wǎng)絡(luò)分區(qū)等,并提供了相應(yīng)的解決方案和代碼實(shí)例,有助于開發(fā)者正確且安全地使用Redis實(shí)現(xiàn)分布式鎖
    2024-09-09
  • 了解Redis常見應(yīng)用場景

    了解Redis常見應(yīng)用場景

    Redis是一個key-value存儲系統(tǒng),現(xiàn)在在各種系統(tǒng)中的使用越來越多,大部分情況下是因?yàn)槠涓咝阅艿奶匦?,被?dāng)做緩存使用,這里介紹下Redis經(jīng)常遇到的使用場景
    2021-06-06
  • redis哈希類型_動力節(jié)點(diǎn)Java學(xué)院整理

    redis哈希類型_動力節(jié)點(diǎn)Java學(xué)院整理

    這篇文章主要介紹了redis哈希類型的常用方法及原理淺析,感興趣的朋友一起看看吧
    2017-08-08
  • 在ssm項(xiàng)目中使用redis緩存查詢數(shù)據(jù)的方法

    在ssm項(xiàng)目中使用redis緩存查詢數(shù)據(jù)的方法

    本文主要簡單的使用Java代碼進(jìn)行redis緩存,即在查詢的時候先在service層從redis緩存中獲取數(shù)據(jù)。如果大家對在ssm項(xiàng)目中使用redis緩存查詢數(shù)據(jù)的相關(guān)知識感興趣的朋友跟隨腳本之家小編一起看看吧
    2018-03-03
  • Redis中過期鍵刪除的三種方法

    Redis中過期鍵刪除的三種方法

    Redis中可以設(shè)置鍵的過期時間,并且通過取出過期字典(expires dict)中鍵的過期時間和當(dāng)前時間比較來判斷是否過期,那么一個過期的鍵是怎么被刪除的呢?本文給大家總結(jié)了三種方法,選了其中兩種給大家詳細(xì)的介紹一下,需要的朋友可以參考下
    2024-05-05

最新評論

沐川县| 湖南省| 蓝田县| 建湖县| 营山县| 峨边| 济宁市| 荃湾区| 凤阳县| 巩义市| 乌拉特中旗| 永州市| 射洪县| 郧西县| 武威市| 怀来县| 天台县| 金山区| 盐山县| 张掖市| 翁源县| 永新县| 滨海县| 沙田区| 曲周县| 拉孜县| 灵川县| 遂宁市| 盖州市| 宁强县| 汉川市| 新龙县| 泗洪县| 丹寨县| 桦甸市| 天峨县| 遂平县| 仙居县| 开远市| 大连市| 珠海市|