Airflow 任务依赖与调度最佳实践:DAG 设计精要

📅 2026/8/2 17:25:09 👁️ 阅读次数 📝 编程学习
Airflow 任务依赖与调度最佳实践:DAG 设计精要

Airflow 任务依赖与调度最佳实践:DAG 设计精要

一、DAG 不是流程图:从依赖本质说起

很多团队刚接触 Airflow 时,把它当成一张普通的流程图来画。结果很快就踩坑。

DAG 的核心不是「画得好看」,而是「依赖关系正确且可重放」。这是调度系统的命门。

如果依赖写错,下游任务会在上游数据还没就绪时就启动。脏数据就这样流向下游。

循环依赖则更隐蔽。Airflow 在解析阶段就会拒绝成环的 DAG,但定位成本高。

还有一些团队喜欢把所有任务塞进一个巨型 DAG。一旦某节点失败,整图都要重算。

合理的做法是按业务边界拆分 DAG。让每个 DAG 内聚、外部通过数据集或传感器解耦。

本文聚焦三件事:DAG 该怎么设计、依赖如何管理、重试与告警怎样才不误伤。

可观测性设计这个环节常被低估。DAG 跑通只是开始,能看懂才关键。

日志要结构化,关键节点要打点。否则一次失败,排障同学要在海量日志里捞线索。

SLA 也要提前约定。哪些任务必须在几点前完成,超时即触发值班,而非等下游投诉。

这些工程习惯,比任何炫技的算子都更能决定调度系统在凌晨是否把你叫醒。

二、依赖管理与触发链路:从上游到下游

依赖管理要区分「硬依赖」与「软依赖」,并善用 Dataset 做跨 DAG 的数据驱动触发。

硬依赖是最常见的一类。任务 B 必须等任务 A 成功,才允许开始执行。这是强约束。

软依赖则更灵活。比如任务 B 等上游数据到达即触发,不必关心上游任务是否成功。

Airflow 2.4 之后引入的 Dataset,让跨 DAG 触发变得声明式,不再依赖人工计时错位。

序列化任务的状态也要管好。失败要有重试,重试耗尽要有告警,告警要能直达负责人。

在 Data-aware 调度模式下,上下游 DAG 通过数据集实现衔接。上游任务完成事实表写入后,会产出特定的 Dataset 对象。该对象作为触发条件,会自动激活下游的指标汇总 DAG 与质量校验 DAG。若指标任务执行失败,系统会执行重试策略;若质检任务失败,则直接触发告警通知负责人。

依赖要尽量窄。一个任务依赖的上游越少,失败时的爆炸半径就越小,排查也越快。

跨 DAG 触发优先用 Dataset,而不是用长轮询传感器。后者既耗资源又容易误判超时。

三、生产级 DAG 与重试告警代码

下面给出一段生产级 DAG 定义。它包含重试策略、超时、告警回调与空数据兜底。

代码强调把失败处理显式化,避免任务静默成功却产出了空结果误导下游。

from datetime import datetime, timedelta from airflow import DAG from airflow.operators.python import PythonOperator from airflow.datasets import Dataset from airflow.utils.email import send_email def _on_failure(context): """任务失败回调:发送告警邮件,附上执行上下文便于快速定位。""" dag_id = context.get("dag").dag_id task_id = context.get("task").task_id ts = context.get("ts") send_email( to=["oncall@corp.com"], subject=f"[Airflow 失败] {dag_id}.{task_id}", html=f"<p>任务在 {ts} 失败,请查看日志。</p>", ) def extract_orders(**kwargs): """抽取订单数据,空结果时显式抛错,避免下游误信空表为正常。""" rows = _query_source() # 假设返回列表 if not rows: raise ValueError("源表返回空,疑似上游延迟,主动失败以便重试") return len(rows) def _query_source(): # 占位:真实场景替换为数据库/数仓查询,返回行列表 raise NotImplementedError("请实现 _query_source 以对接真实数据源") with DAG( dag_id="orders_summary_v2", start_date=datetime(2026, 7, 1), schedule="@daily", catchup=False, # 禁止历史补跑,避免突发负载 default_args={ "retries": 3, # 失败重试三次 "retry_delay": timedelta(minutes=5), "execution_timeout": timedelta(hours=1), # 超时熔断,防挂死 "on_failure_callback": _on_failure, }, tags=["dwd", "orders"], ) as dag: t_extract = PythonOperator( task_id="extract_orders", python_callable=extract_orders, outlets=[Dataset("dataset://fact_orders")], # 声明产出数据集 ) t_extract

重试次数要克制。无限重试会卡住调度器并掩盖真实故障,三到五次通常足够。

执行超时务必设置。否则一个慢查询可能长期占用 worker,拖垮整批任务的并发。

四、边界条件与权衡:何时该信、何时该拦

Airflow 调度也有边界。第一个边界是「任务粒度」。过细会产生海量小任务。

海量小任务会压垮调度器元数据库,让 UI 卡顿、解析变慢,整体稳定性下降。

第二个边界是「跨系统依赖」。Airflow 擅长编排,不擅长做重型数据计算本身。

第三个边界是「时间敏感型触发」。纯定时在上下游波动大时容易错位或空跑。

在 Trade-offs 上,我们主张「声明式 Dataset 优先于隐式计时」。让数据说话。

但也别盲目追求细粒度解耦。拆得太碎,可观测性反而变差,排障成本直线上升。

适用场景:批处理编排、跨系统任务依赖、需要重试与可观测性的 ETL 链路。

禁用场景:毫秒级实时流处理、把 Airflow 当计算引擎硬算 TB 级数据、无监控告警。

还要注意版本演进带来的范式变化。Dataset 触发虽好,但要求调度器版本足够新。

老旧集群若不支持数据感知调度,应继续用传感器或外部编排做跨 DAG 衔接。

不要为了追新特性而强行升级。稳定压倒一切,调度系统是下游所有数据的节拍器。

最后提醒一点:catchup 默认开启时,新 DAG 上线会疯狂补跑历史,务必显式关闭。

五、总结

Airflow 的精髓在依赖正确与可重放,而不在流程图是否花哨。这是调度设计的根本。

依赖要分清硬软,跨 DAG 优先用 Dataset 声明式触发,避免长轮询传感器的资源浪费。

生产 DAG 必须显式配置重试、超时与失败告警,并对空结果主动失败以防误导下游。

粒度要平衡:太粗爆炸半径大,太细压垮调度器。让数据驱动,而不是让计时猜测。