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

Python Airflow 声明DAG时如何获取任务的执行返回结果

实现方法

Airflow的DAG结构默认在文件解析阶段就会被Scheduler固定下来,Task1的返回值是任务实际运行时才会产生的数据,解析阶段无法直接获取,因此直接在DAG定义层写循环读取Task1运行结果的写法无法生效,根据业务场景可以选择以下两种方案实现:

方案1:结果可在解析阶段预知时,提前生成固定任务

如果Task1返回的列表不依赖任务运行时上下文(比如不需要读取运行时才返回的接口数据、库表查询结果),可以把生成结果的逻辑抽为公共函数,在DAG解析阶段直接调用拿到结果,循环生成固定结构的下游任务即可。
示例代码:

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

# 抽离Task1的结果生成逻辑,解析阶段就可执行
def get_task1_result():
    return ["a", "b", "c"]

def Task1():
    return get_task1_result()

def Task2(value):
    print(f"处理值: {value}")
    return

with DAG(
    dag_id="static_dynamic_demo",
    start_date=datetime(2024, 1, 1),
    schedule_interval=None,
    catchup=False
) as dag:
    task1 = PythonOperator(
        task_id="task1_id",
        python_callable=Task1,
    )

    # 解析阶段直接拿到结果,循环生成下游任务
    task1_result = get_task1_result()
    for value in task1_result:
        t = PythonOperator(
            task_id=f"task2_id_{value}",
            python_callable=Task2,
            op_kwargs={"value": value}
        )
        task1 >> t
  • 优点:所有任务结构在DAG解析完成后就完全固定,和普通DAG行为一致,兼容所有Airflow版本
  • 缺点:无法适配运行时结果动态变化的场景,如果解析阶段调用的逻辑依赖外部服务,还可能因为服务不可用导致DAG解析失败。不要尝试在解析阶段拉取历史任务的XCom结果生成任务,这种写法会大幅拖慢Scheduler解析性能,还会因历史数据变动导致DAG结构频繁变化引发调度异常。

方案2:运行时动态结果使用官方动态任务映射(推荐,Airflow 2.3+支持)

如果Task1的返回值是运行时才能确定的动态数据,直接使用Airflow 2.3版本后推出的动态任务映射能力即可,不需要在解析阶段固定下游任务数量,Task1运行完成后,Airflow会自动根据返回值生成对应数量的下游Task实例。
示例代码:

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

def Task1():
    # 运行时执行逻辑返回动态结果,比如查库、调接口拿到的列表
    return ["a", "b", "c"]

def Task2(value):
    print(f"处理值: {value}")
    return

with DAG(
    dag_id="official_dynamic_mapping_demo",
    start_date=datetime(2024, 1, 1),
    schedule_interval=None,
    catchup=False
) as dag:
    task1 = PythonOperator(
        task_id="task1_id",
        python_callable=Task1,
    )

    # partial定义任务公共参数,expand指定要遍历映射的参数
    task2 = PythonOperator.partial(
        task_id="task2_id",
        python_callable=Task2
    ).expand(value=task1.output)

运行时Task1返回["a","b","c"]后,Airflow会自动生成3个Task2实例,每个实例分别接收入参value="a"、value="b"、value="c",和预期的遍历生成下游任务效果一致。如果需要自定义任务后缀而非默认索引,可以通过map_index_template参数配置自定义的任务名渲染规则。

内容的提问来源于stack exchange,提问作者Phạm Hoài Lâm

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 23:09:18