小王算钱

Interview Prep · Data Infra

任务调度
从 cron 到 Airflow 和 Temporal

cron 只会到点跑命令。数据管道还需要另外四样东西:任务之间的依赖、每次运行的状态、失败后的重试和补跑,以及不依赖单台机器。下面七个系统按出现的先后排列,每一个都在补上一个缺的东西,同时带来新的问题。

时间线

从 cron 到 Temporal,每一步在补上一步缺的什么

  1. 出处

    Maxime Beauchemin 2014 年 10 月在 Airbnb 开始写,2015 年 6 月正式公布,2016 年 3 月进入 Apache 孵化器,2019 年 1 月成为顶级项目

    它要解决的问题

    Oozie 的 XML 难写,又离不开 Hadoop;Luigi 没有定时触发,也看不到每次运行的历史。Airflow 想要的是:Python 定义、自带调度,每次运行的状态都能在界面上查到,失败的运行也能在界面上重跑。
    Airflow 3 的组成
    • DAG 文件Python 代码,放在 DAG bundle 里
    • DAG processor解析 DAG 文件,序列化后写入元数据库
    • scheduler(内含 executor)创建 DAG run,挑出可以跑的任务,交给 executor
    • 元数据库PostgreSQL 或 MySQL,存 DAG、运行和任务的状态
    • worker执行任务;Airflow 3 里通过 API server 汇报状态,不直接连数据库
    • API serverREST API 和界面,也接收任务汇报的状态
    • triggerer(可选)在 asyncio 事件循环里等待被挂起的任务

    它怎么工作

    • DAG 写在 Python 文件里。DAG processor 解析这些文件,把结构序列化后存进元数据库。
    • scheduler 循环做三件事:给到期的 DAG 创建 DAG run;检查哪些任务的上游都成功了;在 pool 和并发上限允许的范围内,把任务交给 executor。
    • executor 决定任务在哪里跑:LocalExecutor 在调度器所在的机器上跑,CeleryExecutor 交给常驻的 worker,KubernetesExecutor 每个任务起一个 pod。
    • 按 data interval 调度时(Airflow 2 的默认行为),每个 DAG run 对应一段要处理的时间,logical date 是这段时间的起点,运行在这段时间结束之后才触发。Airflow 3 里直接写 cron 表达式,默认不再划分时间段,logical date 就是计划触发的那一刻。
    • 所有状态都在元数据库(通常是 PostgreSQL 或 MySQL)。把一个任务实例 clear 掉,调度器就会重新跑它。

    定义和提交一个流程

    1. 写一个 Python 文件,定义 DAG:schedule、起始日期和任务。
    2. 用 >> 或 TaskFlow 的函数调用声明任务之间的依赖。
    3. 把文件放进 DAG bundle。
    4. DAG processor 解析它并写入元数据库,之后它就出现在界面上。

    一次运行的过程

    1. 到了触发时间,调度器创建 DAG run 和它的任务实例。
    2. 调度器把上游都成功的任务标成 scheduled。
    3. 在 pool 和并发上限内把它们标成 queued,交给 executor。
    4. worker 执行任务并汇报状态。
    5. 任务成功,调度器在下一轮放行它的下游。
    6. 任务失败且有重试次数,状态变成 up_for_retry,过了重试间隔,再由调度器重新安排。
    7. 所有任务都结束,DAG run 标记为成功或失败。

    优点

    • DAG 是 Python 代码,可以动态生成。
    • 自带调度,每次运行有自己的 logical date,重跑和补跑是现成的。
    • 界面上能看到每个任务每次运行的状态和日志。
    • 现成的 operator 和 provider 很多,executor 可以换。
    • 调度器可以同时跑多个,互为备份。

    缺点

    • 任务不是上游一结束就立刻开始,要等调度器的下一轮循环,所以不适合大量很短的任务。
    • 元数据库是所有组件的中心,它慢了,整个系统都慢。
    • 任务之间传数据靠 XCom,只适合传小的值。大数据要自己放到外部存储。
    • DAG 是代码。在 Airflow 3 之前,改了 DAG 文件,历史运行在界面上会按新结构显示。

    用在哪里,现在还在哪里用

    按时间表运行、有明确开始和结束的批处理管道:ETL、报表、模型的定期训练。

    被什么取代,为什么

    没有被取代。Dagster 这类系统在「按数据而不是按任务来调度」上走得更远,Airflow 3 也在往这个方向加功能。

    面试时,一句话

    Python 定义 DAG,调度器按时间段创建运行,状态存在数据库,执行方式可插拔。

    面试时,两分钟

    Airflow 的核心是 DAG processor、scheduler、元数据库和 API server,再加上执行任务的 worker。DAG processor 解析 Python 文件,把 DAG 的结构存进元数据库。scheduler 循环做三件事:给到期的 DAG 创建运行,找出上游都完成的任务,在并发限制内交给 executor。executor 是 scheduler 里的一项配置,决定任务在哪跑,可以是本机、Celery worker 或者每任务一个 pod。元数据库存所有状态。最关键的概念是 logical date:每次运行有一个固定的逻辑时间,重跑时不变,所以重跑和补跑才说得清。按 data interval 调度时,一次运行处理一个确定的时间段,在这个时间段结束后才触发,所以每天的 DAG 在 10 月 2 日跑的那次,logical date 是 10 月 1 日。Airflow 3 改了默认值:直接写 cron 表达式时不再划分时间段,logical date 等于触发时间。它的弱点是调度有延迟,不适合大量短任务,并且重度依赖元数据库。2.0 加了多调度器,3.0 让 worker 通过 API 汇报而不直接连数据库,还加了 DAG 版本。

    值得补一句的

    最容易被追问的是 logical date。按 data interval 调度的每天的 DAG,10 月 2 日零点触发的那次运行,logical date 是 10 月 1 日,因为它处理的是 10 月 1 日的数据。Airflow 3 改了 cron 表达式的默认行为,同一次运行的 logical date 是 10 月 2 日。回答前先说清是哪个版本、哪种 timetable。

