2.1 DAG定义与配置


2.1 DAG 定义与配置

本节摘要:DAG 配置是一组时空约束:start_date 给出原点,schedule 切出数据窗口,catchup 决定要不要清算历史,max_active_runs 限制同时开工的窗口数,default_args 把重试与责任人写成默认宪法。配错其中一项,图看起来仍正确,时钟已经违约。

先说结论

阅读完本节,你应当能够:

  1. 用数据区间而不是墙上时钟解释一次 DagRun 处理哪一段数据
  2. 为生产给出 catchup 与 backfill 的分工
  3. 判断何时开 depends_on_past,何时改用对外部任务的显式等待
  4. 描述动态 DAG 工厂要守住的命名与无副作用解析

一、窗口比闹钟更重要

把调度写成 cron 表达式,容易让人以为 Airflow 只是带界面的闹钟。真正被调度器消费的是窗口。日调度意味着每一个 DagRun 承诺处理长度为一天的区间。旧文档里的 execution_date 常被讲成窗口左端;后续版本用数据区间起止把这件事说得更白:运行实例绑定的是一段左闭右开的数据,而不是“函数被调用的那一秒”。

原文给过一个容易踩空的对齐例子:start_date 设在 6 月 1 日 2 点、间隔一天,第一个窗口从 2 点对齐,而不是自动回到当天零点。原点是偏移,不是午夜魔法。时区再叠加一层:naive 时间按 Airflow 默认时区解释,跨夏令时的调度要用带时区的原点,否则会在换日那天多一次或少一次。

with DAG( dag_id="orders_daily", start_date=datetime(2024, 6, 1, 2, 0), schedule="0 2 * * *", catchup=False, max_active_runs=1, default_args={ "owner": "order_data", "retries": 2, "retry_delay": timedelta(minutes=5), "email_on_failure": False, }, tags=["etl", "prod"], doc_md="每日订单事实表,窗口按活动日切分,失败不自动补历史。", ) as dag: ...

tags 给界面过滤和告警路由用。doc_md 常被当成注释垃圾,原文却把它连到治理:说明数据质量规则、敏感字段、业务影响。OpenLineage 一类集成也会从描述里取业务术语。配置因此不只是调度旋钮,也是给下一个值班的人留下的合同附件。

图 窗口、补数与并发上限如何咬合

图 窗口、补数与并发上限如何咬合

二、默认契约与容易被全局打开的开关

default_args 会落到未单独覆盖的任务上。ownerretriesretry_delay 是最低配置。email_on_failure 在没接好邮件后端时只会制造误导。execution_timeout 给单次执行套上限,防止传感器或外部 API 把 Worker 焊死。sla 是记过不是杀进程:超时后继续跑,但记 SLA miss,给监控用。

depends_on_past=True 要求上一窗口同一任务成功。移动平均、余额滚存需要它。代价是昨天失败则今天永远排队,直到有人清状态或补跑。原文倾向:能不用全局 depends_on_past 就不用,改成对特定上游 DAG 或任务的 ExternalTaskSensor,故障隔离更好。wait_for_downstream 更狠,会等下游也完成,几乎只留给强一致账本,日常 ETL 不要开。

开关 打开的正当理由 默认建议
catchup 新 DAG 必须补一段已知历史 生产 False,补数走 backfill
max_active_runs=1 窗口之间写同一物理分区 写冲突场景保持 1
depends_on_past 跨窗口滚动状态 默认 False
wait_for_downstream 下游未完成不准出下一窗口 默认 False
is_paused_upon_creation 新 DAG 先挂起等人审 生产建议 True

动态生成是配置的另一面。租户很多时,用工厂函数读配置表,循环构造 DAG 对象,dag_id 带命名空间。原文提醒两点:解析阶段就会执行这段 Python,所以工厂里禁止连生产写库、禁止不可预期的网络;调试比静态文件难,要用列出 DAG、列出导入错误,并给工厂做 dry-run 打印拓扑。某出行公司把两百多个静态文件收成一个工厂,新增区域只加配置行——那是治理成熟以后的收益,不是第一周的作业。

