本文作者:DurkBlue

java当中使用Redis ZSet + Lua 脚本自研延迟队列来实现定时到期执行业务处理功能推荐

DurkBlue 今天 62
java当中使用Redis ZSet + Lua 脚本自研延迟队列来实现定时到期执行业务处理功能摘要:       定时到期任务需求在生产中是个常见的需求,比如,生成的待支付订单在5分钟内用户没有支付,需要自动将订单取消,如果单纯靠轮询数据库来定时查询5...

      定时到期任务需求在生产中是个常见的需求,比如,生成的待支付订单在5分钟内用户没有支付,需要自动将订单取消,如果单纯靠轮询数据库来定时查询5分钟内用户没有支付的订单,会造成大量无效的数据库查询产生,下面介绍一些可行的解决方案:Redis ZSet + Lua 脚本自研延迟队列来实现定时到期执行业务处理功能

    1.设置任务(指定多久之后触发通知)

    2.高频续期(原子更新到期时间)

    3.删除任务

    4.后台轮询消费到期任务

    5.异常捕获、日志、防止任务重复执行

    6.支持任务业务对象存储

    


    maven引入相关依赖

    

    

<!--        Spring Boot的Redis数据访问模块。-->
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-data-redis</artifactId>
</dependency>
<!-- redis分布式锁 -->
<dependency>
    <groupId>org.redisson</groupId>
    <artifactId>redisson</artifactId>
    <version>3.16.6</version> <!-- 使用最新版本 -->
</dependency>

   分享redis和redisson配置文件

package com.tehn.zhfljobsapi.config;

import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.data.redis.connection.RedisConnectionFactory;
import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.data.redis.serializer.StringRedisSerializer;

@Configuration
public class RedisConfig {

    @Bean
    public RedisTemplate<String, Object> redisTemplate(RedisConnectionFactory connectionFactory) {
        RedisTemplate<String, Object> redisTemplate = new RedisTemplate<>();
        redisTemplate.setConnectionFactory(connectionFactory);
        redisTemplate.setKeySerializer(new StringRedisSerializer());


        return redisTemplate;
    }

    @Bean
    public RedisTemplate<String,Integer> stringIntegerRedisTemplate(RedisConnectionFactory redisConnectionFactory){
        RedisTemplate<String,Integer> template = new RedisTemplate();
        template.setConnectionFactory(redisConnectionFactory);

        StringRedisSerializer stringRedisSerializer = new StringRedisSerializer();
        template.setKeySerializer(stringRedisSerializer);
        template.setHashKeySerializer(stringRedisSerializer);

        return template;
    }


}


package com.tehn.zhfljobsapi.config;


import com.fasterxml.jackson.annotation.JsonAutoDetect;
import com.fasterxml.jackson.annotation.PropertyAccessor;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.redisson.Redisson;
import org.redisson.api.RedissonClient;
import org.redisson.codec.JsonJacksonCodec;
import org.redisson.config.Config;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

/**
 * Notes:
 * Author: DurkBlue
 * Date: 2026/4/8/00810:08
 **/
@Configuration
public class RedissonConfig {
    @Value("${spring.redis.host}")
    private String redisHost;
    @Value("${spring.redis.port}")
    private Integer redisPort;
    @Value("${spring.redis.password}")
    private String redisPassword;
    @Value("${spring.redis.database}")
    private Integer database;
    @Bean(destroyMethod="shutdown")
    public RedissonClient redisson() {
        Config config = new Config();
        config.useSingleServer()
                .setAddress("redis://" + redisHost + ":" + redisPort)
                .setPassword(redisPassword) // 如果需要密码
                .setDatabase(database);
        ObjectMapper objectMapper = new ObjectMapper();
        objectMapper.setVisibility(PropertyAccessor.ALL, JsonAutoDetect.Visibility.ANY);
        config.setCodec(new JsonJacksonCodec(objectMapper));
        return Redisson.create(config);
    }
}


    1.任务实体

    

package com.tehn.zhfljobsapi.pojo;

import lombok.Data;

import java.io.Serializable;
/**
 * Notes:
 * Author: DurkBlue
 * Date: 2026/8/14/01415:18
 **/
