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

FastAPI+aiosqlite+SQLAlchemy并发请求InvalidRequestError修复

解决FastAPI+aiosqlite+SQLAlchemy高并发下session.refresh失败问题

问题核心分析

高并发下await session.refresh(new)抛出InvalidRequestError: Could not refresh instance,本质是SQLite的文件锁机制、连接池配置不合理,以及会话生命周期管理不当共同导致的:

  • SQLite默认日志模式下,写操作会独占数据库锁,高并发场景下锁竞争会导致部分写入操作未完成就执行refresh
  • SingletonThreadPool是单连接池,高并发下多个会话复用同一连接,会引发事务操作交织,导致实例未被正确持久化
  • 会话中refresh的时机与事务提交/flush的时序不匹配,导致实例未写入数据库就尝试刷新

针对性解决方案

1. 切换为异步安全的连接池

替换SingletonThreadPool为AsyncAdaptedQueuePool,配置合理的连接池参数,确保高并发下每个会话能获取独立连接:

from sqlalchemy.ext.asyncio import create_async_engine
from sqlalchemy.pool import AsyncAdaptedQueuePool

engine = create_async_engine(
    "sqlite+aiosqlite:///./your_db.db",
    poolclass=AsyncAdaptedQueuePool,
    pool_size=15,  # 基于并发量调整,建议大于单批次并发请求数的1/2
    max_overflow=30,  # 超出pool_size时临时创建的连接数
    pool_recycle=300,  # 连接自动回收时间,避免长期占用
    pool_pre_ping=True  # 每次获取连接前检查可用性
)

2. 启用SQLite WAL日志模式

WAL模式支持多读单写,大幅降低写操作的锁冲突概率,是提升SQLite并发性能的关键:

engine = create_async_engine(
    "sqlite+aiosqlite:///./your_db.db?journal_mode=WAL",
    # 搭配上述连接池参数
)

3. 调整会话操作时序

确保refresh在flush之后执行,若依赖自增ID等数据库生成字段,flush已足够将实例写入数据库(无需等待commit);同时避免在未完成持久化时执行refresh:

# 正确的端点示例
from fastapi import Depends, FastAPI
from sqlalchemy.ext.asyncio import AsyncSession

app = FastAPI()

async def get_db():
    async with AsyncSession(engine) as session:
        yield session

@app.post("/save")
async def save_data(db: AsyncSession = Depends(get_db)):
    new_instance = YourDBModel(...)
    db.add(new_instance)
    # 先flush将实例写入数据库(事务未提交,但已生成数据库层面的实例)
    await db.flush()
    # 此时再执行refresh获取最新状态
    await db.refresh(new_instance)
    # 最后提交事务
    await db.commit()
    return {"id": new_instance.id}

4. 添加重试机制处理临时锁冲突

高并发下偶尔的锁冲突属于正常情况,通过重试机制可自动恢复:

from tenacity import retry, stop_after_attempt, wait_exponential, retry_if_exception_type
import sqlalchemy.exc

@retry(
    stop=stop_after_attempt(3),  # 最多重试3次
    wait=wait_exponential(multiplier=1, min=1, max=3),  # 指数退避等待
    retry=retry_if_exception_type(sqlalchemy.exc.InvalidRequestError)
)
async def safe_refresh(session, instance):
    await session.refresh(instance)

# 在端点中使用
@app.post("/save")
async def save_data(db: AsyncSession = Depends(get_db)):
    new_instance = YourDBModel(...)
    db.add(new_instance)
    await db.flush()
    await safe_refresh(db, new_instance)
    await db.commit()
    return {"id": new_instance.id}

5. 确保会话的请求隔离

通过FastAPI的依赖注入管理会话生命周期,保证每个请求拥有独立的会话实例,避免跨请求的会话共享:

# 依赖注入函数,每个请求生成新的会话
async def get_db():
    async with AsyncSession(engine) as session:
        yield session

# 端点中通过Depends获取会话
@app.post("/save")
async def save_data(db: AsyncSession = Depends(get_db)):
    # 业务逻辑

内容的提问来源于stack exchange,提问作者Andreas

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 11:06:33