调度、回填与补数

本文系统讲解工作流的调度、回填与补数,回答 cron 与间隔调度该怎么选、上游未就绪如何等待、回填为什么容易重跑出错、时区与夏令时有哪些陷阱。覆盖数据区间语义、锚点对齐与调度漂移、依赖感知与条件门控、三种回填模式与并发控制、补数与对账、幂等前提、SLA 定义与告警、大规模调度的容量规划,并给出可落地的配置片段与排查方法。

引言

调度看起来是最简单的一环:写一个时间表达式,到点执行。真正投入生产后才会发现,几乎所有调度事故都不来自「表达式写错」,而来自三个更深的问题:这次执行到底处理哪一段时间的数据、上游还没好怎么办、漏跑了几天的数据怎么补。

第一个问题决定了回填能否安全运行。如果任务用「当前时刻」而不是「所属的数据区间」来决定处理范围,那么回填 30 天就会变成「用今天的条件跑 30 次」,结果 30 份重复数据。这是数据领域最经典的事故,也是「数据区间语义」必须被显式设计的原因。

第二个问题决定了流程的可靠性。定时触发假设「到点上游一定好了」,但上游可能延迟、失败、或者根本还没产出。依赖感知触发(等数据就绪而不是等时间)是更正确的模型,代价是链路变长、排查变难。

第三个问题决定了运维的日常体验。线上一定会出现「某天数据没跑」「某天跑错了要重来」「新加了指标要补历史」。回填与补数因此不是异常操作,而是常态操作,必须作为一等公民来设计,而不是临时写个脚本。

本文按「调度语义 → 依赖与条件 → 回填与补数 → 时区与 SLA → 容量与运维」的顺序展开。Airflow 侧的具体 API 与 DAG 写法参见 Airflow DAG 调度体系 ,事件触发与调度的边界参见事件驱动触发相关章节。

目录

  1. 调度的三个基本问题
  2. cron 表达式的语义细节
  3. 间隔调度与锚点对齐
  4. 数据区间与执行时刻的分离
  5. 依赖感知触发
  6. 条件触发与门控
  7. 回填的本质与三种模式
  8. 回填的并发与资源控制
  9. 补数:与回填的区别
  10. 幂等:回填能安全运行的前提
  11. 时区陷阱
  12. 夏令时与时间跳跃
  13. 调度漂移与错过执行
  14. SLA 的定义与度量
  15. SLA 监控与告警
  16. 优先级与抢占
  17. 调度器的可观测与自愈
  18. 大规模调度的容量规划
  19. 落地路线图
  20. 权衡取舍
  21. 常见坑清单
  22. 小结

1. 调度的三个基本问题

任何调度系统都必须明确回答三个问题:什么时候触发(时间表达式)、这次触发处理什么数据(数据区间)、不满足条件时怎么办(等待 / 跳过 / 失败重试)。这三个答案决定了它的能力边界。

问题二是最容易被忽略、也最致命的。多数调度器(Airflow、Dagster、Prefect)都内置了「数据区间」的概念,但它的语义需要使用者理解并正确使用。核心规则是:一次执行对应一个确定的数据区间,任务的所有读写都以这个区间为准,而不是以执行时刻为准。

问题三决定系统的健壮性。「等待」适合上游延迟是常态的场景,「跳过」适合「错过就放弃」的场景(比如每小时抓一次行情),「失败并重试」适合必须成功的场景。三种策略要按任务性质分别选择,而不是全局统一。

2. cron 表达式的语义细节

cron 是事实标准,但它的细节比大多数人以为的多。五个字段依次是分钟、小时、日、月、星期,其中最容易踩坑的有四点:

一是「日」与「星期」是或的关系。0 0 1 * 1 表示「每月 1 号或每周一」,而不是「每月 1 号且是周一」。想要「且」必须写成 0 0 1 * 1 在标准 cron 里做不到,需要用脚本判断。

二是步长表达式的起点。*/5 在分钟位表示 0、5、10…,起点是 0 而不是当前时间。10/5 表示从 10 开始每 5 分钟,即 10、15、20…。混淆这两个会让「每 5 分钟」的预期变成「从某个偏移开始每 5 分钟」。