@Data
public class ExpireNoticeTask implements Serializable {
    /**
     * 唯一任务ID,核心主键(deviceId/业务唯一标识)
     */
    private String taskId;
    /**
     * 业务类型,区分不同通知逻辑
     */
    private String bizType;
    /**
     * 业务参数
     */
    private String bizData;
    /**
     * 任务到期时间戳(ms)
     */
    private long expireTime;
}

    2.Redis 常量

    

package com.tehn.zhfljobsapi.pojo;

/**
 * Notes: Redis 常量
 * Author: DurkBlue
 * Date: 2026/8/14/01415:19
 **/
public class DelayConstant {
    /** zset 存储到期时间戳 */
    public static final String DELAY_ZSET_KEY = "tehn_zhfl_message_center_device_status:expire:zset";
    /** hash存储任务详情 */
    public static final String DELAY_HASH_KEY = "tehn_zhfl_message_center_device_status:expire:hash";
    /** 轮询间隔 单位ms, */
    public static final long POLL_INTERVAL = 1000;
    /** 每次拉取最大数量,防止一次取出巨量任务 */
    public static final int BATCH_SIZE = 200;
}

3.核心延迟服务(Lua 原子续期、新增、删除)

package com.tehn.zhfljobsapi.service.impl;

import com.alibaba.fastjson.JSON;
import com.tehn.zhfljobsapi.pojo.DelayConstant;
import com.tehn.zhfljobsapi.pojo.ExpireNoticeTask;
import lombok.extern.slf4j.Slf4j;
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.data.redis.core.script.DefaultRedisScript;
import org.springframework.stereotype.Service;

import javax.annotation.Resource;
import java.util.Arrays;
import java.util.List;

/**
 * Notes: 核心延迟服务(Lua 原子续期、新增、删除)
 * Author: DurkBlue
 * Date: 2026/8/14/01415:20
 **/

@Slf4j
@Service
public class ZSetDelayNoticeService {

    @Resource
    private StringRedisTemplate stringRedisTemplate;

    /**
     * Lua脚本:新增/更新任务
     * KEY[1] zset key
     * KEY[2] hash key
     * ARGV[1] taskId
     * ARGV[2] 到期时间戳
     * ARGV[3] json任务内容
     */
    private static final String UPSERT_TASK_LUA =
            "redis.call('ZADD', KEYS[1], ARGV[2], ARGV[1]) " +
                    "redis.call('HSET', KEYS[2], ARGV[1], ARGV[3]) " +
                    "return 1";

    /**
     * Lua脚本:删除任务
     */
    private static final String REMOVE_TASK_LUA =
            "redis.call('ZREM', KEYS[1], ARGV[1]) " +
                    "redis.call('HDEL', KEYS[2], ARGV[1]) " +
                    "return 1";

    /**
     * 设置/刷新任务存活时间(新增任务 OR 续期任务统一入口)
     * @param task 任务
     */
    public void setOrRefreshTask(ExpireNoticeTask task) {
        DefaultRedisScript<Long> script = new DefaultRedisScript<>();
        script.setScriptText(UPSERT_TASK_LUA);
        script.setResultType(Long.class);

        List<String> keys = Arrays.asList(DelayConstant.DELAY_ZSET_KEY, DelayConstant.DELAY_HASH_KEY);
        String taskJson = JSON.toJSONString(task);

        stringRedisTemplate.execute(script, keys,
                task.getTaskId(),
                String.valueOf(task.getExpireTime()),
                taskJson
        );
        log.debug("【延迟任务】更新/续期成功 taskId:{} expireTime:{}", task.getTaskId(), task.getExpireTime());
    }

    /**
     * 删除任务(主动取消,未过期不再触发通知)
     */
    public void removeTask(String taskId) {
        DefaultRedisScript<Long> script = new DefaultRedisScript<>();
        script.setScriptText(REMOVE_TASK_LUA);
        script.setResultType(Long.class);

        List<String> keys = Arrays.asList(DelayConstant.DELAY_ZSET_KEY, DelayConstant.DELAY_HASH_KEY);
        stringRedisTemplate.execute(script, keys, taskId);
        log.debug("【延迟任务】移除任务 taskId:{}", taskId);
    }

    /**
     * 获取单个任务详情
     */
    public ExpireNoticeTask getTask(String taskId) {
        Object obj = stringRedisTemplate.opsForHash().get(DelayConstant.DELAY_HASH_KEY, taskId);
        if (obj == null) {
            return null;
        }
        return JSON.parseObject(obj.toString(), ExpireNoticeTask.class);
    }
}

