如何在Celery任务中使用SQLAlchemy的AsyncSession?报错解决
问题解决方案
错误根源
async_scoped_session的scopefunc参数要求传入一个可调用对象,但你写的asyncio.current_task()会直接执行函数并返回当前的Task对象,而非传递函数本身,这就导致后续调用时出现“Task对象不可调用”的错误。
修正步骤
1. 修复scoped_session的scopefunc参数
把asyncio.current_task()改成asyncio.current_task,传递函数引用而非执行结果,同时修正finally块的会话清理逻辑:
@asynccontextmanager async def scoped_session(): scoped_factory = async_scoped_session( async_session, scopefunc=asyncio.current_task # 去掉括号,传递函数本身 ) try: async with scoped_factory() as s: yield s finally: # 直接调用scoped_factory.remove(),无需重复实例化 await scoped_factory.remove()
2. 适配Celery同步环境运行异步代码
Celery是同步任务框架,直接用asyncio.run()可能引发事件循环冲突,建议用asgiref.sync.async_to_sync包装异步逻辑:
from asgiref.sync import async_to_sync @celery.task(name='is_event_done', bind=True, ignore_result=True) def is_event_done(self): async_to_sync(logic)()
3. 修正时区与结果处理问题
datetime.now()使用本地时区,建议改用UTC时间和数据库保持一致;同时fetchall()返回元组,需用scalars()提取Event对象:
from datetime import datetime, timezone async def logic(): async with scoped_session() as session: stmt = select(event.models.Event).where( event.models.Event.end_time <= datetime.now(timezone.utc) # 使用UTC时间 ) results = await session.execute(stmt) # 提取Event对象实例 for event in results.scalars().all(): print(event.is_event_done)
完整修正后代码
数据库会话定义
engine = create_async_engine(settings.db_url, echo=True) async_session = sessionmaker( engine, class_=AsyncSession, expire_on_commit=False ) @asynccontextmanager async def scoped_session(): scoped_factory = async_scoped_session( async_session, scopefunc=asyncio.current_task ) try: async with scoped_factory() as s: yield s finally: await scoped_factory.remove()
业务逻辑与Celery任务
from datetime import datetime, timezone from asgiref.sync import async_to_sync async def logic(): async with scoped_session() as session: stmt = select(event.models.Event).where( event.models.Event.end_time <= datetime.now(timezone.utc) ) results = await session.execute(stmt) for event in results.scalars().all(): print(event.is_event_done) @celery.task(name='is_event_done', bind=True, ignore_result=True) def is_event_done(self): async_to_sync(logic)()
额外说明
async_scoped_session的scopefunc用于标识作用域,必须传入可调用对象返回当前作用域标识。- 用
async_to_sync比直接asyncio.run()更适配Celery同步环境,避免事件循环复用问题。 - 统一用UTC时间处理数据库时间字段,可避免跨时区环境下的时间判断错误。
内容的提问来源于stack exchange,提问作者jabajke
相关产品推荐
相关产品推荐

