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

FastAPI集成SQLAlchemy异步查询偶发连接并发操作不允许错误

问题

在基于FastAPI的服务中使用SQLAlchemy结合asyncio做异步数据库查询(包含读操作)时,出现偶发错误:
sqlalchemy.exc.InvalidRequestError: This session is provisioning a new connection; concurrent operations are not permitted
服务正常运行一段时间后会突然触发该问题。

使用版本

  • Python ^3.12
  • sqlalchemy ^2.0.29
  • asyncpg ^0.29.0
  • greenlet ^3.0.3

相关代码

connection.py

from typing import AsyncIterable
from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker, create_async_engine
from sqlalchemy import NullPool

SessionFactoryType = async_sessionmaker[AsyncSession]

def create_session_factory(connection_url: str, **kwargs: dict[str, ...]) -> async_sessionmaker[AsyncSession]:
    async_engine = create_async_engine(url=connection_url, **kwargs, poolclass=NullPool)
    return async_sessionmaker(async_engine, autoflush=False, expire_on_commit=False)

async def create_async_session(session_factory: SessionFactoryType) -> AsyncIterable[AsyncSession]:
    async with session_factory() as session:
        yield session

unit_of_work.py

import contextlib
from sqlalchemy.exc import SQLAlchemyError
from sqlalchemy.ext.asyncio import AsyncSession, AsyncSessionTransaction
from src.common.interfaces.unit_of_work import AbstractUnitOfWork

class SQLAlchemyUnitOfWork(AbstractUnitOfWork[AsyncSession, AsyncSessionTransaction]):
    async def commit(self) -> None:
        with contextlib.suppress(SQLAlchemyError):
            await self.session.commit()

    async def rollback(self) -> None:
        with contextlib.suppress(SQLAlchemyError):
            await self.session.rollback()

    async def create_transaction(self) -> None:
        if not self.session.in_transaction() and self.session.is_active:
            self._transaction = await self.session.begin()

    async def close_transaction(self) -> None:
        if self.session.is_active and not self.session.in_transaction():
            await self.session.close()

def unit_of_work_factory(session: AsyncSession) -> SQLAlchemyUnitOfWork:
    return SQLAlchemyUnitOfWork(session=session)

crud.py

from typing import Any, Mapping, Optional, Sequence, TypeVar
from iotsota_models.orm import BaseSQLAlchemyModel
from sqlalchemy import ColumnExpressionArgument, insert, select, update
from sqlalchemy.ext.asyncio import AsyncSession
from src.common.interfaces.crud import AbstractCRUDRepository

ModelType = TypeVar("ModelType", bound=BaseSQLAlchemyModel)

class SQLAlchemyCRUDRepository(AbstractCRUDRepository[ModelType, ColumnExpressionArgument]):
    __slots__ = ("_session",)

    def __init__(self, session: AsyncSession, model: type[ModelType]) -> None:
        super().__init__(model)
        self._session = session

    async def select(self, *clauses: ColumnExpressionArgument) -> ModelType | None:
        stmt = select(self.model).where(*clauses)
        return (await self._session.execute(stmt)).scalars().first()

    async def select_many(
        self, *clauses: ColumnExpressionArgument, offset: Optional[int] = None, limit: Optional[int] = None
    ) -> Sequence[ModelType]:
        stmt = select(self.model).where(*clauses).offset(offset).limit(limit)
        return (await self._session.execute(stmt)).scalars().all()

    async def create(self, **values: Mapping[str, Any]) -> ModelType | None:
        stmt = insert(self.model).values(**values).returning(self.model)
        return (await self._session.execute(stmt)).scalars().first()

    async def update(self, *clauses: ColumnExpressionArgument, **values: Mapping[str, Any]) -> list[ModelType] | None:
        stmt = update(self.model).where(*clauses).values(**values).returning(self.model)
        return (await self._session.execute(stmt)).unique().scalars().all()

错误堆栈

