Airflow 2如何从DAG Run配置动态生成任务?
基于DAG Run配置动态生成任务的解决方案
你遇到的核心问题是:DAG定义阶段无法直接访问dag_run.conf(因为此时还没有具体的DAG Run实例),所以不能直接用配置值循环生成Operator。以下是两种可行的替代方案,完全避免使用全局Variables,同时保留每个DAG Run的配置可追溯性:
方案1:使用Dynamic Task Mapping(Airflow 2.3+ 推荐)
这是Airflow官方支持的动态任务生成方式,在运行时根据DAG Run配置动态创建任务实例,无需提前定义所有任务。
代码示例
from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime def fetch_target_ids(**context): # 从当前DAG Run的配置中获取ids,无配置时使用默认值[1,2,3] return context["dag_run"].conf.get("ids", [1,2,3]) with DAG( dag_id="dynamic_task_demo", start_date=datetime(2024, 1, 1), schedule_interval=None, # 手动触发的DAG catchup=False ) as dag: # 第一步:获取要处理的ids列表 get_ids_task = PythonOperator( task_id="fetch_target_ids", python_callable=fetch_target_ids, provide_context=True ) # 第二步:基于ids列表动态生成处理任务 process_single_id = PythonOperator.partial( task_id="process_id", python_callable=lambda id: print(f"Processing ID: {id}"), provide_context=True ).expand( op_args=get_ids_task.output, task_id_template="process_id_{{ args[0] }}" # 自定义任务ID格式(Airflow 2.4+支持) ) # 设置任务依赖 get_ids_task >> process_single_id
说明
- 运行时触发DAG时,在配置中传入
{"ids": [4,5,6]},就会自动生成process_id_4、process_id_5、process_id_6三个任务 - 每个DAG Run的配置会被Airflow存储,可在UI的"DAG Runs"页面查看具体的参数,历史执行记录清晰
- 如果不需要自定义任务ID,可省略
task_id_template,任务会自动命名为process_id__0、process_id__1等
方案2:Airflow 2.3以下版本兼容方案
如果你的Airflow版本低于2.3,可以通过主DAG触发子DAG的方式实现:
- 主DAG中读取
dag_run.conf的ids,通过TriggerDagRunOperator将ids作为配置传给子DAG - 子DAG中使用
dag_run.conf的ids,结合TaskGroup生成动态任务(每个子DAG Run对应独立配置)
代码示例(主DAG)
from airflow import DAG from airflow.operators.trigger_dagrun import TriggerDagRunOperator from airflow.operators.python import PythonOperator from datetime import datetime def pass_ids_to_subdag(**context): ids = context["dag_run"].conf.get("ids", [1,2,3]) return {"ids": ids} with DAG( dag_id="main_dag", start_date=datetime(2024,1,1), schedule_interval=None, catchup=False ) as dag: get_ids = PythonOperator( task_id="get_target_ids", python_callable=pass_ids_to_subdag, provide_context=True ) trigger_subdag = TriggerDagRunOperator( task_id="trigger_subdag", trigger_dag_id="subdag_process_ids", conf="{{ ti.xcom_pull(task_ids='get_target_ids') }}", wait_for_completion=True ) get_ids >> trigger_subdag
代码示例(子DAG)
from airflow import DAG from airflow.operators.python import PythonOperator from airflow.utils.task_group import TaskGroup from datetime import datetime def process_id(id): print(f"Processing ID: {id}") with DAG( dag_id="subdag_process_ids", start_date=datetime(2024,1,1), schedule_interval=None, catchup=False ) as dag: # 从当前子DAG Run的配置中获取ids ids = dag.params.get("ids", [1,2,3]) with TaskGroup(group_id="process_ids_group") as process_group: for id in ids: PythonOperator( task_id=f"process_id_{id}", python_callable=process_id, op_args=[id] )
说明
- 主DAG负责读取配置并传递给子DAG,子DAG根据传入的配置生成对应任务
- 每个子DAG Run的配置独立,可在UI中查看历史执行的参数
内容的提问来源于stack exchange,提问作者Roman Guru
相关产品推荐
相关产品推荐

