FastAPI中WebSocket连接5分钟后触发ConnectionClosedOK(1000)问题排查
WebSocket连接5分钟后断开问题排查与解决
问题描述
我基于FastAPI官方WebSocket示例实现了3个客户端定期向服务器发送数据的功能:服务器通过ConnectionManager管理WebSocket连接,客户端使用websockets库周期性发送数据并接收响应。但扩展到5个客户端后,运行约5分钟会触发websockets.exceptions.ConnectionClosedOK: received 1000 (OK)错误,无法定位根源。
服务器代码
class ConnectionManager: def __init__(self): self.active_connections: list[WebSocket] = [] async def connect(self, websocket: WebSocket): await websocket.accept() self.active_connections.append(websocket) def disconnect(self, websocket: WebSocket): self.active_connections.remove(websocket) async def send_response(self, message: dict, websocket: WebSocket): await websocket.send_json(message) async def broadcast(self, message: str): for connection in self.active_connections: await connection.send_text(message) manager = ConnectionManager() @app.websocket("/chargeStationState/{client_id}") async def websocket_endpoint(websocket: WebSocket, client_id: int, db: Session = Depends(deps.get_db)): await manager.connect(websocket) try: while True: message = await websocket.receive_json() logging.info(message) ## read data from db response = { "stations": "repsonse", "timestamp": int(time.time()) } await manager.send_response(response, websocket) #await manager.broadcast(f"Client #{client_id} says: {data}") except WebSocketDisconnect: manager.disconnect(websocket)
客户端代码
async for websocket in websockets.connect("ws://127.0.0.1:8001/chargeStationState/1"): message = {'name':'station1'} await websocket.send(json.dumps(message)) p = await asyncio.wait_for(websocket.recv(), timeout=10) print(p) await asyncio.sleep(2)
可能原因与解决方案
1. 服务器空闲超时断开
Uvicorn(FastAPI默认服务器)默认WebSocket空闲超时较短,长时间无数据交互会主动断开连接。即使客户端每2秒发一次数据,网络延迟或服务器配置也可能触发超时。
解决办法:启动Uvicorn时增加ping相关参数,通过心跳帧维持连接:
uvicorn main:app --host 0.0.0.0 --port 8001 --ws-ping-interval 30 --ws-ping-timeout 60
--ws-ping-interval:服务器每隔30秒发送一次ping帧--ws-ping-timeout:60秒内未收到pong响应则断开连接
2. 客户端缺乏异常处理与重连机制
原客户端代码未处理连接断开、接收超时等异常,一旦连接关闭就会停止运行,且没有自动重连逻辑。
解决办法:修改客户端代码,添加异常捕获和重连逻辑:
import asyncio import json import websockets from websockets.exceptions import ConnectionClosedOK, ConnectionClosedError async def client_loop(client_id): uri = f"ws://127.0.0.1:8001/chargeStationState/{client_id}" while True: try: async with websockets.connect(uri) as websocket: while True: message = {'name': f'station{client_id}'} await websocket.send(json.dumps(message)) try: p = await asyncio.wait_for(websocket.recv(), timeout=10) print(f"客户端{client_id}收到: {p}") except asyncio.TimeoutError: print(f"客户端{client_id}接收超时,发送ping维持连接...") await websocket.ping() await asyncio.sleep(2) except (ConnectionClosedOK, ConnectionClosedError): print(f"客户端{client_id}连接已关闭,正在重连...") await asyncio.sleep(5) except Exception as e: print(f"客户端{client_id}出现错误: {e},正在重连...") await asyncio.sleep(5) # 启动5个客户端任务 async def main(): tasks = [client_loop(i) for i in range(1, 6)] await asyncio.gather(*tasks) if __name__ == "__main__": asyncio.run(main())
3. ConnectionManager线程不安全
原ConnectionManager使用普通列表管理连接,多客户端并发连接/断开时会出现竞态条件,导致连接管理混乱,引发异常。
解决办法:给ConnectionManager添加异步锁,确保列表操作线程安全:
from asyncio import Lock class ConnectionManager: def __init__(self): self.active_connections: list[WebSocket] = [] self.lock = Lock() async def connect(self, websocket: WebSocket): await websocket.accept() async with self.lock: self.active_connections.append(websocket) async def disconnect(self, websocket: WebSocket): async with self.lock: try: self.active_connections.remove(websocket) except ValueError: # 连接可能已被移除,忽略该错误 pass async def send_response(self, message: dict, websocket: WebSocket): await websocket.send_json(message) async def broadcast(self, message: str): async with self.lock: # 复制连接列表,避免遍历过程中原列表被修改 connections = self.active_connections.copy() for connection in connections: await connection.send_text(message)
4. 数据库会话泄漏
服务器端通过Depends获取的数据库Session会在整个WebSocket连接周期内保持,当连接数量增加时,可能耗尽数据库连接池,导致服务器处理能力下降,间接引发连接断开。
解决办法:在每次处理请求时获取新的数据库会话,避免长期占用:
# 假设你使用SQLAlchemy异步会话,示例如下 from sqlalchemy.ext.asyncio import AsyncSession, create_async_engine from sqlalchemy.orm import sessionmaker engine = create_async_engine("sqlite+aiosqlite:///./test.db") async_session = sessionmaker(engine, class_=AsyncSession, expire_on_commit=False) @app.websocket("/chargeStationState/{client_id}") async def websocket_endpoint(websocket: WebSocket, client_id: int): await manager.connect(websocket) try: while True: message = await websocket.receive_json() logging.info(message) # 每次处理请求时创建新的会话 async with async_session() as session: # 执行数据库查询操作 # db_data = await session.query(...).first() response = { "stations": "response", "timestamp": int(time.time()) } await manager.send_response(response, websocket) except WebSocketDisconnect: await manager.disconnect(websocket)
内容的提问来源于stack exchange,提问作者Matej Senožetnik
相关产品推荐
相关产品推荐

