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
相关产品推荐
相关产品推荐

