跳过导航

如何设计可水平扩展的分布式定时任务系统?

约 18 分钟...次浏览
专栏微服务架构与工程治理第 9 篇

在单机上,定时任务可能只是一条 @Scheduled。部署多个实例后,同一任务会被执行多次;加一把分布式锁可以暂时避免重复,却又会引入锁过期、进程暂停、故障接管和单点吞吐问题。当任务数量增长到百万级、执行时间从毫秒到小时不等、还要求租户隔离和可追溯时,“定时触发”已经变成一个完整的分布式系统。

可水平扩展的调度平台需要分别解决三个问题:

  1. 什么时候产生一次执行:计算触发时间,处理时区、漏触发和日历规则;
  2. 由谁执行:任务分片、抢占、租约、容量和故障转移;
  3. 执行结果如何可信:幂等、重试、超时、状态机、审计和人工处置。

一、先定义语义:通常只能做到至少一次

假设 Worker 已经完成扣款,准备把任务标记为成功时进程崩溃。系统无法知道外部副作用是否发生,租约到期后只能重新派发。若为了避免重复而不重试,又可能丢失任务。

因此通用调度系统最现实的交付语义是 at-least-once:任务可能重复,但不会因一次 Worker 故障永久丢失。业务处理器必须幂等,或提供去重与补偿。

调度系统保证:到期执行最终会被领取
业务处理器保证:同一个 executionId 重复到达不会产生重复副作用

所谓 exactly-once 往往只在一个受限事务边界内成立。例如,执行记录和业务写入位于同一个数据库,可以依靠唯一键和本地事务;一旦调用第三方支付、发送邮件或跨库写入,仍要面对“不知道对方成功但响应丢失”的不确定状态。

二、把控制平面和执行平面分开

推荐架构如下:

                 +----------------------+
User / API ----> | Control Plane        |
                 | job definition       |
                 | validation / audit   |
                 +----------+-----------+
                            |
                     job definitions
                            v
                 +----------------------+
                 | Trigger Generators   |
                 | partitioned scanning |
                 +----------+-----------+
                            |
                     ready executions
                            v
                 +----------------------+
                 | Durable Queue / DB   |
                 +----------+-----------+
                            |
                +-----------+-----------+
                v                       v
          Worker Pool A           Worker Pool B
          short / IO              long / CPU

控制平面管理任务定义、权限、暂停、修改、审计和查询;触发器负责把“应该在某时运行”转换成不可变执行实例;执行平面领取实例并运行。三者独立扩缩容,避免一个慢任务堵住 Cron 扫描。

不要让调度节点同步调用业务接口后才扫描下一个任务。调度器应该快速、可重放地产生执行,Worker 才承担长耗时副作用。

三、数据模型:定义与执行必须分离

一个任务定义会产生很多执行实例,两者生命周期不同:

CREATE TABLE scheduled_job (
    job_id             BIGINT       NOT NULL,
    tenant_id          BIGINT       NOT NULL,
    schedule_type      VARCHAR(16)  NOT NULL,
    schedule_expr      VARCHAR(128) NOT NULL,
    time_zone          VARCHAR(64)  NOT NULL,
    next_fire_at       TIMESTAMP(6) NOT NULL,
    misfire_policy     VARCHAR(32)  NOT NULL,
    concurrency_policy VARCHAR(32)  NOT NULL,
    enabled            BOOLEAN      NOT NULL,
    version            BIGINT       NOT NULL,
    PRIMARY KEY (job_id),
    KEY idx_job_due (enabled, next_fire_at, job_id)
);

CREATE TABLE job_execution (
    execution_id       VARCHAR(64)  NOT NULL,
    job_id             BIGINT       NOT NULL,
    tenant_id          BIGINT       NOT NULL,
    scheduled_at       TIMESTAMP(6) NOT NULL,
    status             VARCHAR(24)  NOT NULL,
    attempt             INT         NOT NULL,
    available_at       TIMESTAMP(6) NOT NULL,
    lease_owner        VARCHAR(128) NULL,
    lease_until        TIMESTAMP(6) NULL,
    fencing_token      BIGINT       NOT NULL,
    started_at         TIMESTAMP(6) NULL,
    finished_at        TIMESTAMP(6) NULL,
    last_error_code    VARCHAR(64)  NULL,
    PRIMARY KEY (execution_id),
    UNIQUE KEY uk_job_schedule (job_id, scheduled_at),
    KEY idx_execution_ready (status, available_at, execution_id)
);

