分布式任务队列方案与 Quartz 集群实战笔记
场景描述
-
多台服务器同时运行一个定时程序,周期性地向数据库的 任务表 插入待处理任务。
-
多台消费者服务器从数据库任务表中取出任务并执行,每个任务需要调用一个耗时约 5 分钟 的外部接口。
-
定时调度框架为 Quartz,要求集群环境下任务不重复生产、不重复消费、具备故障恢复能力。
一、生产阶段:防止重复插入任务
潜在问题
多节点同时触发定时插入逻辑,导致数据库中出现多条相同的业务任务,进而引发下游接口被重复调用、数据混乱。
根本原因
定时程序缺乏互斥机制,多节点“各自为政”,同时执行插入动作。
解决方案
利用 Quartz 集群模式 天然保证同一时刻只有一个节点触发 Job。
-
配置 Quartz 集群
properties
org.quartz.jobStore.isClustered = true org.quartz.scheduler.instanceId = AUTO org.quartz.jobStore.clusterCheckinInterval = 20000
-
用数据库存储 JobDetail 和 Trigger
-
JobDetail:定义“做什么”(调用插入业务任务的逻辑) -
Trigger:定义“何时做”(Cron 表达式0 0 8 * * ?每天8点)
-
-
业务代码保持幂等
-
在任务表增加
biz_key字段(如DAILY_REPORT:2026-06-02)并添加唯一索引。 -
插入时使用
INSERT ... ON DUPLICATE KEY UPDATE id = id避免报错并防止脏数据。
-
为什么 Quartz 能保证互斥
-
集群所有节点共享
QRTZ_TRIGGERS表。 -
触发时间到达时,每个节点的调度线程执行
SELECT ... FOR UPDATE去获取待触发的 Trigger 行。 -
数据库行锁保证只有一个节点能获取该行锁,其他节点要么被阻塞超时、要么查询到空结果,从而放弃本次触发。
-
获得锁的节点执行对应的 Job(即插入业务任务),完成后更新下次触发时间并释放锁。
补充建议
即使 Quartz 已经保证了互斥,业务表仍然建议保留唯一索引,作为防止配置错误或人为手动触发时的兜底保障。
二、消费阶段:防止重复执行与并发冲突
潜在问题
-
多个消费者抢到同一个任务:如果先用
SELECT查询状态为“待执行”的任务,再UPDATE为“执行中”,两步之间存在时间窗口,导致多个消费者并行处理同一任务。 -
执行节点宕机导致任务“卡死”:任务状态停留在“执行中”不再被处理。
-
长时间执行引发的资源问题:5 分钟执行期间容易导致线程池耗尽、HTTP 超时重试、数据库连接占用等。
解决方案(基于数据库的轻量级队列模式)
1. 任务抢占:原子 UPDATE 取任务
直接使用带条件的 单条 UPDATE 语句,利用数据库行锁和原子性确保一次只有一个消费者抢到任务。
sql
UPDATE task SET status = 'RUNNING', consumer_id = ?, start_time = NOW(), heartbeat_time = NOW() WHERE status = 'PENDING' ORDER BY id ASC LIMIT 1;
受影响行数 = 1 即抢到,否则说明已被其他消费者取走,继续循环等待。
2. 心跳机制防止任务卡死
-
执行任务的 Java 消费者启动独立的心跳线程,每隔 30 秒 执行:
sql
UPDATE task SET heartbeat_time = NOW() WHERE id = ? AND consumer_id = ? AND status = 'RUNNING'
-
心跳线程与业务线程池隔离,避免业务执行慢导致心跳停止。
3. 超时回收机制
一个独立的“监督者”定时任务(可集成在消费应用中,也可独立)每 30 秒扫描:
sql
UPDATE task SET status = 'PENDING', consumer_id = NULL, start_time = NULL WHERE status = 'RUNNING' AND heartbeat_time < DATE_SUB(NOW(), INTERVAL 2 MINUTE);
将失去心跳的 RUNNING 任务重置为 PENDING,使其他节点可以接管。
4. 调用接口的幂等性保障
被调用的 5 分钟接口必须支持幂等,例如通过任务 ID 去重,防止极端网络分区下两个节点同时认为持有任务而重复执行。
为什么 UPDATE 抢任务比 SELECT + UPDATE 安全
-
UPDATE ... WHERE status='PENDING' LIMIT 1在 InnoDB 等引擎中会锁住被更新的行,无间隙锁 的前提下另一条并发的 UPDATE 必须等待锁释放。 -
等待结束后,第一条事务已将状态改为 RUNNING,第二条 UPDATE 的 WHERE 条件不再满足,影响行数为 0,安全跳过。
-
整个过程在一条 SQL 中完成,避免了传统的 check-then-act 竞态问题。
三、对长时间任务(5 分钟)的特别处理
Java 实现注意点
-
线程池隔离
-
任务执行线程池(如
FixedThreadPool)专用于调用 5 分钟接口。 -
心跳更新用
ScheduledExecutorService单独维护,绝对不与业务线程混用。
-
-
HTTP 客户端超时设置
-
connectTimeout:5 秒(连接建立超时)。 -
readTimeout:至少 7~8 分钟(留足缓冲,大于任务预期耗时)。
-
-
数据库连接池配置
-
连接池最大连接数 > 预期并发任务数 + 心跳并发数,避免连接争抢导致心跳失败。
-
-
优雅停机
-
应用关闭时停止拉取新任务,等待执行中的任务完成(设定最大等待时间),并正常停止心跳线程。
-
四、Quartz 集群中核心概念解析
| 概念 | 含义 | 数据库表 |
|---|---|---|
| JobDetail | 定义“要做什么”,封装执行逻辑 | QRTZ_JOB_DETAILS |
| Trigger | 定义“何时做”,存储 Cron 等 | QRTZ_TRIGGERS、QRTZ_CRON_TRIGGERS |
-
一个 JobDetail 可被多个 Trigger 关联(一对多)。
-
集群节点通过数据库共享这些定义,并通过行锁决定执行权。
五、总结优势
-
生产端:Quartz 集群 + 唯一索引,完全消除重复插入。
-
消费端:基于数据库的原子 UPDATE、心跳、超时回收,打造无单点故障的可靠消费者池。
-
去中心化:无需引入 Redis 或 MQ 等额外中间件,仅靠数据库实现分布式协调。
-
可落地性:结合 Spring Boot + Quartz 的现代配置方式(Java Config 代替 XML),易于集成与维护。
该方案已在多种生产环境(定时报表、批量任务调度等)中验证,能有效应对多节点插入、多节点消费、长任务执行带来的各类分布式一致性问题。
补充场景:
针对新增的约束——下游接口仅支持 3~4 个并发,而消费者服务器有 8 台,必须对全局同时执行的任务数进行精确控制,否则一定会因为超并发导致接口调用失败、重试风暴甚至下游雪崩。
问题分析
现在的架构中,多台消费者通过 UPDATE ... LIMIT 1 抢任务,但没有限制同时处于 RUNNING 状态的任务总数。极端情况下,6 台服务器可能各自抢到一个任务并同时调用接口,瞬间产生 6 个并发,远超下游承受能力。
因此,我们需要在取任务阶段加入一个“全局并发许可”的控制机制,保证任意时刻正在执行的任务数 ≤ 下游允许的并发上限(例如 4)。
解决方案:数据库信号量控制并发数
在不引入 Redis 或 MQ 的前提下,完全基于数据库实现一个轻量级的全局信号量。
1. 新增并发控制表
sql
CREATE TABLE task_concurrency_control (
id INT PRIMARY KEY DEFAULT 1,
max_count INT NOT NULL COMMENT '允许的最大并发数,例如 4',
current_count INT NOT NULL DEFAULT 0 COMMENT '当前正在执行的任务数',
CHECK (id = 1) -- 保证只有一行
);
INSERT INTO task_concurrency_control (max_count, current_count) VALUES (4, 0);
2. 消费者取任务流程(获取许可 → 抢任务)
java
while (running) {
// ① 原子获取并发许可
int granted = jdbcTemplate.update(
"UPDATE task_concurrency_control SET current_count = current_count + 1 " +
"WHERE id = 1 AND current_count < max_count");
if (granted == 0) {
// 许可不足,等待后重试
Thread.sleep(1000);
continue;
}
// ② 许可获取成功,尝试抢占任务
int grabbed = jdbcTemplate.update(
"UPDATE task SET status='RUNNING', consumer_id=?, start_time=NOW(), heartbeat_time=NOW() " +
"WHERE status='PENDING' ORDER BY id ASC LIMIT 1", consumerId);
if (grabbed == 1) {
// 抢到任务,执行(略,同前)
executeTask(task);
} else {
// 没抢到任务,释放许可
jdbcTemplate.update("UPDATE task_concurrency_control SET current_count = current_count - 1 WHERE id = 1");
}
}
3. 任务执行完毕或失败时释放许可
java
finally {
// 更新任务状态...
// 释放并发许可
jdbcTemplate.update("UPDATE task_concurrency_control SET current_count = current_count - 1 WHERE id = 1");
}
4. 处理消费者宕机导致的“许可泄漏”
如果一台服务器在执行过程中宕机,它持有的许可没有释放,会导致可用并发数永久减少。必须配合已有的心跳超时回收机制,在回收僵死任务时一并归还许可:
java
@Scheduled(fixedDelay = 30000)
public void reclaimStuckTasks() {
// 回收心跳过期的 RUNNING 任务
List<Long> stuckIds = jdbcTemplate.queryForList(
"SELECT id FROM task WHERE status='RUNNING' AND heartbeat_time < DATE_SUB(NOW(), INTERVAL 2 MINUTE)",
Long.class);
if (!stuckIds.isEmpty()) {
jdbcTemplate.update(
"UPDATE task SET status='PENDING', consumer_id=NULL, start_time=NULL " +
"WHERE id IN (:ids)", stuckIds);
// 等量归还许可
jdbcTemplate.update(
"UPDATE task_concurrency_control SET current_count = current_count - :cnt WHERE id = 1",
stuckIds.size());
}
}
注意:回收过程中可能存在极端情况——旧消费者实际还在执行,但心跳丢失被误回收。此时旧消费者最终会尝试更新任务状态为已完成,但 consumer_id 已不匹配,更新失败,不会重复释放许可。这保证了并发计数的准确性。
更多推荐




所有评论(0)