先懂这七个概念

后面反复出现的词

DAG
有向无环图。节点是任务,箭头是「谁要等谁」。没有环,所以一定能排出一个先后顺序。
批处理调度器几乎都用它。有环就意味着两个任务互相等,永远开始不了。
logical date / nominal time
一次运行「应该处理哪个时间段的数据」,不是它实际在哪一刻开始跑。Oozie 叫 nominal time,Airflow 叫 logical date。
有了它,重跑和补跑才有意义:重跑 10 月 1 日那次,处理的仍然是 10 月 1 日的数据。
幂等(idempotent)
同一个任务对同一个时间段跑一次和跑多次,结果一样。
没有调度器能保证任务恰好跑一次。重试、补跑、调度器自己出故障,都可能让同一个任务再跑一次。
backfill / catchup
补跑过去的时间段。catchup 是调度器自动补上错过的周期,backfill 是人指定一段历史去跑。Airflow 3 里 catchup 默认是关的。
新上线一条管道要补历史数据,或者逻辑改了要重算,都靠它。任务不幂等就没法补。
executor
决定任务在哪里、以什么形式执行的那一层:本机进程、常驻 worker,还是每个任务一个容器。
它决定了隔离程度和启动延迟。常驻 worker 启动快但互相影响,每任务一个 pod 隔离好但启动慢。
sensor / 数据触发
不按时间,而是等一个条件成立再跑:文件到了、上游的表更新了。
只按时间触发,就得猜上游几点跑完。猜早了读到不完整的数据,猜晚了白等。
durable execution
把代码执行到哪一步持久化下来,进程挂了能接着跑,而不是从头开始。
这是 Temporal 这一类系统和 DAG 调度器的分界线。

展开讲 · 一

Airflow 的一次运行,从头到尾发生了什么

