生产级实战:基于Spring Boot + Redis的分布式延迟队列设计与实现

简介: 本文介绍基于Redis的生产级分布式延迟队列方案:摒弃ZSet简单轮询(避免CPU空转与惊群效应),采用分级时间轮(近实时秒级+远时分钟级)、Sorted Set+List+Pub/Sub组合及Lua原子脚本,结合优雅停机、幂等处理与监控告警,显著降低IO压力,保障高并发下的可靠性与稳定性。

 一、为什么不用 Redis ZSet 的 Simple 方案?

很多教程会教你用 ZADD 添加分数(时间戳),然后用 ZRANGEBYSCORE 轮询。这在生产环境有一个致命问题:CPU空转与惊群效应

  • Simple方案:每隔100ms全量扫描Key,即使没有数据也会执行。
  • 生产级方案:使用 Redis Sorted Set + List + Pub/Sub,结合 Time Wheel (时间轮) 思想,大幅降低Redis IO压力。

二、核心设计:分级时间轮 (Hierarchical Time Wheels)

我们将延迟时间分为两级:

  1. 近实时轮 (Near-time Wheel):处理接下来1小时内的任务,精度秒级。
  2. 远时轮 (Far-time Wheel):处理超过1小时的任务,精度分钟级。

数据结构设计:

  • delay:near:{slot} (Sorted Set): score为时间戳,member为jobId。
  • delay:far:{slot} (List): 存储序列化后的Job数据。
  • delay:bucket (List): 暂存区,用于原子化迁移数据。
  • delay:processing (Hash): 正在消费的消息,用于ACK确认。

三、生产级代码实现

1. Maven依赖

<dependencies>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-data-redis</artifactId>
    </dependency>
    <dependency>
        <groupId>redis.clients</groupId>
        <artifactId>jedis</artifactId>
        <version>4.4.3</version>
    </dependency>
    <dependency>
        <groupId>com.fasterxml.jackson.core</groupId>
        <artifactId>jackson-databind</artifactId>
    </dependency></dependencies>

image.gif

2. 核心配置类 (RedisConfig)

为了生产级性能,必须配置连接池和序列化方式。

@Configurationpublic class RedisDelayQueueConfig {    @Bean
    public RedisTemplate<String, Object> redisTemplate(RedisConnectionFactory factory) {
        RedisTemplate<String, Object> template = new RedisTemplate<>();
        template.setConnectionFactory(factory);        
        // 使用Jackson2JsonRedisSerializer替换默认的JDK序列化
        Jackson2JsonRedisSerializer<Object> serializer = 
            new Jackson2JsonRedisSerializer<>(Object.class);        ObjectMapper mapper = new ObjectMapper();
        mapper.setVisibility(PropertyAccessor.ALL, JsonAutoDetect.Visibility.ANY);
        mapper.activateDefaultTyping(LaissezFaireSubTypeValidator.instance, 
                                     ObjectMapper.DefaultTyping.NON_FINAL);
        serializer.setObjectMapper(mapper);
        
        template.setKeySerializer(new StringRedisSerializer());
        template.setValueSerializer(serializer);
        template.setHashKeySerializer(new StringRedisSerializer());
        template.setHashValueSerializer(serializer);
        template.afterPropertiesSet();        return template;
    }    @Bean
    public RedissonClient redissonClient() {        Config config = new Config();
        config.useSingleServer()
              .setAddress("redis://127.0.0.1:6379")
              .setPassword("your_password")
              .setDatabase(0)
              .setConnectionPoolSize(64)
              .setConnectionMinimumIdleSize(10);        return Redisson.create(config);
    }
}

image.gif

3. 延迟任务实体 (DelayJob)

