分布式定时任务:Quartz、xxl-job 与 ElasticJob

对比 Spring Scheduler、Quartz、xxl-job 与 ElasticJob,掌握分布式环境下的定时任务调度、分片执行与故障转移

定时任务是业务系统中常见的需求,如数据同步、报表生成、缓存预热、订单超时取消等。在单机环境下 Spring 的 @Scheduled 足够使用,但在分布式环境中,必须解决任务重复执行、集群调度、分片处理与故障转移等问题。

一、定时任务方案对比

特性Spring SchedulerQuartzxxl-jobElasticJob
分布式支持不支持支持(JDBC/Redis)原生支持原生支持
管理控制台需自研完善完善
任务分片不支持不支持支持原生支持
弹性扩容不支持不支持支持支持
失败重试不支持支持支持支持
触发类型固定频率/CronCronCron/固定间隔/APICron/API
学习成本
适用场景单体/简单任务企业级调度中小型分布式大数据量分片

二、Spring Scheduler 基础

2.1 快速入门

@Configuration
@EnableScheduling
public class SchedulingConfig {
    
    @Scheduled(fixedRate = 5000)           // 每 5 秒执行(从上一次开始)
    public void fixedRateTask() {
        log.info("Fixed rate task: {}", LocalDateTime.now());
    }
    
    @Scheduled(fixedDelay = 5000)          // 每 5 秒执行(从上一次结束)
    public void fixedDelayTask() {
        log.info("Fixed delay task: {}", LocalDateTime.now());
    }
    
    @Scheduled(cron = "0 0 2 * * ?")       // 每天凌晨 2 点
    public void cronTask() {
        log.info("Daily cleanup task");
    }
    
    @Scheduled(cron = "0 0/10 9-18 * * MON-FRI")  // 工作日 9-18 点每 10 分钟
    public void workHourTask() {
        log.info("Working hours task");
    }
}

2.2 线程池配置

@Configuration
public class AsyncConfig {
    
    @Bean("taskScheduler")
    public TaskScheduler taskScheduler() {
        ThreadPoolTaskScheduler scheduler = new ThreadPoolTaskScheduler();
        scheduler.setPoolSize(10);
        scheduler.setThreadNamePrefix("task-scheduler-");
        scheduler.setAwaitTerminationSeconds(60);
        scheduler.setWaitForTasksToCompleteOnShutdown(true);
        scheduler.setRemoveOnCancelPolicy(true);
        scheduler.setErrorHandler(t -> log.error("调度任务异常", t));
        return scheduler;
    }
}

2.3 动态调度

@Service
public class DynamicSchedulerService {
    
    @Autowired
    private TaskScheduler taskScheduler;
    
    private final Map<String, ScheduledFuture<?>> tasks = new ConcurrentHashMap<>();
    
    public void scheduleTask(String taskId, Runnable task, String cron) {
        cancelTask(taskId);
        
        CronTrigger trigger = new CronTrigger(cron);
        ScheduledFuture<?> future = taskScheduler.schedule(task, trigger);
        tasks.put(taskId, future);
    }
    
    public void cancelTask(String taskId) {
        ScheduledFuture<?> future = tasks.remove(taskId);
        if (future != null) {
            future.cancel(false);
        }
    }
}

三、Quartz 企业级调度

3.1 核心概念

Job(任务)      → 具体执行逻辑,实现 Job 接口
Trigger(触发器) → 定义执行规则(时间、频率)
Scheduler(调度器)→ 管理 Job 和 Trigger 的执行
JobDetail        → Job 的详细定义(含 JobDataMap)

3.2 Spring Boot 集成

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-quartz</artifactId>
</dependency>
spring:
  quartz:
    job-store-type: jdbc            # 持久化到数据库
    jdbc:
      initialize-schema: never      # 手动初始化表结构
    properties:
      org:
        quartz:
          scheduler:
            instanceName: clusteredScheduler
            instanceId: AUTO
          jobStore:
            class: org.quartz.impl.jdbcjobstore.JobStoreTX
            driverDelegateClass: org.quartz.impl.jdbcjobstore.StdJDBCDelegate
            tablePrefix: QRTZ_
            isClustered: true         # 集群模式
            clusterCheckinInterval: 20000
          threadPool:
            class: org.quartz.simpl.SimpleThreadPool
            threadCount: 10

3.3 Job 定义与触发器

@DisallowConcurrentExecution      // 禁止并发执行同一 Job
@PersistJobDataAfterExecution     // 持久化 JobDataMap
public class OrderCleanupJob implements Job {
    
    @Autowired
    private OrderService orderService;
    
    @Override
    public void execute(JobExecutionContext context) throws JobExecutionException {
        JobDataMap data = context.getMergedJobDataMap();
        int timeoutHours = data.getInt("timeoutHours");
        
        log.info("执行订单清理任务,超时时间: {} 小时", timeoutHours);
        int count = orderService.cancelTimeoutOrders(timeoutHours);
        log.info("取消超时订单 {} 个", count);
        
        // 更新执行次数到 JobDataMap
        int execCount = data.getInt("execCount") + 1;
        context.getJobDetail().getJobDataMap().put("execCount", execCount);
    }
}

