下载工作台
FastAPI 开发

WebSocket 与 SSE 实时通信

试读上半部分 · 解锁后可读全文

第 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运营看板、在线用户计数
轮询客户端反复 GETHTTP仅作降级备选

原则命令走 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 日志
首条 authURL 干净需约定协议
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 禁用缓冲
        },
    )

以下内容需解锁后阅读

试读已结束。解锁本章 ¥5.00,或开通年度会员畅读全部教程。
年度会员 ¥199.00/年; 小紫 AI 工作台有效会员 ¥99.00/年

正文仅在服务端鉴权后下发,未付费无法获取下半部分内容。