Quartz 分布式任务调度集群实战:Redis 分布式锁增强 + 故障自动接管

作者:忆笙智云官方 | 发布时间:2026-05-10 09:00 | 更新时间:2026-06-10 09:00

Quartz分布式任务调度集群实现

引言

在企业级应用中,定时任务是不可或缺的基础能力——数据同步、报表生成、缓存刷新、消息推送等场景都依赖定时任务。当应用部署为集群时,如何保证同一个任务在多个节点上只执行一次,如何处理节点故障时的任务接管,成为必须解决的核心问题。

本文将分析Quartz集群模式的工作原理,以及如何结合Redis分布式锁实现更可靠的任务调度方案。

一、Quartz核心概念

1.1 基础架构

graph TB
    A[Scheduler 调度器] --> B[Job 任务]
    A --> C[Trigger 触发器]
    A --> D[JobDetail 任务详情]

    C --> C1[SimpleTrigger<br/>简单触发]
    C --> C2[CronTrigger<br/>Cron表达式触发]

    B --> E[JobExecutionContext<br/>执行上下文]
    D --> F[JobDataMap<br/>任务参数]

    style A fill:#e8f4f8
    style B fill:#d5f9d5
    style C fill:#f9f5d5
组件 说明
Scheduler 调度器,Quartz的核心,负责管理任务和触发器
Job 任务接口,定义任务执行逻辑
JobDetail 任务详情,包含Job类、参数等
Trigger 触发器,定义任务执行时间和频率
JobDataMap 任务参数,可在调度时传递数据

1.2 单节点 vs 集群模式

graph TB
    subgraph 单节点模式
        A1[应用节点] --> B1[内存JobStore]
        B1 --> C1[❌ 节点宕机任务丢失]
    end

    subgraph 集群模式
        D1[节点1] --> E1[数据库JobStore]
        D2[节点2] --> E1
        D3[节点3] --> E1
        E1 --> F1[✅ 故障自动接管]
        E1 --> F2[✅ 任务互斥执行]
        E1 --> F3[✅ 负载均衡]
    end

    style C1 fill:#f9d5d5
    style F1 fill:#d5f9d5
    style F2 fill:#d5f9d5
    style F3 fill:#d5f9d5

二、Quartz集群原理

2.1 数据库表结构

Quartz集群通过共享数据库实现状态同步,核心表包括:

表名 说明
QRTZ_JOB_DETAILS 任务详情
QRTZ_TRIGGERS 触发器信息
QRTZ_CRON_TRIGGERS Cron触发器
QRTZ_FIRED_TRIGGERS 已触发的触发器(正在执行)
QRTZ_SCHEDULER_STATE 调度器状态(节点心跳)
QRTZ_LOCKS 集群锁(保证互斥)

2.2 集群工作流程

sequenceDiagram
    participant N1 as 节点1
    participant N2 as 节点2
    participant DB as 共享数据库

    Note over N1,N2: 心跳检测(每15秒)
    N1->>DB: UPDATE SCHEDULER_STATE (last_heartbeat)
    N2->>DB: UPDATE SCHEDULER_STATE (last_heartbeat)

    Note over N1,N2: 任务触发(每30秒轮询)
    N1->>DB: SELECT FOR UPDATE (获取集群锁)
    DB-->>N1: 锁获取成功
    N1->>DB: 查询待触发的Trigger
    DB-->>N1: 返回Trigger列表
    N1->>DB: INSERT FIRED_TRIGGER (标记为已触发)
    N1->>DB: 释放集群锁
    N1->>N1: 执行任务

    Note over N2: 同时轮询
    N2->>DB: SELECT FOR UPDATE (获取集群锁)
    DB-->>N2: 锁已被N1持有,等待
    N2->>DB: 锁获取成功(N1释放后)
    N2->>DB: 查询待触发的Trigger
    DB-->>N2: 无待触发Trigger(已被N1处理)

    Note over N1: 节点1宕机
    N1-xN1: 停止心跳

    Note over N2: 检测到N1宕机
    N2->>DB: 检查SCHEDULER_STATE
    DB-->>N2: N1心跳超时
    N2->>DB: 恢复N1的FIRED_TRIGGER
    N2->>N2: 重新执行N1未完成的任务