File "/Users/oleksiiyudin/Documents/WORK/authentication-service/src/databases/orm/repositories/user/reader.py", line 13, in select_user
    return await self.repository.crud.select(clauses)
           ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/Users/oleksiiyudin/Documents/WORK/authentication-service/src/databases/orm/repositories/crud.py", line 21, in select
    return (await self._session.execute(stmt)).scalars().first()
            ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/Users/oleksiiyudin/Library/Caches/pypoetry/virtualenvs/auth_service-v_N9SFpi-py3.12/lib/python3.12/site-packages/sqlalchemy/ext/asyncio/session.py", line 461, in execute
    result = await greenlet_spawn(
             ^^^^^^^^^^^^^^^^^^^^^
  File "/Users/oleksiiyudin/Library/Caches/pypoetry/virtualenvs/auth_service-v_N9SFpi-py3.12/lib/python3.12/site-packages/sqlalchemy/util/_concurrency_py3k.py", line 190, in greenlet_spawn
    result = context.switch(*args, **kwargs)
             ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/Users/oleksiiyudin/Library/Caches/pypoetry/virtualenvs/auth_service-v_N9SFpi-py3.12/lib/python3.12/site-packages/sqlalchemy/orm/session.py", line 2306, in execute
    return self._execute_internal(
           ^^^^^^^^^^^^^^^^^^^^^^^
  File "/Users/oleksiiyudin/Library/Caches/pypoetry/virtualenvs/auth_service-v_N9SFpi-py3.12/lib/python3.12/site-packages/sqlalchemy/orm/session.py", line 2181, in _execute_internal
    conn = self._connection_for_bind(bind)
           ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/Users/oleksiiyudin/Library/Caches/pypoetry/virtualenvs/auth_service-v_N9SFpi-py3.12/lib/python3.12/site-packages/sqlalchemy/orm/session.py", line 2050, in _connection_for_bind
    return trans._connection_for_bind(engine, execution_options)
           ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "<string>", line 2, in _connection_for_bind
  File "/Users/oleksiiyudin/Library/Caches/pypoetry/virtualenvs/auth_service-v_N9SFpi-py3.12/lib/python3.12/site-packages/sqlalchemy/orm/state_changes.py", line 103, in _go
    self._raise_for_prerequisite_state(fn.__name__, current_state)
  File "/Users/oleksiiyudin/Library/Caches/pypoetry/virtualenvs/auth_service-v_N9SFpi-py3.12/lib/python3.12/site-packages/sqlalchemy/orm/session.py", line 945, in _raise_for_prerequisite_state
    raise sa_exc.InvalidRequestError(
sqlalchemy.exc.InvalidRequestError: This session is provisioning a new connection; concurrent operations are not permitted
解决方案

这个错误的核心原因是同一个AsyncSession被多个异步任务并发调用。SQLAlchemy的异步Session并非并发安全,同一时间只能有一个操作执行,当Session正在建立新连接时,若有其他操作发起调用就会触发该错误。

触发场景及修复方法

  1. Session被多请求/任务复用

    • 检查FastAPI依赖注入逻辑:如果create_async_session被注册为全局依赖且未设置use_cache=False,会导致同一个Session被多个请求复用。
    • 修复:确保每个请求获取独立Session,注册依赖时添加use_cache=False,或在依赖中每次创建新Session。
  2. NullPool导致连接频繁重建

    • 使用NullPool意味着每次请求都会创建新连接,Session获取连接的开销更大,并发场景下更容易出现连接创建冲突。
    • 修复:改用SQLAlchemy默认的AsyncAdaptedQueuePool,它会复用连接,减少连接创建频率,降低冲突概率。
  3. UOW中Session生命周期管理不当

    • close_transaction方法中,Session不在事务中时直接关闭,若后续代码尝试复用已关闭的Session,会触发重新建立连接的操作,此时并发调用就会报错。
    • 修复:确保Session生命周期与请求/事务严格绑定,一旦关闭就不再复用;仅当事务完成且确实无需再使用Session时才关闭它。
  4. Repository实例被多任务共享

    • 若Repository实例被多个异步任务共享,其内部持有的Session也会被并发调用。
    • 修复:每个请求/任务创建独立的Repository实例,确保每个Repository持有专属Session。

关键代码调整示例

修改connection.py中的连接池

替换NullPool为默认异步连接池:

def create_session_factory(connection_url: str, **kwargs: dict[str, ...]) -> async_sessionmaker[AsyncSession]:
    # 移除poolclass=NullPool,使用默认的AsyncAdaptedQueuePool
    async_engine = create_async_engine(url=connection_url, **kwargs)
    return async_sessionmaker(async_engine, autoflush=False, expire_on_commit=False)

确保FastAPI依赖每次返回新Session

注册依赖时添加use_cache=False:

from fastapi import Depends

async def get_session(session_factory: SessionFactoryType = Depends(get_session_factory)) -> AsyncIterable[AsyncSession]:
    async with session_factory() as session:
        yield session

# 路由中使用
@app.get("/users/{user_id}")
async def get_user(user_id: int, session: AsyncSession = Depends(get_session, use_cache=False)):
    # 业务逻辑

内容的提问来源于stack exchange,提问作者Олексій Юдін

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 05:27:02