Java调度案例

wen java案例 1

本文目录导读:

Java调度案例

  1. 基础定时任务调度(使用ScheduledExecutorService)
  2. Quartz框架完整案例(企业级调度)
  3. 多任务调度案例(含优先级和依赖)
  4. 动态定时任务管理器
  5. 分布式任务调度(使用Redis实现简单分布)
  6. 数据库支持的持久化调度

我将为您提供几个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中任务调度的不同方式:

  1. ScheduledExecutorService - 适合简单的定时任务
  2. Quartz - 企业级调度框架,支持Cron表达式和持久化
  3. CompletableFuture - 支持任务依赖和异步组合
  4. 动态任务管理器 - 支持任务的动态添加、暂停和取消
  5. 分布式调度 - 使用Redis实现分散式任务处理
  6. 数据库持久化 - 保证任务状态不丢失

选择哪种方式取决于您的具体需求:

  • 简单任务:使用ScheduledExecutorService
  • 复杂调度:选择Quartz
  • 微服务架构:考虑分布式方案
  • 任务状态追踪:使用数据库持久化

上一篇Table API案例

下一篇DataStream案例

抱歉,评论功能暂时关闭!