@Data@AllArgsConstructor@NoArgsConstructorpublic class DelayJob implements Serializable {    private static final long serialVersionUID = 1L;    
    /**
     * 任务ID (全局唯一)
     */
    private String jobId;    
    /**
     * 主题 (业务类型,如 ORDER_CLOSE, SMS_SEND)
     */
    private String topic;    
    /**
     * 延迟时间 (时间戳,毫秒)
     */
    private long delayTime;    
    /**
     * 消息体
     */
    private String body;    
    /**
     * 重试次数
     */
    private int retryCount;    
    /**
     * 最大重试次数
     */
    private int maxRetry;
}

image.gif

4. 生产者:投递延迟消息

这里使用了 Lua脚本 保证原子性,这是生产级代码的标配。

@Component@Slf4jpublic class DelayProducer {    
    @Autowired
    private RedisTemplate<String, Object> redisTemplate;    
    private static final String DELAY_NEAR_PREFIX = "delay:near:";    private static final String DELAY_FAR_PREFIX = "delay:far:";    
    // Lua脚本:防止重复投递
    private static final String PUSH_LUA =
        "if redis.call('EXISTS', KEYS[1]) == 0 then " +        "  redis.call('ZADD', KEYS[1], ARGV[1], ARGV[2]) " +        "  return 1 " +        "else " +        "  return 0 " +        "end";    
    private DefaultRedisScript<Long> redisScript;    
    @PostConstruct
    public void init() {
        redisScript = new DefaultRedisScript<>();
        redisScript.setScriptText(PUSH_LUA);
        redisScript.setResultType(Long.class);
    }    
    /**
     * 投递延迟任务
     */
    public boolean send(DelayJob job) {        long now = System.currentTimeMillis();        long delaySeconds = job.getDelayTime() / 1000;
        String key;        
        if (delaySeconds <= now + 3600) { // 1小时内进近时轮
            key = DELAY_NEAR_PREFIX + (delaySeconds % 60); // 按分钟分槽
            Long result = redisTemplate.execute(redisScript, 
                Collections.singletonList(key), 
                String.valueOf(job.getDelayTime()), 
                JSON.toJSONString(job));            return result != null && result == 1;
        } else { // 超过1小时进远时轮
            key = DELAY_FAR_PREFIX + (delaySeconds / 60 % 60); // 按小时分槽
            redisTemplate.opsForList().rightPush(key, job);            return true;
        }
    }
}

image.gif

5. 消费者:时间轮驱动与消息处理

这是最核心的部分。我们需要一个后台线程不断扫描时间轮,并将到期的任务转移到就绪队列。

