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

使用装饰器定义Airflow DAG时依赖关系异常的解决咨询

问题解决:Airflow TaskFlow装饰器依赖异常修复

问题原因

使用@task装饰器时,每次调用被装饰的函数(如task_c())都会创建一个全新的Task实例。你当前代码里task_b() >> task_c()和task_d() >> task_c()中的task_c是两个完全独立的任务节点,并非同一个,所以生成的DAG会出现两个task_c,导致依赖关系不符合预期。而PythonOperator是先实例化单个对象再关联依赖,所以不会有这个问题。

修复方案

先将所有任务的实例赋值给变量,再通过变量定义依赖关系,确保所有分支指向同一个task_c实例,推荐写法如下:

from airflow.decorators import task, dag,task_group
from datetime import datetime , timedelta
from airflow.operators.dummy import DummyOperator

default_args = {
    "owner" : "khanh",
    "retries" : 1,
    "retry_delay" : timedelta(minutes = 2)
}

@dag(
    start_date= datetime(2025,1,1),
    schedule= "@daily",
    catchup= False,
    default_args= default_args
)
def complex_dags():
    # 先定义所有任务函数
    @task
    def task_a():
        print("a")
    @task
    def task_b():
        print("b")
    @task
    def task_c():
        print("c")
    @task
    def wait_all_task():
        print("all task was done")
    @task
    def task_d():
        print("d")

    # 实例化任务到变量,确保每个任务仅创建一次
    ta = task_a()
    tb = task_b()
    tc = task_c()
    td = task_d()

    # 通过变量设置依赖,所有分支指向同一个task_c
    ta >> tb >> tc
    td >> tc

complex_dags = complex_dags()

验证效果

修改后,所有分支都会指向同一个task_c实例,DAG的依赖关系会完全符合预期:task_a→task_b→task_c、task_d→task_c。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 03:32:39