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

redis使用zset實現延時隊列的示例代碼

 更新時間:2023年06月02日 15:31:50   作者:搶老婆酸奶的小肥仔  
本文主要介紹了redis使用zset實現延時隊列的示例代碼,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧

最近在使用redis時,就想能不能用其實現消息隊列?也在網上看了下其他小伙伴寫的實現,結合自身業(yè)務實現了如下消息隊列,希望對大家有用。

廢話不多說,直接開擼。

1、為什么zset可以做消息隊列?

首先我們來看下,設計消息隊列需要考慮的需求:有序性,消息重復性,可靠性。

  • 有序性:zset所有元素可以根據成員關聯的score來進行從低到高的排序,例如,我們可以利用時間戳來進行排序
  • 消息重復性:在zset中每個元素都是唯一的,這也保證了消息的唯一性
  • 可靠性:zset會自動維護元素之間的順序,在添加或刪除元素時無需手動排序,提升操作速度。

2、使用的zset命令

命令描述
zadd將一個給定score的成員添加到有序集合中,返回添加元素的個數
zrange根據元素在有序排序中的位置,從有序集合中獲取多個元素
rank(K key, Object o)獲取指定元素在集合中的索引,索引從0開始

3、代碼實現

使用zset實現消息隊列時,具體的流程,如下:

生產者流程:

  • 用戶獲取消息Id,并封裝消息體
  • 用戶發(fā)送數據到生產者,先獲取鎖
  • 如果獲取到鎖,則校驗該消息體是否已添加到隊列中,已添加則直接返回提醒。
  • 若未添加則調用方法將數據保存到zset集合中,否則等到指定時間后再獲取鎖。
  • 推送數據后,釋放鎖

消費者流程:

  • 調用方法獲取數據
  • 獲取到數據,則直接返回,否則到指定時間后再次獲取數據,直到獲取到數據并返回。

統一返回類:

    /**
     * @Author: jiangjs
     * @Description:
     * @Date: 2021/11/12 15:46
     **/
    @Data
    @Builder
    @NoArgsConstructor
    @AllArgsConstructor
    public class ResultUtil<T> implements Serializable {
        private int code;
        private String msg;
        private T data;
        public static <T> ResultUtil<T> success(){
            return ResultUtil.<T>builder().code(1000).msg("成功").build();
        }
        public static <T> ResultUtil<T> success(T data){
            return ResultUtil.<T>builder().code(1000).msg("成功").data(data).build();
        }
        public static <T> ResultUtil<T> error(String msg){
            return ResultUtil.<T>builder().code(5000).msg(msg).data(null).build();
        }
        public static <T> ResultUtil<T> error(int code,String msg){
            return ResultUtil.<T>builder().code(code).msg(msg).build();
        }
    }

3.1 消息實體

需添加消息Id,主要防止消息重復提交。

    /**
     * @author: jiangjs
     * @description: 消息實體
     * @date: 2023/5/30 11:11
     **/
    @Data
    @Accessors(chain = true)
    public class QueueTask<T> {
        /**
         * 消息Id
         */
        private String taskId;
        /**
         * 任務
         */
        private T task;
    }

3.2 隊列類型

隊列類型可以理解為隊列的名稱,通過枚舉,可以隨意添加隊列名稱。

    /**
     * @author: jiangjs
     * @description: 隊列類型
     * @date: 2023/5/30 10:53
     **/
    public enum QueueTypeEnum {
        /**
         * 訂單
         */
        ORDER("order");
        private final String type;
        QueueTypeEnum(String type){
            this.type = type;
        }
        public String getType(){
            return type;
        }
    }

