You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何基于配置及手动触发参数创建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触发

  1. 进入目标DAG的详情页面
  2. 点击右上角的 Trigger DAG w/ config 按钮
  3. 在弹出的输入框中填写JSON格式的参数,例如:
{"task_list": ["data_cleaning", "report_generation", "send_notification"]}
  1. 点击 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.17 04:15:39