如何在Celery任务中正确调用async异步函数
Celery默认注册的任务为同步普通函数,普通函数内无法直接使用await关键字调用异步协程。直接调用async def定义的函数时,不会执行函数内部逻辑,只会返回未调度的协程对象,这就是调用select_users()无法拿到数据库查询结果的根本原因。
方案1:任务内手动启动事件循环执行协程
这是改造成本最低的方案,无需调整Celery启动配置,借助标准库asyncio即可实现,适合仅少量任务需要调用异步函数的场景。
修改后的tasks.py代码如下:
from .celery import app import asyncio import db @app.task def update_credits(): # 启动事件循环运行异步协程,拿到实际返回结果 users = asyncio.run(db.select_users()) print(users)
注意:不要在同一个线程内嵌套运行多个事件循环,否则会抛出
RuntimeError: This event loop is already running错误。如果你的Celery worker运行在已有事件循环的环境中,asyncio.run()会触发上述错误,此时可替换为asyncio.get_event_loop().run_until_complete(db.select_users());如果仍存在嵌套冲突,可以安装nest_asyncio库,在代码入口执行nest_asyncio.apply()打补丁即可解决。
方案2:使用Celery原生异步任务(Celery 5.0及以上版本支持)
Celery 5.0版本后原生支持异步任务,只要将worker的执行池切换为asyncio模式,就可以直接把任务定义为异步函数,内部直接使用await调用协程,适配全异步逻辑的项目。
- 第一步:启动Celery worker时指定异步池:
celery -A 你的celery应用模块名 worker --pool=asyncio --loglevel=info - 第二步:修改
tasks.py,将任务定义为异步函数:
from .celery import app import db @app.task async def update_credits(): # 直接await调用异步函数即可 users = await db.select_users() print(users)
方案3:封装同步调用包装层
如果有多个Celery任务、或者多个同步场景需要调用异步的数据库方法,可以单独封装同步版本的调用函数,避免重复编写事件循环相关代码。
示例:在db.py中新增同步包装方法
import asyncio # 保留原有的异步select_users方法不变 async def select_users(): sql = "SELECT * FROM Users WHERE " sql, parameters = self.format_args(sql, parameters=kwargs) return await self.execute(sql, *parameters, fetchrow=True) # 新增同步版本方法供同步场景调用 def select_users_sync(): return asyncio.run(select_users())
之后在Celery任务中直接调用同步方法即可:
from .celery import app import db @app.task def update_credits(): users = db.select_users_sync() print(users)
内容的提问来源于stack exchange,提问作者mirodil

