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正在建立新连接时,若有其他操作发起调用就会触发该错误。
触发场景及修复方法
Session被多请求/任务复用
- 检查FastAPI依赖注入逻辑:如果
create_async_session被注册为全局依赖且未设置use_cache=False,会导致同一个Session被多个请求复用。 - 修复:确保每个请求获取独立Session,注册依赖时添加
use_cache=False,或在依赖中每次创建新Session。
- 检查FastAPI依赖注入逻辑:如果
NullPool导致连接频繁重建
- 使用
NullPool意味着每次请求都会创建新连接,Session获取连接的开销更大,并发场景下更容易出现连接创建冲突。 - 修复:改用SQLAlchemy默认的
AsyncAdaptedQueuePool,它会复用连接,减少连接创建频率,降低冲突概率。
- 使用
UOW中Session生命周期管理不当
close_transaction方法中,Session不在事务中时直接关闭,若后续代码尝试复用已关闭的Session,会触发重新建立连接的操作,此时并发调用就会报错。- 修复:确保Session生命周期与请求/事务严格绑定,一旦关闭就不再复用;仅当事务完成且确实无需再使用Session时才关闭它。
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,提问作者Олексій Юдін
相关产品推荐
相关产品推荐

