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

如何基于Airflow任务返回值动态创建新任务?

Airflow动态生成任务的问题解答

不能直接用你提供的写法实现需求,原因很简单:Airflow的DAG结构是在解析阶段(也就是DAG文件被Airflow调度器加载时)确定的,而task_1的返回值(XCom数据)是在DAG运行阶段才会产生的,解析阶段根本拿不到这个值,所以你写的for task in task_1.output循环在解析时不会生成任何任务。

下面提供两种可行的解决方案:

方案一:静态预生成任务(解析阶段确定任务列表)

如果你的任务列表可以在DAG解析时就确定(比如func_test不需要依赖运行时数据),可以直接在DAG定义前执行函数拿到任务列表,再循环创建任务:

from datetime import datetime
from airflow.operators.python import PythonOperator

def func_test():
    return ['task_2', 'task_3']

def another_function(**context):
    # 自定义业务逻辑
    pass

# 解析阶段直接执行函数获取任务列表
task_list = func_test()

with DAG(
    'dag_name',
    schedule_interval="@once",
    start_date=datetime(2022, 4, 19),
    catchup=False,
    default_args= {
        'depends_on_past': False,
        'retries': 0
    }
) as dag:

    task_1 = PythonOperator(
        task_id='func_test',
        python_callable=func_test,
        provide_context=True
    )

    # 循环创建任务并设置依赖
    for task_id in task_list:
        new_task = PythonOperator(
            task_id=task_id,
            python_callable=another_function,
            provide_context=True
        )
        task_1 >> new_task

方案二:使用动态任务映射(运行时生成任务,Airflow 2.2+支持)

如果任务列表必须依赖运行时数据(比如func_test的返回值是从数据库或其他外部服务获取的),推荐使用Airflow官方的Dynamic Task Mapping功能,它会在运行时根据上游任务的返回值自动生成对应数量的任务实例:

from datetime import datetime
from airflow.operators.python import PythonOperator

def func_test():
    return ['task_2', 'task_3']

def another_function(task_name):
    # 接收映射参数,处理对应任务逻辑
    print(f"Executing task: {task_name}")

with DAG(
    'dag_name',
    schedule_interval="@once",
    start_date=datetime(2022, 4, 19),
    catchup=False,
    default_args= {
        'depends_on_past': False,
        'retries': 0
    }
) as dag:

    task_1 = PythonOperator(
        task_id='func_test',
        python_callable=func_test,
        provide_context=True
    )

    # 基于task_1的返回值动态生成任务
    mapped_tasks = PythonOperator.partial(
        task_id='mapped_task',
        python_callable=another_function,
    ).expand(
        op_args=task_1.output
    )

    # 设置上游依赖
    task_1 >> mapped_tasks

这种方式下,Airflow会自动为task_1返回的每个元素生成一个任务实例,任务ID会自动添加后缀(如mapped_task__0、mapped_task__1),无需手动循环创建。

内容的提问来源于stack exchange,提问作者Renan Nogueira

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 04:35:29