@Component@Slf4jpublic class DelayConsumer implements ApplicationRunner {    
    @Autowired
    private RedisTemplate<String, Object> redisTemplate;    
    @Autowired
    private DelayProducer producer;    
    private static final String READY_QUEUE = "delay:ready";    private static final String PROCESSING_HASH = "delay:processing";    private static final String NEAR_PREFIX = "delay:near:";    
    // 定时任务线程池
    private final ScheduledExecutorService scheduler = 
        Executors.newScheduledThreadPool(4);    
    @Override
    public void run(ApplicationArguments args) {        // 启动近时轮扫描器
        scheduler.scheduleAtFixedRate(this::scanNearTimeWheel, 0, 1, TimeUnit.SECONDS);        // 启动远时轮迁移器
        scheduler.scheduleAtFixedRate(this::migrateFarToNear, 0, 30, TimeUnit.SECONDS);        // 启动消息处理器
        scheduler.scheduleAtFixedRate(this::handleReadyMessages, 0, 500, TimeUnit.MILLISECONDS);
        log.info("Delay queue consumers started...");
    }    
    /**
     * 扫描近时轮
     */
    private void scanNearTimeWheel() {        long now = System.currentTimeMillis();        int slot = (int) (now / 1000 % 60);        String key = NEAR_PREFIX + slot;        
        try {            // 获取所有到期的任务
            Set<String> jobs = redisTemplate.opsForZSet()
                .rangeByScore(key, 0, now);            
            if (jobs != null && !jobs.isEmpty()) {                for (String jobStr : jobs) {                    DelayJob job = JSON.parseObject(jobStr, DelayJob.class);                    // 原子化移动到就绪队列
                    moveToReadyQueue(job, key, jobStr);
                }
            }
        } catch (Exception e) {
            log.error("Scan near time wheel error", e);
        }
    }    
    /**
     * 原子化移动任务到就绪队列 (Lua脚本)
     */
    private static final String MOVE_TO_READY_LUA =
        "if redis.call('ZREM', KEYS[1], ARGV[1]) == 1 then " +        "  redis.call('LPUSH', KEYS[2], ARGV[2]) " +        "  return 1 " +        "else " +        "  return 0 " +        "end";    
    private void moveToReadyQueue(DelayJob job, String zsetKey, String jobStr) {
        DefaultRedisScript<Long> script = new DefaultRedisScript<>(MOVE_TO_READY_LUA, Long.class);        Long result = redisTemplate.execute(script, 
            Arrays.asList(zsetKey, READY_QUEUE), 
            jobStr, jobStr);        
        if (result != null && result == 1) {
            log.debug("Job moved to ready queue: {}", job.getJobId());
        }
    }    
    /**
     * 处理就绪队列中的消息
     */
    private void handleReadyMessages() {        try {            // 阻塞式弹出,防止CPU空转
            List<Object> objs = redisTemplate.executePipelined((RedisCallback<Object>) connection -> {
                connection.listCommands().bLPop(2, READY_QUEUE.getBytes());                return null;
            });            
            if (objs != null && !objs.isEmpty()) {                Object raw = objs.get(0);                if (raw instanceof byte[]) {                    DelayJob job = JSON.parseObject((byte[]) raw, DelayJob.class);
                    processJob(job);
                }
            }
        } catch (Exception e) {
            log.error("Handle ready messages error", e);
        }
    }    
    /**
     * 具体的业务逻辑处理
     */
    private void processJob(DelayJob job) {        String lockKey = "lock:delay:" + job.getJobId();        Boolean locked = redisTemplate.opsForValue()
            .setIfAbsent(lockKey, "1", 30, TimeUnit.SECONDS);        
        if (Boolean.TRUE.equals(locked)) {            try {
                log.info("Processing job: {}, Topic: {}", job.getJobId(), job.getTopic());                
                // TODO: 这里是你的具体业务代码
                // 例如:检查订单状态,如果未支付则关闭
                boolean success = doBusinessLogic(job);                
                if (!success && job.getRetryCount() < job.getMaxRetry()) {                    // 重试逻辑:指数退避
                    job.setRetryCount(job.getRetryCount() + 1);                    long nextDelay = (long) (Math.pow(2, job.getRetryCount()) * 1000);
                    job.setDelayTime(System.currentTimeMillis() + nextDelay);
                    producer.send(job);
                }
            } finally {
                redisTemplate.delete(lockKey);
            }
        }
    }    
    private boolean doBusinessLogic(DelayJob job) {        // 模拟业务处理
        return true;
    }    
    /**
     * 将远时轮任务迁移到近时轮
     */
    private void migrateFarToNear() {        // ... 省略具体实现,逻辑类似scanNearTimeWheel,但针对delay:far:*进行操作
        log.debug("Migrating far time wheel tasks...");
    }
}

image.gif

6. 优雅停机与健康检查

生产环境必须考虑JVM退出时的任务回收,否则会导致消息丢失。

@Component@Slf4jpublic class DelayShutdownHook {    
    @PreDestroy
    public void destroy() {
        log.info("Shutting down delay queue consumer...");        // 1. 停止接收新任务
        // 2. 等待当前正在处理的任务完成 (通常需要一个计数器)
        // 3. 将processing中的任务重新放回delay队列
        recoverProcessingJobs();
        log.info("Delay queue shutdown completed.");
    }    
    private void recoverProcessingJobs() {        // 读取 processing hash
        // 遍历并重新投递
    }
}