三、和基础设施咬合的字段

pool 把任务挂到命名配额上。queue 在 Celery 里把任务打到特定 Worker 队列,例如 GPU 机器只消费训练队列。priority_weight 影响同池里谁先被领走。executor_config 在 Kubernetes 执行器下传入资源请求、镜像覆盖。这些字段看起来像边角,其实是把基础设施语义写进业务图。标签同样如此:criticalbackfill 应走不同告警通道。

调度表达式除了 cron,也可以是 timedelta,或更复杂的时间表对象。手动-only 的 DAG 把 schedule 设为 None,只接受触发。REST API 可以创建带配置的运行,把 Airflow 当成被其他系统调用的中枢——原文称为 DAG as API。这要求 dag_id 稳定、参数经 params 或 DAG Run conf 传入,而不是改代码发版。

⚠️ 常见坑:在 DAG 文件顶层读取 Variable 来决定图的结构。Variable 在解析期就会打元数据库,DAG 一多,解析循环会被变量查询拖死。结构用代码或轻量本地配置;Variable 留给运行时。
💡 关键直觉:你能在界面里改的多是运行与暂停,改不了已经写进代码的窗口语义。时钟错误必须改 DAG 再发布,不要指望点几次清状态。

配置评审我用三问:窗口是否与业务日一致;失败重试是否仍能在 SLA 内结束;补数会不会与在线窗口抢池。三问都写进 doc_md,比单独写“请看代码”有用。params 给手动触发时的可选项,例如只跑某个分区;默认值要安全,不能让空参数扫全表。

时序债务是原文提出的判断:任务耗时加上重试等待逼近调度间隔时,表面上天天绿,实际上已经没有缓冲。对策不是把 retries 加到 10,而是拆任务、下推计算、或把间隔拉长。配置改不了物理耗时,只能把承诺写清楚,让违约可见。

四、配置评审会上怎么问

把 DAG 参数当成合同附件,评审不要从 Operator 开始,从窗口开始。问业务:这一窗的数据从哪一秒到哪一秒,失败了是补这一窗还是跳过。答案写成 schedulestart_date 的对齐方式,以及 catchup 是否允许。再问:两个窗口能不能同时写同一张表。不能则 max_active_runs=1。再问:昨天失败,今天算不算。要算滚动则 depends_on_past,否则不要开。最后问 SLA:从窗口闭合到必须交付,中间塞得进几次 retry_delay。塞不进就减少重试或拆任务,而不是假装坚韧。

executor_configpoolqueuepriority_weight 放在第二轮,等第一轮时钟问完。过早讨论 Pod 资源,会让会议变成基础设施吐槽,窗口错误被带过。tagsdoc_md 当场写:关键、环境、层级、失败影响谁。原文把 description/doc_md 连到血缘与治理,评审里空着就等于合同缺附件。

动态工厂单独开会。必须展示:dag_id 如何保证唯一、解析期读什么配置、配置错误时是导入失败还是静默少图、数量上限是多少。没有上限的工厂不准进生产解析目录。dry-run 打印拓扑应能在 CI 跑。原文出行公司把 200+ 文件收成工厂,那是配置表已经可信之后的事;配置表本身要有评审,不能比 DAG 代码更随便。

手动-only DAG(schedule 为空)要写清谁有权触发、conf 有哪些键、默认值会不会扫全表。REST API 触发等于把触发权交给另一个系统,认证与限流在第 5 章,但配置上必须 max_active_runs 仍然有效,防止外部重试把同一窗打十次。

问题:schedule_intervalschedule 哪个才对?

以你锁定的文档为准。旧例子大量用 schedule_interval,新写法常用 schedule。两者都在表达窗口,不要在同一仓库混用两套而不注明。本教程示例用 schedule,读旧代码看到 schedule_interval 时按同一套窗口语义理解,不要以为是另一种时钟。

问题:时序债务怎么量化?

