1. 引言:为什么需要定制定时任务程序?
在现代软件开发中,定时执行任务的需求无处不在:数据同步、报表生成、缓存清理、消息推送、系统监控等。虽然市面上有许多成熟的定时任务框架(如 Quartz、Spring Scheduler、Celery 等),但在某些场景下,我们需要定制自己的定时任务程序:
- 特殊调度需求:非标准的 cron 表达式或复杂的调度逻辑
- 资源限制:轻量级部署,不希望引入重型框架
- 特定业务逻辑:需要与现有系统深度集成
- 学习与理解:掌握定时任务的核心原理
本文将带您从零开始,逐步构建一个可定制、可扩展的定时任务程序。
2. 核心概念与设计思路
2.1 定时任务的基本要素
一个完整的定时任务程序通常包含以下核心组件:
2.2 设计模式选择
- 观察者模式:任务作为观察者,调度器作为被观察者
- 命令模式:将任务封装为命令对象
- 工厂模式:创建不同类型的任务实例
- 策略模式:支持不同的调度策略
3. 基础实现:单机定时任务程序
3.1 环境准备
// 示例:Java 环境下的基础项目结构
// pom.xml 依赖(如果使用 Maven)
<dependencies>
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j–api</artifactId>
<version>2.0.9</version>
</dependency>
<dependency>
<groupId>ch.qos.logback</groupId>
<artifactId>logback–classic</artifactId>
<version>1.4.11</version>
</dependency>
</dependencies>
3.2 定义任务接口
/**
* 任务接口定义
*/
public interface ScheduledTask {
/**
* 获取任务ID
*/
String getTaskId();
/**
* 获取任务名称
*/
String getTaskName();
/**
* 执行任务
* @return 执行结果
*/
TaskResult execute();
/**
* 获取下次执行时间
* @return 下次执行的时间戳(毫秒)
*/
long getNextExecutionTime();
/**
* 是否启用
*/
boolean isEnabled();
}
/**
* 任务执行结果
*/
public class TaskResult {
private boolean success;
private String message;
private long executionTime;
private Object data;
// 构造方法、getter、setter 省略
}
3.3 实现简单的调度器
import java.util.concurrent.*;
/**
* 基于线程池的简单调度器
*/
public class SimpleTaskScheduler {
private final ScheduledExecutorService scheduler;
private final ConcurrentHashMap<String, ScheduledFuture<?>> taskFutures;
public SimpleTaskScheduler(int corePoolSize) {
this.scheduler = Executors.newScheduledThreadPool(corePoolSize);
this.taskFutures = new ConcurrentHashMap<>();
}
/**
* 注册定时任务
* @param task 任务实例
* @param initialDelay 初始延迟(毫秒)
* @param period 执行周期(毫秒)
*/
public void scheduleAtFixedRate(ScheduledTask task, long initialDelay, long period) {
ScheduledFuture<?> future = scheduler.scheduleAtFixedRate(
() -> {
try {
TaskResult result = task.execute();
logExecution(task, result);
} catch (Exception e) {
logError(task, e);
}
},
initialDelay,
period,
TimeUnit.MILLISECONDS
);
taskFutures.put(task.getTaskId(), future);
}
/**
* 注册 cron 风格的任务
* @param task 任务实例
* @param cronExpression cron 表达式
*/
public void scheduleWithCron(ScheduledTask task, String cronExpression) {
// 解析 cron 表达式,计算下次执行时间
long nextExecutionTime = parseCronExpression(cronExpression);
long delay = nextExecutionTime – System.currentTimeMillis();
if (delay > 0) {
scheduleWithDelay(task, delay, cronExpression);
}
}
private void scheduleWithDelay(ScheduledTask task, long delay, String cronExpression) {
ScheduledFuture<?> future = scheduler.schedule(
() -> {
try {
TaskResult result = task.execute();
logExecution(task, result);
// 重新调度下一次执行
scheduleWithCron(task, cronExpression);
} catch (Exception e) {
logError(task, e);
}
},
delay,
TimeUnit.MILLISECONDS
);
taskFutures.put(task.getTaskId(), future);
}
// 其他方法:取消任务、关闭调度器等
}
4. 进阶功能:可配置的任务管理
4.1 任务配置化
/**
* 任务配置类
*/
public class TaskConfig {
private String taskId;
private String taskClass; // 任务类全限定名
private String cronExpression;
private Map<String, Object> parameters;
private boolean enabled;
private String description;
// 支持多种触发类型
private TriggerType triggerType; // CRON, FIXED_RATE, FIXED_DELAY
private long fixedRate; // 固定频率(毫秒)
private long fixedDelay; // 固定延迟(毫秒)
private long initialDelay; // 初始延迟(毫秒)
public enum TriggerType {
CRON, FIXED_RATE, FIXED_DELAY, MANUAL
}
// 构造方法、getter、setter 省略
}
4.2 配置文件示例(YAML)
tasks:
– taskId: "dataSyncTask"
taskClass: "com.example.tasks.DataSyncTask"
triggerType: "CRON"
cronExpression: "0 0 2 * * ?" # 每天凌晨2点执行
enabled: true
parameters:
source: "databaseA"
target: "databaseB"
batchSize: 1000
– taskId: "cacheCleanTask"
taskClass: "com.example.tasks.CacheCleanTask"
triggerType: "FIXED_RATE"
fixedRate: 3600000 # 每小时执行一次
initialDelay: 60000 # 启动后1分钟开始
enabled: true
parameters:
cacheNames: ["userCache", "productCache"]
ttl: 3600
4.3 任务工厂与动态加载
/**
* 任务工厂:根据配置创建任务实例
*/
public class TaskFactory {
private static final Map<String, Class<?>> taskClassCache = new ConcurrentHashMap<>();
public static ScheduledTask createTask(TaskConfig config) throws Exception {
String className = config.getTaskClass();
Class<?> taskClass = taskClassCache.computeIfAbsent(className, key -> {
try {
return Class.forName(key);
} catch (ClassNotFoundException e) {
throw new RuntimeException("Task class not found: " + key, e);
}
});
ScheduledTask task = (ScheduledTask) taskClass.newInstance();
// 通过反射设置参数
if (config.getParameters() != null) {
setTaskParameters(task, config.getParameters());
}
return task;
}
private static void setTaskParameters(Object task, Map<String, Object> parameters)
throws IllegalAccessException {
// 使用反射设置字段值
// 实际实现中应考虑使用更安全的方式,如调用 setter 方法
}
}
5. 高级特性实现
5.1 任务持久化与故障恢复
/**
* 任务状态管理器
*/
public class TaskStateManager {
private final TaskStateStore store;
public TaskStateManager(TaskStateStore store) {
this.store = store;
}
/**
* 保存任务执行状态
*/
public void saveExecutionState(String taskId, TaskResult result) {
TaskExecutionRecord record = new TaskExecutionRecord();
record.setTaskId(taskId);
record.setExecutionTime(System.currentTimeMillis());
record.setSuccess(result.isSuccess());
record.setMessage(result.getMessage());
record.setDuration(result.getExecutionTime());
store.saveRecord(record);
}
/**
* 获取任务历史记录
*/
public List<TaskExecutionRecord> getExecutionHistory(String taskId, int limit) {
return store.getRecords(taskId, limit);
}
/**
* 故障恢复:重启后重新调度未完成的任务
*/
public void recoverTasks(List<TaskConfig> taskConfigs, SimpleTaskScheduler scheduler) {
for (TaskConfig config : taskConfigs) {
if (config.isEnabled()) {
// 检查上次执行状态
TaskExecutionRecord lastRecord = store.getLastRecord(config.getTaskId());
if (lastRecord != null && !lastRecord.isSuccess()) {
// 失败重试逻辑
scheduleRetry(config, scheduler, lastRecord);
} else {
// 正常调度
scheduleTask(config, scheduler);
}
}
}
}
}
5.2 分布式任务调度
/**
* 基于 Redis 的分布式锁实现
*/
public class DistributedTaskScheduler {
private final JedisPool jedisPool;
private final String lockKeyPrefix = "task:lock:";
public boolean tryAcquireLock(String taskId, long expireSeconds) {
try (Jedis jedis = jedisPool.getResource()) {
String lockKey = lockKeyPrefix + taskId;
String requestId = UUID.randomUUID().toString();
String result = jedis.set(lockKey, requestId,
SetParams.setParams().nx().ex(expireSeconds));
return "OK".equals(result);
}
}
public void releaseLock(String taskId) {
try (Jedis jedis = jedisPool.getResource()) {
String lockKey = lockKeyPrefix + taskId;
jedis.del(lockKey);
}
}
/**
* 分布式任务执行
*/
public void executeDistributedTask(ScheduledTask task) {
String taskId = task.getTaskId();
if (tryAcquireLock(taskId, 300)) { // 获取5分钟锁
try {
TaskResult result = task.execute();
logExecution(task, result);
} finally {
releaseLock(taskId);
}
} else {
// 其他节点正在执行此任务
logger.info("Task {} is being executed by another node", taskId);
}
}
}
6. 实战案例:数据备份任务
6.1 具体任务实现
/**
* 数据库备份任务
*/
public class DatabaseBackupTask implements ScheduledTask {
private String taskId = "dbBackupTask";
private String taskName = "数据库备份任务";
// 配置参数
private String databaseUrl;
private String backupPath;
private int retentionDays;
@Override
public String getTaskId() {
return taskId;
}
@Override
public String getTaskName() {
return taskName;
}
@Override
public TaskResult execute() {
long startTime = System.currentTimeMillis();
TaskResult result = new TaskResult();
try {
logger.info("开始执行数据库备份任务: {}", taskName);
// 1. 连接数据库
Connection conn = DriverManager.getConnection(databaseUrl);
// 2. 执行备份(这里以 MySQL 为例)
String backupFile = backupPath + "/backup_" +
LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyyMMdd_HHmmss")) + ".sql";
ProcessBuilder pb = new ProcessBuilder(
"mysqldump",
"-h", "localhost",
"-u", "root",
"–databases", "mydb",
"–result-file=" + backupFile
);
Process process = pb.start();
int exitCode = process.waitFor();
if (exitCode == 0) {
// 3. 清理过期备份
cleanupOldBackups();
result.setSuccess(true);
result.setMessage("数据库备份成功,文件: " + backupFile);
} else {
result.setSuccess(false);
result.setMessage("数据库备份失败,退出码: " + exitCode);
}
conn.close();
} catch (Exception e) {
logger.error("数据库备份任务执行失败", e);
result.setSuccess(false);
result.setMessage("执行异常: " + e.getMessage());
}
result.setExecutionTime(System.currentTimeMillis() – startTime);
return result;
}
private void cleanupOldBackups() {
File backupDir = new File(backupPath);
File[] backupFiles = backupDir.listFiles((dir, name) -> name.startsWith("backup_"));
if (backupFiles != null) {
long cutoffTime = System.currentTimeMillis() – (retentionDays * 24L * 60 * 60 * 1000);
for (File file : backupFiles) {
if (file.lastModified() < cutoffTime) {
if (file.delete()) {
logger.info("删除过期备份文件: {}", file.getName());
}
}
}
}
}
@Override
public long getNextExecutionTime() {
// 每天凌晨3点执行
LocalDateTime now = LocalDateTime.now();
LocalDateTime next = now.toLocalDate().plusDays(1)
.atTime(3, 0, 0);
return next.atZone(ZoneId.systemDefault()).toInstant().toEpochMilli();
}
@Override
public boolean isEnabled() {
return true;
}
// setter 方法
public void setDatabaseUrl(String databaseUrl) {
this.databaseUrl = databaseUrl;
}
public void setBackupPath(String backupPath) {
this.backupPath = backupPath;
}
public void setRetentionDays(int retentionDays) {
this.retentionDays = retentionDays;
}
}
6.2 配置与启动
/**
* 应用启动类
*/
public class TaskSchedulerApplication {
public static void main(String[] args) {
// 1. 加载配置
List<TaskConfig> taskConfigs = loadTaskConfigs("tasks.yaml");
// 2. 创建调度器
SimpleTaskScheduler scheduler = new SimpleTaskScheduler(10);
// 3. 创建任务状态管理器
TaskStateManager stateManager = new TaskStateManager(new FileTaskStateStore());
// 4. 恢复上次未完成的任务
stateManager.recoverTasks(taskConfigs, scheduler);
// 5. 注册并启动所有任务
for (TaskConfig config : taskConfigs) {
if (config.isEnabled()) {
try {
ScheduledTask task = TaskFactory.createTask(config);
switch (config.getTriggerType()) {
case CRON:
scheduler.scheduleWithCron(task, config.getCronExpression());
break;
case FIXED_RATE:
scheduler.scheduleAtFixedRate(task,
config.getInitialDelay(),
config.getFixedRate());
break;
case FIXED_DELAY:
// 实现 fixedDelay 调度
break;
}
logger.info("任务注册成功: {}", config.getTaskId());
} catch (Exception e) {
logger.error("任务注册失败: {}", config.getTaskId(), e);
}
}
}
// 6. 添加关闭钩子
Runtime.getRuntime().addShutdownHook(new Thread(() -> {
scheduler.shutdown();
logger.info("任务调度器已关闭");
}));
logger.info("定时任务程序启动完成");
}
}
7. 监控与运维
7.1 任务执行监控
/**
* 任务监控面板
*/
public class TaskMonitor {
private final Map<String, TaskMetrics> taskMetrics = new ConcurrentHashMap<>();
public void recordExecution(String taskId, TaskResult result) {
TaskMetrics metrics = taskMetrics.computeIfAbsent(taskId,
k -> new TaskMetrics());
metrics.recordExecution(result);
}
public TaskMetrics getMetrics(String taskId) {
return taskMetrics.get(taskId);
}
public Map<String, TaskMetrics> getAllMetrics() {
return new HashMap<>(taskMetrics);
}
/**
* 任务指标统计
*/
public static class TaskMetrics {
private AtomicLong totalExecutions = new AtomicLong(0);
private AtomicLong successfulExecutions = new AtomicLong(0);
private AtomicLong failedExecutions = new AtomicLong(0);
private AtomicLong totalExecutionTime = new AtomicLong(0);
private volatile long lastExecutionTime;
public void recordExecution(TaskResult result) {
totalExecutions.incrementAndGet();
if (result.isSuccess()) {
successfulExecutions.incrementAndGet();
} else {
failedExecutions.incrementAndGet();
}
totalExecutionTime.addAndGet(result.getExecutionTime());
lastExecutionTime = System.currentTimeMillis();
}
// 计算平均执行时间、成功率等指标
public double getAverageExecutionTime() {
long total = totalExecutions.get();
return total > 0 ? (double) totalExecutionTime.get() / total : 0;
}
pu