以一个每天跑一次、按 data interval 调度的 DAG 为例,看 10 月 1 日这一天的数据是怎么被处理的。

  1. 1

    DAG 文件被解析

    DAG processor 读到 DAG 文件,执行它,得到 DAG 的结构,序列化后写进元数据库。调度器之后只读数据库里的这份结构,不碰 DAG 代码。

  2. 2

    一个时间段结束

    每天跑一次的 DAG,10 月 1 日这个 data interval 在 10 月 2 日零点结束。在这之前,这次运行不会被创建。这是按 data interval 调度的行为,也是 Airflow 2 的默认行为。Airflow 3 里直接写 cron 表达式,这次运行同样在 10 月 2 日零点触发,但 logical date 是 10 月 2 日,不带时间段;要原来的行为,得用 CronDataIntervalTimetable。

  3. 3

    调度器创建 DAG run

    调度器发现这个 DAG 该有一次新的运行,就在数据库里插入一条 DAG run,logical date 是 10 月 1 日,同时为 DAG 里的每个任务建一条任务实例。

  4. 4

    挑出可以跑的任务

    调度器检查任务实例:上游都成功了的(这是默认的触发规则),状态改成 scheduled。然后在 pool 和并发上限允许的范围内把它们改成 queued,交给 executor。这一步会对 pool 表的行加锁,所以多个调度器同时运行也不会超出 pool 的上限。

  5. 5

    任务执行并汇报

    worker 拿到任务,在一个子进程里执行。Airflow 3 里,任务通过 API server 汇报状态;Airflow 2 里,任务代码直接连元数据库。

  6. 6

    失败和重试

    任务失败且还有重试次数,状态变成 up_for_retry,等过了重试间隔再回到调度器手里。重试用完,状态是 failed,下游变成 upstream_failed。

  7. 7

    人来重跑

    修好问题后,把失败的任务实例 clear 掉。它的状态被清空,调度器在下一轮就会重新挑中它。logical date 不变,所以它处理的仍然是 10 月 1 日的数据。

展开讲 · 二

调度器给不了「恰好一次」,任务要自己幂等

这是调度系统面试里最常被追问的一点。它和具体用哪个系统无关。

  1. 1

    为什么做不到「恰好一次」

    worker 跑完了任务,在汇报成功之前断了网。调度器没收到结果,分不清「没跑」和「跑了但没汇报」。它只有两个选择:再跑一次,任务可能被执行两次;不再跑,任务可能一次都没完成。配了重试是前一种,叫至少一次;没配重试是后一种,叫至多一次。哪一种都不是恰好一次。

  2. 2

    按时间段覆盖写,不要追加

    任务写的是「10 月 1 日这个分区」,每次都整个覆盖。跑几次结果都一样。如果是往一张表里 INSERT,跑两次就是两份数据。

  3. 3

    用 logical date,不要用当前时间

    任务里写 now() 或者「昨天」,补跑上个月的那次运行时,它处理的却是今天的数据。时间段必须从调度器传进来。

  4. 4

    先写临时位置,再原子地改名

    写到一半失败,留下半个文件,下游会把它当成完整的数据。Luigi 靠输出是否存在来判断完成,所以把这一条写进了设计目标。

  5. 5

    对外部系统的调用带上幂等键

    发邮件、扣款这类操作没法覆盖。给每次调用一个由「流程 ID + 步骤」算出来的键,让对方去重。Temporal 的 Activity 默认无限次重试,文档建议把 Activity 写成幂等的。

横向对比

七个系统放在同一张表里

系统怎么定义流程靠什么触发状态存在哪任务在哪里跑最适合
croncrontab 的一行时间不保存本机进程单机上互不依赖的定时任务
OozieXML(hPDL)时间 + 数据到位关系数据库Hadoop 集群上的作业还在运行的 Hadoop 集群
LuigiPython 类没有,靠 cron 启动输出文件是否存在启动它的那个 worker 进程规模不大、以文件为产出的管道
AirflowPython(DAG)时间、asset 更新、外部事件元数据库由 executor 决定:本机、Celery worker 或 pod按时间表跑的批处理管道
Argo WorkflowsYAML(Kubernetes CRD)提交即运行,定时用 CronWorkflowWorkflow 资源的 status(etcd)每一步一个 pod已经在 Kubernetes 上的容器化任务,如 ML 训练和 CI
DagsterPython(asset)时间、sensor、上游 asset 变化数据库,记录运行和 asset 的事件由 run launcher 和 executor 决定关心数据血缘和新鲜度的数据平台
Temporal普通代码(Workflow 函数)API 调用、signal、定时Event History用户自己部署的 worker长时间运行的业务流程

