本节摘要:Apache Airflow 适合用代码声明、需要可见依赖、允许分钟到小时级延迟的批处理与准调度流水线。它不是消息队列、不是流计算引擎、也不是机器学习实验看板。选型时先问三件事:依赖是否必须显式、失败是否必须按图恢复、延迟能否被业务接受。答不上来就先别迁。
阅读完本节,你应当能够:
官方那句“编写、调度、监控”可以翻译成三条责任。编写:依赖和重试写在 Python 里,能进代码评审。调度:按时间窗口或数据就绪创建运行,而不是有人 SSH 上去敲。监控:每次 TaskInstance 的状态和日志能回放。三条里缺任何一条,cron 加聊天群也能凑合;三条都要,Airflow 才开始划算。
它不承担的事情同样清楚。它不提高 SQL 本身的速度,不替代 Spark 的执行计划,不保证跨任务事务。两个任务之间没有数据库级的 begin/commit。你要的一致性,得靠幂等写入、分区覆盖、以及“失败则下游不跑”的图语义。把 Airflow 当成分布式数据库,会在入库任务里发明出无法回滚的半成品。
适合进 Airflow 的典型形状:每日数仓分层、每周对账、模型按天重训、跨系统的“先抽后转再推”、需要 SLA 和责任人的报表。共同特征是有向无环、窗口清晰、失败要可重放。不适合的形状:每秒数万次的点击校验、需要在一次函数调用里睡三个星期还保持毫秒唤醒的聊天会话、没有边界的交互式分析。后者要么用流系统,要么用专门的长工作流引擎,要么根本不该自动化成图。
对照时不要比 GitHub 星星,要比“核心契约你必须接受什么”。原文在替代方案里强调这一点:Kubeflow 把组件定义成带输入输出路径的容器,执行常委托给 Kubernetes 上的工作流引擎,强项是实验可复现;Prefect 更强调动态工作流与状态;Temporal 把长期运行的函数做成可持久化的状态机,适合订单这种跨天对话;Argo 是集群原生的容器工作流,YAML 即图。Airflow 的契约是:Python 即图、元数据库即账本、调度与执行分离。
| 工具 | 核心契约 | 我更倾向何时用 |
|---|---|---|
| cron / 系统定时 | 到点执行,无图 | 单机单脚本,失败人工看 |
| 纯 Celery 队列 | 消息进,函数出 | 在线请求的异步卸荷,不要窗口语义 |
| Airflow | Python DAG + 元数据账本 | 数据管道、有窗口、要审计 |
| Argo Workflows | 容器步骤在集群里编排 | 已经全盘容器化、步骤即镜像 |
| Kubeflow Pipelines | ML 组件与产物路径 | 实验、训练、产物版本是一等公民 |
| Temporal | 持久化工作流函数 | 跨天业务会话、信号与补偿 |
| 流计算 | 事件连续处理 | 延迟以秒或更短计 |
托管与自建不是第三条产品,而是运维切分。自建意味着你负责元数据库高可用、Web 与调度器进程、Worker 扩缩、日志落地、升级停机。托管(云厂商的托管 Airflow、专业运行时厂商)把调度器扩缩、补丁、部分监控接过去,你仍要写 DAG、管 Connection、对业务失败负责。原文提到编排会往“声明 DAG、运行时隐式提供”的方向走,动手时把它理解成:平台可以托管进程,不能托管你的窗口语义错误。

