SpringBoot使用@Scheduled 和 Redis分布式锁实现分布式定时任务

'# SpringBoot使用@Scheduled 和 Redis分布式锁实现分布式定时任务

一、背景与问题

在微服务架构中,定时任务常常需要在多个服务实例之间协调执行。传统SpringBoot的@Scheduled注解仅适用于单机环境,当部署到分布式集群时会出现以下问题:

  1. 多实例同时执行导致数据重复处理
  2. 任务执行过程中服务宕机导致任务丢失
  3. 任务执行时间过长导致资源争用

为解决这些问题,我们需要引入分布式锁机制。Redis的分布式锁可以保证同一时刻只有一个实例执行任务,结合@Scheduled的定时触发能力,形成可靠的分布式定时任务系统。

二、基本原理

1. @Scheduled 的工作原理

Spring的@Scheduled注解通过TaskScheduler实现定时任务调度。其核心机制是:

  • 使用CronTrigger解析cron表达式
  • 调用TaskScheduler的schedule()方法注册任务
  • 通过Thread管理任务执行线程

在单机环境下,这可以保证任务按计划执行。但在集群环境中,多个实例会同时触发任务。

2. Redis分布式锁原理

Redis分布式锁的核心是通过SETNX命令(或SET的NX选项)实现锁的获取:

String lockKey = "my_lock";
String requestId = UUID.randomUUID().toString();
int expireTime = 30; // 锁过期时间(秒)

// 获取锁
Boolean isLocked = redisTemplate.opsForValue().setIfAbsent(lockKey, requestId, expireTime, TimeUnit.SECONDS);

if (isLocked) {
    try {
        // 执行业务逻辑
    } finally {
        // 释放锁
        if (redisTemplate.opsForValue().get(lockKey).equals(requestId)) {
            redisTemplate.delete(lockKey);
        }
    }
}

关键点:

  • setIfAbsent原子操作保证锁的获取是原子的
  • 设置合理的过期时间防止死锁
  • 释放锁时需要校验请求ID防止误删

三、环境准备

1. 依赖配置

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-data-redis</artifactId>
</dependency>
<dependency>
    <groupId>io.lettuce</groupId>
    <artifactId>lettuce-core</artifactId>
</dependency>

2. Redis配置

spring:
  redis:
    host: 127.0.0.1
    port: 6379
    lettuce:
      pool:
        max-active: 8
        max-idle: 8
        min-idle: 2
        max-wait: 1000ms

四、核心实现

1. 基础定时任务

@Scheduled(cron = "0 0 1 * * ?")
public void basicTask() {
    System.out.println("执行基础定时任务:" + LocalDateTime.now());
}

2. 带分布式锁的定时任务

@Scheduled(cron = "0 0 1 * * ?")
public void lockedTask() {
    String lockKey = "scheduled_task_lock";
    String requestId = UUID.randomUUID().toString();
    int expireTime = 30;
    
    try {
        Boolean isLocked = redisTemplate.opsForValue().setIfAbsent(
            lockKey, requestId, expireTime, TimeUnit.SECONDS);
        
        if (isLocked) {
            try {
                System.out.println("执行带锁的定时任务:" + LocalDateTime.now());
                // 模拟耗时操作
                Thread.sleep(5000);
                // 模拟业务逻辑
                System.out.println("任务完成");
            } catch (Exception e) {
                System.err.println("任务执行异常:" + e.getMessage());
            }
        } else {
            System.out.println("任务已加锁,跳过执行");
        }
    } finally {
        String currentId = redisTemplate.opsForValue().get(lockKey);
        if (currentId != null && currentId.equals(requestId)) {
            redisTemplate.delete(lockKey);
        }
    }
}

3. 优化版锁实现(带重试机制)

public boolean tryLock(String lockKey, String requestId, int expireTime, int retryCount) {
    int retry = 0;
    while (retry < retryCount) {
        Boolean isLocked = redisTemplate.opsForValue().setIfAbsent(
            lockKey, requestId, expireTime, TimeUnit.SECONDS);
        
        if (isLocked) {
            return true;
        }
        
        retry++;
        try {
            Thread.sleep(100); // 短暂等待后重试
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            return false;
        }
    }
    return false;
}

五、完整案例

1. 订单处理定时任务

@Component
public class OrderProcessor {

    @Autowired
    private RedisTemplate<String, String> redisTemplate;
    
    @Scheduled(cron = "0 0 1 * * ?")
    public void processOrders() {
        String lockKey = "order_process_lock";
        String requestId = UUID.randomUUID().toString();
        int expireTime = 60;
        
        if (tryLock(lockKey, requestId, expireTime, 3)) {
            try {
                // 模拟处理订单
                List<String> orderIds = getUnprocessedOrders();
                for (String orderId : orderIds) {
                    processOrder(orderId);
                }
            } catch (Exception e) {
                System.err.println("订单处理异常:" + e.getMessage());
            } finally {
                releaseLock(lockKey, requestId);
            }
        }
    }
    
    private List<String> getUnprocessedOrders() {
        // 模拟从数据库获取未处理订单
        return Arrays.asList("order_1", "order_2", "order_3");
    }
    