@Service
public class QuartzJobService {
    
    @Autowired
    private Scheduler scheduler;
    
    public void scheduleOrderCleanup() throws SchedulerException {
        JobDetail job = JobBuilder.newJob(OrderCleanupJob.class)
            .withIdentity("orderCleanup", "orderGroup")
            .usingJobData("timeoutHours", 24)
            .usingJobData("execCount", 0)
            .build();
        
        Trigger trigger = TriggerBuilder.newTrigger()
            .withIdentity("orderCleanupTrigger", "orderGroup")
            .withSchedule(CronScheduleBuilder
                .cronSchedule("0 0/30 * * * ?")   // 每 30 分钟
                .withMisfireHandlingInstructionFireAndProceed())
            .build();
        
        scheduler.scheduleJob(job, trigger);
    }
}

3.4 监听器

@Component
public class GlobalJobListener implements JobListener {
    
    @Override
    public String getName() {
        return "globalJobListener";
    }
    
    @Override
    public void jobToBeExecuted(JobExecutionContext context) {
        log.info("Job 即将执行: {}", context.getJobDetail().getKey());
    }
    
    @Override
    public void jobExecutionVetoed(JobExecutionContext context) {
        log.warn("Job 被否决: {}", context.getJobDetail().getKey());
    }
    
    @Override
    public void jobWasExecuted(JobExecutionContext context, JobExecutionException jobException) {
        long cost = System.currentTimeMillis() - context.getFireTime().getTime();
        log.info("Job 执行完成: {}, 耗时: {}ms", context.getJobDetail().getKey(), cost);
        
        if (jobException != null) {
            alertService.sendAlert("Job 执行异常", jobException.getMessage());
        }
    }
}

四、xxl-job 分布式任务调度

4.1 架构设计

┌──────────────────────────────────────────┐
│            xxl-job-admin                  │  ← 调度中心(独立部署)
│  - 任务管理、调度触发                      │
│  - 执行器管理、日志查看                    │
│  - 报警通知                               │
└──────────────────────────────────────────┘
                      │
         ┌────────────┼────────────┐
         ▼            ▼            ▼
┌─────────────┐ ┌─────────────┐ ┌─────────────┐
│  Executor-1  │ │  Executor-2  │ │  Executor-3  │  ← 执行器(业务服务)
│  执行分片-1  │ │  执行分片-2  │ │  执行分片-3  │
└─────────────┘ └─────────────┘ └─────────────┘

4.2 执行器集成

xxl:
  job:
    admin:
      addresses: http://xxl-job-admin:8080/xxl-job-admin
    executor:
      appname: order-service-executor
      ip:                                 # 自动获取
      port: 9999                          # 执行器端口
      logpath: /data/applogs/xxl-job
      logretentiondays: 30
    accessToken: ${XXL_JOB_TOKEN:}        # 安全认证
@Component
public class XxlJobHandler {
    
    @Autowired
    private OrderService orderService;
    
    @XxlJob("cancelTimeoutOrder")
    public ReturnT<String> cancelTimeoutOrder() {
        String param = XxlJobHelper.getJobParam();
        int timeout = param != null ? Integer.parseInt(param) : 24;
        
        int count = orderService.cancelTimeoutOrders(timeout);
        XxlJobHelper.log("取消超时订单 {} 个", count);
        
        return ReturnT.SUCCESS;
    }
    
    @XxlJob("syncInventoryJob")
    public ReturnT<String> syncInventory() {
        // 分片参数
        int shardIndex = XxlJobHelper.getShardIndex();
        int shardTotal = XxlJobHelper.getShardTotal();
        
        XxlJobHelper.log("分片参数: index={}, total={}", shardIndex, shardTotal);
        
        // 按分片取模执行
        List<Long> skuIds = inventoryService.getSkuIdsByShard(shardIndex, shardTotal);
        for (Long skuId : skuIds) {
            inventoryService.syncToRedis(skuId);
        }
        
        return ReturnT.SUCCESS;
    }
    
    @XxlJob("dailyReportJob")
    public void dailyReport() {
        // 新版注解支持 void 返回
        String today = LocalDate.now().toString();
        Report report = reportService.generateDailyReport(today);
        
        // 保存执行日志
        XxlJobHelper.log("生成日报完成: {}", report.getFileUrl());
        
        // 设置执行结果
        XxlJobHelper.handleSuccess("报告已生成: " + report.getFileUrl());
    }
}

4.3 分片广播策略

