第 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 协议,之后双向、全双工、长连接。
| 维度 | HTTP | WebSocket |
|---|---|---|
| 通信方向 | 客户端发起,服务端响应 | 双向,任一方可随时发送 |
| 服务端主动推送 | 不支持(需 SSE / 轮询) | 原生支持 |
| 断线感知 | 无状态,天然无感 | 需心跳或 close 帧检测 |
适合:即时聊天、站内通知、协同编辑、实时行情、长任务进度。不适合:一次性数据查询、纯单向推送(SSE 足够)、对代理与防火墙兼容性要求极高的外部开放接口。
15.2 最小可用端点
@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。
❌ 不推荐(断开时刷满异常堆栈,且连接对象从不清理):
@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 三段式清理):
@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
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 完整示例:回声 + 广播聊天室
@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}"})// 浏览器控制台
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 侧同时把超时调大,并确保透传升级头:
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
@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()
@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)会把这次握手直接转成 HTTP403,浏览器看到的是握手失败而不是关闭码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,同时订阅同一频道,收到后再推给本进程内的连接。
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 的取舍
| 维度 | 短轮询 | 长轮询 | SSE | WebSocket |
|---|---|---|---|---|
| 实时性 | 差(取决于间隔) | 较好 | 好 | 最好 |
| 服务端推送 | 不支持 | 变相支持 | 支持(单向) | 支持(双向) |
| 典型场景 | 状态轮询、低频更新 | 兼容性优先的推送 | 通知、进度条、日志流 | 聊天、协同、双向交互 |
选型顺序建议:能用普通 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 |
练习题
- 实现一个支持按「房间」分组广播的
ConnectionManager(dict[str, set[WebSocket]]),并写一个把用户从一个房间移到另一个房间的端点。 - 复现 15.3 的问题:用两个标签页连接
/ws/bad,关掉其中一个,观察另一个是否还能收到广播、日志里出现了什么;再用/ws/good验证修复。 - 为聊天室加上认证:客户端首帧发送
{"token": "..."},服务端校验通过才accept();用过期 token 测试,记录服务端与浏览器分别看到什么。 - 用
--workers 2启动服务,用 4 个客户端验证广播丢失现象,再接入 Redis 版RedisBroadcaster复测。
下一章预告
到这一章为止,功能已经写得差不多了。但「能跑」和「敢改」之间隔着一套测试。下一章用 pytest + httpx 把接口、依赖覆盖和异步测试补齐。