定时任务是业务系统中常见的需求,如数据同步、报表生成、缓存预热、订单超时取消等。在单机环境下 Spring 的 @Scheduled 足够使用,但在分布式环境中,必须解决任务重复执行、集群调度、分片处理与故障转移等问题。
一、定时任务方案对比
| 特性 | Spring Scheduler | Quartz | xxl-job | ElasticJob |
|---|
| 分布式支持 | 不支持 | 支持(JDBC/Redis) | 原生支持 | 原生支持 |
| 管理控制台 | 无 | 需自研 | 完善 | 完善 |
| 任务分片 | 不支持 | 不支持 | 支持 | 原生支持 |
| 弹性扩容 | 不支持 | 不支持 | 支持 | 支持 |
| 失败重试 | 不支持 | 支持 | 支持 | 支持 |
| 触发类型 | 固定频率/Cron | Cron | Cron/固定间隔/API | Cron/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 取模、按范围划分 |
| 可靠性 | 幂等设计 + 超时控制 + 失败重试 |
| 监控运维 | 执行日志 + 指标采集 + 告警通知 |
定时任务看似简单,但在分布式环境中涉及调度一致性、故障转移、任务分片、幂等保障等多个复杂问题。选择合适的调度框架,遵循设计原则,才能确保任务的稳定可靠执行。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。