Apache Airflow 如何生成任务结构可变的动态 DAG?
当你需要让 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_config1 和 dynamic_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_id 和 task_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 生成存在依赖链或副作用的场景。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0629
MiniCPM5-2BMiniCPM5-2B 是一款面向端侧、本地部署和资源受限场景的 2B 稠密 Transformer,能够达到同尺寸开源模型 SOTA 水平。Markdown00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
HivisionIDPhotos⚡️HivisionIDPhotos: a lightweight and efficient AI ID photos tools. 一个轻量级的AI证件照制作算法。Python07
DragonOSDragonOS is an operating system developed from scratch using Rust, with Linux compatibility. It is designed for **Serverless** scenarios. 使用Rust从0自研内核,具有Linux兼容性的操作系统,面向云计算Serverless场景而设计。Rust00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00