业界最新的方向 · 2026 年 10 月核对

这几年在往哪里走

  • Hadoop 时代的调度器退场

    Apache Oozie 在 2025 年 2 月退役,进了 Apache Attic。Apache 的页面没有写原因。一个合理的解释是:计算离开了 Hadoop,一个只为 Hadoop 设计的调度器就没有了位置。

  • Airflow 3 把任务和调度器的数据库隔开

    Airflow 3.0 在 2025 年 4 月 22 日发布。任务不再直接连元数据库,而是通过 API server 汇报;同时加了 DAG 版本、由调度器管理的 backfill 和按 asset 调度。3.0 还改了两个默认值:catchup 默认关闭;直接写 cron 表达式时不再划分 data interval,logical date 等于触发时间。之后 3.1(2025 年 9 月)加了 human-in-the-loop 和 deadline alerts,3.2(2026 年 4 月)加了 asset 分区。查到的最新版本是 3.3.2,2026 年 9 月 17 日发布。

  • 从调度任务到调度数据

    Dagster 把 asset 作为基本单位,Airflow 从 2.4 开始也能按数据的更新来触发 DAG,在 Airflow 3 里这个概念叫 asset。方向是一致的:让调度器知道任务产出了什么。

  • 做编排工具的公司在合并

    2026 年 7 月 13 日,Prefect 宣布收购 Dagster Labs。公告说 Dagster 保留原来的名字和开源许可,合并后的公司用 Prefect 这个名字。对使用者来说,选型时要把项目背后的公司也算进去。

  • 触发条件不再只有时间

    Oozie 的 coordinator 早就能等数据到位。现在的系统把这件事做成了一等功能:上游 asset 更新、外部事件、sensor 都可以触发一次运行。

  • 执行层交给 Kubernetes

    Argo 在 2022 年 12 月从 CNCF 毕业。Airflow 和 Dagster 也都可以把每个任务放进一个 pod 里跑。有了 Kubernetes,调度器不必再自己维护一组常驻的 worker。

  • 持久化执行成了单独的一类

    Temporal 这一类系统不和 DAG 调度器抢数据管道,它们面向的是业务流程。面试里被问到「设计一个可靠的多步骤流程」,要先分清题目说的是哪一类。

常见误解

面试里容易说错的七件事

  • 误解:Airflow 的 logical date 是任务开始跑的时间。

    logical date 是调度器给这次运行定的逻辑时间,重跑时不变,它不是任务实际开始的时刻。按 data interval 调度时(Airflow 2 的默认行为),它是所处理时间段的起点:每天的 DAG 在 10 月 2 日零点触发的那次运行,logical date 是 10 月 1 日。Airflow 3 里直接写 cron 表达式,logical date 默认等于计划触发的时间,也就是 10 月 2 日零点。任务要用调度器传进来的时间决定读哪一天的数据,不要用当前时间。

  • 误解:调度器能保证每个任务恰好执行一次。

    没有调度器能保证这一点。配了重试,得到的是至少一次:worker 跑完任务但没来得及汇报,任务会再跑一次。没配重试,得到的是至多一次。Kubernetes 的文档对 CronJob 说得更直接:可能创建两个 Job,也可能一个都不创建。恰好一次的效果要靠任务自己幂等。

  • 误解:Airflow 是处理数据的引擎。

    Airflow 是编排工具,它决定什么时候跑什么。真正的计算通常发生在别处:Spark、数据仓库、一个容器。它也不是为流处理设计的,它管的是有始有终的批处理。

  • 误解:Temporal 是更新、更好的 Airflow。

    它们解决的不是同一个问题。Airflow 按时间表处理一个个时间段的数据,Temporal 让一个个长时间的业务流程不丢进度。Temporal 没有补跑和 data interval,Airflow 没有逐步持久化。

  • 误解:Luigi 的中心调度器负责执行任务。

    它不执行任何任务。它只保证同一个任务不会有两个实例同时跑,并提供一个界面。任务是 worker 进程跑的,而 worker 要靠 cron 或人来启动。

  • 误解:给 Airflow 配上 KubernetesExecutor,就和 Argo Workflows 一样了。

    只有执行那一层一样,都是每个任务一个 pod。Airflow 仍然有自己的调度器和元数据库,状态存在数据库里;Argo 的状态存在 Workflow 这个 Kubernetes 资源里,调度由一个 controller 完成。

  • 误解:机器重启后,cron 会把错过的任务补上。

    不会。cron 只看当前这一分钟匹配哪些行,过去的不管。能补跑是工作流调度器才有的能力,前提是每次运行都对应一个明确的时间段。