4.到期任务消费调度器

重要:集群部署时,多个实例同时轮询会重复消费!

解决方案:利用你项目已有的 Redisson 分布式锁,同一时刻只有一台实例执行拉取逻辑。


java当中使用Redis ZSet + Lua 脚本自研延迟队列来实现定时到期执行业务处理功能  第1张


package com.tehn.zhfljobsapi.jobs;

/**
 * Notes: 轮询消费器
 * Author: DurkBlue
 * Date: 2026/8/14/01416:22
 **/

import cn.hutool.core.date.DateTime;
import com.tehn.zhfljobsapi.entity.EquipmentInfo;
import com.tehn.zhfljobsapi.entity.EquipmentMaintain;
import com.tehn.zhfljobsapi.enums.RedisKey;
import com.tehn.zhfljobsapi.mapper.EquipmentInfoMapper;
import com.tehn.zhfljobsapi.mapper.EquipmentMaintainMapper;
import com.tehn.zhfljobsapi.pojo.DelayConstant;
import com.tehn.zhfljobsapi.pojo.ExpireNoticeTask;
import com.tehn.zhfljobsapi.service.impl.ZSetDelayNoticeService;
import com.tehn.zhfljobsapi.utils.RedisUtil;
import lombok.extern.slf4j.Slf4j;
import org.redisson.api.RLock;
import org.redisson.api.RedissonClient;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.ApplicationArguments;
import org.springframework.boot.ApplicationRunner;
import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.stereotype.Component;

import javax.annotation.PreDestroy;
import javax.annotation.Resource;
import java.util.Objects;
import java.util.Set;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;

@Slf4j
@Component
public class DelayTaskPollRunner implements ApplicationRunner {

    private static final String POLL_LOCK_KEY = "zhfl:delay:poll:lock";

    @Resource
    private StringRedisTemplate stringRedisTemplate;
    @Resource
    private RedissonClient redissonClient;
    @Resource
    private ZSetDelayNoticeService delayNoticeService;
    @Autowired
    private EquipmentInfoMapper equipmentInfoMapper;
    @Autowired
    private EquipmentMaintainMapper equipmentMaintainMapper;
    @Autowired
    private RedisUtil redisUtil;

    private ScheduledExecutorService scheduler;
    private volatile boolean running = true;

    @Override
    public void run(ApplicationArguments args) {
        scheduler = Executors.newSingleThreadScheduledExecutor();
        // 固定间隔轮询
        scheduler.scheduleAtFixedRate(this::pollExpireTask,
                3,
                DelayConstant.POLL_INTERVAL,
                TimeUnit.MILLISECONDS);
        log.info("ZSet延迟任务轮询器启动成功");
    }

    /**
     * 轮询逻辑
     */
    private void pollExpireTask() {
        if (!running) {
            return;
        }
        RLock lock = redissonClient.getLock(POLL_LOCK_KEY);
        try {
            // 抢占锁,锁自动过期时间略大于轮询间隔,防止死锁
            boolean acquire = lock.tryLock(0, DelayConstant.POLL_INTERVAL + 200, TimeUnit.MILLISECONDS);
            if (!acquire) {
                // 其他节点正在执行,直接跳过
                return;
            }

            long now = System.currentTimeMillis();
            // 查询所有 score <= 当前时间 的任务ID
            Set<String> taskIdSet = stringRedisTemplate.opsForZSet()
                    .rangeByScore(DelayConstant.DELAY_ZSET_KEY, 0, now, 0, DelayConstant.BATCH_SIZE);

            if (taskIdSet == null || taskIdSet.isEmpty()) {
                return;
            }

            for (Object taskIdObj : taskIdSet) {
                String taskId = taskIdObj.toString();
                try {
                    // 查询任务详情
                    ExpireNoticeTask task = delayNoticeService.getTask(taskId);
                    if (task == null) {
                        // 任务已经被手动删除,清理脏数据
                        cleanFinishedTask(taskId);
                        continue;
                    }
                    // 二次校验时间(防止续期后时间被刷新)
                    if (task.getExpireTime() > now) {
                        continue;
                    }
                    // ==========执行业务通知逻辑==========
                    handleExpireNotice(task);
                    // ==================================
                    // 处理完成,清理任务
                    cleanFinishedTask(taskId);
                } catch (Exception e) {
                    log.error("【延迟任务】处理失败 taskId:{}", taskId, e);
                }
            }
        } catch (Exception e) {
            log.error("【延迟任务】轮询异常", e);
        } finally {
            if (lock.isHeldByCurrentThread()) {
                lock.unlock();
            }
        }
    }

