FastAPI WebSocket服务器ASGI应用异常及连接断开问题排查
问题分析与解决方案
核心问题
- 服务器未捕获
WebSocketDisconnect异常:当前服务器仅捕获ConnectionClosedError,但Starlette抛出的WebSocketDisconnect(如错误码1000正常断开)未被处理,导致ASGI应用抛出未捕获异常,连接直接断开。 - 客户端接收逻辑仅执行一次:
receive_messages函数只接收一条消息就结束协程,asyncio.gather完成后触发async with块退出,主动关闭连接,引发服务器端断开异常。 - 同步
input()阻塞事件循环:客户端使用同步input()获取用户输入,会阻塞整个asyncio事件循环,导致WebSocket心跳(ping/pong)无法及时处理,引发1006(异常断开)、1012(服务重启)类超时错误。
修复步骤
1. 服务器端代码修复
- 导入
WebSocketDisconnect并添加到异常捕获列表 - 维护在线客户端列表,实现真正的消息广播
- 优化异常处理逻辑,避免未捕获异常导致ASGI报错
2. 客户端代码修复
- 将
receive_messages改为循环接收,保持连接活跃 - 使用异步输入替代同步
input(),避免阻塞事件循环 - 优化重连逻辑,覆盖所有WebSocket断开场景
修正后的代码
服务器端(server.py)
from fastapi import FastAPI, WebSocket from fastapi.staticfiles import StaticFiles from websockets.exceptions import ConnectionClosedError from starlette.websockets import WebSocketDisconnect # 新增导入 import uvicorn from pathlib import Path app = FastAPI() current_file = Path(__file__) static_root_absolute = current_file.parent.resolve() app.mount("/static", StaticFiles(directory=static_root_absolute / 'static'), name="static") # 维护所有在线客户端连接 connected_clients = [] async def server_handle_message(message, ws): if message == "messageA": await ws.send_text("messageA") print("handle message A") elif message == "messageB": await ws.send_text("messageB") print("handle message B") @app.websocket("/ws") async def websocket_endpoint(websocket: WebSocket): await websocket.accept() connected_clients.append(websocket) print(f"新客户端连接,当前在线:{len(connected_clients)}") try: while True: message = await websocket.receive_text() print(f"收到客户端消息:{message}") await server_handle_message(message, websocket) # 广播消息给其他在线客户端 for client in connected_clients: if client != websocket: await client.send_text(f"[{len(connected_clients)}人在线] 广播消息:{message}") print("已完成消息广播") except (ConnectionClosedError, WebSocketDisconnect) as e: print(f"连接断开,错误码:{e.code}") if websocket in connected_clients: connected_clients.remove(websocket) print(f"客户端已移除,当前在线:{len(connected_clients)}") if __name__ == "__main__": uvicorn.run( app, host="0.0.0.0", port=8000, ws_ping_interval=30, # 缩短ping间隔,保持连接活跃 ws_ping_timeout=10 )
客户端(client.py)
import asyncio import websockets import sys def handle_message(message): if message == "messageA": print("client received msg A") elif message == "messageB": print("client received msg B") else: print(f"client received msg: {message}") async def receive_messages(ws): # 循环持续接收服务器消息 while True: try: message = await ws.recv() print(f"Received message: {message}") handle_message(message) print("message was parsed") except websockets.exceptions.ConnectionClosed: break async def async_input(prompt): # 异步输入,避免阻塞事件循环 loop = asyncio.get_event_loop() return await loop.run_in_executor(None, input, prompt) async def send_messages(ws): while True: try: message = await async_input("Enter message: ") if message == "some text A": print("text A") await ws.send("messageA") elif message == "some text B": print("text B") await ws.send("messageB") else: await ws.send(message) except websockets.exceptions.ConnectionClosed: break async def main(): async with websockets.connect( "ws://localhost:8000/ws", ping_interval=25, # 客户端ping间隔略短于服务器超时时间 ping_timeout=15 ) as websocket: await asyncio.gather(receive_messages(websocket), send_messages(websocket)) if __name__ == "__main__": while True: try: asyncio.run(main()) except (websockets.exceptions.ConnectionClosed, websockets.ConnectionClosed): print("连接已断开,3秒后重连...") await asyncio.sleep(3) except asyncio.exceptions.TimeoutError: print("连接超时,3秒后重连...") await asyncio.sleep(3) except KeyboardInterrupt: print("用户中断程序") sys.exit(0)
关键说明
- 服务器端新增
connected_clients列表维护在线连接,实现真正的跨客户端消息广播。 - 客户端使用
async_input替代同步input(),彻底解决事件循环阻塞问题,避免心跳超时。 - 调整ping参数为更合理的值(服务器ping间隔30s,超时10s;客户端ping间隔25s,超时15s),确保连接稳定。
- 异常捕获覆盖所有WebSocket断开场景,避免未处理异常导致的ASGI错误。
内容的提问来源于stack exchange,提问作者K H Tan
相关产品推荐
相关产品推荐

