如何基于配置及手动触发参数创建Airflow动态任务?
手动触发Airflow DAG并基于传入参数创建动态任务
1. 核心思路
要实现基于触发时传入的参数生成动态任务,Airflow 2.2+ 提供的Dynamic Task Mapping是最优方案,它支持在运行时根据触发配置动态生成任务实例,无需提前硬编码任务数量或参数。
2. 编写支持动态任务的DAG
以下是两种常见写法,选择适合你的风格:
写法一:使用TaskFlow API(推荐,代码更简洁)
from airflow.decorators import dag, task from datetime import datetime @dag( dag_id="dynamic_task_demo", schedule=None, # 禁用自动调度,仅手动触发 start_date=datetime(2024, 1, 1), catchup=False, tags=["dynamic", "manual_trigger"] ) def dynamic_task_workflow(): # 第一步:获取触发时传入的配置参数 @task def fetch_trigger_config(**context): # 从dag_run中提取触发配置,默认值避免空参数报错 trigger_conf = context["dag_run"].conf or {} # 校验参数格式,确保task_list是列表类型 task_list = trigger_conf.get("task_list") if not isinstance(task_list, list): task_list = ["default_task_1", "default_task_2"] return task_list # 第二步:定义动态执行的任务逻辑 @task def execute_dynamic_task(task_name): # 这里替换为你的实际业务逻辑 print(f"正在执行动态任务:{task_name}") # 示例:可以根据task_name执行不同分支逻辑 if task_name.startswith("data"): print("执行数据处理逻辑") elif task_name.startswith("report"): print("执行报表生成逻辑") # 第三步:串联任务,动态映射生成任务实例 task_list = fetch_trigger_config() dynamic_tasks = execute_dynamic_task.expand(task_name=task_list) # 实例化DAG dynamic_task_workflow()
写法二:使用传统Operator(适配习惯旧写法的场景)
from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime def fetch_trigger_config(**context): trigger_conf = context["dag_run"].conf or {} task_list = trigger_conf.get("task_list", ["default_task_1", "default_task_2"]) return task_list def execute_dynamic_task(task_name, **context): print(f"正在执行动态任务:{task_name}") with DAG( dag_id="dynamic_task_demo_traditional", schedule=None, start_date=datetime(2024, 1, 1), catchup=False, tags=["dynamic", "manual_trigger"] ) as dag: fetch_config_task = PythonOperator( task_id="fetch_trigger_config", python_callable=fetch_trigger_config, provide_context=True ) # 利用Dynamic Task Mapping生成动态任务 dynamic_task = PythonOperator.partial( task_id="execute_dynamic_task", python_callable=execute_dynamic_task, provide_context=True ).expand( op_kwargs=[{"task_name": task} for task in fetch_config_task.output] ) fetch_config_task >> dynamic_task
3. 手动触发DAG并传入参数
方式一:Airflow UI触发
- 进入目标DAG的详情页面
- 点击右上角的 Trigger DAG w/ config 按钮
- 在弹出的输入框中填写JSON格式的参数,例如:
{"task_list": ["data_cleaning", "report_generation", "send_notification"]}
- 点击 Trigger 即可触发DAG,系统会根据传入的
task_list生成对应数量的动态任务。
方式二:CLI命令触发
使用airflow dags trigger命令,通过-c参数传入配置:
airflow dags trigger -c '{"task_list": ["data_cleaning", "report_generation", "send_notification"]}' dynamic_task_demo
4. 进阶:复杂动态任务配置
如果需要给每个动态任务传入多组参数,可以使用expand_kwargs:
@task def fetch_complex_config(**context): trigger_conf = context["dag_run"].conf or {} # 传入包含多参数的字典列表 return trigger_conf.get("task_details", [ {"task_id": "t1", "param1": "value1", "param2": 100}, {"task_id": "t2", "param1": "value2", "param2": 200} ]) @task def execute_complex_task(task_id, param1, param2): print(f"执行任务{task_id},参数1:{param1},参数2:{param2}") task_details = fetch_complex_config() execute_complex_task.expand_kwargs(task_details)
触发时传入的配置示例:
{"task_details": [ {"task_id": "data_process", "param1": "raw_data", "param2": 500}, {"task_id": "model_train", "param1": "training_set", "param2": 1000} ]}
内容的提问来源于stack exchange,提问作者romanzdk
相关产品推荐
相关产品推荐