Staff 级面试题

先自己答,再点开看

▶公司现在用 cron 跑 200 个脚本,经常出现下游跑的时候上游还没跑完。你会换成什么?
  • 先分清问题是什么。这里缺的是依赖和状态,不是执行能力。
  • 如果这些脚本是按天、按小时处理数据的批处理,Airflow 这类 DAG 调度器对得上:依赖显式写出来,失败能重试,历史能补跑。
  • 迁移不必一次完成。可以先让调度器原样调用现有脚本,把依赖关系建起来,再逐步把脚本改成幂等的。
  • 要说代价:多了调度器、数据库和 worker 要运维。只有十几个互不依赖的脚本时,cron 加报警可能就够了。
  • 最后取决于脚本之间依赖有多密、失败后人工处理的成本有多高,以及团队有没有人能维护一套新系统。
▶设计一个调度器,要求调度器进程挂掉不影响任务按时触发。
  • 先说状态放哪。DAG run 和任务实例的状态都放在数据库里,调度器进程本身不保存状态,挂了换一个就能接着干。
  • 再说多个调度器怎么不重复触发。两条路:选主(同一时刻只有一个在工作,通常借助 ZooKeeper、etcd 或数据库里的一把锁),或者多活(多个同时工作,靠数据库的行锁保证同一个任务只被一个调度器拿到)。
  • Airflow 2.0 走的是多活:用 SELECT ... FOR UPDATE 加锁,不引入额外的协调服务。Oozie 的高可用则是多个 server 共享数据库,用 ZooKeeper 做锁。
  • 还要处理错过的触发:调度器全挂了十分钟,恢复后要能算出这十分钟里该触发哪些运行。这要求「该不该触发」是从时间表和数据库里的记录算出来的,而不是靠内存里的定时器。
  • 选哪条路取决于你已经有什么:已经有可靠的数据库而不想多运维一个组件,就用行锁;触发频率极高、数据库锁成为瓶颈,才值得考虑分片或专门的协调服务。
▶Airflow 和 Temporal 都叫 workflow 引擎,什么时候用哪个?
  • 看流程是按什么触发的。按时间表、处理一个时间段的数据,是 Airflow 的模型。每来一个请求起一个实例、要等外部事件,是 Temporal 的模型。
  • 看流程的形状。Airflow 的 DAG 在运行前基本确定;Temporal 的流程是代码,可以有循环、可以根据中间结果走不同的分支、可以睡上三十天。
  • 看实例的数量。Airflow 的运行次数由时间表决定;Temporal 面向的是同时存在大量 Workflow 实例的场景。
  • 看你需要什么周边能力。补跑、data interval、数据血缘在 Airflow 和 Dagster 里是现成的;Temporal 没有这些。反过来,Temporal 对代码每一步的持久化,Airflow 没有。
  • 两者经常并存:数据平台用 Airflow,业务后端用 Temporal。归根结底要看你的流程更像哪一种:一张每天重画的时间表,还是一个个各自独立、要存活很久的业务流程。
▶一个每天的任务失败了三天才被发现。你怎么补?补的时候要注意什么?
  • 先确认任务是不是幂等的。不是的话,重跑会产生重复数据,要先修任务或者先清掉这三天写了一半的结果。
  • 按 logical date 补,一天一次运行,而不是跑一次处理三天。这样每次运行的输入输出和平时一样,出问题也容易定位。
  • 注意下游。这三天下游要么没跑,要么读了缺数据的结果,也要一起重跑。调度器知道依赖关系时(Airflow 里 clear 时带上下游),这一步可以自动完成。
  • 注意并发。三天一起补,会不会把数据库或集群压垮?要限制同时运行的数量。
  • 注意上游数据还在不在。有的源只保留几天,过期了就补不回来。
  • 之后要补的是报警:失败三天没人知道,说明缺的是 SLA 或 deadline 告警。具体怎么补,取决于任务幂等的程度、下游的范围和源数据的保留期。
