Airflow DAG导入报错:AttributeError: '_TaskDecorator'无update_relative属性
Airflow DAG导入报错排查:
AttributeError: '_TaskDecorator' object has no attribute 'update_relative' 问题重现
运行以下Airflow DAG代码时出现导入失败:
from airflow.sensors.sql import SqlSensor import pendulum from airflow.decorators import task,dag @dag( dag_id = "database_monitor", schedule_interval = '*/10 * * * *', start_date=pendulum.datetime(2023, 7, 16, 21,0,tz="UTC"), catchup=False,) def Pipeline(): check_db_alive = SqlSensor( task_id="check_db_alive", conn_id="evergreen", sql="SELECT pg_is_in_recovery()", success= lambda x: x == False, poke_interval= 60, #timeout = 60 * 2, mode = "reschedule", ) @task() def alert_of_db_inrecovery(): import requests # result = f"Former primary instance is in recovery, task_instance_key_str: {kwargs['task_instance_key_str']}" data = {"@key":"kkll", "@version" : "alertapi-0.1", "@type":"ALERT", "object" : "Testobject", "severity" : "MINOR", "text" : str("Former primary instance is in recovery") } requests.post('https://httpevents.systems/api/sendAlert',verify=False,data=data) check_db_alive >> alert_of_db_inrecovery dag = Pipeline()
报错信息:
AttributeError: '_TaskDecorator' object has no attribute 'update_relative'
问题原因
报错核心原因是建立任务依赖时,直接引用了@task装饰后的函数对象,而非实际的Task实例。alert_of_db_inrecovery是被@task包装后的装饰器对象,不具备Airflow Task实例的属性和方法,无法与上游的SqlSensor任务建立依赖关系,最终导致DAG解析失败。
修复方案
需要调用被@task装饰的函数,生成实际的Task实例后再建立依赖关系。有两种实现方式:
方式1:先创建Task实例再连接依赖
from airflow.sensors.sql import SqlSensor import pendulum from airflow.decorators import task,dag @dag( dag_id = "database_monitor", schedule_interval = '*/10 * * * *', start_date=pendulum.datetime(2023, 7, 16, 21,0,tz="UTC"), catchup=False,) def Pipeline(): check_db_alive = SqlSensor( task_id="check_db_alive", conn_id="evergreen", sql="SELECT pg_is_in_recovery()", success= lambda x: x == False, poke_interval= 60, #timeout = 60 * 2, mode = "reschedule", ) @task() def alert_of_db_inrecovery(): import requests # result = f"Former primary instance is in recovery, task_instance_key_str: {kwargs['task_instance_key_str']}" data = {"@key":"kkll", "@version" : "alertapi-0.1", "@type":"ALERT", "object" : "Testobject", "severity" : "MINOR", "text" : str("Former primary instance is in recovery") } requests.post('https://httpevents.systems/api/sendAlert',verify=False,data=data) # 生成Task实例 alert_task = alert_of_db_inrecovery() # 建立依赖 check_db_alive >> alert_task dag = Pipeline()
方式2:直接在依赖中调用函数
简化写法,直接在依赖连接时调用函数生成实例:
# 替换原依赖代码行 check_db_alive >> alert_of_db_inrecovery()
关键提示
- 使用Airflow的
@task装饰器时,必须通过调用函数(即加())来生成可用于依赖连接的Task实例 - 如果任务需要接收上游传递的数据,调用时可传入对应参数,确保DAG逻辑的完整性
内容的提问来源于stack exchange,提问作者moth
相关产品推荐
相关产品推荐

