FastAPI下Async SQLAlchemy两种连接方案的压测对比与结论验证
两种Async SQLAlchemy数据库连接方案的优劣分析及结论验证
方案概述
- 方案1:复用单个异步引擎,为每个请求创建新异步会话
- 方案2:为每个请求创建新异步引擎与会话,会话结束后销毁引擎
压测环境与监控方式
Locust压测代码
@task(1) def testFastApitask1(self): headers = {'Accept': 'application/json', 'Content-Type': 'application/json'} resp = self.client.get( url="/data/campaign/channels/no_log_no_auth", auth=None, name='task1' )
PostgreSQL连接状态监控SQL
select max_conn,used,res_for_super,max_conn-used-res_for_super res_for_normal from (select count(*) used from pg_stat_activity) t1, (select setting::int res_for_super from pg_settings where name=$$superuser_reserved_connections$$) t2, (select setting::int max_conn from pg_settings where name=$$max_connections$$) t3;
基准连接数13,总连接数100。
方案细节与压测结果
方案1核心逻辑与结果
核心为复用全局异步引擎,通过会话工厂为每个请求生成会话:
# engine.py核心代码 ConnectionString="{}://{}:{}@{}:{}/{}".format(dialect,DS_POSTGRESQL_USER,DS_POSTGRESQL_PASSWORD,DS_POSTGRESQL_HOST,DS_POSTGRESQL_PORT,DS_POSTGRESQL_DATABASE) async_engine1 = create_async_engine(ConnectionString, future=True, echo=True) async_session = sessionmaker(async_engine1, class_=AsyncSession, expire_on_commit=False) async def run_stmt(stmt, async_session: async_sessionmaker): async with async_session() as session: try: async with session.begin(): 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 else: await session.commit() finally: await session.close()
压测结果:并发用户达5000,错误总数33527次,资源占用更低。
方案2核心逻辑与结果
每个请求创建独立引擎,会话结束后销毁引擎:
# engine.py核心代码 async def create_postgresql_engine(): engine = create_async_engine(ConnectionString, future=True, echo=True) return engine async def run_stmt(stmt, engine): 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 else: await session.commit() finally: await session.close() await engine.dispose()
压测结果:PostgreSQL连接数在90-103波动,错误总数67133次。
优劣分析
方案1
- 优势:
- 引擎复用避免了频繁创建/销毁引擎的开销(包含TCP连接池初始化、数据库协议协商等成本),资源占用更低
- 错误率更低:单个引擎的连接池自动管理连接,避免了方案2中短时间创建大量引擎导致PostgreSQL连接数耗尽的问题
- 劣势:
- 响应时间波动大:并发请求激增时,连接池连接可能被耗尽,新请求需等待空闲连接,导致响应时间不稳定
方案2
- 优势:无实际业务场景下的有效优势(除非极端隔离需求,但无必要)
- 劣势:
- 资源开销极大:每个请求重复初始化引擎,CPU、网络资源浪费严重
- 连接数极易耗尽:每个引擎默认创建多个连接,高并发下直接突破PostgreSQL的
max_connections限制,引发大量连接错误 - 错误率极高:连接耗尽导致的数据库连接失败、超时等错误大幅增加
结论验证
你的结论完全正确:方案1资源占用更少、错误更少,响应时间波动是连接池等待导致的正常现象。
优化建议
针对方案1的响应时间波动问题,可通过调整SQLAlchemy异步引擎的连接池参数优化:
async_engine1 = create_async_engine( ConnectionString, future=True, echo=True, pool_size=80, # 匹配PostgreSQL可用连接数(总100-预留20) max_overflow=10, # 峰值时允许临时扩容连接 pool_recycle=300 # 自动回收闲置连接,避免数据库主动断开 )
内容的提问来源于stack exchange,提问作者Moh-Spark
相关产品推荐
相关产品推荐

