首页
/ Apache Airflow 如何生成任务结构可变的动态 DAG?

Apache Airflow 如何生成任务结构可变的动态 DAG?

2026-09-08 19:38:43作者:明树来

当你需要让 DAG 的任务结构不是手写死的,而是由代码在解析时生成——例如按环境(PROD/DEV)构建不同的任务、按一份外部配置列表批量生成多个 DAG——就可以用 Apache Airflow 的动态 DAG 生成能力。airflow-core/docs/howto/dynamic-dag-generation.rst 描述的做法是:在 DAG 文件的顶层代码中读取结构数据,循环构建 Operator/Task 和依赖关系,Airflow 会在解析阶段自动注册这些 DAG。完成后的效果是:每个配置对应一个独立 DAG,可以各自单独运行。

先明确边界:这套做法生成的 DAG 在每一次 Dag Run 中任务数量不变,变化的是解析时生成的结构本身。如果你需要任务数(或 Task Group 数)随上游任务的运行结果在每次 Run 中动态增减,那是 dynamic-task-mapping 的场景,与本文方法不同。

两条必须先保证的原则

1. 任务生成顺序必须稳定。 文档特别提示:动态生成的 Task 和 Task Group 每次生成的顺序应保持一致,否则每次刷新页面时 Grid View 中任务的排列会跳动。文档给出的实现方式是:数据库查询使用稳定的排序机制,或 Python 中用 sorted() 对待生成列表排序。

2. 顶层代码不要做重活。 调度器会以 dag_processor__min_file_process_interval 为最小间隔反复执行 DAG 文件顶层代码来维持动态调度能力,所以顶层代码中的数据库访问、重计算、网络调用都会拖慢解析并增加数据库负载(见 best-practices.rst 的 "Top level Python Code" 一节)。因此结构数据要么来自环境变量,要么要提前导出到 DAG 目录下的文件中,而不是在顶层代码里去"拉取"远端数据。

让结构数据进入 DAG 文件:三种方式

方式一:环境变量——适合按环境切换结构

文档建议:顶层代码中配置代码行为时,用环境变量而不是 Airflow Variables,因为 Airflow Variables 会连元数据库取数,拖慢解析。例如用 DEPLOYMENT 区分生产与开发环境:

deployment = os.environ.get("DEPLOYMENT", "PROD")
if deployment == "PROD":
    task = Operator(param="prod-param")
elif deployment == "DEV":
    task = Operator(param="dev-param")

DEPLOYMENT 在生产环境设为 PROD、开发环境设为 DEV,同一个 DAG 文件在两套环境下构建出不同结构的任务。

方式二:生成带元数据的可导入 Python 模块

适合"元数据本身就是任务列表"的场景。外部流程动态生成 DAG 目录下的 my_company_utils/common.py

# This file is generated automatically !
ALL_TASKS = ["task1", "task2", "task3"]

DAG 直接导入该常量来构建任务:

from my_company_utils.common import ALL_TASKS

with DAG(
    dag_id="my_dag",
    schedule=None,
    start_date=datetime(2021, 1, 1),
    catchup=False,
):
    for task in ALL_TASKS:
        # create your operators and relations here
        ...