    /**
     * 执行业务通知(你在这里写发送通知逻辑)
     */
    private void handleExpireNotice(ExpireNoticeTask task) {
        log.info("【到期触发通知】taskId:{}, bizType:{} data:{}",
                task.getTaskId(), task.getBizType(), task.getBizData());
        // 获取设备对象
        EquipmentInfo equipmentInfo = equipmentInfoMapper.selectById(task.getTaskId());
        if(equipmentInfo.getStatus().equals("3")){
            equipmentInfoMapper.updateStatusByEquipmentId("4", equipmentInfo.getEquipmentId());
            DateTime now = DateTime.now();
            redisUtil.set(RedisKey.REFRESH_KEY_SUBJECT_ID.getKey() + equipmentInfo.getSubjectId(), now.getTime());
            // 离线需要插入维护表
            EquipmentMaintain maintain = new EquipmentMaintain();
            maintain.setEquipmentId(equipmentInfo.getEquipmentId());
            maintain.setEquipmentName(equipmentInfo.getEquipmentName());
            maintain.setFailureReason("掉线");
            maintain.setFailureTime(now);  // 当前时间为故障时间
            maintain.setSubjectId(equipmentInfo.getSubjectId());
            // 插入维护表
            equipmentMaintainMapper.insert(maintain);
        }
    }

    /**
     * 清理已经处理完成的任务
     */
    private void cleanFinishedTask(String taskId) {
        stringRedisTemplate.opsForZSet().remove(DelayConstant.DELAY_ZSET_KEY, taskId);
        stringRedisTemplate.opsForHash().delete(DelayConstant.DELAY_HASH_KEY, taskId);
    }

    @PreDestroy
    public void stop() {
        running = false;
        if (scheduler != null) {
            scheduler.shutdown();
        }
        log.info("延迟任务轮询器停止");
    }
}

5.使用示例

package com.zhfl.msgserver.controller;

import com.zhfl.msgserver.model.delay.ExpireNoticeTask;
import com.zhfl.msgserver.service.delay.ZSetDelayNoticeService;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
import javax.annotation.Resource;

@RestController
public class DelayTestController {

    @Resource
    private ZSetDelayNoticeService delayNoticeService;

    /**
     * 创建任务,10秒后到期通知
     */
    @GetMapping("/delay/add")
    public String addTask(@RequestParam String taskId) {
        ExpireNoticeTask task = new ExpireNoticeTask();
        task.setTaskId(taskId);
        task.setBizType("DEVICE_OFFLINE");
        task.setBizData("{\"deviceId\":\"" + taskId + "\"}");
        // 当前时间 + 10秒
        task.setExpireTime(System.currentTimeMillis() + 10000);
        delayNoticeService.setOrRefreshTask(task);
        return "任务创建成功";
    }

    /**
     * 刷新存活时间,再次延后10秒
     */
    @GetMapping("/delay/refresh")
    public String refreshTask(@RequestParam String taskId) {
        ExpireNoticeTask task = delayNoticeService.getTask(taskId);
        if (task == null) {
            return "任务不存在";
        }
        // 重置到期时间,延后10秒
        task.setExpireTime(System.currentTimeMillis() + 10000);
        delayNoticeService.setOrRefreshTask(task);
        return "续期刷新成功";
    }

    /**
     * 主动取消任务
     */
    @GetMapping("/delay/remove")
    public String remove(@RequestParam String taskId) {
        delayNoticeService.removeTask(taskId);
        return "任务已移除";
    }
}


此篇文章由DurkBlue发布,麻烦转载请注明来处
文章投稿或转载声明

来源:DurkBlue版权归原作者所有,转载请保留出处。本站文章发布于 今天
温馨提示:文章内容系作者个人观点,不代表DurkBlue博客对其观点赞同或支持。

赞(0)

觉得文章有用就打赏一下文章作者

支付宝扫一扫打赏

微信扫一扫打赏

阅读
分享