2.3 Spring Boot集成Quartz集群

# application.yml - Quartz集群配置
spring:
  quartz:
    job-store-type: jdbc                    # 使用数据库存储
    properties:
      org.quartz.scheduler.instanceName: YsClusterScheduler
      org.quartz.scheduler.instanceId: AUTO  # 自动生成实例ID
      org.quartz.jobStore.class: org.quartz.impl.jdbcjobstore.JobStoreTX
      org.quartz.jobStore.driverDelegateClass: org.quartz.impl.jdbcjobstore.StdJDBCDelegate
      org.quartz.jobStore.tablePrefix: QRTZ_
      org.quartz.jobStore.isClustered: true   # 启用集群模式
      org.quartz.jobStore.clusterCheckinInterval: 15000  # 心跳间隔15秒
      org.quartz.jobStore.misfireThreshold: 60000        # 超时阈值60秒
      org.quartz.threadPool.class: org.quartz.simpl.SimpleThreadPool
      org.quartz.threadPool.threadCount: 10
      org.quartz.threadPool.threadPriority: 5
/**
 * Quartz配置类
 */
@Configuration
public class QuartzConfig {

    @Autowired
    private DataSource dataSource;

    @Bean
    public SchedulerFactoryBean schedulerFactoryBean() {
        SchedulerFactoryBean factory = new SchedulerFactoryBean();
        factory.setDataSource(dataSource);
        factory.setOverwriteExistingJobs(true);
        factory.setAutoStartup(true);
        factory.setStartupDelay(5);  // 启动延迟5秒,等待应用初始化完成

        // 集群配置
        Properties props = new Properties();
        props.setProperty("org.quartz.scheduler.instanceId", "AUTO");
        props.setProperty("org.quartz.jobStore.isClustered", "true");
        props.setProperty("org.quartz.jobStore.clusterCheckinInterval", "15000");
        factory.setQuartzProperties(props);

        return factory;
    }
}

三、Redis分布式锁增强

3.1 Quartz集群的局限

Quartz集群依赖数据库锁实现互斥,存在以下问题:

  1. 性能瓶颈:所有节点竞争同一把数据库行锁,高并发时数据库压力大
  2. 故障恢复慢:依赖心跳超时检测,默认需要30秒以上才能发现节点宕机
  3. 数据库依赖:必须使用JDBC JobStore,无法使用内存存储

3.2 Redis分布式锁方案

/**
 * Redis分布式锁
 * 基于SET NX EX命令实现
 */
@Component
public class RedisDistributedLock {

    @Autowired
    private StringRedisTemplate redisTemplate;

    /** 锁前缀 */
    private static final String LOCK_PREFIX = "quartz:lock:";

    /**
     * 尝试获取锁
     * @param lockKey 锁键
     * @param requestId 请求ID(用于标识锁的持有者)
     * @param expireSeconds 过期时间(秒)
     * @return 是否获取成功
     */
    public boolean tryLock(String lockKey, String requestId, long expireSeconds) {
        String key = LOCK_PREFIX + lockKey;
        Boolean result = redisTemplate.opsForValue()
                .setIfAbsent(key, requestId, Duration.ofSeconds(expireSeconds));
        return Boolean.TRUE.equals(result);
    }

    /**
     * 释放锁(Lua脚本保证原子性)
     * 只有锁的持有者才能释放锁
     */
    public boolean unlock(String lockKey, String requestId) {
        String key = LOCK_PREFIX + lockKey;
        String script = """
            if redis.call('get', KEYS[1]) == ARGV[1] then
                return redis.call('del', KEYS[1])
            else
                return 0
            end
            """;
        Long result = redisTemplate.execute(
                new DefaultRedisScript<>(script, Long.class),
                Collections.singletonList(key),
                requestId);
        return Long.valueOf(1).equals(result);
    }
}

3.3 分布式锁任务包装器

/**
 * 分布式锁任务包装器
 * 确保同一任务在集群中只有一个节点执行
 */
