【Redis】redis高阶-使用zset实现延时队列

2024-06-03 22:20

本文主要是介绍【Redis】redis高阶-使用zset实现延时队列,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

Hi,大家好,我是抢老婆酸奶的小肥仔。

最近在使用redis时,就想能不能用其实现消息队列?也在网上看了下其他小伙伴写的实现,结合自身业务实现了如下消息队列,希望对大家有用。

废话不多说,直接开撸。

1、为什么zset可以做消息队列?

首先我们来看下,设计消息队列需要考虑的需求:有序性,消息重复性,可靠性。

  • 有序性:zset所有元素可以根据成员关联的score来进行从低到高的排序,例如,我们可以利用时间戳来进行排序
  • 消息重复性:在zset中每个元素都是唯一的,这也保证了消息的唯一性
  • 可靠性:zset会自动维护元素之间的顺序,在添加或删除元素时无需手动排序,提升操作速度。

2、使用的zset命令

命令描述
zadd将一个给定score的成员添加到有序集合中,返回添加元素的个数
zrange根据元素在有序排序中的位置,从有序集合中获取多个元素
rank(K key, Object o)获取指定元素在集合中的索引,索引从0开始

3、代码实现

使用zset实现消息队列时,具体的流程,如下:

生产者流程:

  1. 用户获取消息Id,并封装消息体
  2. 用户发送数据到生产者,先获取锁
  3. 如果获取到锁,则校验该消息体是否已添加到队列中,已添加则直接返回提醒。
  4. 若未添加则调用方法将数据保存到zset集合中,否则等到指定时间后再获取锁。
  5. 推送数据后,释放锁

消费者流程:

  1. 调用方法获取数据
  2. 获取到数据,则直接返回,否则到指定时间后再次获取数据,直到获取到数据并返回。

统一返回类:

/*** @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 创建消息工具

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 返回加锁状态*/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");}
}

在消息工具类中,创建消息任务时添加了锁,只有在获取锁的前提下才能添加消息任务。

提供获取消息Id的方法是为了让提交消息任务前,先获取Id,即使在提交时网络发生问题,提交的Id还是同一个,再进行消息消费时,可以根据这个Id来进行判断该消息任务是否已被消费,被消费则直接丢弃。

3.4 消费消息

/*** @author: jiangjs* @description: 启动消费* @date: 2023/5/30 14:27**/
@Component
public class CustomerTaskLineRunner implements CommandLineRunner {@Resourceprivate RedisTemplate<String,QueueTask<String>> redisTemplate;private final static String QUEUE_TYPE = QueueTypeEnum.ORDER.getType();private final static Logger log = LoggerFactory.getLogger(CustomerTaskLineRunner.class);@Overridepublic 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分钟内未被消费,则直接将该消息从队列中删除。在消息任务被消费后,将数据从队列中移除。

执行结果:

谢谢大家,今天的分享就到这,不合理的地方希望大家指正。

这篇关于【Redis】redis高阶-使用zset实现延时队列的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/1028278

相关文章

Java中流式并行操作parallelStream的原理和使用方法

《Java中流式并行操作parallelStream的原理和使用方法》本文详细介绍了Java中的并行流(parallelStream)的原理、正确使用方法以及在实际业务中的应用案例,并指出在使用并行流... 目录Java中流式并行操作parallelStream0. 问题的产生1. 什么是parallelS

C++中unordered_set哈希集合的实现

《C++中unordered_set哈希集合的实现》std::unordered_set是C++标准库中的无序关联容器,基于哈希表实现,具有元素唯一性和无序性特点,本文就来详细的介绍一下unorder... 目录一、概述二、头文件与命名空间三、常用方法与示例1. 构造与析构2. 迭代器与遍历3. 容量相关4

Linux join命令的使用及说明

《Linuxjoin命令的使用及说明》`join`命令用于在Linux中按字段将两个文件进行连接,类似于SQL的JOIN,它需要两个文件按用于匹配的字段排序,并且第一个文件的换行符必须是LF,`jo... 目录一. 基本语法二. 数据准备三. 指定文件的连接key四.-a输出指定文件的所有行五.-o指定输出

Linux jq命令的使用解读

《Linuxjq命令的使用解读》jq是一个强大的命令行工具,用于处理JSON数据,它可以用来查看、过滤、修改、格式化JSON数据,通过使用各种选项和过滤器,可以实现复杂的JSON处理任务... 目录一. 简介二. 选项2.1.2.2-c2.3-r2.4-R三. 字段提取3.1 普通字段3.2 数组字段四.

C++中悬垂引用(Dangling Reference) 的实现

《C++中悬垂引用(DanglingReference)的实现》C++中的悬垂引用指引用绑定的对象被销毁后引用仍存在的情况,会导致访问无效内存,下面就来详细的介绍一下产生的原因以及如何避免,感兴趣... 目录悬垂引用的产生原因1. 引用绑定到局部变量,变量超出作用域后销毁2. 引用绑定到动态分配的对象,对象

Linux kill正在执行的后台任务 kill进程组使用详解

《Linuxkill正在执行的后台任务kill进程组使用详解》文章介绍了两个脚本的功能和区别,以及执行这些脚本时遇到的进程管理问题,通过查看进程树、使用`kill`命令和`lsof`命令,分析了子... 目录零. 用到的命令一. 待执行的脚本二. 执行含子进程的脚本,并kill2.1 进程查看2.2 遇到的

SpringBoot基于注解实现数据库字段回填的完整方案

《SpringBoot基于注解实现数据库字段回填的完整方案》这篇文章主要为大家详细介绍了SpringBoot如何基于注解实现数据库字段回填的相关方法,文中的示例代码讲解详细,感兴趣的小伙伴可以了解... 目录数据库表pom.XMLRelationFieldRelationFieldMapping基础的一些代

Java HashMap的底层实现原理深度解析

《JavaHashMap的底层实现原理深度解析》HashMap基于数组+链表+红黑树结构,通过哈希算法和扩容机制优化性能,负载因子与树化阈值平衡效率,是Java开发必备的高效数据结构,本文给大家介绍... 目录一、概述:HashMap的宏观结构二、核心数据结构解析1. 数组(桶数组)2. 链表节点(Node

Java AOP面向切面编程的概念和实现方式

《JavaAOP面向切面编程的概念和实现方式》AOP是面向切面编程,通过动态代理将横切关注点(如日志、事务)与核心业务逻辑分离,提升代码复用性和可维护性,本文给大家介绍JavaAOP面向切面编程的概... 目录一、AOP 是什么?二、AOP 的核心概念与实现方式核心概念实现方式三、Spring AOP 的关

详解SpringBoot+Ehcache使用示例

《详解SpringBoot+Ehcache使用示例》本文介绍了SpringBoot中配置Ehcache、自定义get/set方式,并实际使用缓存的过程,文中通过示例代码介绍的非常详细,对大家的学习或者... 目录摘要概念内存与磁盘持久化存储:配置灵活性:编码示例引入依赖:配置ehcache.XML文件:配置