本节摘要:TaskFlow 用装饰器把函数返回值变成 XCom,少写 Operator 样板。动态任务映射把一个列表在解析或运行时展开成一组同构子任务,每个子任务有独立状态和日志。两者让图更贴近“对每个分区跑一次校验”,但映射长度受解析成本与调度器压力约束,不能当成普通 Python for 循环的无限放大。
阅读完本节,你应当能够:
TaskFlow 把函数标成任务,return 自动推 XCom,下游函数参数按名字注入。读起来像普通 Python,跑起来仍是多个 TaskInstance。超时、重试、池,该写的还要写,只是换了装饰器参数位。
from airflow.decorators import dag, task from datetime import datetime @dag(start_date=datetime(2024, 6, 1), schedule="@daily", catchup=False) def dim_refresh(): @task def list_tables(): return ["user_dim", "item_dim"] @task def checksum(name: str) -> str: return f"{name}:ok" @task def summarize(rows: list) -> str: return ",".join(rows) summarize(checksum.expand(name=list_tables())) dim_refresh()
expand 是映射入口:list_tables 返回两张表名,checksum 展开成两个子任务。失败一个,另一个仍可成功——这是映射相对“一个 Python 任务里 for 两张表”的核心收益:隔离重试与日志。for 循环里第二张表失败,整任务重跑可能把第一张表再写一遍,除非你自己做检查点。
TaskFlow 不是新引擎。调度器看见的仍是任务图。过长的装饰器链会让新手忘记:函数体在 Worker 上执行,闭包里不要捕获笔记本路径、不要捕获未序列化的数据库句柄。参数必须能进 XCom。传入一个活连接对象,展开时就会在另一进程里摔碎。
映射有两种常见来源。一是上游任务返回列表。二是用常量列表 expand(name=["a","b"])。运行时映射的长度可以随窗口变化:今天 10 个分区,明天 200 个。这正是它比静态复制任务强的地方,也是风险:200 个子任务意味着 200 行实例、200 份日志、调度器 200 次状态推演。原文讨论动态 DAG 时强调命名与解析成本;映射把成本从“很多 DAG 文件”挪到“一次运行里很多实例”。
部分失败策略要事先说清。默认下游若依赖整个映射结果列表,可能等全部结束。你想“失败的分区告警、成功的照装”,就要把装载也映射,并在汇聚处用合适的触发规则。map_index 出现在 UI 和日志里,回放时必须带上它,否则你说不清是哪张表失败。
| 写法 | 隔离 | 适用 |
|---|---|---|
| 单任务内部 for | 无,一次失败整段重跑 | 极轻、无外部副作用 |
| 静态复制任务 | 有,但表数写死 | 表数长期稳定且很少 |
| 动态映射 | 每个元素一个实例 | 列表随窗口变化 |
| 动态 DAG 工厂 | 每个租户一张图 | 调度与权限要按租户拆 |
映射不是环。图在展开后仍必须无环。不要幻想在运行时根据结果再映射出依赖自己的下一步无限生长——那是普通程序,不是 DAG。需要 while 的业务,要么拆到任务内部循环并自带退出条件,要么换长工作流引擎。
与传统 Operator 混编时,TaskFlow 下游要 xcom_pull 的话注意默认 return 键。专用 Operator 的输出字段各不相同,不要假设都叫 return_value。混编是现实:SQL 仍用 PostgresOperator,Python 胶水用 TaskFlow。保持边界比追求全装饰器纯洁更重要。