public abstract class DistributedJob implements Job {

    @Autowired
    private RedisDistributedLock distributedLock;

    /** 锁持有时间(秒),应大于任务最大执行时间 */
    protected abstract long getLockExpireSeconds();

    /** 获取锁的键名 */
    protected abstract String getLockKey();

    @Override
    public void execute(JobExecutionContext context) throws JobExecutionException {
        String requestId = UUID.randomUUID().toString();
        String lockKey = getLockKey();

        // 尝试获取分布式锁
        boolean locked = distributedLock.tryLock(lockKey, requestId, getLockExpireSeconds());
        if (!locked) {
            // 其他节点正在执行,跳过
            return;
        }

        try {
            // 执行业务逻辑
            doExecute(context);
        } catch (Exception e) {
            throw new JobExecutionException("任务执行异常", e);
        } finally {
            // 释放锁
            distributedLock.unlock(lockKey, requestId);
        }
    }

    /**
     * 子类实现具体的任务逻辑
     */
    protected abstract void doExecute(JobExecutionContext context) throws Exception;
}

3.4 具体任务示例

/**
 * 数据同步定时任务
 * 使用分布式锁保证集群中只有一个节点执行
 */
@Component
@DisallowConcurrentExecution  // 禁止同一节点并发执行
public class DataSyncJob extends DistributedJob {

    @Autowired
    private DataSyncService dataSyncService;

    @Override
    protected String getLockKey() {
        return "data_sync_job";
    }

    @Override
    protected long getLockExpireSeconds() {
        return 300;  // 锁过期时间5分钟
    }

    @Override
    protected void doExecute(JobExecutionContext context) throws Exception {
        // 执行数据同步逻辑
        dataSyncService.syncFromExternalSystem();
    }
}

四、超时控制与重试机制

4.1 任务超时控制

/**
 * 带超时控制的任务执行器
 * 使用CompletableFuture实现任务超时中断
 */
@Component
public class JobTimeoutExecutor {

    private final ScheduledExecutorService scheduler =
            Executors.newScheduledThreadPool(2);

    /**
     * 带超时执行任务
     * @param task 任务逻辑
     * @param timeoutSeconds 超时时间(秒)
     */
    public void executeWithTimeout(Runnable task, long timeoutSeconds) {
        CompletableFuture<Void> future = CompletableFuture.runAsync(task);

        // 设置超时
        scheduler.schedule(() -> {
            if (!future.isDone()) {
                future.cancel(true);
            }
        }, timeoutSeconds, TimeUnit.SECONDS);

        try {
            future.get(timeoutSeconds, TimeUnit.SECONDS);
        } catch (TimeoutException e) {
            throw new RuntimeException("任务执行超时(" + timeoutSeconds + "秒)");
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            throw new RuntimeException("任务被中断");
        } catch (ExecutionException e) {
            throw new RuntimeException("任务执行异常", e.getCause());
        }
    }
}

4.2 任务重试机制

/**
 * 可重试的任务基类
 * 支持自定义重试次数和重试间隔
 */
public abstract class RetryableJob extends DistributedJob {

    /** 最大重试次数 */
    protected int getMaxRetries() {
        return 3;
    }

    /** 重试间隔(毫秒) */
    protected long getRetryInterval() {
        return 5000;
    }

    @Override
    protected void doExecute(JobExecutionContext context) throws Exception {
        int retryCount = 0;
        Exception lastException = null;

        while (retryCount <= getMaxRetries()) {
            try {
                executeWithRetry(context);
                return;  // 执行成功,直接返回
            } catch (Exception e) {
                lastException = e;
                retryCount++;

                if (retryCount <= getMaxRetries()) {
                    // 等待重试间隔
                    Thread.sleep(getRetryInterval() * retryCount);  // 指数退避
                }
            }
        }

        // 所有重试都失败
        throw lastException;
    }

    /**
     * 子类实现具体的任务逻辑(可重试)
     */
    protected abstract void executeWithRetry(JobExecutionContext context) throws Exception;
}

五、集群监控与管理

5.1 任务执行状态监控

/**
 * 任务执行记录服务
 * 记录每次任务执行的状态、耗时、结果
 */
