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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 18:23:10