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

Airflow报错AttributeError: 'coroutine'对象无update_relative属性求助

解决Airflow DAG中coroutine对象的AttributeError问题

问题代码

用户的DAG代码如下:

from datetime import timedelta
import pendulum
from airflow.decorators import dag
from stage import stage_data
from table_async_pg import main

@dag(
dag_id = "data-sync",
schedule_interval = "1 * * * *",
start_date=pendulum.datetime(2023, 3, 8, tz="Asia/Hong_Kong"),
catchup=False,
dagrun_timeout=timedelta(minutes=20),
)
def Pipeline():
    a = stage_data()
    b = main()

    a >> b

pipeline = Pipeline()

错误现象

仅保留a = stage_data()时功能正常,加入b=main()和a>>b后触发以下错误:

/opt/airflow/dags/airflow_sched.py | Traceback (most recent call last):                                                                       
                                   |   File "/opt/.venv/lib/python3.9/site-packages/airflow/models/taskmixin.py", line 230, in set_downstream
                                   |     self._set_relatives(task_or_task_list, upstream=False, edge_modifier=edge_modifier)                 
                                   |   File "/opt/.venv/lib/python3.9/site-packages/airflow/models/taskmixin.py", line 175, in _set_relatives
                                   |     task_object.update_relative(self, not upstream)                                                      
                                   | AttributeError: 'coroutine' object has no attribute 'update_relative'         

问题原因

main()是异步函数,直接调用会返回coroutine对象,而Airflow的任务依赖关系(>>)只能在Airflow官方的Task对象之间建立,coroutine对象不具备Airflow任务的update_relative方法,因此报错。

解决方案

需要将异步函数包装为Airflow可识别的异步任务,使用Airflow 2.2+支持的@task.async装饰器即可实现。

方案1:在DAG中包装异步函数

修改DAG代码,给异步函数main()添加包装:

from datetime import timedelta
import pendulum
from airflow.decorators import dag, task
from stage import stage_data
from table_async_pg import main

@dag(
dag_id = "data-sync",
schedule_interval = "1 * * * *",
start_date=pendulum.datetime(2023, 3, 8, tz="Asia/Hong_Kong"),
catchup=False,
dagrun_timeout=timedelta(minutes=20),
)
def Pipeline():
    a = stage_data()
    
    # 用@task.async包装异步函数,转换为Airflow异步任务
    @task.async
    async def async_main_task():
        await main()
    
    b = async_main_task()
    a >> b

pipeline = Pipeline()

方案2:直接修改异步函数文件

如果有权限修改table_async_pg.py,可以直接在函数上添加装饰器:

# table_async_pg.py
from airflow.decorators import task

@task.async
async def main():
    # 原异步逻辑代码
    ...

之后DAG中无需额外包装,直接调用b = main()即可正常建立依赖关系。

原理说明

@task.async装饰器会将普通异步函数转换为Airflow的AsyncTask对象,该对象继承了Airflow任务的核心方法(包括update_relative),因此可以正常参与任务依赖的构建。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 15:43:33