我见过的失败迁入,多半不是选错品牌,而是把错误的形状塞进 DAG。
用 DAG 每分钟空转去“近似实时”。调度器解析、写库、排队的开销会被放大;元数据库变成心跳盘。真有秒级需求,把计算放进流作业,Airflow 只负责每小时对齐一次检查点或每天重建一次表。
用 XCom 传整份数据帧。元数据库会在高峰期锁成一团。正确形状是上游把数据写进对象存储或表,XCom 只传地址。
在任务里起一个永不退出的服务进程。Worker 槽位被永久占用,心跳与重试语义全部失效。长服务交给部署平台,DAG 只做发布或健康检查。
把所有租户写进一张静态巨图,每天改代码发版。原文对动态 DAG 工厂的描述是:元数据驱动生成 dag_id。那是规模到了以后的事;第一刀仍然是一条可评审的静态 DAG,先证明窗口与幂等。
第一套执行器怎么选,只给直觉,细节在第 3 章。开发用 Sequential 或 Local。单机小团队、任务不密集,Local 可以撑一段时间。多机、要隔离、要排队,再上 Celery。任务镜像差异大、希望一次任务一个容器,再上 Kubernetes 执行器。不要因为“生产感”一上来就上 K8s——镜像拉取的八秒延迟,会让你误以为 Airflow 很慢;用 Local 对照同一 DAG,才能把锅分给集群调度而不是 DAG 逻辑。
⚠️ 常见坑:用 Airflow 当通用异步任务队列,把用户点击产生的作业丢进 DAG 触发。DAG 解析与运行创建的粒度配不上请求峰值,界面也会被一次性运行刷屏。
💡 关键直觉:Airflow 卖的是“可审计的因果图 + 时间窗口”,不是“更快的 CPU”。快来自下游引擎;图来自你的声明。
团队规模也影响选型。三个人、十条 DAG,最贵的是把时间语义搞错,不是没上高可用。三十个人、几百条 DAG,没有规范、没有测试闸门、没有池隔离,Scheduler 会在解析和锁上先倒下——那是第 6 章的主题,但边界意识现在就要有:工具能表达契约,不能代替契约评审。
和业务方沟通时,把 SLA 翻译成 DAG 参数:九点前出报表,对应调度时刻、超时、重试、失败回调。原文举过市场部日报的例子:重试两次、间隔五分钟、失败通知到指定渠道、owner 写团队名。这比会议纪要更像可执行合同。若业务其实要的是“用户一点击立刻出结果”,那份合同就不该签给 Airflow。
把候选流水线按五项打分,每项 0 或 1,低于 3 分先别迁。依赖是否必须显式(多步骤且失败要停下游)记 1。是否按窗口重放(补昨天与跑今天必须同一套代码)记 1。延迟是否允许分钟级记 1。是否需要责任人、SLA、审计日志记 1。是否主要是批处理而不是请求级记 1。实时特征校验、点击级风控、长对话订单补偿,往往前几项就拿不满。
和 cron 的分界不是“有没有界面”,而是“有没有图”。只有一个脚本、失败了重跑整个脚本可接受,cron 更简单,少一个元数据库。和纯队列的分界是有没有数据区间:队列消费“这条消息”,Airflow 消费“这个窗口的账”。和 Argo 的分界是步骤是否已经全部容器化、团队是否以 YAML 为图;Python 团队用 Airflow 更顺。和 Kubeflow 的分界是产物与实验是否一等公民;纯数仓分层不必上 ML 平台。和 Temporal 的分界是要不要在工作流里睡很多天等人类信号。
托管怎么选,只给运维直觉。没有 DBA、没有人值班 Broker,优先托管运行时。已有 Kubernetes 平台团队、DAG 量极大、要按任务隔离,再评估 K8s 执行器或在集群里养 Celery。不要因为采购清单上有“云原生”三个字就淘汰 Local 对照——对照永远便宜。
能,而且常常应该。Airflow 管有图的链,cron 管单机清理磁盘这类无依赖杂务。硬把杂务塞进 DAG,只会增加解析量和无意义告警。反过来,把有依赖的五步链拆成五个 cron 再靠睡眠等待,是本教程反对的起点。边界在链,不在工具崇拜。
能。流作业持续写表或对象存储,Airflow 按小时或按 Dataset 做对齐、对账、导出、训练。Airflow 不替代流作业,只在“需要一个闭合窗口的账”时接手。不要用 DAG 去逐条处理流里的事件。
第一刀是静态可评审 DAG 加 Local 或托管起步。第二刀才是池、队列、远程日志。第三刀才是 Celery 或 K8s。第四刀才是动态工厂。把第四刀当第一刀的团队,会在命名空间和解析成本里淹死,还以为是“平台不成熟”。演进是责任增加,不是名词升级。若业务形态变了——从日批变成要秒级——应退出 Airflow 主路径,而不是给 DAG 打激素。退出也是选型能力。本教程不把 Airflow 写成万能城,1.3 的价值就是允许你说不。说不比错迁更便宜。
不要汇报“我们上了业界标准编排器”。汇报三句:哪些链必须有图所以进 Airflow;哪些请求级作业留在队列;哪些训练实验若产物是一等公民则另看 ML 流水线。预算对应的是元数据库与 Worker,不是界面授权数。托管与自建的差别用“谁半夜起来切主库”一句话说清。说不清,就还没选型,只是装上了。装上了之后的错迁成本,高于选错品牌:把秒级链路塞进 DAG,会逼团队把扫描间隔调到荒唐,然后宣布产品不行。产品行不行,在 1.3 的打分表上就已经决定。把打分表附在立项文档里,比附一张架构彩虹图更能保护后续的人。彩虹图不能阻止错误形状。打分表能。能说不的选型,才是动手优先而不是工具优先。工具优先会把所有自动化都称作 DAG,于是名词贬值,真正有窗口的账也被淹没。
立项附件只保留五项 0/1 与总分,以及“不进 Airflow 的理由”一句话。理由允许是“延迟不够分钟级”或“单脚本足够”。允许的理由会在一年后保护你,当有人质问为何没用编排器。质问常常来自工具崇拜。崇拜不看形状。形状在打分表。表在附件。附件比记忆可靠。可靠才能在人员变动后仍说不。说不是 1.3 的核心技能。技能要有档案。档案不要写成论文。论文没人读。五行分数加一句理由,够了。够了就可以进入第 2 章写图。写图的前提是这张图该存在。该存在,才配得上后面所有字段。不该存在,字段再正确也是错工具上的正确。错工具上的正确,维护成本会惩罚所有正确的人。惩罚表现为扫描间隔被调到荒唐、元数据库被心跳盘占满、团队宣布产品不行。产品行。形状不行。形状由打分表拦住。拦住了,教程才有资格继续。继续去写窗口。窗口只属于该进 Airflow 的链。
下一章开始把最小图画成可配置、可集成、可通信的生产 DAG。