    private void processOrder(String orderId) {
        // 模拟处理逻辑
        System.out.println("处理订单:" + orderId);
        try {
            Thread.sleep(1000);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
    
    private boolean tryLock(String lockKey, String requestId, int expireTime, int retryCount) {
        int retry = 0;
        while (retry < retryCount) {
            Boolean isLocked = redisTemplate.opsForValue().setIfAbsent(
                lockKey, requestId, expireTime, TimeUnit.SECONDS);
            
            if (isLocked) {
                return true;
            }
            
            retry++;
            try {
                Thread.sleep(100);
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
                return false;
            }
        }
        return false;
    }
    
    private void releaseLock(String lockKey, String requestId) {
        String currentId = redisTemplate.opsForValue().get(lockKey);
        if (currentId != null && currentId.equals(requestId)) {
            redisTemplate.delete(lockKey);
        }
    }
}

六、源码解析

1. Redis锁的原子性保证

Redis的SETNX命令是原子操作,确保在多线程环境下:

Boolean isLocked = redisTemplate.opsForValue().setIfAbsent(
    lockKey, requestId, expireTime, TimeUnit.SECONDS);

这个操作会同时完成:

  1. 检查键是否存在
  2. 如果不存在则设置键值
  3. 自动设置过期时间

2. 锁释放的条件判断

String currentId = redisTemplate.opsForValue().get(lockKey);
if (currentId != null && currentId.equals(requestId)) {
    redisTemplate.delete(lockKey);
}

必须校验当前锁的请求ID是否与当前线程的ID一致,否则可能误删其他线程的锁。

七、进阶使用

1. 动态锁失效时间

int expireTime = Math.min(30, (int) (System.currentTimeMillis() / 1000) + 60);

根据当前时间动态计算锁的过期时间,防止任务执行时间过长导致锁提前失效。

2. 多级锁机制

String lockKey = "order_process_lock_" + orderId;

对不同订单使用不同的锁,提高并发度。

3. 带重试机制的锁获取

public boolean tryLockWithRetry(String lockKey, String requestId, int expireTime, int retryCount) {
    int retry = 0;
    while (retry < retryCount) {
        Boolean isLocked = redisTemplate.opsForValue().setIfAbsent(
            lockKey, requestId, expireTime, TimeUnit.SECONDS);
        
        if (isLocked) {
            return true;
        }
        
        retry++;
        try {
            Thread.sleep(100);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            return false;
        }
    }
    return false;
}

八、性能与工程实践

1. 性能优化

  • 锁粒度控制:避免锁范围过大,如使用业务ID作为锁键
  • 锁过期时间:设置合理的过期时间(建议30-60秒)
  • 异步处理:将耗时操作放入队列异步处理
  • 缓存结果:对重复任务进行结果缓存

2. 异常处理

catch (RedisException e) {
    System.err.println("Redis连接异常:" + e.getMessage());
    // 尝试重连
}

3. 安全考虑

  • 防止锁误删:严格校验请求ID
  • 防止锁竞争:设置合理的锁等待时间
  • 防止死锁:设置锁过期时间

九、常见问题与踩坑

1. 锁未释放

错误代码:

redisTemplate.delete(lockKey);

问题:未校验锁的请求ID

解决:必须校验当前锁的请求ID是否与当前线程一致

2. 锁失效时间设置不当

错误场景:任务执行时间超过锁的过期时间

解决:动态计算锁的过期时间或延长锁的过期时间

3. 锁竞争导致任务堆积

错误场景:多个实例同时获取锁,导致任务执行顺序混乱

解决:使用更细粒度的锁,或采用队列调度机制

4. Redis连接问题

错误场景:Redis服务宕机导致锁失效

解决:设置合理的重试机制和断线处理逻辑

十、最佳实践

  1. 锁粒度控制:按业务实体划分锁,避免全局锁
  2. 锁过期时间:设置合理的过期时间(建议30-60秒)
  3. 异常处理:添加完善的异常捕获和重试机制
  4. 日志记录:记录锁的获取和释放日志,便于排查问题
  5. 监控报警:对锁的获取失败情况设置监控报警

十一、总结

通过结合SpringBoot的@Scheduled定时任务和Redis分布式锁,可以构建可靠的分布式定时任务系统。在实际开发中需要注意:

  • 适用场景:需要精确控制任务执行顺序、处理关键业务的场景
  • 不适用场景:任务执行时间极短、对一致性要求不高的场景

本方案通过Redis的原子操作和锁管理,有效解决了分布式环境下的定时任务问题。在实现时需要注意锁的粒度、过期时间设置、异常处理等关键点,通过合理的架构设计和代码实现,可以构建稳定可靠的分布式定时任务系统。

评论已关闭

推荐阅读

AIGC实战——Transformer模型
2024年12月01日
Socket TCP 和 UDP 编程基础(Python)
2024年11月30日
python , tcp , udp
如何使用 ChatGPT 进行学术润色?你需要这些指令
2024年12月01日
AI
最新 Python 调用 OpenAi 详细教程实现问答、图像合成、图像理解、语音合成、语音识别(详细教程)
2024年11月24日
ChatGPT 和 DALL·E 2 配合生成故事绘本
2024年12月01日
omegaconf,一个超强的 Python 库!
2024年11月24日
【视觉AIGC识别】误差特征、人脸伪造检测、其他类型假图检测
2024年12月01日
[超级详细]如何在深度学习训练模型过程中使用 GPU 加速
2024年11月29日
Python 物理引擎pymunk最完整教程
2024年11月27日
MediaPipe 人体姿态与手指关键点检测教程
2024年11月27日
深入了解 Taipy:Python 打造 Web 应用的全面教程
2024年11月26日
基于Transformer的时间序列预测模型
2024年11月25日
Python在金融大数据分析中的AI应用(股价分析、量化交易)实战
2024年11月25日
AIGC Gradio系列学习教程之Components
2024年12月01日
Python3 `asyncio` — 异步 I/O,事件循环和并发工具
2024年11月30日
llama-factory SFT系列教程:大模型在自定义数据集 LoRA 训练与部署
2024年12月01日
Python 多线程和多进程用法
2024年11月24日
Python socket详解,全网最全教程
2024年11月27日
python之plot()和subplot()画图
2024年11月26日
理解 DALL·E 2、Stable Diffusion 和 Midjourney 工作原理
2024年12月01日