如何设计可水平扩展的分布式定时任务系统?
在单机上,定时任务可能只是一条 @Scheduled。部署多个实例后,同一任务会被执行多次;加一把分布式锁可以暂时避免重复,却又会引入锁过期、进程暂停、故障接管和单点吞吐问题。当任务数量增长到百万级、执行时间从毫秒到小时不等、还要求租户隔离和可追溯时,“定时触发”已经变成一个完整的分布式系统。
可水平扩展的调度平台需要分别解决三个问题:
- 什么时候产生一次执行:计算触发时间,处理时区、漏触发和日历规则;
- 由谁执行:任务分片、抢占、租约、容量和故障转移;
- 执行结果如何可信:幂等、重试、超时、状态机、审计和人工处置。
一、先定义语义:通常只能做到至少一次
假设 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/Shanghai 或 America/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,再为时区、漏触发、重试和背压定义明确策略。系统允许重复、延迟和接管,却不允许这些现实被隐藏;只有把不确定性变成状态机、指标和审计,调度平台才真正具备水平扩展与生产可信度。