3.3 創(chuàng)建消息工具

    package com.jiashn.springbootproject.redis.utils;
    import com.jiashn.springbootproject.redis.domain.QueueTask;
    import com.jiashn.springbootproject.utils.ResultUtil;
    import org.slf4j.Logger;
    import org.slf4j.LoggerFactory;
    import org.springframework.data.redis.core.RedisTemplate;
    import javax.annotation.Resource;
    import java.time.LocalDateTime;
    import java.util.Objects;
    import java.util.Set;
    import java.util.UUID;
    import java.util.concurrent.TimeUnit;
    /**
     * @author: jiangjs
     * @description: redis實現消息隊列
     * @date: 2023/5/30 10:51
     **/
    public class RedisQueueUtil<T> {
        private static final Logger log = LoggerFactory.getLogger(RedisQueueUtil.class);
        private RedisTemplate<String,QueueTask<T>> redisTemplate;
        /**
         * 隊列類型,即名稱
         */
        private final QueueTypeEnum typeEnum;
        public RedisQueueUtil(QueueTypeEnum typeEnum,RedisTemplate<String,QueueTask<T>> redisTemplate){
            this.typeEnum = typeEnum;
            this.redisTemplate = redisTemplate;
        }
        /**
         * 添加消息數據
         * @param queueTask 消息
         * @param time 延遲時間,單位s
         */
        public ResultUtil<String> sendQueueTask(QueueTask<T> queueTask, long time){
            //加鎖
            if (getLock()){
                try {
                    Long rank = redisTemplate.opsForZSet().rank(typeEnum.getType(), queueTask);
                    if (Objects.nonNull(rank)){
                        return ResultUtil.error(6000,"消息數據已經存在,不予添加......");
                    }
                    Boolean result = redisTemplate.opsForZSet().add(typeEnum.getType(), queueTask, System.currentTimeMillis() + time*1000);
                    if (Objects.nonNull(result) && result){
                        log.info("添加消息數據成功:" + queueTask + ",添加時間:" + LocalDateTime.now());
                        return ResultUtil.success("添加消息數據成功");
                    }
                    return ResultUtil.error("添加消息數據失敗");
                }finally {
                    //釋放鎖
                    releaseLock();
                }
            } else {
                log.info("未獲取到鎖,稍后再試");
                return ResultUtil.error("未獲取到鎖,稍后再試");
            }
        }
        /**
         * 獲取zset前count數據
         * @param count 數據數
         * @return 返回獲取到數據
         */
        public Set<QueueTask<T>> loopGetTask(int count) {
                //rangeByScore,根據score順序獲取zset數據的值
                return redisTemplate.opsForZSet().rangeByScore(typeEnum.getType(), 0, System.currentTimeMillis(), 0, count-1);
        }
        /**
         * 注銷消息隊列
         * @param typeEnum 消息隊列名稱
         */
        public void destroy(QueueTypeEnum typeEnum){
            redisTemplate.opsForZSet().remove(typeEnum.getType());
        }
        /**
         * 獲取任務Id
         * @return 返回消息Id
         */
        public String getTaskId(){
           return typeEnum.getType() + "_" + UUID.randomUUID().toString().replace("-","");
        }
        /**
         * 獲取鎖
         * @return 返回加鎖狀態(tài)
         */
        private boolean getLock(){
            Boolean absent = redisTemplate.opsForValue().setIfAbsent(typeEnum.getType() + "_Locked", null, 30L, TimeUnit.MINUTES);
            return Objects.nonNull(absent) ? absent : false;
        }
        /**
         * 釋放鎖
         */
        public void releaseLock(){
            redisTemplate.delete(typeEnum.getType() + "_Locked");
        }
    }

在消息工具類中,創(chuàng)建消息任務時添加了鎖,只有在獲取鎖的前提下才能添加消息任務。

提供獲取消息Id的方法是為了讓提交消息任務前,先獲取Id,即使在提交時網絡發(fā)生問題,提交的Id還是同一個,再進行消息消費時,可以根據這個Id來進行判斷該消息任務是否已被消費,被消費則直接丟棄。

3.4 消費消息

    /**
     * @author: jiangjs
     * @description: 啟動消費
     * @date: 2023/5/30 14:27
     **/
    @Component
    public class CustomerTaskLineRunner implements CommandLineRunner {
        @Resource
        private RedisTemplate<String,QueueTask<String>> redisTemplate;
        private final static String QUEUE_TYPE = QueueTypeEnum.ORDER.getType();
        private final static Logger log = LoggerFactory.getLogger(CustomerTaskLineRunner.class);
        @Override
        public void run(String... args) throws Exception {
            RedisQueueUtil<String> queueUtil = new RedisQueueUtil<>(QueueTypeEnum.ORDER,redisTemplate);
            while (true){
                Set<QueueTask<String>> queueTasks = queueUtil.loopGetTask(10);
                if (CollectionUtils.isNotEmpty(queueTasks)){
                    for (QueueTask<String> queueTask : queueTasks) {
                        //校驗當前消息是否已消費,主要防止網絡延時,導致多次提交同一任務 存在
                        QueueTask<String> stringQueueTask = redisTemplate.opsForValue().get(QUEUE_TYPE + "_" + queueTask.getTaskId());
                        if (Objects.nonNull(stringQueueTask)){
                            log.info("該任務已經消費,不能重復消費");
                            redisTemplate.opsForZSet().remove(QUEUE_TYPE,queueTask);
                            continue;
                        }
                        Long removeNum = redisTemplate.opsForZSet().remove(QUEUE_TYPE,queueTask);
                        if (Objects.nonNull(removeNum) && removeNum > 0){
                            String task = queueTask.getTask();
                            log.info("消費任務數據:" + task);
                            //設置過期時間,10分鐘內則默認是重復提交
                            redisTemplate.opsForValue().set(QUEUE_TYPE + "_" + queueTask.getTaskId(),queueTask,10L, TimeUnit.MINUTES);
                        }
                    }
                }
                log.info("------1分鐘后再次獲取------");
                Thread.sleep(60000);
            }
        }
    }

