本文目录导读:

- 基础定时任务调度(使用ScheduledExecutorService)
- Quartz框架完整案例(企业级调度)
- 多任务调度案例(含优先级和依赖)
- 动态定时任务管理器
- 分布式任务调度(使用Redis实现简单分布)
- 数据库支持的持久化调度
我将为您提供几个Java调度(任务调度)的实用案例,涵盖从简单到复杂的场景。
基础定时任务调度(使用ScheduledExecutorService)
import java.util.concurrent.*;
import java.time.LocalTime;
public class BasicScheduler {
public static void main(String[] args) throws InterruptedException {
ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(2);
// 1. 延迟执行任务
scheduler.schedule(() -> {
System.out.println("[延迟任务] 5秒后执行: " + LocalTime.now());
}, 5, TimeUnit.SECONDS);
// 2. 固定频率执行(固定周期)
scheduler.scheduleAtFixedRate(() -> {
System.out.println("[固定频率] 每2秒执行: " + LocalTime.now());
}, 0, 2, TimeUnit.SECONDS);
// 3. 固定延迟执行(上次执行完毕后延迟)
scheduler.scheduleWithFixedDelay(() -> {
System.out.println("[固定延迟] 每次执行间隔3秒: " + LocalTime.now());
try {
Thread.sleep(1000); // 模拟任务执行耗时
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}, 0, 3, TimeUnit.SECONDS);
// 运行10秒后关闭
Thread.sleep(10000);
scheduler.shutdown();
System.out.println("调度器已关闭");
}
}
Quartz框架完整案例(企业级调度)
// 需要引入依赖:org.quartz-scheduler:quartz
import org.quartz.*;
import org.quartz.impl.StdSchedulerFactory;
import java.util.Date;
// 1. 自定义任务类
public class EmailJob implements Job {
@Override
public void execute(JobExecutionContext context) throws JobExecutionException {
JobDataMap dataMap = context.getJobDetail().getJobDataMap();
String email = dataMap.getString("email");
System.out.println("[" + new Date() + "] 发送邮件到: " + email);
}
}
// 2. 调度配置类
public class QuartzSchedulerExample {
public static void main(String[] args) throws SchedulerException {
// 创建调度器
Scheduler scheduler = StdSchedulerFactory.getDefaultScheduler();
// 创建JobDetail
JobDetail jobDetail = JobBuilder.newJob(EmailJob.class)
.withIdentity("emailJob", "group1")
.usingJobData("email", "user@example.com") // 传递参数
.build();
// 创建Trigger - 每天上午10点执行
Trigger trigger = TriggerBuilder.newTrigger()
.withIdentity("emailTrigger", "group1")
.startNow()
.withSchedule(CronScheduleBuilder.cronSchedule("0 0 10 * * ?"))
.build();
// 注册任务和触发器
scheduler.scheduleJob(jobDetail, trigger);
// 启动调度器
scheduler.start();
System.out.println("Quartz调度器已启动");
// 如果需要停止,可以添加shutdown hook
Runtime.getRuntime().addShutdownHook(new Thread(() -> {
try {
scheduler.shutdown(true);
} catch (SchedulerException e) {
e.printStackTrace();
}
}));
}
}
多任务调度案例(含优先级和依赖)
import java.util.*;
import java.util.concurrent.*;
public class ComplexTaskScheduler {
// 任务定义
static class Task {
String name;
int priority;
long duration;
Runnable action;
Task(String name, int priority, long duration, Runnable action) {
this.name = name;
this.priority = priority;
this.duration = duration;
this.action = action;
}
}
public static void main(String[] args) throws InterruptedException {
// 优先级队列
ScheduledExecutorService executor = Executors.newScheduledThreadPool(5);
Map<String, CompletableFuture<Void>> taskDependencies = new HashMap<>();
// 任务1:基础数据准备
CompletableFuture<Void> task1 = CompletableFuture.runAsync(() -> {
System.out.println("任务1: 准备基础数据...");
sleep(2000);
System.out.println("任务1: 完成");
}, executor);
// 任务2:依赖任务1的处理
CompletableFuture<Void> task2 = task1.thenRun(() -> {
System.out.println("任务2: 处理任务1的数据...");
sleep(1500);
System.out.println("任务2: 完成");
});
// 任务3:周期性数据备份(独立任务)
ScheduledFuture<?> task3 = executor.scheduleAtFixedRate(() -> {
System.out.println("任务3: 执行数据备份...");
sleep(800);
}, 0, 10, TimeUnit.SECONDS);
// 任务4:延迟任务 - 1分钟后执行清理
executor.schedule(() -> {
System.out.println("任务4: 执行系统清理");
}, 1, TimeUnit.MINUTES);
// 等待任务2完成
task2.join();
System.out.println("所有依赖任务执行完毕");
// 在程序结束前停止调度器
executor.shutdown();
}
private static void sleep(long millis) {
try {
Thread.sleep(millis);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
}
动态定时任务管理器
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicBoolean;
public class DynamicTaskManager {
private final ScheduledExecutorService scheduler;
private final Map<String, ScheduledFuture<?>> tasks = new ConcurrentHashMap<>();
private final AtomicBoolean running = new AtomicBoolean(false);
public DynamicTaskManager(int threadPoolSize) {
this.scheduler = Executors.newScheduledThreadPool(threadPoolSize);
}
// 添加固定频率任务
public String addFixedRateTask(String name, Runnable task, long interval, TimeUnit unit) {
validateTaskName(name);
ScheduledFuture<?> future = scheduler.scheduleAtFixedRate(
wrapTask(name, task), 0, interval, unit
);
tasks.put(name, future);
return name;
}
// 添加Cron式任务(简化版)
public String addPeriodicTask(String name, Runnable task, long delay, long period, TimeUnit unit) {
validateTaskName(name);
ScheduledFuture<?> future = scheduler.scheduleWithFixedDelay(
wrapTask(name, task), delay, period, unit
);
tasks.put(name, future);
return name;
}
// 取消任务
public boolean cancelTask(String name) {
ScheduledFuture<?> future = tasks.remove(name);
if (future != null) {
return future.cancel(true);
}
return false;
}
// 暂停任务
public boolean pauseTask(String name) {
ScheduledFuture<?> future = tasks.get(name);
return future != null && future.cancel(true);
}
// 检查任务状态
public boolean isTaskRunning(String name) {
ScheduledFuture<?> future = tasks.get(name);
return future != null && !future.isDone();
}
// 获取所有任务名称
public Set<String> getTaskNames() {
return tasks.keySet();
}
// 关闭所有任务
public void shutdown() {
running.set(false);
tasks.values().forEach(future -> future.cancel(true));
scheduler.shutdownNow();
System.out.println("所有任务已停止");
}
// 包装任务,添加日志
private Runnable wrapTask(String name, Runnable task) {
return () -> {
System.out.println("[任务启动] " + name + " @ " +
SimpleDateFormat.getTimeInstance().format(new Date()));
task.run();
System.out.println("[任务结束] " + name);
};
}
private void validateTaskName(String name) {
if (tasks.containsKey(name)) {
throw new IllegalArgumentException("任务已存在: " + name);
}
}
// 测试用例
public static void main(String[] args) throws InterruptedException {
DynamicTaskManager manager = new DynamicTaskManager(2);
// 添加任务
manager.addFixedRateTask("backup", () -> {
System.out.println("执行数据备份");
}, 3, TimeUnit.SECONDS);
manager.addPeriodicTask("cleanup", () -> {
System.out.println("执行清理");
}, 1, 5, TimeUnit.SECONDS);
Thread.sleep(10000);
// 暂停备份任务
manager.pauseTask("backup");
System.out.println("备份任务已暂停");
Thread.sleep(5000);
// 重新添加任务
manager.addFixedRateTask("report", () -> {
System.out.println("生成报告");
}, 2, TimeUnit.SECONDS);
Thread.sleep(6000);
// 关闭管理器
manager.shutdown();
}
}
分布式任务调度(使用Redis实现简单分布)
import redis.clients.jedis.Jedis;
import java.util.UUID;
import java.util.concurrent.*;
public class DistributedTaskScheduler {
private static final String LOCK_PREFIX = "task:lock:";
private static final String QUEUE_PREFIX = "task:queue:";
private final Jedis redis;
private final String instanceId;
private final ScheduledExecutorService scheduler;
public DistributedTaskScheduler(String redisHost, int redisPort) {
this.redis = new Jedis(redisHost, redisPort);
this.instanceId = UUID.randomUUID().toString();
this.scheduler = Executors.newScheduledThreadPool(3);
System.out.println("分布式调度器实例启动: " + instanceId);
}
// 发送定时任务到队列
public void enqueueTask(String taskType, String payload, int delaySeconds) {
String taskId = UUID.randomUUID().toString();
String key = QUEUE_PREFIX + taskType;
// 使用Redis List作为任务队列
redis.lpush(key, taskId + ":" + payload);
// 设置任务延迟检查
scheduler.schedule(() -> {
processTask(taskType, taskId, payload);
}, delaySeconds, TimeUnit.SECONDS);
}
// 处理任务(带分布式锁)
private void processTask(String taskType, String taskId, String payload) {
String lockKey = LOCK_PREFIX + taskId;
String lockValue = UUID.randomUUID().toString();
// 尝试获取分布式锁(SET NX EX)
String result = redis.set(lockKey, lockValue, "NX", "EX", 30);
if ("OK".equals(result)) {
try {
System.out.println("[" + instanceId + "] 开始处理任务: " +
taskType + " - " + payload);
// 模拟任务执行
Thread.sleep(1000);
System.out.println("[" + instanceId + "] 任务完成: " + taskType);
} catch (Exception e) {
System.err.println("任务执行失败: " + e.getMessage());
} finally {
// 释放锁
String currentValue = redis.get(lockKey);
if (lockValue.equals(currentValue)) {
redis.del(lockKey);
}
}
} else {
System.out.println("[" + instanceId + "] 任务已被其他实例处理");
}
}
public void shutdown() {
scheduler.shutdown();
redis.close();
}
public static void main(String[] args) throws InterruptedException {
// 模拟两个分布式实例
DistributedTaskScheduler instance1 = new DistributedTaskScheduler("localhost", 6379);
DistributedTaskScheduler instance2 = new DistributedTaskScheduler("localhost", 6379);
// 发送任务
instance1.enqueueTask("EMAIL", "user1@example.com", 5);
instance2.enqueueTask("SMS", "13800138000", 3);
instance1.enqueueTask("PUSH", "device123", 1);
Thread.sleep(10000);
instance1.shutdown();
instance2.shutdown();
}
}
数据库支持的持久化调度
// 使用MySQL存储任务状态
import java.sql.*;
import java.time.LocalDateTime;
import java.util.concurrent.*;
public class DatabaseBackedScheduler {
// 数据库连接
private Connection connection;
private ScheduledExecutorService executor;
public DatabaseBackedScheduler(String jdbcUrl, String username, String password)
throws SQLException {
this.connection = DriverManager.getConnection(jdbcUrl, username, password);
this.executor = Executors.newScheduledThreadPool(2);
initDatabase();
}
// 初始化数据库表
private void initDatabase() throws SQLException {
String sql = """
CREATE TABLE IF NOT EXISTS tasks (
id INT AUTO_INCREMENT PRIMARY KEY,
name VARCHAR(100) NOT NULL,
status VARCHAR(20) DEFAULT 'PENDING',
scheduled_time TIMESTAMP,
last_executed TIMESTAMP,
result TEXT,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
)
""";
try (Statement stmt = connection.createStatement()) {
stmt.execute(sql);
}
}
// 添加调度任务
public void scheduleTask(String name, LocalDateTime triggerTime) throws SQLException {
String sql = """
INSERT INTO tasks (name, scheduled_time, status)
VALUES (?, ?, 'PENDING')
""";
try (PreparedStatement ps = connection.prepareStatement(sql)) {
ps.setString(1, name);
ps.setTimestamp(2, Timestamp.valueOf(triggerTime));
ps.executeUpdate();
}
// 调度执行
long delay = Duration.between(LocalDateTime.now(), triggerTime).toMillis();
executor.schedule(() -> {
try {
executeTask(name);
} catch (Exception e) {
e.printStackTrace();
}
}, delay, TimeUnit.MILLISECONDS);
}
// 执行任务
private void executeTask(String taskName) throws SQLException {
// 更新任务状态为运行中
String updateSQL = "UPDATE tasks SET status = 'RUNNING' WHERE name = ?";
try (PreparedStatement ps = connection.prepareStatement(updateSQL)) {
ps.setString(1, taskName);
ps.executeUpdate();
}
System.out.println("[" + LocalDateTime.now() + "] 执行任务: " + taskName);
// 模拟任务执行
try {
Thread.sleep(2000);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
// 完成任务
String completeSQL = """
UPDATE tasks SET status = 'COMPLETED',
last_executed = ?, result = ? WHERE name = ?
""";
try (PreparedStatement ps = connection.prepareStatement(completeSQL)) {
ps.setTimestamp(1, Timestamp.valueOf(LocalDateTime.now()));
ps.setString(2, "任务执行成功");
ps.setString(3, taskName);
ps.executeUpdate();
}
}
// 查询未执行的任务
public void checkPendingTasks() throws SQLException {
String sql = "SELECT name FROM tasks WHERE status = 'PENDING'";
try (Statement stmt = connection.createStatement();
ResultSet rs = stmt.executeQuery(sql)) {
while (rs.next()) {
System.out.println("待执行任务: " + rs.getString("name"));
}
}
}
public static void main(String[] args) throws Exception {
DatabaseBackedScheduler scheduler = new DatabaseBackedScheduler(
"jdbc:mysql://localhost:3306/jobs",
"root",
"password"
);
// 添加任务
scheduler.scheduleTask("每日数据备份", LocalDateTime.now().plusSeconds(5));
scheduler.scheduleTask("清理临时文件", LocalDateTime.now().plusSeconds(10));
// 查询任务
scheduler.checkPendingTasks();
Thread.sleep(15000);
}
}
这些案例展示了Java中任务调度的不同方式:
- ScheduledExecutorService - 适合简单的定时任务
- Quartz - 企业级调度框架,支持Cron表达式和持久化
- CompletableFuture - 支持任务依赖和异步组合
- 动态任务管理器 - 支持任务的动态添加、暂停和取消
- 分布式调度 - 使用Redis实现分散式任务处理
- 数据库持久化 - 保证任务状态不丢失
选择哪种方式取决于您的具体需求:
- 简单任务:使用ScheduledExecutorService
- 复杂调度:选择Quartz
- 微服务架构:考虑分布式方案
- 任务状态追踪:使用数据库持久化