UNIQUE(job_id, scheduled_at) 是关键:即使两个触发节点同时认为任务到期,也只会生成一个逻辑执行。execution_id 可由 job_id + scheduled_at + definition_version 确定性生成,方便下游幂等。

任务定义更新使用乐观锁 version,避免两个管理员互相覆盖。执行实例保存当时的任务参数快照或参数版本,否则事后无法解释一次历史执行到底使用了什么配置。

四、如何水平扫描到期任务

方案一:数据库竞争领取

对于中等规模系统,可以让多个触发节点使用 FOR UPDATE SKIP LOCKED 并发领取。它不是公平队列:长期锁住的行会被跳过,热点任务可能延迟;事务必须短小,并用“最老待处理年龄”监控饥饿:

BEGIN;

SELECT job_id, next_fire_at, version
FROM scheduled_job
WHERE enabled = TRUE
  AND next_fire_at <= CURRENT_TIMESTAMP(6)
ORDER BY next_fire_at, job_id
LIMIT 200
FOR UPDATE SKIP LOCKED;

-- 为每条任务插入执行实例,并推进 next_fire_at

COMMIT;

优点是实现简单、由数据库保证互斥;缺点是高频扫描会给索引和主库带来压力,超大规模下数据库成为协调热点。

方案二:固定逻辑分片

用稳定哈希把 job 分到大量逻辑分片:

shard = hash(job_id) mod 4096

调度节点通过协调服务租用若干分片,只扫描自己负责的范围。逻辑分片数远大于节点数,扩容时只移动部分分片;不要直接 mod worker_count,节点数变化会让几乎所有任务重新映射。

SELECT job_id, next_fire_at
FROM scheduled_job
WHERE shard_id = :shard
  AND enabled = TRUE
  AND next_fire_at <= :scan_horizon
ORDER BY next_fire_at
LIMIT :batch;

分片租约只能决定“谁应该扫描”,最终仍依靠执行唯一键消除切换期间的重复触发。

方案三:时间桶与分层时间轮

当任务数量极大且大部分离当前触发时间很远,可按分钟/小时建立时间桶,近期任务进入内存时间轮,远期任务保存在数据库或对象存储。时间轮降低全局索引扫描,但节点故障后必须能从持久化状态重建,不能只依赖内存。

方案优点主要限制适用规模
DB + SKIP LOCKED简单、事务清晰主库扫描与锁竞争中小规模、低到中频率
逻辑分片扫描易水平扩展、故障域清楚分片租约与再平衡复杂大量持久化任务
消息延迟队列投递与消费解耦超长延迟、修改取消能力依赖产品短中期一次性任务
时间轮/时间桶近期触发效率高恢复、持久化和迁移复杂超大规模高频调度

五、租约解决接管,fencing token 解决旧持有者复活

Worker 领取执行时写入租约:

UPDATE job_execution
SET status = 'RUNNING',
    lease_owner = :workerId,
    lease_until = CURRENT_TIMESTAMP(6) + INTERVAL 30 SECOND,
    fencing_token = fencing_token + 1,
    started_at = COALESCE(started_at, :now)
WHERE execution_id = :executionId
  AND status IN ('READY', 'RETRY_WAIT')
  AND available_at <= CURRENT_TIMESTAMP(6)
  AND (lease_until IS NULL OR lease_until < CURRENT_TIMESTAMP(6));

执行时间超过租约时,Worker 定期续租。如果 Worker 崩溃,其他节点在租约过期后接管。但租约并不能让旧 Worker 立即停止:一次长 GC、网络分区或进程挂起后,它可能恢复并继续写入,此时新 Worker 已经取得任务。

因此每次领取都递增 fencing_token,受保护的资源拒绝旧 token:

UPDATE report_output
SET object_key = :key,
    last_fencing_token = :token
WHERE execution_id = :executionId
  AND last_fencing_token < :token;

分布式锁只能说明某一时刻谁持有锁,fencing token 才能让下游识别过期持有者。若第三方系统无法校验 token,就必须依赖幂等键、状态查询或业务补偿。

租约时间不能随意设置。太短会因 GC 或瞬时网络抖动频繁误接管,太长会延迟故障恢复。应基于心跳间隔、P99 暂停时间和恢复目标设置,并使用数据库时间或一致时间源,避免各节点本地时钟偏差决定锁是否过期。

六、执行状态机必须单向且可审计

推荐状态:

SCHEDULED -> READY -> RUNNING -> SUCCEEDED
                         |  \-> RETRY_WAIT -> READY
                         \----> FAILED / DEAD
SCHEDULED / READY -> CANCELED

状态转换使用条件更新,避免旧 Worker 覆盖新状态:

UPDATE job_execution
SET status = 'SUCCEEDED',
    finished_at = :now,
    lease_until = NULL
WHERE execution_id = :executionId
  AND status = 'RUNNING'
  AND lease_owner = :workerId
  AND fencing_token = :token;

若更新行数为 0,说明任务已被接管或取消,当前 Worker 的结果不能继续发布。执行日志与状态事件应保留 actor、时间、旧状态、新状态、attempt 和错误码,以支持事故还原。

七、幂等要覆盖真正的业务副作用

“任务表只执行一次”并不能保证邮件、账单或转账只发生一次。常见策略如下:

数据库唯一键

INSERT INTO monthly_invoice(tenant_id, customer_id, billing_month, amount)
VALUES (:tenant, :customer, :month, :amount)
ON DUPLICATE KEY UPDATE amount = amount;

这里是 MySQL 示例,前提是 (tenant_id, customer_id, billing_month) 已建立唯一约束,并且应用区分真正插入与重复命中;无操作更新仍可能触发审计字段或触发器,更清晰的生产实现通常是捕获重复键错误后读取既有记录。若使用 PostgreSQL,则可使用 ON CONFLICT (...) DO NOTHING。业务唯一键比随机请求 ID 更可靠,因为它表达了“同一客户同一账期只能有一张账单”。

幂等收件箱

BEGIN;

INSERT INTO job_inbox(execution_id, handler)
VALUES (:executionId, 'invoice-generator')
ON CONFLICT DO NOTHING;

-- 只有首次插入成功才执行业务写入

COMMIT;

第三方幂等键

调用支付或邮件供应商时传递稳定的 execution_id。超时后先按幂等键查询状态,再决定重试,不要盲目创建新请求。

若副作用天然无法幂等,例如“调用旧设备执行一次动作”,任务平台只能提供去重尽力而为,并把不确定状态交给人工确认或业务补偿。

八、Cron、时区和 DST 是业务语义

0 0 9 * * ? 并不能完整表达“每天上午九点”,还需要 IANA 时区,如 Asia/ShanghaiAmerica/New_York。不要只保存固定 UTC 偏移,因为夏令时规则会变化。

在夏令时切换日,本地时间可能不存在或重复:

  • 春季跳时:02:30 可能根本不存在;
  • 秋季回拨:01:30 可能出现两次。

平台必须定义策略,例如不存在时跳过或顺延到下一个有效时间,重复时只执行一次或两次都执行。这个选择是产品语义,不能交给不同语言库的默认行为。

next_fire_at 建议持久化为 UTC 时间点,任务定义保留原始时区和表达式。每次触发后从上一次计划时间计算下一次,而不是从实际完成时间计算,否则延迟会不断漂移:

错误:next = actual_finish + 24h
正确:next = cron.next(previous_scheduled_time, time_zone)

时区数据库升级也可能改变未来结果,应记录 tzdb 版本并对关键日历任务做回归测试。

九、漏触发策略必须显式

调度集群停机两小时后恢复,一个每分钟任务积累了 120 次计划执行。不同业务需要不同处理:

策略行为适用场景
SKIP丢弃错过的触发,从未来继续高频刷新、过期即无价值
FIRE_ONCE_NOW立即补一次,合并所有错过周期缓存重建、状态同步
CATCH_UP_ALL按顺序补齐全部实例账务、按周期生成不可缺失记录
CATCH_UP_LIMITED(n)最多补最近 n 次控制恢复风暴
MARK_FOR_MANUAL暂停并等待人工判断高风险资金或外部副作用

补跑时必须设置全局和租户级速率,避免恢复瞬间制造“惊群”。对于可合并任务,可让处理器接收时间范围 [last_success, now),一次处理多段,而不是创建几千个小任务。

十、并发策略决定同一任务能否重叠

当上一次执行尚未完成、下一次计划时间又到了,可提供:

  • ALLOW:允许重叠,适合彼此独立的分区任务;
  • FORBID:跳过或等待,适合全量同步;
  • REPLACE:取消旧执行并启动新执行,前提是处理器支持协作取消;
  • SERIALIZE:每个计划都保留,但同一 job 串行执行。

“取消”通常只是设置标志或发送中断信号,无法强制撤销已经提交的外部副作用。Handler 应在安全点检查取消令牌,并明确哪些阶段不可取消。

并发限制不只在 job 维度,还可能按租户、任务类型、资源池和下游依赖设置。例如每个租户最多 5 个报表任务,全平台最多 20 个访问某旧数据库的任务。

十一、重试需要错误分类、退避与预算

不是所有失败都值得重试:

错误处理
网络超时、临时 503指数退避 + jitter
限流 429尊重 Retry-After,并降低并发
参数非法、权限不足直接失败,不重试
下游长时间故障熔断、延迟重试或暂停队列
结果未知先查询幂等键状态,再决定

退避示例:

delay = min(maxDelay, base × 2^attempt) × random(0.5, 1.5)

任务要同时受到尝试次数、最大存活时间和总执行时长预算约束,不能无限重试。达到上限后进入 DEAD,保留输入摘要、错误码、最后堆栈、版本和处置记录。重放死信必须生成审计事件,并继续使用原业务幂等键。

十二、隔离不同工作负载并实现背压

短小任务与长时间报表共用一个队列和线程池,会产生队头阻塞。应按资源特征建立执行池:

fast-io        P99 < 1s
external-api   有严格下游配额
cpu-heavy      CPU 配额隔离
long-running   分钟到小时,可检查点

Worker 拉取任务前根据自身空闲 slot 决定批量,不能一次领取数千条后在本地排队,导致租约过期和不公平。队列深度高时,平台应限制新建、降低低优先级任务速率或延后非关键任务,而不是不断扩容直到压垮数据库和下游。

多租户平台还需要加权公平:一个租户提交十万任务时,其他租户仍能获得执行机会。可以按 tenant 维护令牌桶,调度时使用 deficit round robin 或分层队列。

十三、长任务需要检查点,而不是无限续租

小时级任务若失败后从头开始,会浪费大量资源。Handler 可定期保存 checkpoint:

{
  "executionId": "exec-20260712-001",
  "lastProcessedId": 5839200,
  "outputParts": 37,
  "fencingToken": 8
}

新 Worker 接管后从已提交检查点继续,并用 fencing token 阻止旧 Worker 覆盖。检查点必须与输出提交顺序协调:先生成临时结果,再原子登记 manifest,避免状态指向尚未完整写入的文件。

对于 Map/Reduce 类任务,可以把大执行拆成可重试子任务,父执行只聚合完成状态。拆分粒度要平衡调度开销与失败重算成本。

十四、修改、暂停和删除任务的语义

编辑 Cron 时,需要明确已经生成的未来执行如何处理。推荐让任务定义版本化:新版本只影响切换点之后的触发,历史执行仍引用旧版本。暂停通常停止生成新执行,不自动终止正在运行的任务;若要取消运行中任务,应单独操作并审计。

删除任务采用软删除或 disabled,先停止触发并经过保留期,再清理定义。执行历史按合规策略归档,不能因为删除定义而失去事故证据。

十五、可观测性与 SLO

调度系统至少要回答:该执行是否按时产生、等待多久、运行多久、是否成功、当前由谁持有、为什么重试。