三是不同实现的方言差异。Linux cron 不支持秒,Spring 的 @Scheduled 是 6 位(秒在最前),Quartz 是 7 位(秒 + 年),Kubernetes CronJob 是标准 5 位且从 1.27 起支持时区。跨系统复制表达式时务必确认位数与方言。

四是 @daily 这类别名的实际时间。@daily 等价于 0 0 * * *,是本地时区的 0 点。在多时区部署的环境里,不同节点上的 @daily 触发时刻不同,会造成重复触发或漏触发。解决方式是显式指定时区。

# Kubernetes CronJob:显式指定时区,避免节点时区差异
apiVersion: batch/v1
kind: CronJob
spec:
  schedule: "0 2 * * *"
  timeZone: "Asia/Shanghai"     # 1.27+ 支持,低版本必须靠容器内 TZ 环境变量

3. 间隔调度与锚点对齐

不是所有调度都适合 cron。对于「每 15 分钟跑一次」这类需求,间隔调度(interval)比 cron 更自然,但引入了一个新问题:锚点从哪算起。

方式 A:从任务启动时刻算起
  启动于 10:07 -> 10:22 -> 10:37 -> 10:52
  问题:重启后锚点变化,区间边界漂移,无法与下游对齐

方式 B:从固定锚点算起(epoch 对齐)
  固定对齐到整 15 分钟 -> 10:00 -> 10:15 -> 10:30 -> 10:45
  优点:区间边界确定,重启不影响,下游可预测

生产系统必须用方式 B。锚点对齐的实现是把「当前时间向下取整到间隔的整数倍」,而不是「上次执行时间 + 间隔」:

