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
相关产品推荐
相关产品推荐

