本节摘要:依赖边声明谁必须先结束,触发规则声明结束成什么样才轮到自己,XCom 在元数据里传小纸条。三者一起才构成可执行因果:只有边没有规则,失败清理跑不起来;只有 XCom 没有边,调度器不会等上游写完再读。大结果进对象存储,纸条只传地址。
阅读完本节,你应当能够:
>> 表达结构依赖,并用 trigger_rule 表达失败后仍要跑的清理结构依赖是图纸上的箭头:pull >> validate >> load。调度器在推进某个 TaskInstance 前,先看上游实例在同一 DagRun 里的状态。默认 all_success:任一上游失败,下游进入 upstream_failed,不一定执行业务代码。这正是最小链想要的。
业务很快会超出直线。校验失败要告警,入库成功与否都要删临时目录,多路并行后要汇聚。这时不是再画一条边就完了,而要改触发规则。one_failed 让告警任务在有人失败时启动。all_done 让清理在上游结束(无论成败)后启动。none_failed 允许上游跳过但不能失败。none_failed_min_one_success 更严一点。规则选错,会出现“告警任务永远 skipped”或“清理在上游还在跑时就删了目录”。
alert = PythonOperator( task_id="alert_schema", python_callable=notify, trigger_rule="one_failed", ) cleanup = PythonOperator( task_id="cleanup_tmp", python_callable=remove_tmp, trigger_rule="all_done", ) validate >> [load, alert] [load, alert] >> cleanup
校验 ──► 入库 │ │ └─失败告警 │ ▼ 清理 规则 all_done
原文把同一 DAG 在不同日期走出不同路径,归因于规则与状态:某一天校验失败,告警被点亮;下一天全成功,告警跳过。图是活的。排障要看那一天的实例状态,不要只看静态箭头。
| 规则 | 何时启动自己 | 典型用途 |
|---|---|---|
| all_success | 上游全成功 | 默认 ETL 链 |
| all_done | 上游全结束 | 清理、最终通知 |
| one_failed | 至少一个失败 | 告警 |
| none_failed | 无失败,允许跳过 | 可选分支后的汇聚 |
| none_skipped | 无人跳过 | 强制所有分支都走 |
| always | 不管上游 | 极少用,易掩盖故障 |
任务可以通过 xcom_push 或 TaskFlow 的 return 把值写入元数据库,下游 xcom_pull。默认后端是数据库表。原文定位:轻量、去中心化的信使。适合:对象存储路径、影响行数、分区名、布尔开关。不适合:查询结果集、模型权重、图片。
体积问题有两层。一层是表膨胀,调度器和 Web 查询变慢。一层是序列化:不能 pickle 的对象、过大的对象会在推送时失败,或在另一台 Worker 上反序列化失败。Celery 跨机器时尤其明显。约定:大于几 KB 的东西写存储,XCom 只留 URI。
def write_partition(**ctx): path = f"s3://bucket/dt={ctx['ds']}/part.parquet" # 实际写入省略 return path def load_from_path(**ctx): ti = ctx["ti"] path = ti.xcom_pull(task_ids="write_partition") # 按 path 装载
跨任务读 XCom 要写明 task_ids,必要时写 map_index。动态映射后每个子任务都有自己的条子,pull 错 index 会拿到别人的路径。include_prior_dates 能读历史窗口,方便但危险:昨天的路径今天再用,可能装错分区。默认不要开。
密钥不要进 XCom。日志、UI、数据库备份都会留下。连接密码走 Connection;短期令牌如果必须传,考虑加密后端或根本不传、让下游自己去保险库取。
TaskGroup 是界面与结构上的文件夹,仍在同一个 DAG 里,调度器看见的还是平铺任务,只是 id 带组前缀。用它降低画布噪音,不要用它做隔离。SubDAG 曾经流行,另起一个 DAG 当任务跑。现代实践谨慎:调度嵌套、UI 卡顿、失败重试语义绕人。能用 TaskGroup 就不要上 SubDAG。
跨 DAG 依赖用 ExternalTaskSensor 等待另一个 dag_id 的某个 task_id 在对齐的窗口成功。窗口对齐是坑:两个 DAG 的 schedule 不同,execution 对不上,传感器会等到超时。写清楚双方的数据区间约定,或改用 Dataset:上游任务声明产出数据集,下游 DAG 按数据就绪触发,减少“对钟”。原文把 Dataset 出现在 2.4 之后,作为从任务触发升到数据就绪的一层。两者可以并存,但一条链上不要又传感器又数据集还 depends_on_past,否则没人说得清为什么没跑。
⚠️ 常见坑:汇聚任务仍用 all_success,导致可选分支 skipped 时汇聚永远不跑。可选分支后的汇聚用 none_failed。
💡 关键直觉:先问是控制流还是数据流。控制流用边和规则;数据流用存储加短 XCom。用 XCom 冒充控制流(下游靠拉到空值决定跳过)会让状态机撒谎。
触发规则与分支算子配合时,被切走的一侧是 skipped。下游若要求 all_success,会被 skipped 卡住。这是分支后第一件要改的事。BranchPythonOperator 返回下一个 task_id 列表,未选中的路径跳过,不是失败。
并发边可以扇出:一个校验后面挂多个互不依赖的装载。扇出加大 Worker 压力,用池限制重任务。扇入必须选对规则。画 ASCII 时把规则写在汇聚箭头上,评审比只看代码里的关键字快。
代码里的 trigger_rule 不容易在画布上被看见。评审时用 ASCII 把规则标在汇聚处,比只看关键字少漏。典型漏项:分支后仍 all_success;可选路径 skipped 导致汇聚不跑;清理任务用了默认规则,失败时临时目录永不删;告警任务挂在成功边上,永远 skipped。把这四条写成评审口头禅。
XCom 评审只问体积和秘密。返回值是路径还是整帧?有没有 token?跨机器能否 pickle?include_prior_dates 开了没有?映射后 pull 有没有写 map_index?任一问含糊,就假定会撑库或串窗。自定义 XCom 后端(对象存储)是规模上来之后的选项,小团队先纪律,再换后端。纪律比换后端便宜。
跨 DAG 传感器评审问对齐:两边 schedule 是否同长度、逻辑日期是否同一套业务日、对方失败时传感器是超时失败还是无限 poke。改用 Dataset 时问:谁更新、更新几次、下游超时谁告警。TaskGroup 只问前缀稳不稳,SubDAG 默认否决,除非能讲出 TaskGroup 做不到的隔离理由——大多数讲不出。
原文把依赖分成结构、规则、外部三层。评审按三层打勾,不要只看有没有 >>。有边无规则,失败清理不跑;有规则无边,调度器不认为你该等;有外部等待无对齐,等于随机超时。
all_done 会不会把失败当成成功传给更下游?all_done 只让清理或通知启动,清理自己成功不代表业务成功。更下游若还依赖清理,且用 all_success,会在“清理成功、业务失败”时继续跑,把坏窗口送进报表。正确形状是:业务链用 all_success,清理旁路 all_done,报表仍依赖业务链而不是依赖清理。旁路不要接回主干,除非你明确要“无论成败都出一封邮件然后结束”,那封邮件不是数据产物。
可以,若只有几个键且无密钥。表名、分区、行数、布尔开关都合适。一旦字典里开始塞抽样行、explain 计划、整段日志,就停。小字典也会在映射下复制成 N 份,N 很大时同样伤库。能改成只传分区名就只传分区名。
复盘先贴 ASCII 图,标出当天实际走的边与规则,再贴实例状态表。不要先贴 Python。状态表能看出是 skipped 卡住还是失败未通知。然后才改规则。改完补结构测试:给定上游 skipped,汇聚必须启动或必须失败,写成断言,防止下周又改回 all_success。XCom 事故复盘则先量体积:那一次 push 了多少字节、表增长多少、调度延迟是否同步变差。能对应上,才能立“禁止传帧”的规矩。对应不上,规矩会被当成风格偏好而反弹。跨 DAG 等待事故复盘必须写出两边窗口换算,只写“传感器超时”等于没复盘。
主干从左到右表示业务成功路径,全部默认 all_success。旁路向下表示告警与清理,规则写在箭头旁。禁止旁路再接回主干,除非主干明确接受“通知也算产物”。接回是汇聚事故的温床。TaskGroup 只包主干或只包旁路,不要一组里混两种规则,画布折叠后规则看不见。跨 DAG 等待画在主干左侧虚线,标注对方 dag_id 与窗口换算。Dataset 画成主干上方的到位灯,灯亮才发车。一张图上虚线与灯不要同时成为主触发。纪律看起来像美术,实际是排障速度:出事时眼睛先找主干红点,再看旁路有没有亮。旁路该亮不亮,是规则错误;主干绿旁路亮,可能是误报。没有画布纪律,规则再正确也会在折叠的组里失踪。失踪的规则等于没有规则。没有规则的边,只是装饰箭头。
CI 可以做不到称重每一次 return,但代码评审必须问“这是路径还是载荷”。载荷超几 KB 就打回。打回理由写体积,不写风格。风格会被争论。体积有数。有数的门禁能活。活着的门禁会逼人把数据放进表或对象存储。放对地方,元数据库才能继续当账本而不是当湖。当湖的账本,3.1 的锁等待会来。锁等待来了,6.2 会让你减 XCom。减的时候业务已依赖大纸条,减不动。减不动就只能加盘。加盘是认输。认输记录在性能债。债来自 2.3 没有门禁。门禁放在本节,是为了让债不要出生。不出生比还债便宜。便宜的门禁只是一句问话。问话要烦。烦到作者自动只 return 分区名。自动才可持续。可持续才不必每次评审都吵架。吵架会让门禁被绕过。绕过一次,范本就被污染。污染的范本会复制载荷。复制是利息。利息在实例行数上滚。映射会把利息变成本金。本金在 2.4。2.4 放大 2.3 的选择。选择在问话。问话在评审。评审在合并前。合并前问:路径还是载荷。答载荷且大,打回。打回是爱护调度器。调度器不会感谢你。它只会在你不打回时默默变慢。变慢以后,有人会加 Worker。加 Worker 是 6.2 反对的第一反应。第一反应的根,常常是一张过大的纸条。纸条在 2.3。门禁也在 2.3。把门禁用上。用上就少加机器。少加机器是性能。性能是纪律。纪律是问话。问话只要五个字:是不是路径。
评审画布时用手指沿主干走到底,再看旁路有没有箭头指回来。指回来就问产物是什么。答不出数据产物,必须拆开。拆开才能让 all_done 的清理不污染报表。污染是静默错账。错账全绿。全绿比全红难发现。难发现来自接回。接回在本节禁止。禁止要检查。检查用手指。手指比工具快。快才能在评审里做。做了,2.3 的门禁才完整:体积问路径,画布问接回。两问都过,才允许映射放大。放大在 2.4。2.4 会复制接回错误。复制前先拆。拆完结束 2.3 最后一问。问完结束本节。本节结束于两问。两问结束事故温床。温床结束于手指。手指结束评审。评审结束合并。合并结束。
下一节用 TaskFlow 和动态映射减少样板,并把“对列表里每一项跑一次”变成图上的一组工位。