如何实现SQLAlchemy AsyncSession生成器?异常时自动回滚关闭
解决SQLAlchemy AsyncSession生成器报错及异常自动回滚问题
问题原因
原代码中get_session使用yield返回会话,导致该方法成为异步生成器。直接赋值session = db.get_session()得到的是生成器对象,而非实际的AsyncSession实例,因此调用session.scalar()会触发AttributeError。同时原代码未实现异常时自动回滚、关闭会话的逻辑。
修复后的代码
修改DatabaseHelper的get_session方法,将其实现为支持异步上下文管理器的生成器,同时添加异常回滚、会话关闭的逻辑:
from sqlalchemy.ext.asyncio import create_async_engine, async_sessionmaker, AsyncSession from sqlalchemy import select class DatabaseHelper: def __init__( self, url: str, echo: bool = False, ) -> None: self.engine = create_async_engine( url=url, echo=echo, ) self.session_factory = async_sessionmaker( bind=self.engine, expire_on_commit=False # 可选:避免提交后实例过期 ) async def get_session(self) -> AsyncSession: session = self.session_factory() try: yield session # 正常执行时提交会话 await session.commit() except Exception: # 发生异常时回滚会话 await session.rollback() # 重新抛出异常,让上层处理 raise finally: # 无论是否异常,最终关闭会话 await session.close()
正确调用方式
使用async with语法获取会话,自动处理生成器迭代,拿到实际的AsyncSession实例:
async def main(): db = DatabaseHelper(url="your-database-url") # 使用async with获取会话 async with db.get_session() as session: group = await session.scalar(select(Group).where(Group.tg_id == 1)) # 执行其他数据库操作...
逻辑说明
- 正常流程:代码块执行完毕后,自动提交会话,最后关闭会话。
- 异常流程:代码块抛出异常时,先回滚会话,重新抛出异常供上层捕获处理,最后关闭会话。
expire_on_commit=False是可选配置,避免提交会话后数据库实例被标记为过期,方便后续使用实例属性。
内容的提问来源于stack exchange,提问作者Pitoni
相关产品推荐
相关产品推荐