▶调度器要支撑每天一百万个任务,瓶颈会在哪里?
  • 先看任务的时长分布。一百万个各跑一小时的任务,和一百万个各跑两秒的任务,是两个问题。
  • 对短任务,瓶颈是每个任务的固定开销:调度器轮询的间隔、起一个 pod 的时间、每次状态变化写一次数据库。这时该考虑把多个小任务合并成一个,或者换成常驻 worker。
  • 对调度器本身,瓶颈通常是数据库:每一轮都要查「哪些任务可以跑」,任务实例表越大越慢。办法是加索引、清理历史、多个调度器分担。
  • 对 Argo 这类把状态放在 Kubernetes 里的系统,瓶颈是 API server 和 etcd,以及单个 Workflow 资源 1 MB 的上限。
  • 解析也是瓶颈:DAG 是代码,几千个 DAG 文件每次都要执行一遍才能知道结构。
  • 先量再动。答案取决于任务时长、是否需要强隔离,以及能接受多大的调度延迟。
▶为什么新一代调度器要从「任务」转向「数据(asset)」?值得迁移吗?
  • 按任务调度时,调度器只知道任务成功了。它不知道任务写了哪张表,所以答不出「这张表是不是最新的」「它坏了会影响谁」。
  • 按 asset 调度时,依赖是在数据之间声明的。调度器因此能做三件事:上游更新了才触发下游、只重算受影响的部分、直接给出血缘。
  • 代价是要重新建模。不产出数据的任务套不进去,团队也要换一种写法。
  • 不一定要换系统。Airflow 2.4 起有 dataset,Airflow 3 改名叫 asset 并支持按它来调度,可以在原有 DAG 上逐步加。
  • 值不值得,取决于你的痛点是不是「不知道数据的状态」。如果痛点只是任务偶尔失败,按任务调度加上好的告警就够了。

动手练习 · 约 5 分钟

四十多行的迷你调度器:依赖、状态、重跑

只用 Python 标准库(3.9 以上),不用装任何东西。把下面的代码存成 mini_scheduler.py。它按日期跑一个四步的 DAG,把每个任务的状态记在 state.json 里。其中 10 月 2 日的 transform 第一次会失败。

import json
import os
from graphlib import TopologicalSorter

# task -> the tasks it depends on
DAG = {"extract": [], "transform": ["extract"], "load": ["transform"], "report": ["load"]}
DATES = ["2026-10-01", "2026-10-02", "2026-10-03"]
STATE_FILE = "state.json"  # plays the part of the metadata database


def run_task(task, date, attempt):
    # transform fails the first time it runs for 10-02, like an upstream outage.
    if task == "transform" and date == "2026-10-02" and attempt == 1:
        raise RuntimeError("upstream timeout")
    if task == "load":
        # Idempotent write: one file per date, overwritten, never appended.
        with open(f"out_{date}.txt", "w") as f:
            f.write(f"rows for {date}\n")


def schedule():
    state = json.load(open(STATE_FILE)) if os.path.exists(STATE_FILE) else {}
    for date in DATES:
        for task in TopologicalSorter(DAG).static_order():
            key = f"{date}/{task}"
            record = state.setdefault(key, {"status": "none", "tries": 0})
            if record["status"] == "success":
                print(f"{key:24} skip (already success)")
                continue
            if any(state[f"{date}/{up}"]["status"] != "success" for up in DAG[task]):
                record["status"] = "upstream_failed"
                print(f"{key:24} upstream_failed")
                continue
            record["tries"] += 1
            try:
                run_task(task, date, record["tries"])
                record["status"] = "success"
            except RuntimeError:
                record["status"] = "failed"
            print(f"{key:24} {record['status']} (try {record['tries']})")
    json.dump(state, open(STATE_FILE, "w"), indent=1)