校驗重復消息,若消息重復且在10分鐘內未被消費,則直接將該消息從隊列中刪除。在消息任務被消費后,將數據從隊列中移除。

執(zhí)行結果:

到此這篇關于redis使用zset實現延時隊列的示例代碼的文章就介紹到這了,更多相關redis zset延時隊列內容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關文章希望大家以后多多支持腳本之家!

相關文章

  • redis?手機驗證碼實現示例

    redis?手機驗證碼實現示例

    本文主要介紹了redis?手機驗證碼實現示例,文中通過示例代碼介紹的非常詳細,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2021-11-11
  • Redis 旁路緩存深度解析

    Redis 旁路緩存深度解析

    旁路緩存是 Redis 最常用的緩存策略,本文主要介紹了Redis旁路緩存深度解析,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧
    2026-04-04
  • 詳解如何使用Redis實現分布式鎖

    詳解如何使用Redis實現分布式鎖

    Redis 作為一個獨立的三方系統,其天生的優(yōu)勢就是可以作為一個分布式系統來使用,因此使用 Redis 實現的鎖都是分布式鎖,所以本文就給大家講講如何使用Redis實現分布式鎖,感興趣的小伙伴跟著小編來看看吧
    2023-08-08
  • redis實現主從模式(1主2從)

    redis實現主從模式(1主2從)

    本文主要介紹了在Windows環(huán)境下搭建和測試Redis的主從復制模式,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧
    2024-12-12
  • redis實現計數器-防止刷單方法介紹

    redis實現計數器-防止刷單方法介紹

    本文主要向大家介紹了redis實現計數器防止刷單的方法和有關代碼,具有一定參考價值,需要的朋友可以了解下。
    2017-11-11
  • Redis之ZipList壓縮列表的使用

    Redis之ZipList壓縮列表的使用

    這篇文章主要介紹了Redis之ZipList壓縮列表的使用,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教
    2025-06-06
  • Redis內存碎片處理實例詳解

    Redis內存碎片處理實例詳解

    內存碎片是redis服務中分配器分配存儲對象內存的時產生的,下面這篇文章主要給大家介紹了關于Redis內存碎片處理的相關資料,文中通過實例代碼介紹的非常詳細,需要的朋友可以參考下
    2022-05-05
  • Redis常見分布鎖的原理和實現

    Redis常見分布鎖的原理和實現

    這篇文章主要介紹了Redis常見分布鎖的原理和實現,文章圍繞主題展開詳細的內容介紹,具有一定的參考價值,需要的小伙伴可以參考一下
    2022-08-08
  • redis 解決key的亂碼問題,并清理詳解

    redis 解決key的亂碼問題,并清理詳解

    這篇文章主要介紹了redis 解決key的亂碼問題,并清理詳解,具有很好的參考價值,希望對大家有所幫助。一起跟隨小編過來看看吧
    2020-07-07
  • Redis 單機安裝和哨兵模式集群安裝的實現

    Redis 單機安裝和哨兵模式集群安裝的實現

    本文主要介紹了Redis 單機安裝和哨兵模式集群安裝的實現,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧
    2022-07-07

最新評論

黄陵县| 崇仁县| 法库县| 临朐县| 高唐县| 长宁县| 钦州市| 舞钢市| 葫芦岛市| 滕州市| 南安市| 瑞昌市| 平和县| 清河县| 合作市| 伊吾县| 乐至县| 富顺县| 盖州市| 周至县| 琼海市| 朔州市| 全州县| 沭阳县| 临漳县| 乐业县| 专栏| 彝良县| 濮阳市| 黄大仙区| 环江| 封开县| 阿鲁科尔沁旗| 剑河县| 交城县| 衡东县| 隆化县| 惠州市| 庄浪县| 建始县| 峨山|