Skip to content

第 15 章:WebSocket 实时通信 ​

HTTP 是「客户端问、服务端答」的单向模型,轮询能凑合但代价高。聊天、站内通知、协同编辑、长任务进度这些场景需要服务端主动推数据——本章把 FastAPI 的 WebSocket 端点、连接管理、心跳与多进程广播讲完整。


学习目标 ​

  • 理解 WebSocket 与 HTTP 的协议差异,能判断一个需求该不该用长连接
  • 掌握 @app.websocket 端点的完整生命周期:accept / receive_* / send_* / close
  • 能用 ConnectionManager 实现广播,并正确处理断开连接导致的发送异常
  • 掌握应用层心跳与代理层空闲超时的关系,能实现 WebSocket 认证
  • 理解多 worker 下广播失效的根因,并知道用 Redis Pub/Sub 解决

15.1 WebSocket 与 HTTP 的区别 ​

WebSocket 连接由一个普通 HTTP 请求发起:客户端发出带 Upgrade: websocket 头的握手请求,服务端返回 101 Switching Protocols 后,同一条 TCP 连接切换到 WebSocket 协议,之后双向、全双工、长连接。

维度HTTPWebSocket
通信方向客户端发起,服务端响应双向,任一方可随时发送
服务端主动推送不支持(需 SSE / 轮询)原生支持
断线感知无状态,天然无感需心跳或 close 帧检测

适合:即时聊天、站内通知、协同编辑、实时行情、长任务进度。不适合:一次性数据查询、纯单向推送(SSE 足够)、对代理与防火墙兼容性要求极高的外部开放接口。


15.2 最小可用端点 ​

python
@app.websocket("/ws/echo")
async def echo(websocket: WebSocket) -> None:
    await websocket.accept()                      # 完成握手,之后可以收发
    while True:
        await websocket.send_text(f"echo: {await websocket.receive_text()}")
方法作用
await websocket.accept()接受握手;不调用则连接挂起
await websocket.send_json(obj) / receive_json()JSON 编解码的便捷方法
await websocket.close(code=1000)主动关闭,可携带关闭码

用 websocat ws://127.0.0.1:8000/ws/echo 或浏览器控制台验证。WebSocket 端点不会出现在 Swagger UI 里,也不参与 Pydantic 响应校验;但在 APIRouter 上同样可以用 @router.websocket 声明。


15.3 断连处理:必须捕获 WebSocketDisconnect ​

客户端关闭标签页、切网络、进程被杀,都会让服务端 receive_* 抛 WebSocketDisconnect。

❌ 不推荐(断开时刷满异常堆栈,且连接对象从不清理):

python
@app.websocket("/ws/bad")
async def bad(websocket: WebSocket) -> None:
    await websocket.accept()
    manager.active_connections.append(websocket)
    while True:
        # ❌ 客户端一断开就抛 WebSocketDisconnect,堆栈打进日志,
        #    连接也永远留在 active_connections 里
        await manager.broadcast({"text": await websocket.receive_text()})

✅ 推荐(try / except / finally 三段式清理):

python
@app.websocket("/ws/good")
async def good(websocket: WebSocket) -> None:
    await manager.connect(websocket)
    try:
        while True:
            await manager.broadcast({"text": await websocket.receive_text()})
    except WebSocketDisconnect:
        pass                                  # ✅ 正常断开,不视为错误
    finally:
        await manager.disconnect(websocket)   # ✅ 无论如何都清理登记

🔑 finally 是必须的。除 WebSocketDisconnect 外,任何业务异常都会跳出循环;只在 except 里清理会让异常路径上的连接永久泄漏。


15.4 连接管理器 ConnectionManager ​

python
class ConnectionManager:
    def __init__(self) -> None:
        self.active_connections: list[WebSocket] = []

    async def connect(self, websocket: WebSocket) -> None:
        await websocket.accept()
        self.active_connections.append(websocket)

    async def disconnect(self, websocket: WebSocket) -> None:
        if websocket in self.active_connections:
            self.active_connections.remove(websocket)

    async def broadcast(self, payload: dict[str, object]) -> None:
        dead: list[WebSocket] = []
        for connection in list(self.active_connections):   # 浅拷贝,避免遍历中被修改
            try:
                await connection.send_json(payload)
            except Exception:
                dead.append(connection)      # ✅ 只标记,不向外抛
        for connection in dead:
            await self.disconnect(connection)
manager = ConnectionManager()