@Service
public class JobExecutionLogService {

    @Autowired
    private JobExecutionLogMapper logMapper;

    /**
     * 记录任务执行
     */
    public void logExecution(String jobName, String jobGroup,
                              boolean success, long duration, String errorMsg) {
        JobExecutionLog log = new JobExecutionLog();
        log.setJobName(jobName);
        log.setJobGroup(jobGroup);
        log.setExecuteNode(getNodeIdentifier());
        log.setSuccess(success);
        log.setDuration(duration);
        log.setErrorMsg(errorMsg);
        log.setExecuteTime(LocalDateTime.now());
        logMapper.insert(log);
    }

    /**
     * 获取节点标识(IP+进程ID)
     */
    private String getNodeIdentifier() {
        try {
            String hostAddress = InetAddress.getLocalHost().getHostAddress();
            String processId = ManagementFactory.getRuntimeMXBean().getName();
            return hostAddress + ":" + processId;
        } catch (Exception e) {
            return "unknown";
        }
    }
}

5.2 集群节点状态

flowchart TB
    subgraph 集群监控面板
        A[节点列表] --> A1[节点1: 192.168.1.10 ✅ 运行中]
        A --> A2[节点2: 192.168.1.11 ✅ 运行中]
        A --> A3[节点3: 192.168.1.12 ❌ 已离线]

        B[任务统计] --> B1[今日执行: 156次]
        B --> B2[成功: 148次]
        B --> B3[失败: 8次]
        B --> B4[运行中: 3个]

        C[告警信息] --> C1[⚠️ 节点3心跳超时]
        C --> C2[❌ 报表生成任务连续3次失败]
    end

    style A1 fill:#d5f9d5
    style A2 fill:#d5f9d5
    style A3 fill:#f9d5d5
    style C1 fill:#f9f5d5
    style C2 fill:#f9d5d5

六、完整架构

flowchart TB
    A[任务调度请求] --> B[Quartz Scheduler]
    B --> C{集群模式}

    C -->|节点1| D1[获取Redis分布式锁]
    C -->|节点2| D2[获取Redis分布式锁]
    C -->|节点3| D3[获取Redis分布式锁]

    D1 -->|成功| E1[执行任务]
    D2 -->|失败| E2[跳过本次执行]
    D3 -->|失败| E3[跳过本次执行]

    E1 --> F[超时控制]
    F --> G{执行成功?}
    G -->|是| H[记录执行日志]
    G -->|否| I{重试次数未超?}
    I -->|是| J[等待重试间隔]
    J --> E1
    I -->|否| K[记录失败日志+告警]

    subgraph 故障恢复
        L[节点1宕机] --> M[Redis锁自动过期]
        M --> N[下次调度由其他节点获取锁]
        N --> O[Quartz恢复未完成任务]
    end

    style E1 fill:#d5f9d5
    style E2 fill:#f9f5d5
    style E3 fill:#f9f5d5
    style K fill:#f9d5d5

结论与建议

核心要点

  1. Quartz集群通过共享数据库实现状态同步和故障恢复,是成熟的集群调度方案
  2. Redis分布式锁作为补充,提供更快速的互斥控制和故障检测
  3. 超时控制重试机制是保证任务可靠执行的关键

实践建议

  1. 锁过期时间:应大于任务最大执行时间,但不宜过长(建议任务执行时间的2-3倍)
  2. 心跳间隔:Quartz集群心跳间隔建议15秒,过短增加数据库压力,过长影响故障恢复速度
  3. 任务幂等性:所有定时任务必须保证幂等性,因为网络抖动可能导致重复执行
  4. 监控告警:对任务执行失败、超时、节点离线等异常建立完善的告警机制
  5. 数据库选型:Quartz集群必须使用JDBC JobStore,建议使用独立的数据源,避免影响业务数据库

方案选型

场景 推荐方案
单节点简单调度 Quartz内存模式
多节点互斥执行 Quartz集群 + Redis锁
大规模分布式调度 XXL-Job / ElasticJob
云原生K8s环境 Kubernetes CronJob

相关资源