image.gif

四、生产环境优化建议 (大牛视角)

  1. 幂等性设计
  • 所有的消费者逻辑必须支持幂等。因为网络抖动或重试机制,同一条消息可能会被处理两次。建议使用 jobId 作为唯一键,配合数据库的唯一索引或Redis的 SETNX
  1. 内存与持久化
  • AOF: 务必开启 appendfsync everysec,防止Redis宕机导致大量延迟任务丢失。
  • Key过期: 对于 processing Hash中的字段,设置合理的TTL,防止死信堆积。
  1. 监控告警
  • 监控 delay:ready 的长度。如果长度持续增长,说明消费能力不足,需要扩容消费者。
  • 监控 delay:processing 的长度。如果长度过大,说明有大量任务处理超时或被卡住。
  1. 集群模式下的分布式锁
  • 本文使用的是单Redis实例的简单锁。在生产集群环境中,强烈建议使用 RedissonRLock(红锁)或 Zookeeper 来实现分布式选主,确保同一时间只有一个消费者线程在执行 scanNearTimeWheel,避免资源浪费和惊群效应。

五、总结

这套方案相比简单的 ZRANGEBYSCORE,引入了时间轮分槽双缓冲队列,极大地降低了Redis的CPU负载。同时,通过 Lua脚本 保证了操作的原子性,通过 分布式锁重试机制 保证了消息的可靠性。

在高并发场景下(如QPS 10k+),此架构已在多个生产项目中验证稳定。如果你面临更高的吞吐需求,可以考虑将 Ready Queue 替换为 Kafka,由 Redis 负责调度,Kafka 负责削峰填谷。

本文由 摸鱼不慌 发布,转载请注明出处。

文章链接:生产级实战:基于Spring Boot + Redis的分布式延迟队列设计与实现 - 摸鱼不慌

