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

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

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

最近在使用redis時(shí),就想能不能用其實(shí)現(xiàn)消息隊(duì)列?也在網(wǎng)上看了下其他小伙伴寫的實(shí)現(xiàn),結(jié)合自身業(yè)務(wù)實(shí)現(xiàn)了如下消息隊(duì)列,希望對(duì)大家有用。

廢話不多說,直接開擼。

1、為什么zset可以做消息隊(duì)列?

首先我們來看下,設(shè)計(jì)消息隊(duì)列需要考慮的需求:有序性,消息重復(fù)性,可靠性。

  • 有序性:zset所有元素可以根據(jù)成員關(guān)聯(lián)的score來進(jìn)行從低到高的排序,例如,我們可以利用時(shí)間戳來進(jìn)行排序
  • 消息重復(fù)性:在zset中每個(gè)元素都是唯一的,這也保證了消息的唯一性
  • 可靠性:zset會(huì)自動(dòng)維護(hù)元素之間的順序,在添加或刪除元素時(shí)無需手動(dòng)排序,提升操作速度。

2、使用的zset命令

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

3、代碼實(shí)現(xiàn)

使用zset實(shí)現(xiàn)消息隊(duì)列時(shí),具體的流程,如下:

生產(chǎn)者流程:

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

消費(fèi)者流程:

  • 調(diào)用方法獲取數(shù)據(jù)
  • 獲取到數(shù)據(jù),則直接返回,否則到指定時(shí)間后再次獲取數(shù)據(jù),直到獲取到數(shù)據(jù)并返回。

統(tǒng)一返回類:

    /**
     * @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 消息實(shí)體

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

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

3.2 隊(duì)列類型

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

    /**
     * @author: jiangjs
     * @description: 隊(duì)列類型
     * @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實(shí)現(xiàn)消息隊(duì)列
     * @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;
        /**
         * 隊(duì)列類型,即名稱
         */
        private final QueueTypeEnum typeEnum;
        public RedisQueueUtil(QueueTypeEnum typeEnum,RedisTemplate<String,QueueTask<T>> redisTemplate){
            this.typeEnum = typeEnum;
            this.redisTemplate = redisTemplate;
        }
        /**
         * 添加消息數(shù)據(jù)
         * @param queueTask 消息
         * @param time 延遲時(shí)間,單位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,"消息數(shù)據(jù)已經(jīng)存在,不予添加......");
                    }
                    Boolean result = redisTemplate.opsForZSet().add(typeEnum.getType(), queueTask, System.currentTimeMillis() + time*1000);
                    if (Objects.nonNull(result) && result){
                        log.info("添加消息數(shù)據(jù)成功:" + queueTask + ",添加時(shí)間:" + LocalDateTime.now());
                        return ResultUtil.success("添加消息數(shù)據(jù)成功");
                    }
                    return ResultUtil.error("添加消息數(shù)據(jù)失敗");
                }finally {
                    //釋放鎖
                    releaseLock();
                }
            } else {
                log.info("未獲取到鎖,稍后再試");
                return ResultUtil.error("未獲取到鎖,稍后再試");
            }
        }
        /**
         * 獲取zset前count數(shù)據(jù)
         * @param count 數(shù)據(jù)數(shù)
         * @return 返回獲取到數(shù)據(jù)
         */
        public Set<QueueTask<T>> loopGetTask(int count) {
                //rangeByScore,根據(jù)score順序獲取zset數(shù)據(jù)的值
                return redisTemplate.opsForZSet().rangeByScore(typeEnum.getType(), 0, System.currentTimeMillis(), 0, count-1);
        }
        /**
         * 注銷消息隊(duì)列
         * @param typeEnum 消息隊(duì)列名稱
         */
        public void destroy(QueueTypeEnum typeEnum){
            redisTemplate.opsForZSet().remove(typeEnum.getType());
        }
        /**
         * 獲取任務(wù)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)建消息任務(wù)時(shí)添加了鎖,只有在獲取鎖的前提下才能添加消息任務(wù)。

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

3.4 消費(fèi)消息

    /**
     * @author: jiangjs
     * @description: 啟動(dòng)消費(fèi)
     * @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) {
                        //校驗(yàn)當(dāng)前消息是否已消費(fèi),主要防止網(wǎng)絡(luò)延時(shí),導(dǎo)致多次提交同一任務(wù) 存在
                        QueueTask<String> stringQueueTask = redisTemplate.opsForValue().get(QUEUE_TYPE + "_" + queueTask.getTaskId());
                        if (Objects.nonNull(stringQueueTask)){
                            log.info("該任務(wù)已經(jīng)消費(fèi),不能重復(fù)消費(fèi)");
                            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("消費(fèi)任務(wù)數(shù)據(jù):" + task);
                            //設(shè)置過期時(shí)間,10分鐘內(nèi)則默認(rèn)是重復(fù)提交
                            redisTemplate.opsForValue().set(QUEUE_TYPE + "_" + queueTask.getTaskId(),queueTask,10L, TimeUnit.MINUTES);
                        }
                    }
                }
                log.info("------1分鐘后再次獲取------");
                Thread.sleep(60000);
            }
        }
    }

校驗(yàn)重復(fù)消息,若消息重復(fù)且在10分鐘內(nèi)未被消費(fèi),則直接將該消息從隊(duì)列中刪除。在消息任務(wù)被消費(fèi)后,將數(shù)據(jù)從隊(duì)列中移除。

執(zhí)行結(jié)果:

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

相關(guān)文章

  • redis?手機(jī)驗(yàn)證碼實(shí)現(xiàn)示例

    redis?手機(jī)驗(yàn)證碼實(shí)現(xiàn)示例

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

    Redis 旁路緩存深度解析

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

    詳解如何使用Redis實(shí)現(xiàn)分布式鎖

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

    redis實(shí)現(xiàn)主從模式(1主2從)

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

    redis實(shí)現(xiàn)計(jì)數(shù)器-防止刷單方法介紹

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

    Redis之ZipList壓縮列表的使用

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

    Redis內(nèi)存碎片處理實(shí)例詳解

    內(nèi)存碎片是redis服務(wù)中分配器分配存儲(chǔ)對(duì)象內(nèi)存的時(shí)產(chǎn)生的,下面這篇文章主要給大家介紹了關(guān)于Redis內(nèi)存碎片處理的相關(guān)資料,文中通過實(shí)例代碼介紹的非常詳細(xì),需要的朋友可以參考下
    2022-05-05
  • Redis常見分布鎖的原理和實(shí)現(xiàn)

    Redis常見分布鎖的原理和實(shí)現(xiàn)

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

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

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

    Redis 單機(jī)安裝和哨兵模式集群安裝的實(shí)現(xiàn)

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

最新評(píng)論

米脂县| 浦江县| 万年县| 高淳县| 台中县| 井研县| 江安县| 巧家县| 西安市| 迁安市| 武安市| 新乐市| 图木舒克市| 济南市| 平舆县| 金门县| 方城县| 淳安县| 临西县| 双柏县| 龙川县| 天门市| 奉化市| 白沙| 武功县| 九龙坡区| 治县。| 依安县| 沙湾县| 临澧县| 安吉县| 维西| 南阳市| 天峻县| 清丰县| 连城县| 陵水| 湖北省| 深水埗区| 郸城县| 林周县|