一、背景

Scheduler 每 10 秒扫一次到期任务
每条任务执行前 tryLock
成功就执行,失败走重试,最后解锁
这个方案在任务量小时没问题,但任务规模和节点规模上来后,问题明显:

逐条抢锁开销高:先查出一批任务,再每条任务单独 UPDATE … AND locked=0 抢锁,数据库往返次数高。
多节点竞争浪费:节点 A/B/C 都查到同一批候选任务,再逐条竞争,很多请求天然是“无效请求”。
并发边界可解释性弱:只靠 lock_node + lock_time 回查“我这轮认领了哪些任务”,在高并发和时间精度边界下不够硬。
所以第一阶段优化目标很明确:

把“逐条抢锁”升级为“批量认领”,并且让每轮认领有唯一身份可追踪。

二、优化目标

把调度执行权竞争,从“每条任务竞争一次”收敛为“每轮批次竞争一次”,并用 claim_token 让认领结果可证明。

三、优化后的链路

Scheduler 扫描
-> 批量 claim(一次 UPDATE 认领一批)
-> 按 claim_token 回查本轮任务
-> 线程池并发执行
-> 成功/失败更新 next_run_time 或 retry
-> finally 解锁并清理 claim_token
注意两点:

调度线程只负责发现与派发,执行在线程池里做。
同一任务仍然通过 locked=0 -> 1 保证多节点互斥执行。

四、第一部分:从逐条抢锁改为批量认领(Claim)

1)核心思路
过去是:
findDueJobs() 找到候选
对每条候选执行一次 tryLock(jobId)

现在改成:
一条 SQL 直接把“到期且未锁定”的前 N 条任务标记为当前节点认领
再查回“本轮认领成功”的任务列表去执行
这样做的收益非常直接:
数据库交互次数下降、无效竞争减少、吞吐更稳定。

2)核心 SQL:批量认领

@Update("""
    UPDATE inspection_job
    SET locked = 1,
        lock_node = #{nodeId},
        lock_time = NOW(6),
        claim_token = #{claimToken}
    WHERE status = 'RUNNING'
    AND next_run_time <= NOW(6)
    AND locked = 0
    ORDER BY next_run_time ASC
    LIMIT #{limit}
    """)
int claimDueJobs(
        @Param("nodeId") String nodeId,
        @Param("claimToken") String claimToken,
        @Param("limit") int limit
);

这段 SQL 的关键点:
locked = 0 作为并发门闩(谁先改成功谁拿执行权)
ORDER BY next_run_time ASC LIMIT 控制调度公平性和批量上限
NOW(6) 统一数据库时间,避免节点本地时钟偏差

3)调度入口改造:认领后再并发执行

@Scheduled(fixedDelayString = "${inspection.scheduling.scan-interval-ms:10000}")
public void schedule() {
    cleanLocks();
    List<InspectionJob> jobs = inspectionJobService.claimDueJobs();
    for (InspectionJob job : jobs) {
        InspectionJob snapshot = job;
        inspectionJobExecutor.execute(() -> inspectionJobService.runInspectionJob(snapshot));
    }
}

五、第二部分:并发边界保护(claim_token + DB 时间)

第一部分解决了性能问题,第二部分解决“认领结果可证明”的工程问题。

1)为什么需要 claim_token
如果只按 lock_node + lock_time 回查本轮任务,极端并发场景下语义不够强。
我希望做到:
“这轮认领出来的任务,就是这轮,不多不少,可追踪。”

所以每轮认领生成一个 UUID,写入 claim_token,回查时按 token 精确过滤。

2)服务层:每轮 token 唯一

public List<InspectionJob> claimDueJobs() {
    String nodeId = resolveNodeId();
    int limit = inspectionProperties.getScheduling().getBatchSize();
    String claimToken = UUID.randomUUID().toString();
    int claimed = inspectionJobMapper.claimDueJobs(nodeId, claimToken, limit);
    if (claimed <= 0) {
        return Collections.emptyList();
    }
    return inspectionJobMapper.findClaimedJobs(nodeId, claimToken, limit);
}

3)回查 SQL:按 token 精确取本轮结果

@Select("""
    SELECT *
    FROM inspection_job
    WHERE locked = 1
    AND lock_node = #{nodeId}
    AND claim_token = #{claimToken}
    ORDER BY next_run_time ASC
    LIMIT #{limit}
    """)
List<InspectionJob> findClaimedJobs(
        @Param("nodeId") String nodeId,
        @Param("claimToken") String claimToken,
        @Param("limit") int limit
);

4)解锁与超时清理:同步清空 token

@Update("""
    UPDATE inspection_job
    SET locked = 0,
        claim_token = NULL
    WHERE job_id = #{jobId}
    AND lock_node = #{nodeId}
""")
int unlock(Long jobId, String nodeId);
@Update("""
    UPDATE inspection_job
    SET locked = 0,
        claim_token = NULL
    WHERE locked = 1
    AND lock_time < NOW() - INTERVAL 5 MINUTE
""")
void releaseTimeoutLocks();

这一步很重要:
claim_token 是“本轮认领批次身份”,任务完成或超时回收后必须清理,避免脏状态污染下一轮。

六、实体与表结构补充

1)实体字段
/** 本轮批量认领的唯一标识,避免并发场景误取任务 */
private String claimToken;
2)数据库迁移
ALTER TABLE inspection_job
ADD COLUMN max_retry INT NOT NULL DEFAULT 3,
ADD COLUMN retry_count INT NOT NULL DEFAULT 0,
ADD COLUMN retry_interval_seconds INT NOT NULL DEFAULT 60,
ADD COLUMN claim_token VARCHAR(64) NULL;
– 建议索引
– CREATE INDEX idx_job_due ON inspection_job (status, next_run_time);
– CREATE INDEX idx_job_claim ON inspection_job (lock_node, claim_token);

七、失败重试与 token 的关系(常见疑问)

很多人会问:任务失败重试后,claim_token 会不会对不上?

答案是:不会,这是预期行为。

claim_token 表示“本轮认领批次”
一次执行完成(无论成功失败)最终都会解锁并清 token
重试进入下一轮调度,再认领时会生成新的 token
也就是说:

job_id 是任务长期身份
claim_token 是短生命周期的批次身份

八、优化前后对比

维度 优化前(逐条抢锁) 优化后(批量认领 + token)
竞争粒度 每条任务竞争一次 每轮批次集中竞争
DB 往返
并发解释性 弱(回查语义不够硬) 强(claim_token 可追踪)
时钟一致性 受节点本地时间影响 统一使用 DB NOW(6)
故障排查 较难定位“本轮到底认领了什么” 可按 token 精确定位

九、这阶段最大的工程收获

这次优化让我更清楚一件事:

调度系统的难点不是“能不能执行任务”,而是“多节点并发下,执行权如何高效且可证明地分配”。

从“逐条抢锁”到“批量认领”,解决了效率问题;
从“普通回查”到“claim_token 批次身份”,解决了正确性和可解释性问题。
这两个点合在一起,才是真正可上线的第一阶段优化。

Logo

AtomGit 是由开放原子开源基金会联合 CSDN 等生态伙伴共同推出的新一代开源与人工智能协作平台。平台坚持“开放、中立、公益”的理念,把代码托管、模型共享、数据集托管、智能体开发体验和算力服务整合在一起,为开发者提供从开发、训练到部署的一站式体验。

更多推荐