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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 03:34:51