场景描述

  • 多台服务器同时运行一个定时程序,周期性地向数据库的 任务表 插入待处理任务。

  • 多台消费者服务器从数据库任务表中取出任务并执行,每个任务需要调用一个耗时约 5 分钟 的外部接口。

  • 定时调度框架为 Quartz,要求集群环境下任务不重复生产、不重复消费、具备故障恢复能力。


一、生产阶段:防止重复插入任务

潜在问题

多节点同时触发定时插入逻辑,导致数据库中出现多条相同的业务任务,进而引发下游接口被重复调用、数据混乱。

根本原因

定时程序缺乏互斥机制,多节点“各自为政”,同时执行插入动作。

解决方案

利用 Quartz 集群模式 天然保证同一时刻只有一个节点触发 Job。

  1. 配置 Quartz 集群

    properties

    org.quartz.jobStore.isClustered = true
    org.quartz.scheduler.instanceId = AUTO
    org.quartz.jobStore.clusterCheckinInterval = 20000
  2. 用数据库存储 JobDetail 和 Trigger

    • JobDetail:定义“做什么”(调用插入业务任务的逻辑)

    • Trigger:定义“何时做”(Cron 表达式 0 0 8 * * ? 每天8点)

  3. 业务代码保持幂等

    • 在任务表增加 biz_key 字段(如 DAILY_REPORT:2026-06-02)并添加唯一索引。

    • 插入时使用 INSERT ... ON DUPLICATE KEY UPDATE id = id 避免报错并防止脏数据。

为什么 Quartz 能保证互斥

  • 集群所有节点共享 QRTZ_TRIGGERS 表。

  • 触发时间到达时,每个节点的调度线程执行 SELECT ... FOR UPDATE 去获取待触发的 Trigger 行。

  • 数据库行锁保证只有一个节点能获取该行锁,其他节点要么被阻塞超时、要么查询到空结果,从而放弃本次触发。

  • 获得锁的节点执行对应的 Job(即插入业务任务),完成后更新下次触发时间并释放锁。

补充建议

即使 Quartz 已经保证了互斥,业务表仍然建议保留唯一索引,作为防止配置错误或人为手动触发时的兜底保障。


二、消费阶段:防止重复执行与并发冲突

潜在问题

  1. 多个消费者抢到同一个任务:如果先用 SELECT 查询状态为“待执行”的任务,再 UPDATE 为“执行中”,两步之间存在时间窗口,导致多个消费者并行处理同一任务。

  2. 执行节点宕机导致任务“卡死”:任务状态停留在“执行中”不再被处理。

  3. 长时间执行引发的资源问题: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 实现注意点

  1. 线程池隔离

    • 任务执行线程池(如 FixedThreadPool)专用于调用 5 分钟接口。

    • 心跳更新用 ScheduledExecutorService 单独维护,绝对不与业务线程混用。

  2. HTTP 客户端超时设置

    • connectTimeout:5 秒(连接建立超时)。

    • readTimeout:至少 7~8 分钟(留足缓冲,大于任务预期耗时)。

  3. 数据库连接池配置

    • 连接池最大连接数 > 预期并发任务数 + 心跳并发数,避免连接争抢导致心跳失败。

  4. 优雅停机

    • 应用关闭时停止拉取新任务,等待执行中的任务完成(设定最大等待时间),并正常停止心跳线程。


四、Quartz 集群中核心概念解析

概念 含义 数据库表
JobDetail 定义“要做什么”,封装执行逻辑 QRTZ_JOB_DETAILS
Trigger 定义“何时做”,存储 Cron 等 QRTZ_TRIGGERSQRTZ_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 已不匹配,更新失败,不会重复释放许可。这保证了并发计数的准确性。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