Quartz 分布式任务调度集群实战:Redis 分布式锁增强 + 故障自动接管
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集群依赖数据库锁实现互斥,存在以下问题:
- 性能瓶颈:所有节点竞争同一把数据库行锁,高并发时数据库压力大
- 故障恢复慢:依赖心跳超时检测,默认需要30秒以上才能发现节点宕机
- 数据库依赖:必须使用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
结论与建议
核心要点
- Quartz集群通过共享数据库实现状态同步和故障恢复,是成熟的集群调度方案
- Redis分布式锁作为补充,提供更快速的互斥控制和故障检测
- 超时控制和重试机制是保证任务可靠执行的关键
实践建议
- 锁过期时间:应大于任务最大执行时间,但不宜过长(建议任务执行时间的2-3倍)
- 心跳间隔:Quartz集群心跳间隔建议15秒,过短增加数据库压力,过长影响故障恢复速度
- 任务幂等性:所有定时任务必须保证幂等性,因为网络抖动可能导致重复执行
- 监控告警:对任务执行失败、超时、节点离线等异常建立完善的告警机制
- 数据库选型:Quartz集群必须使用JDBC JobStore,建议使用独立的数据源,避免影响业务数据库
方案选型
| 场景 | 推荐方案 |
|---|---|
| 单节点简单调度 | Quartz内存模式 |
| 多节点互斥执行 | Quartz集群 + Redis锁 |
| 大规模分布式调度 | XXL-Job / ElasticJob |
| 云原生K8s环境 | Kubernetes CronJob |