You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

FastAPI WebSocket服务器ASGI应用异常及连接断开问题排查

问题分析与解决方案

核心问题

  1. 服务器未捕获WebSocketDisconnect异常:当前服务器仅捕获ConnectionClosedError,但Starlette抛出的WebSocketDisconnect(如错误码1000正常断开)未被处理,导致ASGI应用抛出未捕获异常,连接直接断开。
  2. 客户端接收逻辑仅执行一次:receive_messages函数只接收一条消息就结束协程,asyncio.gather完成后触发async with块退出,主动关闭连接,引发服务器端断开异常。
  3. 同步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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.22 11:52:32