Async SQLAlchemy执行SQL语句的正确方式及FastAPI端点卡顿排查
问题分析与修复方案
核心问题
- 每次请求创建并销毁数据库引擎:每个端点调用都会新建
AsyncEngine,且run_stmt最后执行engine.dispose()销毁引擎。引擎自带连接池,频繁创建销毁会导致每次请求都要重新建立数据库连接,完全失去连接池的复用优势。一旦有耗时SQL执行,后续请求都要等待新引擎初始化+连接建立,引发全局卡顿。 - 冗余的Session操作:
async with AsyncSession(...)会自动管理Session生命周期,手动调用await session.close()属于多余操作。- 只读SQL操作不需要执行
session.commit(),该操作会触发不必要的事务提交流程,增加额外开销。
修复后的代码
1. 全局复用数据库引擎
将引擎初始化改为全局单例,避免每次请求重复创建:
# 全局初始化引擎,仅创建一次 engine = create_async_engine(ConnectionString, future=True, echo=True) async def create_postgresql_engine(): # 返回全局引擎实例,不再每次新建 return engine
2. 精简run_stmt函数
移除冗余操作,保留核心逻辑:
import datetime import pandas as pd from sqlalchemy.exc import SQLAlchemyError from sqlalchemy.ext.asyncio import AsyncSession async def run_stmt(stmt, engine): print("+++++------- time is:", datetime.now() , " and statement is: ", stmt, "-------++++++++++") df = pd.DataFrame() async with AsyncSession(engine) as session: try: result = await session.execute(stmt) rows = result.fetchall() columns = result.keys() df = pd.DataFrame(rows, columns=columns) if len(rows) > 0 else pd.DataFrame(columns=columns) df = df.rename(columns=str.lower) except SQLAlchemyError as e: error = str(e.__cause__) await session.rollback() raise RuntimeError(error) from e return df
3. 端点调用逻辑调整
确保复用全局引擎,避免重复创建:
async def camp_channel_unique(engine): # 构造你的SQL语句 stmt = ... df = await run_stmt(stmt, engine) return df async def get_channels_in_camp(): df = await camp_channel_unique(engine=await create_postgresql_engine()) return df
额外优化建议
- 配置连接池参数:根据并发量调整连接池大小,避免连接耗尽:
engine = create_async_engine( ConnectionString, future=True, echo=True, pool_size=10, # 常驻连接数 max_overflow=20, # 临时扩容连接数 pool_recycle=300 # 自动回收闲置连接,防止数据库断开 ) - 优化耗时SQL:为查询字段添加索引、拆分大查询或使用分页,从根源减少长耗时请求的影响。
- 避免同步阻塞操作:确保异步函数内无同步阻塞代码(当前pandas内存操作无影响,若后续涉及IO需保持异步)。
内容的提问来源于stack exchange,提问作者Moh-Spark
相关产品推荐
相关产品推荐

