FastAPI WebSocket结合数据库时服务冻结问题求助
核心原因:数据库操作阻塞ASGI事件循环 + 连接池耗尽
你的问题本质是同步数据库操作与ASGI异步模型的冲突,加上连接池配置不合理导致的服务假死,具体拆解:
同步SQLAlchemy阻塞事件循环
FastAPI的WebSocket依赖ASGI异步事件循环,如果你用的是普通的同步SQLAlchemy(非asyncio版本),每条消息的DB写入都是阻塞IO操作。当50+连接同时发送消息,大量阻塞操作会占满事件循环,导致新的WebSocket连接请求、其他路由请求都无法被处理,直接表现为服务无响应。即使停止发送消息,已经阻塞的任务可能还在排队,或者连接池资源没释放,所以无法自动恢复。数据库连接池耗尽
SQLAlchemy默认连接池大小很小(通常是5),当50+并发请求同时发起DB操作,连接池瞬间被占满,后续请求会一直等待空闲连接。如果代码里没有正确释放连接(比如事务未提交/回滚、会话没关闭),这些被占用的连接不会归还到池里,即使用户断开WebSocket,连接也一直被挂着,导致服务彻底无法获取新连接,只能重启。未正确处理事务/会话
如果每个消息写入都开启了会话,但没在操作完成后及时提交、回滚并关闭会话,连接会被持续占用,连接池很快枯竭,后续所有依赖DB的操作都会卡住,进而拖垮整个服务。
针对性解决方案
切换到异步SQLAlchemy + 异步驱动
这是根治问题的最优解,用SQLAlchemy的异步版本(sqlalchemy.ext.asyncio)配合异步数据库驱动(比如PostgreSQL用asyncpg,MySQL用aiomysql),所有DB操作都是非阻塞的,不会占用ASGI事件循环。示例代码框架:from sqlalchemy.ext.asyncio import AsyncSession, create_async_engine from sqlalchemy.orm import sessionmaker engine = create_async_engine("postgresql+asyncpg://user:pass@localhost/dbname") AsyncSessionLocal = sessionmaker(engine, class_=AsyncSession, expire_on_commit=False) async def save_message(db: AsyncSession, message_data): # 异步DB操作 new_msg = Message(**message_data) db.add(new_msg) await db.commit()调整连接池参数
如果你暂时不想改异步,至少要增大连接池的容量,同时确保数据库本身允许足够的连接数。在创建SQLAlchemy引擎时配置:from sqlalchemy import create_engine engine = create_engine( "postgresql://user:pass@localhost/dbname", pool_size=20, # 常规连接数 max_overflow=30, # 临时溢出连接数 pool_recycle=3600 # 定期回收闲置连接,避免连接失效 )注意:
pool_size + max_overflow不能超过数据库的max_connections配置(比如PostgreSQL默认是100)。强制释放数据库连接
用同步SQLAlchemy时,必须确保每个DB操作完成后都关闭会话,用try/finally块兜底:def save_message(message_data): db = SessionLocal() try: new_msg = Message(**message_data) db.add(new_msg) db.commit() db.refresh(new_msg) except Exception as e: db.rollback() raise e finally: db.close() # 必须关闭,归还连接到池用线程池隔离同步操作
如果必须用同步DB操作,要把DB任务放到独立的线程池里,避免阻塞ASGI事件循环。可以自定义一个更大的线程池:from concurrent.futures import ThreadPoolExecutor from fastapi import BackgroundTasks executor = ThreadPoolExecutor(max_workers=30) # 根据并发量调整 async def handle_websocket_message(websocket, message_data, background_tasks: BackgroundTasks): # 把DB操作丢到线程池,不阻塞事件循环 background_tasks.add_task(executor.submit, save_message, message_data) # 分发消息给其他客户端...批量写入优化
如果消息量极大,可以攒一批消息再批量插入,减少DB连接的调用次数:async def batch_save_messages(db: AsyncSession, messages_list): db.add_all([Message(**data) for data in messages_list]) await db.commit()
内容的提问来源于stack exchange,提问作者Мария Мамонова