用“从窗口右端到全部成功”的时长,去除以调度间隔。靠近 1 或超过 1,就是债务。连续一周都在 0.8 以上,要拆任务或降频,不要加 retries。retries 只增加承诺的分子。原文把任务完成时间加上重试等待必须塞进间隔减去调度开销,写成约束而不是感受。把这个数画进监控,比把“感觉变慢”写进周报有用。

五、参数变更如何发版

改 start_date、schedule、catchup,都视为时钟变更,必须写清:已有 Run 怎么处理、要不要 backfill、max_active_runs 是否临时降低。只改 retries 或 timeout,视为契约变更,监控的 SLA 图要一起改。只改 doc_md 与 tags,视为文档变更,可随日常发。三类混在一次发布里,复盘时说不清哪次改动导致补数爆炸。动态工厂的配置表变更按时钟变更管理,因为多出来的 dag_id 会立刻进入解析。配置表没有发版记录,等于 DAG 没有 Git。把配置表纳入同一套评审,是工厂模式的隐藏成本,原文用“元数据成为活文档”描述收益,成本是活文档也要治理。

六、default_args 模板里每一项的否决权

owner 无法路由到团队频道,否决。retries 大于团队上限且无书面理由,否决。retry_delay 之和明显塞不进 SLA,否决。execution_timeout 缺失而任务会调外部 API,否决。depends_on_past 为真但没有滚动状态说明,否决。catchup 为真但没有补数范围,否决。email_on_failure 为真但邮件后端未验收,否决。这张否决表比“请按规范编写”有效。模板提供推荐值,否决表提供牙齿。动态工厂生成的 DAG 必须继承同一否决表,不能因为是机器写的就豁免。机器写的错,错得更快。start_date 写“今天”的工厂,会让每天的原点漂移,窗口无法重放。工厂的原点应来自配置表里的固定日期,而不是生成时刻的日期。生成时刻是墙钟,墙钟进不了合同。合同要经得起上周四被重新执行。评审时把这一句问出口。

七、时钟变更的发布说明模板

必须出现:旧窗口定义、新窗口定义、已存在 Run 是否保留、是否 backfill、backfill 范围、max_active_runs 临时值、回退后如何恢复。缺任一项,变更视为未说明。未说明的时钟变更不准合并。合并机器人看不了自然语言是否完整,但可以看发布说明文件是否存在。存在性检查加人工读模板。人工读的时候只核这七项,不评文采。文采会掩盖范围。范围写“适当补数”等于没写。必须写日期。日期是区间。区间才能被调度器执行。调度器不执行形容词。形容词是会议语言。会议语言进了发布说明,补数就会长成金融案例那种一百多个 Run。案例已经发生过,不必再演。不演的方法是模板。模板强迫写日期。日期强迫人思考。思考比 catchup 开关更安全。开关只是布尔。布尔太容易被复制成 True。True 加上遥远原点,就是事故配方。配方在 2.1 被拆开。拆开后,只允许显式 backfill。显式才有范围。有范围才有池。有池才有在线 SLA。SLA 比方便重要。方便是 catchup。重要是日历。日历在模板里。模板在合并前。合并前读七项。读完才能点合并。点合并之后,调度器会按你写的日期办事。办事没有情绪。情绪必须在模板阶段用完。用完就执行。执行就是合同。

重点提炼

  • DagRun 绑定数据区间,start_date 是原点偏移,不是开机键
  • 生产默认关闭 catchup,历史用范围明确的 backfill
  • max_active_runs 管峰值,catchup 管总量,两者一起才构成安全阀
  • depends_on_past 会级联阻塞,跨 DAG 等待更可隔离
  • 解析必须无副作用,动态工厂尤其不能在 import 时写生产
  • pool、queue、tags、doc_md 是基础设施与治理接口,不是装饰

下一节把节点从 Python 函数换成对外部系统的可治理集成。


作者与出处
原作者: 灏天文库
来源:灏天文库
整理: 灏天文库整理
由灏天文库平台收录,内容或由平台用户上传,仅供学习交流
发布者: 作者: 灏天文库 转发
评论区 (0)
U