如何在FastAPI+SQLAlchemy仓储模式应用中避免事务空闲
解决FastAPI+SQLAlchemy中事务空闲资源浪费问题
当前基于FastAPI+SQLAlchemy的架构能正常运行,但存在明显的数据库资源浪费:在create_user流程中,事务会在第一个数据库请求时开启,直到get_session依赖结束才关闭。中间调用Keycloak的HTTP请求耗时约0.5秒,这段时间事务处于Idle in transaction状态,数据库连接被无效占用。
现有核心代码回顾
基础会话与仓储定义
from sqlalchemy.ext.asyncio import create_async_engine, async_sessionmaker engine = create_async_engine(settings.DB.async_dns, **settings.DB.OPTIONS) Session = async_sessionmaker(bind=engine, expire_on_commit=False) # 事务管理依赖 async def get_session(): session = Session() try: yield session await session.commit() except Exception as e: await session.rollback() raise e finally: await session.close() class BaseRepository: def __init__(self, session: AsyncSession = Depends(get_session)): self._session = session # 查询与写操作方法混合定义 async def count(self, stmt: Select) -> int: stmt = stmt.with_only_columns(func.count(literal_column("1")), maintain_column_froms=True) return await self._session.scalar(stmt) async def create(self, **kwargs): instance = self.model(**kwargs) self._session.add(instance) return instance
服务层核心流程
class Service: async def create_user(self, input: InputSchema): # 1. 查询邮箱是否已注册(事务开启) if await self.repository.email_registered(input.email): raise HTTPException(status_code=400, detail="User already registered") # 2. 调用Keycloak接口(事务Idle,连接被占用) if not await self.user_exists_in_keycloak(input.email): raise HTTPException(status_code=400, detail="User does not exist in keycloak") # 3. 查询用户类型(事务继续Idle) user_type = await self.user_types_repository.get_by_pk(input.type_id) if not user_type: raise HTTPException(status_code=400, detail="User type not found") # 4. 执行写操作 await self.repository.some_method_that_modifies_something() user = await self.users_repository.create(**input.dict()) return user
最优解决方案:拆分只读/事务性会话
核心思路是将查询操作与写操作分离到不同的会话中:
- 只读会话:开启
autocommit模式,执行完查询立即释放连接,不保持事务 - 事务性会话:仅用于写操作和需要原子性的逻辑,事务仅在必要阶段开启
1. 定义两种会话依赖
from sqlalchemy.ext.asyncio import create_async_engine, async_sessionmaker engine = create_async_engine(settings.DB.async_dns, **settings.DB.OPTIONS) # 只读会话:自动提交,适合纯查询场景 ReadOnlySession = async_sessionmaker(bind=engine, expire_on_commit=False, autocommit=True) # 事务性会话:用于写操作,保证原子性 TransactionalSession = async_sessionmaker(bind=engine, expire_on_commit=False) async def get_readonly_session(): session = ReadOnlySession() try: yield session finally: await session.close() async def get_transactional_session(): session = TransactionalSession() try: yield session await session.commit() except Exception as e: await session.rollback() raise e finally: await session.close()
2. 拆分仓储类
将查询和写操作分别归属到不同的仓储基类:
# 只读仓储:仅包含查询方法 class BaseReadOnlyRepository: def __init__(self, session: AsyncSession = Depends(get_readonly_session)): self._session = session async def count(self, stmt: Select) -> int: stmt = stmt.with_only_columns(func.count(literal_column("1")), maintain_column_froms=True) return await self._session.scalar(stmt) async def scalar_or_raise(self, stmt, error: str = None): if result := await self._session.scalar(stmt): return result raise NotFoundError(error or "Resource not found") async def get_by_pk(self, pk): if instance := await self._session.get(self.model, pk): return instance raise NotFoundError # 事务性仓储:继承只读方法,添加写操作 class BaseTransactionalRepository(BaseReadOnlyRepository): def __init__(self, session: AsyncSession = Depends(get_transactional_session)): super().__init__(session) async def update_instance(self, instance, update_data: dict): for field, value in update_data.items(): setattr(instance, field, value) await self._session.flush() return instance async def create(self, **kwargs): instance = self.model(**kwargs) self._session.add(instance) return instance
3. 调整模型仓储实现
# 用户只读仓储(仅查询) class UsersReadOnlyRepository(BaseReadOnlyRepository): model = User # 用户事务仓储(写操作) class UsersRepository(BaseTransactionalRepository): model = User # 用户类型只读仓储 class UserTypesReadOnlyRepository(BaseReadOnlyRepository): model = UserType async def get_by_code(self, code): stmt = select(self.model).where(self.model.code == code) return await self.scalar_or_raise(stmt, "User type not found") # 端点专用只读仓储 class CreateUserReadOnlyRepository(BaseReadOnlyRepository): async def email_registered(self, email): stmt = select("1").where(User.email == email) return await self._session.scalar(stmt) is not None # 端点专用事务仓储 class CreateUserTransactionalRepository(BaseTransactionalRepository): async def some_method_that_modifies_something(self): stmt = update(...) await self._session.execute(stmt)
4. 重构服务层流程
将查询与写操作彻底分离,确保事务仅在最后写阶段开启:
class Service: def __init__( self, readonly_repo: CreateUserReadOnlyRepository = Depends(), transactional_repo: CreateUserTransactionalRepository = Depends(), users_repo: UsersRepository = Depends(), user_types_repo: UserTypesReadOnlyRepository = Depends(), ): self.readonly_repo = readonly_repo self.transactional_repo = transactional_repo self.users_repo = users_repo self.user_types_repo = user_types_repo async def create_user(self, input: InputSchema): # 1. 只读查询:邮箱是否已注册(无事务,查询完立即释放连接) if await self.readonly_repo.email_registered(input.email): raise HTTPException(status_code=400, detail="User already registered") # 2. 调用Keycloak接口:此时无任何数据库事务/连接占用 if not await self.user_exists_in_keycloak(input.email): raise HTTPException(status_code=400, detail="User does not exist in keycloak") # 3. 只读查询:用户类型是否存在(无事务) user_type = await self.user_types_repo.get_by_pk(input.type_id) # 4. 事务性操作:仅在此阶段开启事务,执行所有写逻辑 await self.transactional_repo.some_method_that_modifies_something() user = await self.users_repo.create( first_name=input.first_name, last_name=input.last_name, type_id=input.type_id ) return user async def user_exists_in_keycloak(self, email) -> bool: # 外部HTTP请求逻辑 pass
方案优势对比
- 对比「手动commit查询」:避免了频繁开启/关闭事务的开销,也不会出现查询间数据不一致的问题,维护成本更低
- 对比「随意混用两个会话」:明确拆分只读/写职责,从架构层面保证资源高效利用,同时保留写操作的原子性
内容的提问来源于stack exchange,提问作者Альберт Александров
相关产品推荐
相关产品推荐