def next_aligned(interval_seconds, now=None):
    now = now if now is not None else time.time()
    current = int(now // interval_seconds) * interval_seconds   # 当前区间起点
    return current + interval_seconds                           # 下一区间起点

对齐还有一个容易忽略的细节:间隔必须能整除一天(比如 15 分钟、1 小时、6 小时),否则每天的区间数量会变化,按「每天 N 个区间」做校验的逻辑会失效。

4. 数据区间与执行时刻的分离

这是整篇文章最重要的概念。用两个时间点来描述一次执行:数据区间 [interval_start, interval_end) 是这次执行要处理的数据范围(业务语义上的时间),执行时刻 run_at 是实际开始的时间(系统语义上的时间),两者满足 run_at >= interval_end,因为区间结束后才能处理完整数据。

一个 0 2 * * * 的任务,在 10 月 8 日 02:00 执行时,处理的数据区间是 [10-07 00:00, 10-08 00:00)。这个「滞后一天」的关系是刻意的:要处理 10 月 7 日的完整数据,必须等 10 月 7 日结束。由此推出两条硬规则:

规则一:任务里禁止使用「当前时间」做业务判断。WHERE created_at >= NOW() - INTERVAL '1 day' 这样的查询在正常调度下碰巧正确(因为执行时刻接近区间结束),但在回填时完全错误——回填 9 月 1 日的区间,查询条件却是「最近一天」,处理的是今天的数据。

规则二:所有查询条件都从数据区间派生。

@task
def extract(data_interval_start=None, data_interval_end=None):
    # 正确:条件来自数据区间
    return run_query(
        "SELECT * FROM ods.orders WHERE created_at >= %s AND created_at < %s",
        (data_interval_start, data_interval_end),
    )

例外情况是「窗口跨越多个区间」的指标(比如「近 7 日均值」)。这类任务应该用 interval_end - 7 天 到 interval_end 的滚动窗口,且这个窗口必须从 interval_end 反推,而不是从「现在」反推。

5. 依赖感知触发

定时触发的根本缺陷是「它假设时间到了,数据就绪了」。当上游延迟 2 小时,下游用时间触发会读到不完整的数据——而且是静默地读到,产出的报表看起来正常但数字是错的。这比任务失败更危险。

依赖感知触发把触发条件从「时间到了」换成「依赖就绪了」。三种实现方式,抽象层次递增:传感器等待(轮询检查依赖是否就绪,实现简单但占资源)、资产/数据驱动(上游声明产出、下游声明依赖,解耦但触发粒度粗)、事件通知(上游完成后发消息,延迟低但需要消息基础设施)。

传感器是最直接的实现,但要注意两个参数:timeout(等多久放弃)与 mode(等待期间是否占槽位)。没有 timeout 的传感器会永久等待,把流程挂死;mode="poke" 会占满 Worker 槽位,应该用 mode="reschedule" 或可延迟模式。

# 等待上游分区就绪,最多等 6 小时,等待期间不占 Worker 槽位
wait_upstream = ExternalTaskSensor(
    task_id="wait_dwd_orders",
    external_dag_id="dwd_orders_pipeline",
    external_task_id="load",
    execution_delta=timedelta(hours=0),   # 对齐两个 DAG 的数据区间
    timeout=6 * 3600, mode="reschedule", poke_interval=300,
)

execution_delta 是最容易出错的参数:它定义「我的哪个区间对应上游的哪个区间」。如果两个 DAG 的调度周期不同(一个每天、一个每小时),必须用 execution_delta 或 execution_date_fn 做映射,否则会永远等待一个不存在的上游实例。

6. 条件触发与门控

除了「上游完成」,还有一类触发条件是业务性的:「只有在有数据时才跑」「只有在工作日才跑」「只有开关打开时才跑」。这类条件叫门控(gate),三种处理方式:

跳过(skip):条件不满足时标记为 skipped,下游按 trigger_rule 决定是否继续
短路(short-circuit):条件不满足时让整个分支不执行
阻塞(block):条件不满足时等待,直到满足或超时

跳过是最常用的,但它的陷阱在于下游的触发规则。默认的 all_success 在遇到上游 skipped 时的行为需要实测确认:在 Airflow 里 skipped 被视为「成功」的变体,下游会继续执行——这通常是你想要的,但也可能造成「数据没产出却往下走」的问题。

门控的一个实践原则是**「条件判断要可观测」**:条件不满足时不仅要跳过,还要记录「为什么跳过」(数据为空、非工作日、开关关闭),否则运维看到一堆 skipped 任务会无从判断是正常还是异常。

7. 回填的本质与三种模式

回填(backfill)是为历史区间创建并执行运行实例。它的本质是**「用相同的逻辑跑不同的数据区间」**,因此它对任务的要求是「逻辑必须与区间无关,只有数据范围随区间变化」。

三种回填模式,区别在于「如何决定回填哪些区间」:

模式一:区间枚举(catchup)  从 start_date 到当前时间逐区间创建,适用首次上线补齐历史
模式二:显式指定范围        给定 [start, end],适用修复特定几天的数据
模式三:从状态反推          查询「哪些区间没有成功记录」,适用日常补数与故障恢复

模式三是运维中最实用的,因为它幂等且自愈:不管漏跑的原因是什么,只要查一下「预期区间 vs 已成功区间」的差集就能补。实现方式是维护一张「预期区间表」或按调度表达式生成区间列表,再与运行记录做差:

-- 找出最近 30 天里没有成功记录的日期
WITH expected AS (
    SELECT generate_series(CURRENT_DATE - 30, CURRENT_DATE - 1, 1)::date AS biz_date
)
SELECT e.biz_date
FROM expected e
LEFT JOIN flow_run r ON r.flow_name = 'orders_daily'
   AND r.biz_date = e.biz_date AND r.status = 'SUCCESS'
WHERE r.id IS NULL ORDER BY e.biz_date;

8. 回填的并发与资源控制

回填最大的风险是资源争抢。回填 30 天意味着 30 个运行实例同时开始,它们会争抢数据库连接、计算资源和下游系统的配额,把正常调度的任务挤掉。

控制手段有四层,从粗到细:

max_active_runs       限制同一流程同时运行的区间数(最有效)
pool / 队列           限制某一类资源的总并发
回填限速             回填任务之间加间隔,避免瞬间起量
下游配额             保护外部系统的调用配额
# 回填时显式限制并发,从默认值降到 2
airflow dags backfill orders_daily \
  --start-date 2026-09-01 --end-date 2026-09-30 \
  --max-active-runs 2 --yes

经验值是**「回填并发不超过正常调度并发的 30%」**,这样即使回填把资源占满,正常调度仍有 70% 的余量。回填的调度时机也要避开业务高峰:数据仓库类任务适合在凌晨正常调度结束后启动回填,而不是和正常调度同时跑。

另一个细节是回填的执行顺序。按区间从早到晚回填符合数据依赖方向,避免下游先跑到「依赖的上游还没回填」的区间而失败,但顺序回填耗时更长。折中方案是「按天顺序、天内并行」。

9. 补数:与回填的区别

「回填」与「补数」常被混用,但它们的语义有明确区别:

维度回填(backfill)补数(repair / catch-up)
触发原因主动补齐历史修复缺失或错误的数据
数据来源与正常调度相同可能来自对账、上游重发
范围通常是连续区间通常是离散的少数区间
是否可预期是(计划内)否(故障驱动)
关键要求幂等、可控并发幂等、可追溯、可对账

补数的难点在**「如何发现缺了什么」**。三种手段:对账(与上游或独立数据源比对总量与校验和)、预期区间差集(第 7 节的 SQL)、业务校验(关键指标为空或异常偏低时告警)。

对账是最可靠的,因为它不依赖「自己是否记得跑过」,而是与外部事实比对:

-- 对账:源库与目标表的行数与金额必须一致,任何一项不等都要告警
SELECT
    (SELECT COUNT(*) FROM src.orders WHERE created_at::date = :biz_date) AS src_cnt,
    (SELECT COUNT(*) FROM dws.orders_daily WHERE biz_date = :biz_date)   AS dst_cnt,
    (SELECT SUM(amount) FROM src.orders WHERE created_at::date = :biz_date) AS src_amt,
    (SELECT SUM(amount) FROM dws.orders_daily WHERE biz_date = :biz_date)   AS dst_amt;

补数执行后必须回写对账结果,形成闭环。没有闭环的补数系统会陷入「补了但不知道补对没有」的状态,反复补同一批数据。

10. 幂等:回填能安全运行的前提

回填与补数能安全运行的前提只有一个:任务幂等。同一区间的任务跑一次与跑十次,结果必须相同。

四种幂等模式的适用场景:

-- 模式一:按区间先删后插(最通用)
DELETE FROM dws.orders_daily WHERE biz_date = :biz_date;
INSERT INTO dws.orders_daily SELECT ... WHERE created_at::date = :biz_date;

-- 模式二:分区覆盖(数仓首选,原子性最好)
INSERT OVERWRITE TABLE dws.orders_daily PARTITION (dt = :biz_date)
SELECT ... WHERE created_at::date = :biz_date;

-- 模式三:按主键 upsert
MERGE INTO dws.orders_daily t USING staging s ON t.id = s.id
WHEN MATCHED THEN UPDATE SET ... WHEN NOT MATCHED THEN INSERT ...;
# 模式四:用业务标识去重(适用于有副作用的动作)
@task
def send_report(biz_date):
    if already_done("send_report", biz_date):   # 用 biz_date 而非 run_id
        return
    send(...); mark_done("send_report", biz_date)

模式一是最通用的,但要注意**「先删后插」之间有一个空窗口**:如果此时下游来读,会读到空数据,解决方式是用事务包住(同一事务内删除与插入)。

模式四适用于无法用 SQL 表达的动作(发邮件、调外部接口)。它的关键是去重键的粒度:用 biz_date + 动作名 而不是 run_id,因为回填时 run_id 会变但业务动作只应执行一次。

11. 时区陷阱

时区是调度系统里最容易被低估的复杂来源。核心原则是**「存储用 UTC,展示与业务判断用本地时区,且本地时区必须显式声明」**。

三类常见错误:调度时区与数据时区不一致(调度按 UTC 0 点触发,业务却认为「昨天」是北京时间昨天,处理的数据区间比业务预期晚 8 小时);服务器时区漂移(容器默认 UTC、物理机默认本地时区,迁移后触发时刻整体偏移且没人发现);跨时区的业务日定义不清(全球业务里「一天」的边界是 UTC 日、北京日还是各地本地日,导致同一份报表在不同区域看到不同数字)。

治理方式是在所有地方显式声明时区,不留任何默认值:

BIZ_TZ = ZoneInfo("Asia/Shanghai")
def biz_date_of(ts):                  # 业务日边界用业务时区定义
    return ts.astimezone(BIZ_TZ).strftime("%Y-%m-%d")
def to_storage(ts):                   # 存储一律用 UTC
    return ts.astimezone(ZoneInfo("UTC"))

对于跨时区的业务,标准做法是**「统一用一个业务时区定义日边界,其余时区只做展示转换」**。让每个区域用本地日会导致同一笔跨境交易在两个区域的报表里落在不同的日期,对账永远对不上。

12. 夏令时与时间跳跃

夏令时(DST)会让「本地时间」出现重复或缺失,是调度系统最经典的坑:

春季开始(时钟前跳):02:00 -> 03:00,本地时间 02:00~02:59 不存在
  影响:cron 表达式 0 2 * * * 在当天不触发(或触发时刻不定)

秋季结束(时钟回拨):02:00 -> 01:00,本地时间 01:00~01:59 出现两次
  影响:cron 表达式 0 1 * * * 在当天触发两次

两种后果都会破坏「一天一次」的假设:春季漏跑一天,秋季重复跑一次。使用 UTC 调度的系统天然免疫这个问题(UTC 没有夏令时),但业务日边界仍需处理。

治理方案按优先级:用 UTC 做调度,只在业务展示层转本地时区,这是最彻底的方案;如果业务坚持「每天凌晨 2 点」的本地语义,则必须在 DST 切换日加补偿任务检查漏跑与重复;自研调度器要把「下一次触发时间」的计算交给时区库,而不是自己加 86400 秒。

# 用 aware datetime 做加法,时区库会自动处理 DST 跳跃
nxt = dt.datetime(2026, 3, 8, 1, 30, tzinfo=ZoneInfo("America/New_York")) \
        + dt.timedelta(hours=1)          # 结果是 03:30 EDT,跳过不存在的 02:30

13. 调度漂移与错过执行

调度漂移指「实际执行时刻逐渐偏离预期时刻」。三种成因:

成因一:任务堆积
  上一个区间的任务还没跑完,下一个区间的任务排队等待
  后果:执行时刻越来越晚,最终「追上」甚至跨区间

成因二:调度器重启
  重启期间的触发点被跳过
  后果:漏跑若干区间

成因三:时钟不同步
  多节点调度器之间时钟偏差,导致触发时刻不一致
  后果:重复触发或漏触发

成因一最普遍,处理方式是「限制并发 + 设置错过执行(misfire)策略」。多数调度器提供「立即补跑 / 跳过 / 只跑最后一次」三个选项,选择依据是业务语义——每小时抓行情的任务应该只跑最后一次(中间的小时没有意义),而每天结算的任务必须立即补跑。

成因二需要「启动时补跑」机制:调度器启动时检查「最近 N 个应该触发的区间是否有记录」,缺了就补。这个机制同时也解决了「调度器故障期间漏跑」的问题,是调度系统必备的自愈能力。

成因三靠 NTP 解决:所有调度节点必须同步时钟(NTP 偏差控制在 100 ms 内),且「谁来决定触发」应该是单点(主调度器)而不是多节点各自判断。

14. SLA 的定义与度量

SLA(服务等级协议)在调度场景里有两个不同的含义,必须区分:

含义一:完成时限(deadline)
  「每天 8 点前必须产出报表」,衡量的是「绝对时间点」

含义二:执行时长(duration)
  「这个任务必须在 30 分钟内跑完」,衡量的是「相对耗时」

多数调度器两者都支持(Airflow 的 sla 参数实际上是「从区间开始到任务完成的时限」,更接近含义一)。定义 SLA 时要明确用哪个,并注意SLA 的起点:是从数据区间结束算起,还是从任务实际开始算起?前者反映「业务可用性」,后者反映「任务性能」。两者都有用,但告警的处置方式完全不同。

SLA 的度量方式建议用「达成率」而不是「单次是否超时」:

SLA 达成率 = 达成 SLA 的运行次数 / 总运行次数
告警规则:连续 3 次未达成,或滚动 7 天达成率低于 95%

用达成率的好处是容忍偶发抖动:一次因上游延迟导致的超时不该触发告警,但如果连续发生或形成趋势,就说明系统性问题。

15. SLA 监控与告警

SLA 监控的关键难点是**「没跑」比「跑失败」更难发现**。任务失败会有明确的失败状态与告警,而任务根本没被创建(调度器故障、区间被跳过)时,什么都不会发生,直到有人发现报表是空的。

因此必须有「预期运行存在性」检查:

-- 检查昨天的关键任务是否有成功记录(而不是检查它是否失败)
SELECT flow_name
FROM flow_expectation e
LEFT JOIN flow_run r
    ON r.flow_name = e.flow_name
   AND r.biz_date = CURRENT_DATE - 1
   AND r.status = 'SUCCESS'
WHERE e.is_critical = TRUE AND r.id IS NULL;
-- 有结果就告警:说明「应该跑了但没跑」

告警分层设计:

层级触发条件通知对象响应时间
P1关键任务未按预期运行值班工程师15 分钟
P2任务失败且重试耗尽流程负责人1 小时
P3SLA 达成率下降团队频道次日

P1 必须直接呼人,不能只发消息。因为「关键任务没跑」的影响是持续的(下游数据全错),而消息通知在夜间很容易被忽略。P2 及以上都应该有明确的「谁负责」,而不是发到一个没人看的群里。

16. 优先级与抢占

当资源有限而任务众多时,必须区分优先级。三个维度决定优先级:

业务紧急度:影响线上的任务 > 影响报表的任务
SLA 紧迫度:临近 deadline 的任务 > 时间宽裕的任务
资源占用:短任务优先(减少队列等待时间)

「短任务优先」是一个常被忽略但收益很大的策略。在队列里,一个跑 2 小时的任务会阻塞后面所有任务;如果先跑那些 5 分钟的任务,整体平均等待时间会显著下降。这就是调度理论里的「最短作业优先」,代价是长任务可能饿死——因此需要配合「老化」机制(等待越久优先级越高)。

抢占(preemption)只在特定场景下可行:可中断且可恢复的任务。批处理任务通常可以抢占(保存检查点后暂停),但有外部副作用的任务不能(不能中途打断一次支付)。抢占的实现复杂度很高,多数团队应该先用「队列优先级 + 资源配额」解决问题,而不是上抢占。

17. 调度器的可观测与自愈

调度器本身是最需要被监控的组件,因为它一旦故障,所有任务静默停止。四个必备监控项:

指标含义告警阈值
调度心跳调度器循环的存活信号超过 3 个周期无心跳
待调度队列长度已到点但未派发的任务数持续增长
触发延迟从预期触发到实际触发的时差P99 超过 60 秒
区间缺口预期存在但缺失的运行大于 0

「区间缺口」是最重要的指标,它是唯一能发现「静默漏跑」的信号。实现方式是把「预期区间」当成一个虚拟的检查清单,定期与实际运行记录做差集(第 7 节的 SQL 就是这个逻辑的运维化)。

自愈能力包括三项:调度器启动时补跑错过的区间、任务失败后按策略自动重试、检测到区间缺口后自动触发补数。第三项要谨慎——自动补数可能掩盖真正的问题(比如上游持续故障),建议先自动补一次,失败则告警人工介入。调度器的可观测体系与 工作流可观测与调试 里讲的实例级观测是互补的:一个看「系统是否在运转」,一个看「实例是否正常」。

18. 大规模调度的容量规划

当流程数达到几千、任务数达到几十万时,调度器本身的性能会成为瓶颈。三个量级参考:小规模(少于 500 流程 / 1 万任务实例每天)单调度器加 PostgreSQL 足够;中规模(5005000 流程 / 1 万50 万实例每天)需要独立调度进程与数据库读写分离,元数据库要分区与定期清理;大规模(超过 5000 流程 / 50 万实例每天)需要多调度器按流程哈希分片、事件驱动替代轮询、元数据冷热分离。

调度的性能瓶颈通常在三处。一是 DAG 解析:每次解析都要执行 DAG 文件的顶层代码,流程多时解析时间线性增长(Airflow 的解法是把解析拆到独立进程并控制解析频率)。二是元数据库写入:每次状态变化都写库,吞吐上限由数据库决定(解法是批量写入 + 历史归档)。三是调度循环的扫描成本:每个周期扫描所有流程的触发条件,流程多时单次扫描变慢(解法是分片 + 索引优化)。

容量规划的经验法则是**「为峰值预留 3 倍余量」**:日常 30% 负载、峰值(月末、大促)可能到 90%,如果按日常负载规划容量,峰值时必然堆积。

19. 落地路线图

  • 第 1 周:梳理所有任务的触发方式与数据区间定义,找出用「当前时间」做业务判断的任务(这是回填事故的根源)。
  • 第 2 周:为关键任务补齐幂等(按区间覆盖或 upsert),用「同一区间跑两次」验证结果一致。
  • 第 3 周:实现「预期区间差集」查询,把它做成日常巡检项,接上告警。
  • 第 4 周:把调度时区统一为 UTC,业务日边界显式声明为业务时区,检查是否存在 DST 影响的任务。
  • 第 5 周:做一次回填演练(回填 7 天),观察并发控制与资源占用,调整 max_active_runs 与池配额。
  • 第 6 周:定义关键任务的 SLA(明确是时限还是时长),建立分层告警与责任人。

顺序上「数据区间正确」与「幂等」必须最先做,因为它们是回填与补数能安全运行的前提,也是后续所有运维动作的基础。

20. 权衡取舍

选择收益代价
定时触发实现简单、可预测上游延迟时静默读到不完整数据
依赖感知触发数据正确性有保证链路变长、排查变难、需要超时策略
catchup 自动补齐上线即补齐历史首次上线可能瞬间创建大量运行
显式回填影响面可控需要人工判断回填范围
差集补数幂等自愈、可自动化需要维护「预期区间」的定义
按区间先删后插通用、易理解存在空窗口,需要事务保护
分区覆盖原子性最好依赖数仓支持
UTC 调度免疫 DST 问题业务方需要心智转换
本地时区调度贴合业务直觉必须显式处理 DST
立即补跑 missed不丢数据可能引发堆积与资源争抢
跳过 missed避免堆积有数据缺口,需要事后补数
短任务优先平均等待时间短长任务可能饿死,需要老化机制

21. 常见坑清单

  1. 任务用 NOW() 而不是数据区间做查询条件,回填时所有区间都处理今天的数据。
  2. catchup=True 配上一年前的 start_date,上线瞬间创建 365 个运行实例打爆资源。
  3. 回填时不限制并发,30 天的区间同时起跑,把正常调度的任务挤掉。
  4. 容器默认 UTC 而物理机是本地时区,迁移后所有触发时刻整体偏移且无人发现。
  5. 用本地时区调度且未处理 DST,春季漏跑一天、秋季重复跑一次。
  6. 传感器不设 timeout,上游永久故障时流程永久挂起。
  7. ExternalTaskSensor 的 execution_delta 未对齐两个 DAG 的周期,永远等待不存在的上游实例。
  8. 依赖「跳过」的语义做判断,但没实测 skipped 在下游触发规则里的行为。
  9. 回填的任务不幂等(用 INSERT 而非覆盖),回填一个月数据翻倍。
  10. 只监控「任务失败」,不监控「任务未按预期运行」,调度器故障静默数小时。
  11. 补数没有对账闭环,补完不知道是否补对,反复补同一批数据。
  12. 调度节点时钟不同步,多节点各自判断触发,产生重复或漏触发。
  13. 间隔调度的锚点用「上次执行时间 + 间隔」,重启后区间边界漂移,与下游对不上。
  14. 抢占机制用在有外部副作用的任务上,中断后产生不一致状态。

22. 小结

调度系统的三条底线是:数据区间语义正确(所有业务判断从区间派生,不用当前时间)、任务幂等(同一区间跑多次结果一致)、能发现漏跑(预期区间的差集巡检)。这三条决定了回填与补数能否安全执行,也决定了故障能否被及时发现。

依赖感知触发比定时触发更正确,但它的成本是链路复杂度。折中的工程实践是「关键路径用依赖感知、非关键路径用定时 + 传感器兜底」,而不是全盘切换。切换前要先解决超时策略与可观测性,否则会把「读错数据」的问题换成「流程挂死」的问题。

时区与 DST 是唯一无法靠工程手段完全规避的问题——它来自人类对时间的定义。最彻底的解法是「调度用 UTC、业务日边界显式声明」,把复杂性集中到一个地方(业务日边界的定义),而不是散布在几百个任务里。回填、补数与幂等的通用设计参见 Airflow DAG 调度体系 与 Dagster 与 Prefect 数据编排 。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「工作流引擎」更多文章

  1. 工作流成本优化
  2. 执行器与资源隔离
  3. 工作流数据传递与 Schema