文档给出两个必须处理的配套步骤,否则该目录本身会被调度器当 DAG 扫描:

  • my_company_utils 目录中放置空的 __init__.py 文件;
  • .airflowignore 文件中加一行 my_company_utils/*(使用默认 glob 语法),让整个文件夹被调度器在查找 DAG 时忽略。

方式三:DAG 目录内的结构化配置文件

元数据比一个列表复杂时,文档建议把数据导出为 DAG 目录内的文件(JSON、YAML 都是候选格式),与 DAG 同包/同目录发布,再用模块的 __file__ 属性定位文件:

my_dir = os.path.dirname(os.path.abspath(__file__))
configuration_file_path = os.path.join(my_dir, "config.yaml")
with open(configuration_file_path) as yaml_file:
    configuration = yaml.safe_load(yaml_file)
# Configuration dict is available here

注意这里的前提是配置文件已经被"推送"到 DAG 目录,而不是顶层代码去远端拉取——原因同前文的顶层代码约束。

在循环中生成并注册多个 DAG

拿到结构数据后,在 @dag 装饰器或 with DAG(..) 上下文管理器中循环生成 DAG 即可,Airflow 会自动注册它们。文档给出的完整示例(from airflow.sdk import dag, task 为文档当前版本使用的导入方式):

from datetime import datetime
from airflow.sdk import dag, task

configs = {
    "config1": {"message": "first Dag will receive this message"},
    "config2": {"message": "second Dag will receive this message"},
}

for config_name, config in configs.items():
    dag_id = f"dynamic_generated_dag_{config_name}"

    @dag(dag_id=dag_id, start_date=datetime(2022, 2, 1))
    def dynamic_generated_dag():
        @task
        def print_message(message):
            print(message)

        print_message(config["message"])

    dynamic_generated_dag()

验证注册结果

文档对这段代码给出的预期结果是:生成 dynamic_generated_dag_config1dynamic_generated_dag_config2 两个 DAG,每个都带着各自的配置独立运行。你可以用同样方式核对:解析完成后,每个 config_name 应能在 UI 中看到对应的 dynamic_generated_dag_<config_name>,并单独触发运行。

两点注册相关的说明:

  • 自 Airflow 2.4 起,通过调用 @dag 装饰函数(或 with DAG(...) 上下文管理器)创建的 DAG 自动注册,不再需要存到全局变量里;
  • 如果不想自动注册,给 DAG 设置 auto_register=False 关闭该行为。

优化:跳过非目标任务执行时的重复解析

从 Airflow 2.4 开始,文档针对"单个 DAG 文件生成大量动态 DAG"的场景给出了一个解析延迟优化:任务执行前 Airflow 会解析其来源 DAG 文件,若文件里循环生成很多 DAG,任务启动前会白白等待这些无关 DAG 的构建。解决方式是用 get_parsing_context() 判断当前解析上下文,只生成当前需要的那一个 DAG。

get_parsing_context() 返回 AirflowParsingContext:执行单个 DAG/任务时其中 dag_idtask_id 字段有值;DAG File Processor 做全量解析时两者为 None(当前版本该函数可从 airflow.sdk 导入,实现见 context.py)。文档给出的用法示例:

from airflow.sdk import DAG
from airflow.sdk import get_parsing_context

current_dag_id = get_parsing_context().dag_id

for thing in list_of_things:
    dag_id = f"generated_dag_{thing}"
    if current_dag_id is not None and current_dag_id != dag_id:
        continue  # skip generation of non-selected Dag

    with DAG(dag_id=dag_id, ...):
        ...

效果文档引用了 Airflow 官方博客 "Magic Loop" 一文的案例:任务执行时的解析时间从 120 秒降到 200 毫秒(文档标注该数字来自 2.4 之前的旧实现案例,仅作参考)。同时文档明确警告:当后续 DAG 的生成依赖前面的 DAG、或生成过程有副作用时不能使用该优化,必须谨慎使用并充分测试。

限制与边界汇总

  • 本方法保证的是"结构在解析时由代码生成",单次 Dag Run 内的任务数不变;任务数需随上游结果变化的场景请改用 dynamic-task-mapping
  • 动态生成时务必保持任务生成顺序稳定(稳定排序或 sorted()),否则 Grid View 中任务每次刷新都会乱序。
  • 顶层代码不做数据库访问、重计算、网络请求;结构数据通过环境变量、生成模块或 DAG 目录内配置文件提供。
  • 生成可导入模块时,记得 __init__.py.airflowignore 两项配套配置,否则该目录会被当作 DAG 目录扫描。
  • get_parsing_context 优化仅在 Airflow 2.4+ 可用,且不适用于 DAG 生成存在依赖链或副作用的场景。
登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.14 K
2.75 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
857
1.35 K
docsdocs
暂无描述
Markdown
898
5.82 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
531
596
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
921
1.84 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.8 K
1.02 K
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.36 K
1.46 K
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
1.02 K
519
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
548
391