核心指标包括:

  • schedule_lag = created_at - scheduled_at:触发器延迟;
  • queue_wait = started_at - created_at:容量是否不足;
  • 执行耗时、成功率、重试率和死亡率;
  • 租约过期与接管次数;
  • 每分片扫描延迟和再平衡次数;
  • 按任务类型/优先级的队列深度与最老年龄;
  • 漏触发数量和补跑积压;
  • 租户配额使用率与公平性。

高基数的 jobId、executionId 不适合作为指标标签,应放在日志和 Trace 中。控制台按 executionId 展示完整状态时间线,并能关联 Worker 日志、定义版本和操作审计。

SLO 应从业务及时性定义,例如“99.9% 的分钟级任务在计划时间后 30 秒内开始”,而不是仅监控调度节点存活。

十六、常见失败模式

所有实例抢一把全局锁

它把系统退化成单节点吞吐,锁服务故障还会阻止全部调度。应让任务或逻辑分片成为并行单位,并用唯一键兜底。

把锁过期等同于旧 Worker 已停止

网络分区和长暂停后旧进程仍可能继续执行。必须使用 fencing token 或业务幂等机制保护副作用。

任务先执行,最后才创建执行记录

崩溃后既无法重试,也无法审计。应先持久化执行意图,再由 Worker 领取。

使用节点本地时钟决定所有权

时钟偏差会让多个节点同时认为租约到期。租约比较使用数据库/协调服务时间,并监控 NTP 偏移。

恢复后补跑所有任务

这会把两小时故障变成恢复风暴。每个任务都要定义 misfire policy 和补跑速率。

一个线程池执行所有任务

长任务、CPU 任务和外部 API 任务互相阻塞,也无法按下游容量背压。应分类路由并设置独立预算。

把任务参数无限放进数据库

大文件和敏感正文增加数据库压力与泄露风险。参数保存受控引用,内容放对象存储并加密、鉴权和设置保留期。

十七、生产检查清单

  • 是否明确系统提供至少一次语义,并要求 Handler 幂等或可补偿?
  • 任务定义、触发实例和执行状态是否分离且版本化?
  • 多调度节点是否通过唯一键防止同一计划时间生成多个逻辑执行?
  • 扫描是否可按逻辑分片水平扩展,并能在节点故障后安全再平衡?
  • 领取是否使用租约,旧 Worker 的副作用是否由 fencing token 或幂等键阻止?
  • 状态转换是否带 owner/token 条件且保留完整审计时间线?
  • Cron 是否保存 IANA 时区,并明确 DST 重复与缺失时间策略?
  • 是否为每个任务定义漏触发、并发、超时、取消和重试策略?
  • 重试是否分类错误、使用 jitter,并受次数与总时长预算限制?
  • 短任务、长任务、CPU 和外部依赖任务是否使用隔离资源池?
  • 是否实现全局、租户、任务类型和下游依赖维度的背压与公平调度?
  • 长任务是否支持检查点,接管后能避免重复发布输出?
  • 编辑、暂停、删除和死信重放是否有明确语义和操作审计?
  • 是否监控 schedule lag、queue wait、租约接管、漏触发、队列年龄和业务 SLO?
  • 是否演练调度节点全停、数据库抖动、Worker 长 GC、时钟偏差和积压恢复?

总结

分布式定时任务系统的难点从来不是解析 Cron,而是在故障和并发下维护可解释的执行状态。让触发与执行解耦,用唯一键容忍重复调度,用租约完成接管,用 fencing token 与业务幂等抵御旧 Worker,再为时区、漏触发、重试和背压定义明确策略。系统允许重复、延迟和接管,却不允许这些现实被隐藏;只有把不确定性变成状态机、指标和审计,调度平台才真正具备水平扩展与生产可信度。

分享:
文章作者:狼码纪
版权声明:本博客所有文章除特别声明外,均采用 CC BY-NC-SA 4.0 许可协议。文章可能参考了其他优秀文章,如有侵权请联系删除。

Morty Proxy This is a proxified and sanitized view of the page, visit original site.