@XxlJob("dataMigrationJob")
public ReturnT<String> dataMigration() {
    int shardIndex = XxlJobHelper.getShardIndex();
    int shardTotal = XxlJobHelper.getShardTotal();
    
    // 数据分片策略:按 ID 取模
    // shardIndex=0 处理 id % 3 == 0 的数据
    // shardIndex=1 处理 id % 3 == 1 的数据
    // shardIndex=2 处理 id % 3 == 2 的数据
    
    int pageSize = 1000;
    int pageNum = 0;
    
    while (true) {
        List<Order> orders = orderDao.findByShard(
            shardIndex, shardTotal, pageNum * pageSize, pageSize);
        
        if (orders.isEmpty()) break;
        
        for (Order order : orders) {
            migrateService.migrate(order);
        }
        
        pageNum++;
        XxlJobHelper.log("分片 {} 处理第 {} 页", shardIndex, pageNum);
    }
    
    return ReturnT.SUCCESS;
}

4.4 任务触发方式

触发方式说明适用场景
Cron定时触发周期性任务
固定间隔固定间隔触发间隔执行
API 触发通过 admin API手动触发、事件触发
子任务父任务完成后触发任务流水线
分片广播广播到所有执行器分片处理

五、ElasticJob 弹性调度

5.1 核心特性

  • 弹性扩容:任务分片随实例数自动调整
  • 高可用:失效转移、错过任务重执行
  • 作业类型:Simple、Dataflow、Script
  • 分片策略:平均、轮询、哈希、自定义

5.2 Spring Boot Starter

<dependency>
    <groupId>org.apache.shardingsphere.elasticjob</groupId>
    <artifactId>elasticjob-lite-spring-boot-starter</artifactId>
    <version>3.0.3</version>
</dependency>
elasticjob:
  regCenter:
    serverLists: localhost:2181
    namespace: elasticjob
  jobs:
    dataSyncJob:
      elasticJobClass: com.example.job.DataSyncJob
      cron: "0/5 * * * * ?"
      shardingTotalCount: 3
      shardingItemParameters: 0=beijing,1=shanghai,2=guangzhou
@Component
public class DataSyncJob implements SimpleJob {
    
    @Override
    public void execute(ShardingContext context) {
        int shard = context.getShardingItem();
        String city = context.getShardingParameter();
        
        log.info("分片 {} 处理 {} 数据", shard, city);
        
        List<Data> dataList = dataService.fetchByCity(city);
        for (Data data : dataList) {
            syncService.sync(data);
        }
    }
}

六、定时任务设计原则

6.1 幂等性保障

@XxlJob("sendSmsJob")
public ReturnT<String> sendSms() {
    List<SmsTask> tasks = smsService.getPendingTasks(100);
    
    for (SmsTask task : tasks) {
        // 分布式锁保证幂等
        String lockKey = "sms:" + task.getId();
        boolean locked = redisLock.tryLock(lockKey, 30);
        
        if (!locked) continue;  // 已被其他节点处理
        
        try {
            if (task.getStatus() == TaskStatus.PENDING) {
                smsService.send(task);
                task.setStatus(TaskStatus.SENT);
                smsService.update(task);
            }
        } finally {
            redisLock.unlock(lockKey);
        }
    }
    
    return ReturnT.SUCCESS;
}

6.2 超时与熔断

@XxlJob("heavyComputeJob")
public ReturnT<String> heavyCompute() {
    ExecutorService executor = Executors.newSingleThreadExecutor();
    Future<?> future = executor.submit(() -> {
        // 耗时计算
        computeService.process();
    });
    
    try {
        future.get(5, TimeUnit.MINUTES);  // 最多执行 5 分钟
    } catch (TimeoutException e) {
        future.cancel(true);
        XxlJobHelper.log("任务执行超时,已取消");
        return ReturnT.FAIL;
    } finally {
        executor.shutdown();
    }
    
    return ReturnT.SUCCESS;
}

6.3 监控指标

指标说明告警阈值
任务执行成功率成功次数 / 总次数< 95%
平均执行时间任务耗时平均值> 历史均值 2 倍
任务堆积数待执行的任务数> 1000
错过触发次数Misfire 次数> 5 次/小时

七、方案选型建议

场景推荐方案理由
单体应用,简单定时任务Spring Scheduler零依赖,简单够用
集群调度,少量定时任务Quartz + JDBC成熟稳定,Spring 原生支持
分布式系统,可视化管理xxl-job开箱即用,运维友好
大数据量分片,弹性伸缩ElasticJob分片策略丰富,弹性好

八、总结

维度关键要点
单机调度Spring Scheduler + 自定义线程池
集群调度Quartz JDBC 集群模式
分布式调度xxl-job / ElasticJob
分片策略按 ID 取模、按范围划分
可靠性幂等设计 + 超时控制 + 失败重试
监控运维执行日志 + 指标采集 + 告警通知

定时任务看似简单,但在分布式环境中涉及调度一致性、故障转移、任务分片、幂等保障等多个复杂问题。选择合适的调度框架,遵循设计原则,才能确保任务的稳定可靠执行。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「java-enterprise」更多文章

  1. 限流算法深度解析:令牌桶、漏桶与滑动窗口计数
  2. Java 代码质量:SonarQube、Checkstyle 与 SpotBugs 工程化实践
  3. Spring IoC 容器与依赖注入原理深度剖析