使用装饰器定义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
相关产品推荐
相关产品推荐

