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

Airflow构建DAG时>>运算符、set_downstream及chain方法报错如何解决

Airflow DAG结构搭建问题排查

预期要实现的DAG结构如下:
预期DAG结构


报错复现过程

第一次编写:位运算符定义依赖报错

代码:

with DAG(
    dag_id="dag_example",
    schedule_interval="@once",
        default_args={
        "owner": "airflow",
        "start_date": datetime(2020, 11, 1),
        "retries": 1,
        "retry_delay": timedelta(seconds=10),
    },
    catchup=False,
) as dag:
    tA = PythonOperator(
        task_id="A",
        python_callable=A,
        provide_context=True,
    ),
    tB = PythonOperator(
        task_id="B",
        python_callable=B,
        provide_context=True,
    ),
    taS = PythonOperator(
        task_id="aS",
        python_callable=aS,
        provide_context=True,
    ),
    taF = PythonOperator(
        task_id="aF",
        python_callable=aF,
        provide_context=True,
    ),
    tbS = PythonOperator(
        task_id="bS",
        python_callable=A,
        provide_context=True,
    ),
    tbF = PythonOperator(
        task_id="bF",
        python_callable=A,
        provide_context=True,
    )
    tA >> tB 
    tA >> taS 
    tA >> taF 
    tB >> tbF 
    tB >> tbS

运行报错:TypeError: unsupported operand type(s) for >>: 'tuple' and 'tuple'

改用set_downstream方法调整依赖定义:

tA.set_downstream(tB)
tA.set_downstream(taS)
tA.set_downstream(taF)
tB.set_downstream(tbS)
tB.set_downstream(tbF)

新报错:AttributeError: 'tuple' object has no attribute 'set_downstream'


第二次编写:chain方法嵌套列表报错

代码:

with DAG(
    dag_id="dag_example",
    schedule_interval="@once",
        default_args={
        "owner": "airflow",
        "start_date": datetime(2020, 11, 1),
        "retries": 1,
        "retry_delay": timedelta(seconds=10),
        # "on_failure_callback": snow_fail,
        # "on_success_callback": snow_succ,
    },
    catchup=False,
) as dag:
    chain(
        PythonOperator(
            task_id="A",
            python_callable=A,
            provide_context=True,
        ),
        [
            PythonOperator(
                task_id="aF",
                python_callable=aF,
                provide_context=True,
                trigger_rule=TriggerRule.ONE_FAILED,
            ),
            PythonOperator(
                task_id="B",
                python_callable=B,
                provide_context=True,
            ),
            [
                PythonOperator(
                    task_id="bF",
                    python_callable=bF,
                    provide_context=True,
                    trigger_rule=TriggerRule.ONE_FAILED,
            ),
                PythonOperator(
                    task_id="bS",
                    python_callable=bS,
                    provide_context=True,
                    trigger_rule=TriggerRule.ONE_SUCCESS,
            ),        
            ],
            PythonOperator(
                task_id="aS",
                python_callable=aS,
                provide_context=True,
                trigger_rule=TriggerRule.ONE_SUCCESS,
            ),
        ],
    )

报错:AttributeError: 'list' object has no attribute 'update_relative'


第三次编写:调整chain参数后结构不符预期

调整后代码:

with DAG(
    dag_id="dag_example",
    schedule_interval="@once",
        default_args={
        "owner": "airflow",
        "start_date": datetime(2020, 11, 1),
        "retries": 1,
        "retry_delay": timedelta(seconds=10),
        # "on_failure_callback": snow_fail,
        # "on_success_callback": snow_succ,
    },
    catchup=False,
) as dag:
    chain(
        PythonOperator(
            task_id="A",
            python_callable=A,
            provide_context=True,
        ),
        [
            PythonOperator(
                task_id="aF",
                python_callable=aF,
                provide_context=True,
                trigger_rule=TriggerRule.ONE_FAILED,
            ),
            PythonOperator(
                task_id="aS",
                python_callable=aS,
                provide_context=True,
                trigger_rule=TriggerRule.ONE_SUCCESS,
            ),
        ],
        PythonOperator(
            task_id="B",
            python_callable=B,
            provide_context=True,
        ),
        [
            PythonOperator(
                task_id="bF",
                python_callable=bF,
                provide_context=True,
                trigger_rule=TriggerRule.ONE_FAILED,
            ),
            PythonOperator(
                task_id="bS",
                python_callable=bS,
                provide_context=True,
                trigger_rule=TriggerRule.ONE_SUCCESS,
            ),        
        ],
    )

生成的DAG结构和预期不符,实际结构如下:
实际生成DAG结构


最终解决方案

核心错误原因

第一次编写时,每个PythonOperator变量定义末尾多写了逗号,导致tA、tB等变量被识别为单元素元组,而非Operator对象,因此无法调用set_downstream方法,也无法使用>>运算符定义依赖。
另外chain方法不支持嵌套列表传参,仅支持同层级的任务列表,因此嵌套写法会报错。

修正后的完整代码

from datetime import datetime, timedelta
from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.utils.trigger_rule import TriggerRule

# 提前定义对应业务逻辑函数
def A(**context):
    pass
def B(**context):
    pass
def aS(**context):
    pass
def aF(**context):
    pass
def bS(**context):
    pass
def bF(**context):
    pass

with DAG(
    dag_id="dag_example",
    schedule_interval="@once",
    default_args={
        "owner": "airflow",
        "start_date": datetime(2020, 11, 1),
        "retries": 1,
        "retry_delay": timedelta(seconds=10),
    },
    catchup=False,
) as dag:
    tA = PythonOperator(
        task_id="A",
        python_callable=A,
        provide_context=True,
    )
    tB = PythonOperator(
        task_id="B",
        python_callable=B,
        provide_context=True,
    )
    taS = PythonOperator(
        task_id="aS",
        python_callable=aS,
        provide_context=True,
        trigger_rule=TriggerRule.ONE_SUCCESS
    )
    taF = PythonOperator(
        task_id="aF",
        python_callable=aF,
        provide_context=True,
        trigger_rule=TriggerRule.ONE_FAILED
    )
    tbS = PythonOperator(
        task_id="bS",
        python_callable=bS,
        provide_context=True,
        trigger_rule=TriggerRule.ONE_SUCCESS
    )
    tbF = PythonOperator(
        task_id="bF",
        python_callable=bF,
        provide_context=True,
        trigger_rule=TriggerRule.ONE_FAILED
    )
    # 按预期定义依赖
    tA >> [tB, taS, taF]
    tB >> [tbS, tbF]

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.27 01:36:03