目录
相关文章
|
Python
Python 游戏开发实战:从入门到精通
本文介绍利用Python与Pygame库进行游戏开发的基础知识。Pygame是专为游戏设计的Python库,提供了丰富的功能简化游戏开发流程。文中首先指导读者完成Pygame库的安装,并通过示例代码演示了游戏窗口创建、基本图形绘制及用户输入处理等核心概念。此外,还展示了如何通过定义类来组织游戏对象,帮助读者更高效地管理游戏代码。适合初学者入门Python游戏开发。
1335 2
|
6天前
|
缓存 人工智能 API
阿里云百炼deepseek-v4-flash模型介绍:模型特点、适用场景、最新优惠及部署流程参考
本文全面解析了阿里云百炼平台托管的DeepSeek-V4-Flash大模型的核心参数与使用指南。这款总参284B、激活13B的轻量化MoE模型,原生支持百万级超长上下文,最大输出长度可达39万+Tokens,推理速度快、调用成本低,适配日常对话、批量文案处理、基础RAG等高并发普惠场景。文章同步梳理了北京、新加坡、法兰克福等全球5大部署节点的能力支持情况、分区域计费标准与限流规则,同时标注了预览版与2026年7月31日正式稳定版的版本差异,帮助开发者快速完成选型与API集成。
|
6天前
|
人工智能 编解码 缓存
Timeline Studio 重大升级:安装一个 Skill,让 AI 在一杯咖啡的时间里完成视频剪辑,并交付可编辑工程
Timeline Studio 是一款AI视频编辑工具,支持从图片/视频/网址自动生成专业视频,并交付可编辑的`.timeline`工程文件。它兼顾自动剪辑与人工精修,实现“一杯咖啡成片,随时局部修改”。开源免费,浏览器本地运行。
|
8月前
|
监控 安全 Unix
iOS 崩溃排查不再靠猜!这份分层捕获指南请收好
从 Mach 内核异常到 NSException,从堆栈遍历到僵尸对象检测,阿里云 RUM iOS SDK 基于 KSCrash 构建了一套完整、异步安全、生产可用的崩溃捕获体系,让每一个线上崩溃都能被精准定位。
2453 158
|
23天前
|
人工智能
Qwen3.8抢先体验!正式版即将发布并开源!
千问Qwen3.8即将开源,参数达2.4T,进化速度以“天”计,实力媲美Fable 5。预览版Qwen3.8-Max已上线阿里Token Plan等平台,限时优惠:日间Credits低至1折,夜间更优,个人/团队版月付仅35元起!
2017 130
|
13天前
|
人工智能 自然语言处理 文字识别
阿里云通义千问大模型最新功能介绍
阿里云通义千问是通义实验室打造的全栈式大模型家族,覆盖通用对话、长文本处理、多模态理解、代码生成、智能体执行等全场景能力,形成从入门到旗舰的完整产品矩阵。最新迭代的Qwen3系列,以**百万级上下文、原生多模态、全链路智能体**为核心标签,彻底突破传统大模型“对话工具”的局限,转向**自主思考、工具调用、任务执行**的智能体时代。
3094 1
|
17天前
|
人工智能 自然语言处理 开发工具
最新版通义千问(Qwen3.8-Max-Preview)功能介绍
通义千问Qwen3.8-Max-Preview是通义千问团队推出的旗舰级预览版大模型,以2.4万亿参数的MoE混合专家架构为核心,实现了原生多模态融合、超长上下文处理、全栈代码工程、多智能体协同等能力的跨越式升级,成为面向复杂生产场景的全域生产力模型。该模型不仅在参数规模上实现突破,更通过架构优化、能力重构,解决了传统大模型在长文本处理、复杂推理、工程落地中的诸多痛点,为开发者、企业用户提供了更强大、更高效、更灵活的AI能力支撑。
10994 4
|
7天前
|
编解码 人工智能 安全
2核4G/4核8G/8核16G阿里云服务器如何选择实例?经济型e、通用算力型u2i与计算型c9i选哪个?
本文介绍了阿里云2核4G、4核8G、8核16G三档主流配置下经济型e、通用算力型u2i和计算型c9i三种实例的最新活动价格与适用场景。同配置下三者价差显著,以2核4G为例,经济型e低至599.93元/年,计算型c9i则高达1742.08元/年。文章详细解析了各实例的性能定位:经济型e适合轻负载入门场景,u2i兼顾稳定算力与性价比,c9i凭借第9代至强处理器与芯片级安全能力支撑高性能业务。同时提示用户可叠加满减优惠券享受折上折,建议根据业务负载与预算综合决策。
538 112
|
11天前
|
存储 人工智能 关系型数据库
阿里云AI产品与云产品最新组合套餐:Token Plan、AI coding及云服务器和建站等组合优惠价
阿里云推出全新“算力+模型+应用”一站式云与AI组合套餐活动,覆盖从个人开发者到中大型企业的全场景需求。核心亮点为分三档定价的Token Plan订阅服务,支持Qwen3.8-Max-Preview大模型调用,错峰时段最低可享0.2折优惠。活动同步推出AI Coding、智能体部署、云电脑托管、0代码建站等十余类场景化组合,搭配99元/年的普惠云服务器、88元/年的入门数据库等经典特惠产品,还为企业提供1V1定制化AI转型方案,大幅降低了不同用户群体拥抱AI的技术门槛与采购成本。
730 111
|
5天前
|
存储 运维 安全
从0到1搭建企业网盘:架构设计、技术选型与落地实践全解析
本文从CTO视角系统解析企业网盘建设全链路:涵盖需求分析、四层架构(接入/服务/存储/检索)、RBAC权限、异构存储、混合检索、RAG智能增强及安全合规等核心环节,兼顾技术深度与落地实践。(239字)
52 3