第 15 章 · WebSocket 与 SSE 实时通信
本章目标:在 svc-demo 中实现 WebSocket 长连接与 SSE(Server-Sent Events) 单向推送;编写 ConnectionManager 管理多客户端连接、房间与广播;完成实时通知(商品上下架、系统公告)演示;理解与 REST 混用的职责划分、鉴权方式与部署注意事项(Nginx 代理、超时)。
学时建议:5~6 小时(含 2 小时浏览器与 curl 联调)
前置:本模块 ch01~ch14;重点复习 ch07 JWT、ch08 中间件、ch12 svc-demo 路由结构。
15.1 何时用 WebSocket、SSE 还是 REST
实时性需求
│
┌───────────────┼───────────────┐
▼ ▼ ▼
低频轮询 SSE 单向 WebSocket 双向
GET /notifications 服务端推送 聊天、协作编辑
(简单但浪费) 通知、行情 游戏、遥控
│ │ │
└───────────────┴───────────────┘
│
REST 仍负责 CRUD 与命令
(创建订单、改库存)
| 方式 | 方向 | 协议 | svc-demo 场景 |
|---|---|---|---|
| REST | 请求-响应 | HTTP | 商品 CRUD、登录 |
| SSE | 服务端→客户端 | HTTP 长连接 | 通知流、上下架事件 |
| WebSocket | 双向 | WS | 运营看板、在线用户计数 |
| 轮询 | 客户端反复 GET | HTTP | 仅作降级备选 |
原则:命令走 REST(有明确 HTTP 状态码与幂等);事件走 WS/SSE(推送状态变化)。
示例项目:svc-demo;前端:user-demo;禁止泄露内部消息中间件或生产域名。
15.2 WebSocket 基础端点
svc_demo/api/v1/ws.py:
from fastapi import APIRouter, WebSocket, WebSocketDisconnect
router = APIRouter(tags=["realtime"])
@router.websocket("/ws/echo")
async def ws_echo(websocket: WebSocket):
await websocket.accept()
try:
while True:
data = await websocket.receive_text()
await websocket.send_text(f"echo: {data}")
except WebSocketDisconnect:
pass
挂载:
# main.py 或 v1/router.py
from svc_demo.api.v1.ws import router as ws_router
app.include_router(ws_router)
测试(websocat 或浏览器控制台):
websocat ws://127.0.0.1:8000/ws/echo
# 输入 hello → echo: hello
15.3 ConnectionManager 设计
svc_demo/realtime/manager.py:
from collections import defaultdict
from fastapi import WebSocket
class ConnectionManager:
"""管理 WebSocket 连接:按 room 分组广播。"""
def __init__(self) -> None:
self._rooms: dict[str, set[WebSocket]] = defaultdict(set)
async def connect(self, websocket: WebSocket, room: str = "global") -> None:
await websocket.accept()
self._rooms[room].add(websocket)
def disconnect(self, websocket: WebSocket, room: str = "global") -> None:
self._rooms[room].discard(websocket)
if not self._rooms[room]:
del self._rooms[room]
async def send_personal(self, message: str, websocket: WebSocket) -> None:
await websocket.send_text(message)
async def broadcast(self, message: str, room: str = "global") -> None:
dead: list[WebSocket] = []
for ws in self._rooms.get(room, set()):
try:
await ws.send_text(message)
except Exception:
dead.append(ws)
for ws in dead:
self.disconnect(ws, room)
def count(self, room: str = "global") -> int:
return len(self._rooms.get(room, set()))
manager = ConnectionManager()
| 方法 | 用途 |
|---|---|
connect / disconnect | 入退房间 |
broadcast | 同房间全员推送 |
send_personal | 单聊 |
count | 在线数统计 |
生产环境可扩展:用户 id 索引、心跳 ping、最大连接数。
15.4 带房间的 WebSocket 通知
import json
from fastapi import WebSocket, WebSocketDisconnect, Depends
from svc_demo.realtime.manager import manager
@router.websocket("/ws/notifications")
async def ws_notifications(websocket: WebSocket, room: str = "products"):
await manager.connect(websocket, room)
await manager.send_personal(
json.dumps({"type": "connected", "room": room}),
websocket,
)
try:
while True:
# 客户端可发 ping 或订阅指令(选修)
data = await websocket.receive_text()
if data == "ping":
await manager.send_personal(
json.dumps({"type": "pong"}), websocket
)
except WebSocketDisconnect:
manager.disconnect(websocket, room)
REST 触发广播(商品上架后):
# products.py 上架成功后
import json
from svc_demo.realtime.manager import manager
event = {
"type": "product.published",
"data": {"id": product.id, "name": product.name},
}
await manager.broadcast(json.dumps(event), room="products")
15.5 WebSocket 鉴权
WebSocket 不能像普通路由一样简单用 Depends(get_current_user) 依赖 Header(握手阶段 Query/Header 传 token)。
方式 A:Query 参数 token
from svc_demo.core.security import decode_access_token
@router.websocket("/ws/notifications")
async def ws_notifications(
websocket: WebSocket,
token: str = "",
room: str = "products",
):
payload = decode_access_token(token)
if not payload:
await websocket.close(code=1008)
return
user_id = int(payload["sub"])
await manager.connect(websocket, room)
...
客户端:ws://127.0.0.1:8000/ws/notifications?token=eyJ...
方式 B:首条消息鉴权
连接后 5 秒内必须发 {"type":"auth","token":"..."},否则断开。
| 方式 | 优点 | 缺点 |
|---|---|---|
| Query token | 实现简单 | token 可能进 URL 日志 |
| 首条 auth | URL 干净 | 需约定协议 |
| Cookie | 同源 Web 方便 | 跨域 SPA 麻烦 |
教学推荐:Query + 短过期 access token;生产配合 WSS。
15.6 SSE(Server-Sent Events)
SSE 适合服务端单向推送:通知列表、进度条、行情。
svc_demo/api/v1/sse.py:
import asyncio
import json
from fastapi import APIRouter
from fastapi.responses import StreamingResponse
router = APIRouter(prefix="/sse", tags=["realtime"])
async def notification_stream():
"""模拟每 3 秒推送一条通知(教学用)。"""
n = 0
while n < 10:
n += 1
payload = {"type": "heartbeat", "seq": n, "message": f"通知 #{n}"}
yield f"data: {json.dumps(payload, ensure_ascii=False)}\n\n"
await asyncio.sleep(3)
@router.get("/notifications")
async def sse_notifications():
return StreamingResponse(
notification_stream(),
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache",
"Connection": "keep-alive",
"X-Accel-Buffering": "no", # Nginx 禁用缓冲
},
)