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

如何在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,提问作者Альберт Александров

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 08:29:52