给映射一个心理上限:开发环境几十,生产若经常上千,先问是否该拆成按域划分的多条 DAG,或把计算下推到能原生并行的引擎,只让 Airflow 提交一次作业。调度器对实例行数敏感,原文在性能章指出状态推演复杂度随节点上升,映射是节点数的放大器。
max_active_tis_per_dag 一类限制可防止一个映射打满集群。池仍然有用:每个子任务都去打同一数据库时,隔离了失败不等于隔离了连接数。
⚠️ 常见坑:在 list 任务里返回巨大字典列表,每个元素带整行业务数据。映射参数本身会进 XCom。只返回 id 或分区名,数据留在表里。
💡 关键直觉:映射换来的是实例级重试与可观测,付出的是账本行数。划不来就回到单任务循环,但必须自己做幂等检查点。
测试映射要覆盖空列表:空展开后下游汇聚是否 skipped。还要覆盖单元素,避免只在多元素时才暴露的 index 问题。CI 里可用 DAG 解析测试断言任务数在展开前的模板结构,展开后的数量放到集成测试,用短列表。
TaskFlow 与 Dataset 组合时,注意产出数据集的任务是映射还是单任务。若每个分区都声明产出,下游可能被触发多次。数据就绪语义要设计成“整窗完成才更新一个数据集”,还是“每分区都通知”,两者差很远。默认我选整窗一个数据集,分区细节留在表里。
用生产会遇到的最大列表在预发跑一次,看 task_instance 行数、调度循环耗时、日志存储、下游系统连接数。四个数字任一不可接受,就拆 DAG 或下推计算。不要用“映射很现代”当理由。实例级隔离的收益,必须大于账本行数的成本。
与 for 循环的选择可以写成规则:有外部副作用且需要单元素重试,映射;纯内存计算且 naturally 幂等,单任务循环;表数长期三张且几乎不变,静态任务更易读。不要为了展示 TaskFlow 把三张表写成映射。可读性也是性能:排障时能说出 index 对应哪张表,比装饰器链漂亮更重要。
混编时列一张表:哪些任务是装饰器,哪些是 PostgresOperator,返回值键叫什么。下游 pull 写错键是映射故障里最窝囊的一种,因为状态全绿、数据全空。TaskFlow 默认 return_value,专用算子各有各的。测试里断言 XCom 键存在。
空列表:expand 得到零个子任务,汇聚往往 skipped。业务上“今天没有分区”是合法还是事故,要事先定义。合法就让汇聚 skipped 并让更下游用 none_failed;事故就应在 list 任务里失败,而不是静默展开成空。这和 Dataset 静默不跑是同一类治理问题:没发生不等于成功。
可以做 zip 类展开,但复杂度上升很快,UI 上 index 对不齐会让人崩溃。能在一个任务里把配对算完再映射“每一对的 id”,通常更清晰。图上表达的应是隔离单位,不是数学上的笛卡尔积。笛卡尔积很容易把 20×20 变成 400 个实例,而业务只是“每个门店每个品类”,其实可以下推成一次 SQL。
没有。SQL、Bash、云厂商专用动作仍然是专用算子更清晰,凭证走 conn_id,审计更好搜。TaskFlow 擅长 Python 胶水与映射。纯度不是目标。混编是成年团队的常态。追求全装饰器往往会把 SQL 塞进 Python 字符串,既失去专用算子的模板约定,又失去 SQL 评审工具。
映射解决“一次运行里对列表的每一项隔离重试”。工厂解决“每个租户一张图、权限与调度分开”。不要用映射模拟租户:租户要独立暂停、独立 SLA、独立 owner,那是多 DAG。不要用工厂模拟分区:分区同构、同窗、同责任人,那是映射。分工画错,会出现一张巨图无法单租户暂停,或一个租户几千个 DAG 把解析打死。上限策略也不同:映射限 max_active_tis 与池;工厂限配置表行数与 dag_id 前缀配额。TaskFlow 只是这两种结构的写法,不是第三种结构。写法服务于结构,结构服务于谁能独立失败。
超过六七个装饰器任务且没有组,评审会放弃读代码,开始凭测试。测试很重要,但放弃读代码意味着偏离默认的规则没人看见。用 TaskGroup 按阶段切开:列出、校验、装载、汇聚。每个组内任务数保持能在一屏看完。映射任务在组内只出现一次蓝图节点,文档写清 index 含义。闭包捕获的配置必须可序列化,在注释里写“此函数禁止捕获连接对象”。可读性上限是为了保护 2.3 的规则可见性,不是为了审美。TaskFlow 让代码像脚本,脚本过长时,图的优势消失,只剩脚本的缺点。到了那一步,宁可拆 DAG。拆 DAG 比拆装饰器更有利于权限与 SLA。映射再长,也要能在汇聚处说出成功个数与失败 index 列表。说不出,可观测性已经输给了 for 循环加日志。输了就改回去,不要为了新写法保面子。
测试必须覆盖长度为 0 与长度为团队上限。0 验证汇聚 skipped 是否符合业务定义。上限验证池与调度是否仍能呼吸。中间随便一个长度只能证明高兴。高兴不是契约。契约是两端。两端夹具写进 6.3 之前,先在本节成为习惯。习惯是:每次改 list 函数,先想空与满。空是合法还是事故。满是多少。多少写进 doc_md。doc_md 不写上限,评审当没设计。没设计的映射不准上。上了会在活动日变成上千实例。上千实例的日志会把采集打满。采集打满,4.3 的“没日志”会响。响的是映射,不是采集。锅要分对。分对靠上限。上限靠夹具。夹具靠习惯。习惯靠本节把两端写进学习目标的肌肉。肌肉记得:映射换隔离,付行数。行数有顶。顶在文档。文档在评审。评审在合并。合并后活动日才不会意外。意外是没顶。没顶是没设计。没设计不是 TaskFlow 的错。错在把展开当成无限 for。无限 for 不是 DAG。DAG 要能画完。画不完的展开,不是图,是循环伪装。伪装会在账本里现形。现形为行数。行数会说话。让夹具先说话。先说话就不必让生产说话。生产说话很贵。贵在窗口已经闭合。闭合的窗口要补。补要日历。日历在 2.1。一切又绕回合同。合同包含映射上限。上限在本节写清。写清才允许 expand。不允许口头“应该不多”。应该不是数。数才是上限。上限进文档。文档进夹具。夹具进 CI。CI 进 6.3。本节先把两端变成习惯。习惯是 CI 的前置。前置不做,CI 会漏测空列表。漏测会在节假日无分区时静默。静默是 5.1 也讨厌的东西。讨厌要前后一致。一致从空列表开始。开始在 2.4。
下一章打开引擎盖:谁解析、谁排队、谁真正执行。