这是最容易被写错的一段代码。 若 broadcast 写成 for conn in self.active_connections: await conn.send_json(...),只要有一个连接处于「半开」状态(TCP 还没感知到对方已断),send_json 就会抛异常,整个循环中断,后面的客户端全部收不到这条消息,异常还会冒泡到广播发起者,让它跟着断开。


15.5 完整示例:回声 + 广播聊天室 ​

python
@app.websocket("/ws/chat/{room}")
async def chat_room(websocket: WebSocket, room: str, nickname: str = "匿名") -> None:
    await manager.connect(websocket)
    await manager.broadcast({"type": "system", "text": f"{nickname} 加入了 {room}"})
    try:
        while True:
            data = await websocket.receive_json()
            if data.get("type") == "ping":
                await websocket.send_json({"type": "pong"})   # 心跳,不广播
            else:
                await manager.broadcast(
                    {"type": "message", "room": room, "from": nickname, "text": data.get("text", "")}
                )
    except WebSocketDisconnect:
        pass
    finally:
        await manager.disconnect(websocket)
        await manager.broadcast({"type": "system", "text": f"{nickname} 离开了 {room}"})
text
// 浏览器控制台
const ws = new WebSocket("ws://127.0.0.1:8000/ws/chat/general?nickname=moqian");
ws.onopen = () => {
  ws.send(JSON.stringify({ type: "message", text: "大家好" }));
  setInterval(() => ws.send(JSON.stringify({ type: "ping" })), 30000);   // 应用层心跳
};

广播时序:


15.6 心跳与保活 ​

长连接不活跃时会被中间设备回收:Nginx 的 proxy_read_timeout 默认 60 秒无数据即关闭,云负载均衡的空闲超时通常 60 秒到 900 秒,表现都是连接静默消失、前端触发 onclose。

Starlette 已内置协议层 ping/pong(由 ASGI 服务器负责),但协议层 ping 不产生应用数据,很多代理只看「是否有数据帧」判断空闲。因此生产环境通常再加一层应用层心跳:客户端每 30 秒发 {"type": "ping"},服务端立刻回 pong(不广播、不写日志),客户端连续 2~3 次没收到就重连。Nginx 侧同时把超时调大,并确保透传升级头:

text
proxy_http_version 1.1;
proxy_set_header Upgrade $http_upgrade;
proxy_set_header Connection "upgrade";
proxy_read_timeout 3600s;

15.7 认证:浏览器拿不到自定义 Header ​

浏览器的 WebSocket 构造函数不支持自定义请求头,所以不能像 HTTP 那样带 Authorization: Bearer ...。两种可行方案:

方案 A:query 参数带 token

python
@app.websocket("/ws/secure")
async def secure_ws(websocket: WebSocket, token: str = Query(...)) -> None:
    try:
        user_id = decode_access_token(token)
    except ValueError:
        await websocket.close(code=1008)      # 1008 = Policy Violation
        return

方案 B:首帧发 token,校验通过后再 accept()

python
@app.websocket("/ws/auth")
async def authenticated_ws(websocket: WebSocket) -> None:
    try:
        raw = await asyncio.wait_for(websocket.receive_text(), timeout=5.0)
        user_id = decode_access_token(json.loads(raw).get("token", ""))
    except (asyncio.TimeoutError, ValueError, json.JSONDecodeError):
        await websocket.close(code=1008)      # ✅ 校验失败,先拒绝再返回
        return
    await websocket.accept()                  # ✅ 校验通过后才接受连接
    try:
        while True:
            await handle(websocket, user_id, await websocket.receive_text())
    except WebSocketDisconnect:
        pass
    finally:
        await manager.disconnect(websocket)

⚠️ 握手未 accept() 时调用 close(),多数 ASGI 服务器(含 Uvicorn)会把这次握手直接转成 HTTP 403,浏览器看到的是握手失败而不是关闭码 1008;方案 A 会把 token 写进 URL 日志与 Referer,尽量换成短期一次性 ticket。绝不要先 accept() 再校验,那样未认证的客户端已经进入连接池、可能被广播到。


15.8 多进程部署:广播为什么丢了 ​

ConnectionManager 是进程内内存状态。uvicorn --workers 4 之后:客户端 A 连 worker 1、B 连 worker 2,A 发消息时 worker 1 只遍历自己进程里的连接,B 收不到任何消息——而且不会有任何报错。

解决思路是把「广播」从进程内提升为跨进程消息总线:每个 worker 把消息 publish 到 Redis,同时订阅同一频道,收到后再推给本进程内的连接。