schedule()
print("output files:", sorted(f for f in os.listdir(".") if f.startswith("out_")))

连着运行两次:

echo "=== run 1"; python3 mini_scheduler.py
echo "=== run 2"; python3 mini_scheduler.py

两次运行的真实输出:

=== run 1
2026-10-01/extract       success (try 1)
2026-10-01/transform     success (try 1)
2026-10-01/load          success (try 1)
2026-10-01/report        success (try 1)
2026-10-02/extract       success (try 1)
2026-10-02/transform     failed (try 1)
2026-10-02/load          upstream_failed
2026-10-02/report        upstream_failed
2026-10-03/extract       success (try 1)
2026-10-03/transform     success (try 1)
2026-10-03/load          success (try 1)
2026-10-03/report        success (try 1)
output files: ['out_2026-10-01.txt', 'out_2026-10-03.txt']
=== run 2
2026-10-01/extract       skip (already success)
2026-10-01/transform     skip (already success)
2026-10-01/load          skip (already success)
2026-10-01/report        skip (already success)
2026-10-02/extract       skip (already success)
2026-10-02/transform     success (try 2)
2026-10-02/load          success (try 1)
2026-10-02/report        success (try 1)
2026-10-03/extract       skip (already success)
2026-10-03/transform     skip (already success)
2026-10-03/load          skip (already success)
2026-10-03/report        skip (already success)
output files: ['out_2026-10-01.txt', 'out_2026-10-02.txt', 'out_2026-10-03.txt']
  • 第一次运行,10 月 2 日的 transform 失败了。它的下游 load 和 report 没有执行,状态是 upstream_failed。10 月 1 日和 10 月 3 日不受影响,因为每个日期是一次独立的运行。
  • 第二次运行相当于重跑。已经成功的任务都被跳过,只有 10 月 2 日失败的 transform 和它的下游真正执行了。这就是「状态存下来」的用处。
  • transform 第二次的尝试次数是 2,而 load 是 1。load 第一次根本没有被执行,所以不算一次尝试。
  • load 写文件用的是覆盖。把 state.json 删掉再连跑两遍,所有任务都会重新执行,但输出文件的内容不变。把 "w" 改成 "a" 再试一次,就能看到不幂等的任务在重跑时会多出重复的行。
  • 这四十多行缺的东西,正好是真实调度器的主要工作:按时间触发、并行执行、多个调度器同时跑时的加锁,以及任务在别的机器上跑时怎么汇报状态。

速查表

考前十分钟过一遍

cron
到点跑命令。没有依赖、状态、重试和补跑,只在一台机器上。
Oozie
Hadoop 上的调度器。XML 写 DAG,coordinator 按时间加数据到位触发。2025 年 2 月退役。
Luigi
Python 版的 Make。输出存在就算完成。没有内置定时触发。
Airflow
Python 定义 DAG,调度器按 data interval 创建运行,状态在元数据库,executor 可换。
Argo Workflows
工作流是 Kubernetes CRD,每步一个 pod,状态存在资源的 status 里。
Dagster
按 asset 编排。声明有哪些数据、各自怎么算,系统推出依赖,自带血缘。
Temporal
持久化执行。记录每一步,靠重放 Event History 恢复。Workflow 要确定性,Activity 要幂等。
logical date
一次运行的逻辑时间,重跑时不变,不是任务实际开始的时刻。按 data interval 调度时是时间段的起点;Airflow 3 的 cron 默认等于触发时间。
幂等
没有调度器能保证恰好一次,重试和补跑都会让任务再跑。按时间段覆盖写,用调度器给的时间,不用当前时间。
调度器高可用
状态放数据库,调度器无状态。多个调度器靠选主或行锁避免重复触发。
怎么选
按时间表处理数据用 DAG 调度器;每个请求一个长流程用持久化执行;任务已经是容器就看 Kubernetes 原生的。

来源

版本号、发布日期和项目状态在 2026 年 10 月核对过,之后可能已经变化。

留言

说说你的看法

有想法或问题都可以写在这里。留言会立刻显示;想收到回复通知再填邮箱。

还没有留言,你可以写第一条。