分布式定时任务防重设计与实践

分布式定时任务防重设计与实践 1. 分布式定时任务重复执行的本质问题我第一次在生产环境遇到Scheduled重复执行问题时整个团队排查到凌晨三点。当时我们的电商促销定时任务在三个节点上同时触发导致优惠券被重复发放直接造成数十万元损失。这个惨痛教训让我深刻认识到在分布式环境中单纯依赖Spring的Scheduled注解就像在高速公路上骑自行车——迟早要出事。定时任务重复执行的本质是多个实例对共享资源如数据库记录、文件状态的竞态访问。当多个服务实例的定时线程同时触发时如果没有协调机制就会形成三头六臂的工作状态。我曾用Arthas监控过一个典型场景三个Pod上的Scheduled方法在10毫秒内相继启动每个都认为自己应该处理当天的数据归档。关键认知误区很多开发者认为只要控制好cron表达式就不会重复。实际上在K8s滚动更新时哪怕只有5秒的重叠期也足够产生重复操作。2. Scheduled在分布式环境中的五大致命陷阱2.1 无状态陷阱自以为是的单机思维最典型的错误就是在定时方法里直接写业务逻辑Scheduled(cron 0 0 3 * * ?) public void generateDailyReport() { // 查询昨天数据 LocalDate yesterday LocalDate.now().minusDays(1); ListOrder orders orderRepo.findByDate(yesterday); // 生成报告 Report report buildReport(orders); reportService.save(report); }这段代码在单机运行时完美工作但在分布式环境下每个实例都会执行查询每个实例都会生成报告数据库最终存入N份相同报告我见过最离谱的案例是某个财务系统每天生成7份相同的报表直到审计时才发现问题。2.2 时间漂移陷阱你以为的同时其实不同步即使所有节点配置相同的cron表达式实际触发时间也可能存在差异Scheduled(cron 0 0/5 * * * ?) // 每5分钟执行 public void syncInventory() { inventoryService.syncFromERP(); }实测数据表明节点A在00:00:00.123触发节点B在00:00:00.456触发节点C在00:00:01.002触发这种微妙的时间差会导致多个节点几乎同时拉取ERP库存每个节点基于不同时间点的数据做计算最终写入结果相互覆盖2.3 异常处理陷阱失败重试变重复执行没有正确处理异常的场景Scheduled(fixedRate 300000) public void processPendingOrders() { try { ListOrder orders orderRepo.findPending(); orders.forEach(this::fulfillOrder); } catch (Exception e) { // 仅打印日志 log.error(处理订单失败, e); } }当数据库连接闪断时节点A获取到10条待处理订单处理到第3条时连接中断节点B立即启动相同流程最终前3条订单被重复处理2.4 持久化陷阱内存标记在重启后失效常见的伪解决方案Scheduled(cron 0 0 1 * * ?) public void archiveOldData() { if (!MemoryCache.get(archive_running)) { MemoryCache.set(archive_running, true); // 执行归档逻辑... MemoryCache.set(archive_running, false); } }这种方案有三个致命缺陷内存状态在应用重启后丢失多个实例的内存缓存不共享没有处理进程崩溃导致的死锁2.5 锁竞争陷阱分布式锁的错误实现看似正确的分布式锁方案Scheduled(fixedDelay 60000) public void sendReminders() { String lockKey reminder_lock; try { if (redisTemplate.opsForValue().setIfAbsent(lockKey, 1, 30, TimeUnit.SECONDS)) { // 发送提醒逻辑... } } finally { redisTemplate.delete(lockKey); } }实际存在的问题任务执行超过30秒会导致锁自动释放多个实例同时获得锁finally块中的删除操作可能误删其他实例的锁3. 工业级解决方案设计与实现3.1 基于ShedLock的防重方案ShedLock是目前最成熟的解决方案之一。这是我们的生产配置// 1. 添加依赖 implementation net.javacrumbs.shedlock:shedlock-spring:4.42.0 implementation net.javacrumbs.shedlock:shedlock-provider-jdbc-template:4.42.0 // 2. 配置LockProvider Bean public LockProvider lockProvider(DataSource dataSource) { return new JdbcTemplateLockProvider( JdbcTemplateLockProvider.Configuration.builder() .withJdbcTemplate(new JdbcTemplate(dataSource)) .usingDbTime() // 使用数据库时间避免时钟漂移 .build() ); } // 3. 注解使用 Scheduled(cron 0 0 2 * * ?) SchedulerLock(name financial_report, lockAtLeastFor 10m, lockAtMostFor 30m) public void generateFinancialReport() { // 复杂的报表生成逻辑 }关键参数说明lockAtLeastFor最短持有时间防止任务执行过快导致锁提前释放lockAtMostFor最大持有时间防止进程崩溃导致死锁我们在金融系统中实测发现锁表记录增加约3ms的额外开销相比重复执行的风险这点开销完全可以接受3.2 基于Redis的原子锁方案对于无法使用JDBC的环境Redis方案更合适Scheduled(fixedRate 300000) public void syncProductPrices() { String lockKey price_sync_lock; String requestId UUID.randomUUID().toString(); try { // 尝试获取锁 Boolean locked redisTemplate.execute( new RedisCallbackBoolean() { Override public Boolean doInRedis(RedisConnection connection) { return connection.set( lockKey.getBytes(), requestId.getBytes(), Expiration.seconds(300), RedisStringCommands.SetOption.SET_IF_ABSENT ); } } ); if (locked ! null locked) { // 真正的业务逻辑 priceService.syncFromSupplier(); } } finally { // 只删除自己设置的锁 String script if redis.call(get, KEYS[1]) ARGV[1] then return redis.call(del, KEYS[1]) else return 0 end; redisTemplate.execute( new DefaultRedisScriptLong(script, Long.class), Collections.singletonList(lockKey), requestId ); } }这个方案的关键改进使用UUID作为请求标识避免误删其他实例的锁采用Lua脚本保证原子性设置合理的过期时间建议比任务周期长20%3.3 数据库乐观锁方案对于数据驱动的定时任务可以结合版本号控制Scheduled(cron 0 0 4 * * ?) Transactional public void calculateStatistics() { // 1. 获取任务记录 OptionalScheduledTask taskOpt taskRepo.findByName(stats_calculation); // 2. 检查状态 if (taskOpt.isPresent() taskOpt.get().getStatus() TaskStatus.RUNNING) { log.warn(任务已在其他节点运行); return; } // 3. 标记为运行中 ScheduledTask task taskOpt.orElse(new ScheduledTask(stats_calculation)); task.setStatus(TaskStatus.RUNNING); task.setStartedAt(LocalDateTime.now()); taskRepo.save(task); try { // 4. 执行业务逻辑 statsService.calculateAll(); // 5. 标记为完成 task.setStatus(TaskStatus.COMPLETED); task.setFinishedAt(LocalDateTime.now()); taskRepo.save(task); } catch (Exception e) { // 6. 标记为失败 task.setStatus(TaskStatus.FAILED); taskRepo.save(task); throw e; } }这个方案的优点不需要额外中间件天然支持任务状态追踪可以通过数据库记录分析历史执行情况4. 生产环境中的进阶实践4.1 任务分片策略当单个任务需要处理大量数据时我们可以结合分片和分布式锁Scheduled(cron 0 0 1 * * ?) public void processBigData() { // 获取当前实例编号(通过K8s环境变量或启动参数) int instanceId Integer.parseInt(System.getenv(POD_INSTANCE_ID)); int totalInstances Integer.parseInt(System.getenv(TOTAL_INSTANCES)); // 获取分布式锁 if (acquireLock(big_data_processing)) { try { // 查询总数据量 long totalCount dataRepo.countUnprocessed(); // 计算分片范围 long chunkSize totalCount / totalInstances; long start instanceId * chunkSize; long end (instanceId totalInstances - 1) ? totalCount : start chunkSize; // 处理分片数据 dataRepo.findUnprocessed(start, end).forEach(this::processItem); } finally { releaseLock(big_data_processing); } } }这种模式特别适合每日用户行为分析大规模数据迁移全量缓存预热4.2 补偿任务设计对于关键任务我们需要实现补偿机制Scheduled(fixedDelay 60000) public void checkStuckTasks() { // 查找运行超过1小时的任务 ListScheduledTask stuckTasks taskRepo.findByStatusAndStartedAtBefore( TaskStatus.RUNNING, LocalDateTime.now().minusHours(1) ); stuckTasks.forEach(task - { log.warn(发现卡住的任务: {}, task.getName()); // 释放锁 if (task.requiresLock()) { lockManager.release(task.getLockName()); } // 更新状态 task.setStatus(TaskStatus.FAILED); task.setErrorMessage(超时自动终止); taskRepo.save(task); // 触发告警 alertService.notifyAdmin(task); }); }4.3 监控与告警配置完善的监控体系应该包括Prometheus指标采集Scheduled(cron 0 * * * * ?) SchedulerLock(name metrics_collection) public void collectMetrics() { // 记录任务执行次数 metrics.counter(scheduled.tasks.execution.count).increment(); // 记录执行时间 Timer.Sample sample Timer.start(); try { // 实际采集逻辑 collectSystemMetrics(); } finally { sample.stop(metrics.timer(scheduled.tasks.duration)); } }Grafana监控看板应包含任务执行成功率平均耗时分布锁等待时间失败任务排行关键告警规则连续3次任务失败任务执行时间超过阈值锁竞争率过高5. 血泪教训我们踩过的那些坑5.1 时钟同步问题曾经有个生产事故我们所有节点都配置了NTP服务但某台物理机的BIOS电池没电了导致系统时间比实际慢10分钟。结果是该节点上的定时任务总是延迟触发当它终于执行时其他节点已经释放了锁最终数据被重复处理解决方案# 在所有节点上配置强制时间同步 sudo timedatectl set-ntp true sudo systemctl restart systemd-timesyncd # 在K8s中配置NTP spec: template: spec: containers: - name: ntp image: cturra/ntp5.2 锁粒度太粗早期我们为整个报表系统使用同一个锁SchedulerLock(name report_system_lock) public void generateAllReports() { // 生成10种不同的报表 }这导致即使报表之间没有依赖也必须串行执行总执行时间超过1小时锁过期导致部分报表重复生成改进后的方案public void generateReport(String reportType) { SchedulerLock(name report_lock_ reportType) void doGenerate() { // 生成单个报表 } // 并行触发不同类型报表 ForkJoinPool.commonPool().submit(() - doGenerate(reportType)); }5.3 未考虑网络分区某次机房网络故障导致Redis主从切换期间出现两个Master节点。结果是节点A在旧Master上获取锁成功节点B在新Master上获取相同的锁也成功两个节点同时处理相同数据最终我们引入了RedLock算法RedissonClient redisson Redisson.create(config); RLock lock redisson.getLock(my_lock); try { // 等待锁最多100秒获得锁后300秒自动释放 if (lock.tryLock(100, 300, TimeUnit.SECONDS)) { // 处理业务 } } finally { lock.unlock(); }5.4 任务幂等性缺失即使有分布式锁也必须实现幂等处理。我们曾遇到任务执行中途JVM崩溃锁自动释放新实例重新获取锁部分数据被处理两次现在的标准做法void processOrder(Order order) { // 先检查处理状态 if (order.getStatus() ProcessStatus.COMPLETED) { return; } // 用乐观锁控制 int updated orderRepo.updateStatus( order.getId(), ProcessStatus.PENDING, ProcessStatus.PROCESSING ); if (updated 0) { return; // 已被其他进程处理 } // 实际业务处理 doRealWork(order); // 标记完成 order.setStatus(ProcessStatus.COMPLETED); orderRepo.save(order); }定时任务在分布式系统中的正确实现远不止添加几个注解那么简单。它需要考虑时钟同步、网络分区、故障恢复等复杂场景。经过多年实践我的建议是小系统可以用数据库锁方案中等规模推荐ShedLock复杂场景考虑专业的任务调度中间件最后记住任何定时任务都必须实现幂等性这是最后的防线。就像我们团队现在的信条——锁可能会失效但业务数据必须正确。