python
class RedisBroadcaster:
    def __init__(self, url: str, channel: str = "ws:broadcast") -> None:
        self._client = redis.from_url(url, decode_responses=True)   # redis.asyncio
        self._channel = channel
        self._task: asyncio.Task[None] | None = None

    async def start(self, manager: ConnectionManager) -> None:
        self._pubsub = self._client.pubsub()
        await self._pubsub.subscribe(self._channel)
        self._task = asyncio.create_task(self._listen(manager))   # 后台消费

    async def _listen(self, manager: ConnectionManager) -> None:
        async for message in self._pubsub.listen():
            if message["type"] == "message":
                await manager.broadcast(json.loads(message["data"]))

    async def publish(self, payload: dict[str, object]) -> None:
        # 业务侧只 publish,各 worker(含自己)订阅后统一广播给本地连接
        await self._client.publish(self._channel, json.dumps(payload))

stop() 里 cancel 掉 _task 并 aclose() pubsub 与 client,在 lifespan 中成对调用 start / stop;业务处理器把 manager.broadcast(...) 换成 broadcaster.publish(...)。同样模式可用 NATS、RabbitMQ fanout 或 Kafka 实现。


15.9 与轮询、SSE 的取舍 ​

维度短轮询长轮询SSEWebSocket
实时性差(取决于间隔)较好好最好
服务端推送不支持变相支持支持(单向)支持(双向)
典型场景状态轮询、低频更新兼容性优先的推送通知、进度条、日志流聊天、协同、双向交互

选型顺序建议:能用普通 HTTP 就别上长连接 → 单向推送优先 SSE → 需要双向交互才上 WebSocket。WebSocket 的运维成本(心跳、断线重连、跨进程广播、压测)明显更高,不要因为「听起来更实时」就选它。


常见坑与排查 ​

现象原因解决
日志里刷满 WebSocketDisconnect 异常堆栈未捕获该异常用 try / except WebSocketDisconnect / finally 包住收发循环
一个客户端断开后,其他客户端也收不到广播broadcast 中单个 send 抛异常中断了整个循环每个连接单独 try,失败的收集到 dead 列表后统一清理
连接数持续增长,内存缓慢上涨异常路径未从 active_connections 移除清理逻辑放 finally,而不是只放 except WebSocketDisconnect
浏览器报握手失败,Nginx 返回 400未透传 Upgrade / Connection 头,或未用 HTTP/1.1加 proxy_http_version 1.1、Upgrade $http_upgrade、Connection "upgrade"
多 worker 部署后广播只能到达部分客户端ConnectionManager 是进程内状态引入 Redis Pub/Sub(或 NATS / MQ)做跨进程广播
客户端空闲几分钟后连接被静默关闭代理空闲超时(Nginx 默认 60 s)应用层心跳 + 调大 proxy_read_timeout
WebSocket 处理器里 raise HTTPException 无效WebSocket 不是请求—响应模型用 await websocket.close(code=1008),或先 send_json 再 close

本章小结 ​

要点说明
断连处理必须捕获 WebSocketDisconnect,清理逻辑放 finally
广播要点逐连接 try + 收集失败连接 + 统一清理,别让一个死连接中断全局广播
心跳代理空闲超时会掐断连接,用 ping / pong 加 Nginx 超时配置兜底
认证浏览器不能带自定义 Header;用 query token 或首帧 token 后 accept(),失败 close(1008)
多进程连接管理器不跨进程,生产需 Redis Pub/Sub 做跨进程广播
选型单向推送优先 SSE,只有需要双向交互才上 WebSocket

练习题 ​

  1. 实现一个支持按「房间」分组广播的 ConnectionManager(dict[str, set[WebSocket]]),并写一个把用户从一个房间移到另一个房间的端点。
  2. 复现 15.3 的问题:用两个标签页连接 /ws/bad,关掉其中一个,观察另一个是否还能收到广播、日志里出现了什么;再用 /ws/good 验证修复。
  3. 为聊天室加上认证:客户端首帧发送 {"token": "..."},服务端校验通过才 accept();用过期 token 测试,记录服务端与浏览器分别看到什么。
  4. 用 --workers 2 启动服务,用 4 个客户端验证广播丢失现象,再接入 Redis 版 RedisBroadcaster 复测。

下一章预告 ​

到这一章为止,功能已经写得差不多了。但「能跑」和「敢改」之间隔着一套测试。下一章用 pytest + httpx 把接口、依赖覆盖和异步测试补齐。

👉 第 16 章:测试体系(pytest + httpx)

📖本文阅读--次|📊全站访问--次|👥访客--人