Python 基础体系 · 第 90/112 篇。示例统一以 Python 3.14 为语言基线;第三方库使用与其兼容的现代稳定版本,版本敏感行为会单独说明。
Python 实时通信:WebSocket、SSE、心跳、广播和背压
实时通信不是“把一个接口改成长连接”这么简单。一个可维护的实时服务至少要回答五个问题:
- 客户端和服务器是否需要双向通信?
- 连接断开后,双方如何发现?
- 一个事件如何发送给多个连接?
- 某个客户端读取很慢时,是否会拖住其他客户端?
- 应用关闭、进程重启或网络故障时,连接和任务如何收尾?
WebSocket、SSE、心跳、广播和背压分别解决这些问题中的不同部分。理解它们之前,需要先建立 ASGI 和 asyncio 的连接模型。
一、实时通信的共同基础:长连接与异步事件
传统 HTTP 请求通常经历如下过程:
客户端发送请求
│
▼
服务器处理请求
│
▼
服务器发送完整响应
│
▼
连接结束,或交给 HTTP 持久连接复用
实时通信则需要让服务器在请求之后继续发送数据:
客户端建立连接
│
▼
服务器保持连接打开
│
├── 客户端发送数据
├── 服务器发送事件
├── 定期发送心跳
└── 任意一方关闭连接
ASGI 将应用抽象成一个异步调用:
async def application(scope, receive, send):
...
其中:
scope描述连接类型、路径、请求头等元数据;receive()等待来自客户端或服务器的协议事件;send(message)向客户端或服务器发送协议事件。
ASGI 的关键变化是:应用不再只是“接收一次请求并返回一次响应”,而是通过事件持续参与连接生命周期。HTTP 请求通常对应一个请求作用域,而 WebSocket 作用域会持续到 WebSocket 连接关闭。(asgi.readthedocs.io)
FastAPI 建立在 Starlette 和 ASGI 之上,因此 WebSocket 路由、连接断开异常和流式响应最终都要落到 ASGI 的事件模型上。FastAPI 官方提供的 WebSocket 对象也是对底层能力的封装。(fastapi.tiangolo.com)
二、WebSocket:真正的双向长连接
2.1 WebSocket 解决什么问题
WebSocket 是一种全双工通信协议:
- 客户端可以发送消息;
- 服务器也可以主动发送消息;
- 双方不需要为每条消息重新建立 HTTP 请求;
- 一条连接可以承载多个文本消息或二进制消息。
“全双工”意味着发送方向和接收方向相互独立。服务器不能只在收到客户端消息后才发送数据,否则它只是一个“请求驱动的回声服务”,还没有发挥 WebSocket 的价值。
WebSocket 连接最初仍然通过 HTTP 请求发起握手,但握手完成后,后续数据使用 WebSocket 帧传输。ASGI 会把协议服务器解析出的完整消息交给应用,而不是要求应用自己处理 TCP 分片或 WebSocket 帧分片。ASGI 规范还规定,协议服务器应处理 WebSocket 的 PING/PONG 和消息分片。(asgi.readthedocs.io)
2.2 WebSocket 的 ASGI 状态转换
一个典型连接可以表示为:
stateDiagram-v2
[*] --> Handshake
Handshake --> Connected: websocket.connect\n+ websocket.accept
Handshake --> Closed: websocket.close
Connected --> Connected: websocket.receive
Connected --> Connected: websocket.send
Connected --> Closed: websocket.disconnect
Connected --> Closed: websocket.close
Closed --> [*]
关键路径是:
- ASGI 服务器向应用提供
websocket.connect; - 应用发送
websocket.accept,握手完成; - 客户端消息以
websocket.receive到达; - 应用消息以
websocket.send发出; - 任意一方关闭后,应用收到
websocket.disconnect。
在接受连接之前,应用必须在 websocket.accept 和 websocket.close 之间作出选择。若握手阶段直接关闭,协议服务器应拒绝这次升级,通常表现为 HTTP 403。(asgi.readthedocs.io)
2.3 最小 WebSocket 示例
# websocket_minimal.py
from fastapi import FastAPI, WebSocket, WebSocketDisconnect
app = FastAPI()
@app.websocket("/ws")
async def websocket_endpoint(websocket: WebSocket) -> None:
await websocket.accept()
try:
while True:
message = await websocket.receive_text()
await websocket.send_text(f"echo: {message}")
except WebSocketDisconnect:
print("client disconnected")
运行:
python -m pip install fastapi uvicorn
uvicorn websocket_minimal:app --reload
浏览器控制台中可以测试:
const ws = new WebSocket("ws://127.0.0.1:8000/ws");
ws.onopen = () => ws.send("hello");
ws.onmessage = (event) => console.log(event.data);
预期输出:
echo: hello
receive_text() 是一个等待点。当客户端没有消息时,当前协程会挂起,事件循环可以运行其他连接的任务。客户端关闭连接后,接收操作会抛出 WebSocketDisconnect,应用应捕获它并执行清理。FastAPI 官方文档也采用这种方式处理多客户端断开。(fastapi.tiangolo.com)
2.4 为什么不能只写一个发送循环
下面的代码只能实现“收到消息后回复”,不能实现服务器主动推送:
while True:
message = await websocket.receive_text()
await websocket.send_text("reply")
如果服务器有一个外部事件源:
async def publish_event(event: str) -> None:
...
那么 websocket_endpoint() 如果一直阻塞在 receive_text(),就无法及时执行发送逻辑。双向通信通常需要两个并发任务:
接收任务:客户端 ──────> 应用
发送任务:应用 ──────> 客户端
但是,发送任务不能随意并发调用同一个 WebSocket 的发送方法。很多协议实现要求同一连接上的消息按顺序发送;即使底层库允许并发调用,也难以保证业务顺序。因此更可靠的结构是:
业务事件
│
▼
每个连接一个有界队列
│
▼
每个连接一个唯一发送任务
│
▼
WebSocket.send_*
这样,业务生产者只负责放入事件,连接发送任务负责串行写出。
三、SSE:单向的 HTTP 事件流
3.1 SSE 与 WebSocket 的本质区别
SSE,即 Server-Sent Events,是建立在 HTTP 响应之上的服务器到客户端事件流。
它的方向是:
服务器 ───────────────> 客户端
客户端不能通过同一条 SSE 连接向服务器发送业务消息。客户端如果需要提交数据,仍然使用普通 HTTP 请求,例如 POST /messages。
因此:
| 能力 | WebSocket | SSE |
|---|---|---|
| 服务器推送 | 支持 | 支持 |
| 客户端通过同一连接发送 | 支持 | 不支持 |
| 数据基础 | WebSocket 消息 | HTTP 响应文本流 |
| 浏览器原生 API | WebSocket |
EventSource |
| 适合方向 | 双向交互 | 服务端事件、通知、日志 |
| 代理兼容性 | 需要支持 WebSocket 升级 | 通常按 HTTP 流处理 |
SSE 响应通常使用:
Content-Type: text/event-stream
Cache-Control: no-cache
Connection: keep-alive
每个事件由若干字段组成,事件结束使用一个空行表示:
event: message
id: 42
data: {"text":"hello"}
其中:
event是事件名称,可选;id是事件 ID,可选;data是事件数据,可以出现多次;- 空行表示一个事件结束。
SSE 的心跳通常不是业务事件,而是注释行:
: heartbeat
浏览器会忽略注释,但它能让中间代理和网络设备看到连接仍在产生数据。
3.2 使用 StreamingResponse 实现 SSE
FastAPI 的 StreamingResponse 可以接收异步生成器,并逐段发送响应体。异步生成器必须包含可取消的 await,否则长时间运行的流可能无法及时响应任务取消。(fastapi.tiangolo.com)
# sse_minimal.py
import asyncio
import json
from collections.abc import AsyncIterator
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
app = FastAPI()
async def event_stream() -> AsyncIterator[str]:
for number in range(3):
payload = json.dumps({"number": number}, ensure_ascii=False)
yield f"event: number\ndata: {payload}\n\n"
await asyncio.sleep(1)
yield "event: done\ndata: {}\n\n"
@app.get("/events")
async def events() -> StreamingResponse:
return StreamingResponse(
event_stream(),
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache",
"Connection": "keep-alive",
"X-Accel-Buffering": "no",
},
)
浏览器端:
const source = new EventSource("/events");
source.addEventListener("number", (event) => {
console.log(JSON.parse(event.data));
});
source.addEventListener("done", () => {
source.close();
});
X-Accel-Buffering: no 是针对某些 Nginx 配置的常见部署措施,不是 SSE 协议本身的要求。真正需要验证的是:从客户端视角看,事件是否逐条到达,而不是等响应结束后一次性出现。
ASGI 对流式 HTTP 响应的抽象是多次发送:
await send({
"type": "http.response.start",
"status": 200,
"headers": [...],
})
await send({
"type": "http.response.body",
"body": b"data: hello\n\n",
"more_body": True,
})
await send({
"type": "http.response.body",
"body": b"",
"more_body": False,
})
more_body=True 表示后面仍有内容;False 表示响应结束。ASGI 还要求服务器在处理响应体发送时刷新相应数据,因此流式响应可以逐段向客户端传输。(asgi.readthedocs.io)
3.3 SSE 的重连与事件 ID
浏览器的 EventSource 在连接断开后通常会尝试重连。客户端可能在重连请求中发送:
Last-Event-ID: 42
服务器可以据此从事件日志中继续发送 ID 大于 42 的事件。
但这里有一个重要边界:SSE 的自动重连不等于消息可靠投递。
如果服务器只在内存中保存最新事件:
latest_event = event
那么客户端断开期间错过的事件无法恢复。要提供恢复能力,至少需要:
- 事件具有单调递增或可比较的 ID;
- 服务器保留一段事件历史;
- 客户端发送
Last-Event-ID; - 服务器能够根据 ID 查询补发范围;
- 历史过期时返回明确的“需要全量同步”结果。
如果业务要求严格的消费确认、重放和持久化,SSE 本身不是消息队列,应结合数据库或消息系统设计。
四、心跳:发现失效连接,而不是证明业务可用
4.1 三种容易混淆的“心跳”
实时系统中至少有三种不同机制:
WebSocket 协议级 PING/PONG
WebSocket 协议支持 PING/PONG。ASGI 规范要求协议服务器处理 PING/PONG,并在需要时发送 PING 以确认连接存活。应用通常不需要自己拼接 WebSocket 控制帧。(asgi.readthedocs.io)
应用级心跳
应用自己发送一条普通消息:
{"type":"heartbeat","ts":1730000000}
客户端回复:
{"type":"heartbeat_ack","ts":1730000000}
这可以确认:
- 应用任务仍在运行;
- 客户端确实收到了业务层消息;
- 客户端的业务代码仍能处理消息。
SSE 注释心跳
SSE 没有 WebSocket 那样的双向控制帧,服务器通常发送:
: ping
这只能证明服务器仍在写 HTTP 响应流,不能证明客户端业务逻辑已经处理了它。
4.2 心跳检测的形式化条件
设:
t_now为服务器当前时间;t_last为最近一次收到客户端有效消息或确认的时间;T_timeout为超时时间。
连接被判定为失活的条件可以写成:
如果检查任务每隔 T_check 秒运行一次,那么实际断开时间并不是精确的 T_timeout,而是大致落在:
这说明:
- 检查间隔越小,失效发现越及时;
- 检查间隔越小,定时器和调度开销越高;
- 超时时间不能小于正常网络抖动、代理空闲时间和客户端调度延迟的总和。
心跳间隔也不应直接等于超时时间。例如:
heartbeat_interval = 20 秒
timeout = 60 秒
意味着客户端最多可能错过多次心跳后才被清理。
4.3 心跳的失败路径
一个可靠的连接任务至少要处理以下情况:
发送心跳
│
├── 发送成功
│ └── 等待下一次心跳
│
├── 发送抛出连接错误
│ └── 移除连接并取消相关任务
│
└── 客户端未确认
└── 超时关闭连接
不要把“发送函数没有立即抛异常”理解为“客户端在线”。网络内核缓冲区可能暂时接受数据,但客户端已经不可用。只有协议层或应用层的确认,才提供更强的存活证据。
五、广播:从一个事件复制到多个连接
5.1 广播的定义
广播是将一个事件发送给一组订阅者:
event
│
├──> client A
├──> client B
├──> client C
└──> client D
最简单的实现是:
for websocket in clients:
await websocket.send_json(event)
这段代码存在一个严重问题:发送是串行的。如果客户端 B 的网络很慢,服务器会在 B 的 await 上停留,客户端 C 和 D 也无法及时收到事件。
这不是广播本身的问题,而是把所有客户端的发送进度绑定到了同一个协程。
5.2 每个客户端独立发送队列
更合理的结构是:
┌──> queue A ──> sender A ──> client A
event producer ──┼──> queue B ──> sender B ──> client B
├──> queue C ──> sender C ──> client C
└──> queue D ──> sender D ──> client D
广播过程只做入队:
for client in clients:
client.queue.put_nowait(event)
每个连接自己的发送任务再从队列取出数据:
while True:
event = await client.queue.get()
await websocket.send_json(event)
这样一个慢客户端最多拖慢它自己的发送任务,不会直接阻塞其他连接。
5.3 连接管理中的竞态
连接集合并不只是一个 set[WebSocket]。连接的创建和销毁可能与广播同时发生:
任务 A:正在广播
任务 B:客户端刚好断开
任务 C:清理任务正在删除连接
如果广播直接遍历一个会被修改的集合,可能出现:
RuntimeError: Set changed size during iteration
一种简单处理方式是复制快照:
for client in tuple(self._clients):
...
更重要的是,清理操作必须幂等:
self._clients.discard(client)
discard() 在元素不存在时不会报错,因此重复清理不会制造新的异常。
六、背压:当生产速度超过消费速度
6.1 背压的定义
背压(backpressure)是消费者通过阻塞、限速或拒绝的方式,反向约束生产者速度。
设:
λ:事件生产速率,单位为事件/秒;μ:某个客户端的消费速率;Q(t):时刻t的待发送队列长度。
当:
队列通常不会持续增长。
当:
队列长度会近似以:
的速度增长。
例如:
- 生产者每秒产生 100 个事件;
- 慢客户端每秒只能处理 20 个事件;
- 初始队列为 0;
- 队列上限为 1000。
那么每秒净增长:
队列大约在:
后达到上限。
如果队列是无界的,问题会从“延迟增加”演变为“内存持续增长”。如果每条消息平均占用 4 KiB,100 万条排队消息就可能占用约 3.8 GiB,还未计算 Python 对象本身的额外开销。
6.2 asyncio.Queue 的背压
有界队列可以表达明确容量:
queue: asyncio.Queue[dict] = asyncio.Queue(maxsize=100)
生产者有两种典型写法。
阻塞生产者
await queue.put(event)
队列满时,生产者等待消费者取走数据。这是严格背压:生产者速度会被最慢的消费者限制。
非阻塞丢弃
try:
queue.put_nowait(event)
except asyncio.QueueFull:
# 丢弃、合并,或断开这个客户端
...
这不是“没有背压”,而是选择了另一种背压策略:队列达到上限后拒绝更多数据。
6.3 广播中的背压策略
广播系统不能脱离业务语义讨论“队列满了怎么办”。常见策略有四类:
阻塞整个广播
await asyncio.gather(
*(client.queue.put(event) for client in clients)
)
优点是尽量不丢消息,缺点是一个慢客户端可能阻塞整个发布路径。
丢弃最新消息
队列满时不放入新消息。适用于状态更新,例如鼠标位置、实时指标、进度百分比,因为新状态会很快覆盖旧状态。
丢弃最旧消息
先移除旧消息,再加入最新消息:
if queue.full():
queue.get_nowait()
queue.put_nowait(event)
适用于“只关心最近状态”的场景,但要注意并发下需要更严谨的封装,不能假设 full() 检查和 get_nowait() 之间没有其他任务操作队列。
断开慢客户端
当客户端持续积压时关闭连接,要求它重新连接或执行全量同步:
队列满
│
├── 短暂偶发:记录并继续
└── 持续发生:关闭连接,避免无限占用资源
对于聊天消息、订单状态、权限变更等不能静默丢失的事件,通常不能简单丢弃,而应使用持久化事件日志和客户端游标恢复。
七、一个可运行的 FastAPI 广播示例
下面的示例同时提供:
- WebSocket
/ws:双向通信; - SSE
/events:服务端单向事件流; - HTTP
POST /publish:产生广播事件; - 每个连接一个有界队列;
- 慢客户端超过队列容量后被移除;
- 应用关闭时取消后台任务。
# main.py
from __future__ import annotations
import asyncio
import contextlib
import json
import time
from collections.abc import AsyncIterator
from dataclasses import dataclass, field
from typing import Any
from fastapi import FastAPI, Request, WebSocket, WebSocketDisconnect
from fastapi.responses import StreamingResponse
from pydantic import BaseModel
@dataclass(eq=False)
class Client:
queue: asyncio.Queue[dict[str, Any]] = field(
default_factory=lambda: asyncio.Queue(maxsize=100)
)
class PublishRequest(BaseModel):
text: str
class Hub:
def __init__(self) -> None:
self._clients: set[Client] = set()
self._lock = asyncio.Lock()
self._sequence = 0
async def add(self, client: Client) -> None:
async with self._lock:
self._clients.add(client)
async def remove(self, client: Client) -> None:
async with self._lock:
self._clients.discard(client)
async def publish(self, data: dict[str, Any]) -> int:
async with self._lock:
clients = tuple(self._clients)
dropped = 0
for client in clients:
try:
client.queue.put_nowait(data)
except asyncio.QueueFull:
dropped += 1
await self.remove(client)
return dropped
async def next_event(self, text: str) -> dict[str, Any]:
async with self._lock:
self._sequence += 1
sequence = self._sequence
return {
"id": sequence,
"type": "message",
"ts": time.time(),
"text": text,
}
hub = Hub()
app = FastAPI()
@app.post("/publish")
async def publish_message(body: PublishRequest) -> dict[str, int]:
event = await hub.next_event(body.text)
dropped = await hub.publish(event)
return {"event_id": event["id"], "dropped_clients": dropped}
async def websocket_sender(
websocket: WebSocket,
client: Client,
) -> None:
while True:
event = await client.queue.get()
await websocket.send_json(event)
@app.websocket("/ws")
async def websocket_endpoint(websocket: WebSocket) -> None:
await websocket.accept()
client = Client()
await hub.add(client)
sender = asyncio.create_task(websocket_sender(websocket, client))
try:
while True:
message = await websocket.receive_text()
# 这里演示客户端到服务器的业务方向。
event = await hub.next_event(f"client says: {message}")
await hub.publish(event)
except WebSocketDisconnect:
pass
finally:
await hub.remove(client)
sender.cancel()
with contextlib.suppress(asyncio.CancelledError):
await sender
def encode_sse(event: dict[str, Any]) -> str:
event_id = event["id"]
event_type = event["type"]
data = json.dumps(event, ensure_ascii=False)
return (
f"id: {event_id}\n"
f"event: {event_type}\n"
f"data: {data}\n\n"
)
async def sse_stream(request: Request) -> AsyncIterator[str]:
client = Client()
await hub.add(client)
try:
while True:
if await request.is_disconnected():
break
try:
event = await asyncio.wait_for(
client.queue.get(),
timeout=15,
)
yield encode_sse(event)
except asyncio.TimeoutError:
# SSE 注释心跳,不触发客户端业务事件。
yield ": heartbeat\n\n"
finally:
await hub.remove(client)
@app.get("/events")
async def events(request: Request) -> StreamingResponse:
return StreamingResponse(
sse_stream(request),
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache",
"Connection": "keep-alive",
"X-Accel-Buffering": "no",
},
)
运行:
python -m pip install fastapi uvicorn
uvicorn main:app --reload
打开一个 SSE 客户端:
curl -N http://127.0.0.1:8000/events
另一个终端发布事件:
curl -X POST http://127.0.0.1:8000/publish \
-H 'content-type: application/json' \
-d '{"text":"hello"}'
SSE 客户端会看到:
id: 1
event: message
data: {"id": 1, "type": "message", "ts": 1730000000.0, "text": "hello"}
实际时间戳会不同。
7.1 这个示例的数据流
flowchart LR
P["POST /publish"] --> H["Hub.publish"]
H --> Q1["WebSocket queue"]
H --> Q2["SSE queue"]
Q1 --> S1["websocket_sender"]
S1 --> W["WebSocket client"]
Q2 --> G["sse_stream generator"]
G --> E["HTTP response stream"]
E --> C["SSE client"]
Hub.publish() 不直接调用 websocket.send_json(),也不直接 yield SSE 数据,而是把事件放入每个连接自己的队列。这样,广播操作的职责只是分发,具体网络写入由连接所属的发送任务完成。
7.2 为什么 Client 使用 eq=False
这里用 set[Client] 管理连接:
self._clients: set[Client] = set()
@dataclass(eq=False) 让对象按身份进行哈希,而不是根据字段值比较。每个连接对象都代表一个独立订阅者,即使两个连接的队列配置相同,也不能被集合当成同一个客户端。
7.3 为什么要有锁
asyncio 中的协程通常运行在同一个事件循环线程,但这不意味着所有复合操作天然安全:
if client in clients:
clients.remove(client)
如果中间发生切换,连接可能已经被其他任务清理。锁保护的是“读取集合快照”和“修改集合”这些操作的逻辑一致性,而不是保护网络发送本身。
锁的持有范围应尽量短。示例中先在锁内复制连接快照,然后释放锁,再执行入队。否则,如果某个未来实现把队列操作改成阻塞等待,就可能长时间持有全局锁。
八、发送顺序、消息顺序与并发任务
对同一个客户端,应用通常希望满足:
其中 ≺ 表示发送顺序。
如果多个任务直接发送:
asyncio.create_task(websocket.send_json(event1))
asyncio.create_task(websocket.send_json(event2))
就不能仅凭创建任务的先后断言网络上的业务顺序。任务何时运行取决于调度、序列化和底层发送状态。
单写者模型可以恢复顺序:
while True:
event = await queue.get()
await send(event)
对一个连接只允许一个发送任务,队列顺序就成为明确的业务顺序。若需要更强的可靠性,还要给事件分配序号,让客户端能够检测间隙:
收到 41
收到 42
收到 44
客户端发现缺少 43 后,可以请求:
GET /events/replay?after=42
这一步已经超出了传输协议本身,属于应用层可靠性设计。
九、asyncio 流、Reader、Writer 与实时通信的关系
WebSocket 和 SSE 通常由 ASGI 服务器、协议实现和框架处理,不需要直接操作 TCP。理解底层 asyncio 流仍然很重要,因为它揭示了“写入”和“真正发送”之间的差异。
asyncio.open_connection() 返回:
reader, writer = await asyncio.open_connection(host, port)
其中:
StreamReader负责异步读取;StreamWriter负责写入;writer.write()把数据交给底层写缓冲;await writer.drain()根据写缓冲水位等待是否可以继续。
Python 3.14 文档明确说明,write() 可能只是把数据放入内部写缓冲;drain() 才是与底层流控交互的等待点。当写缓冲达到高水位时,drain() 会等待缓冲下降到低水位。(docs.python.org)
示例:
import asyncio
async def send_lines(writer: asyncio.StreamWriter) -> None:
for number in range(1000):
writer.write(f"{number}\n".encode())
await writer.drain()
但 drain() 不是万能的公平调度器。Python 文档指出,当写缓冲低于高水位时,drain() 可能立即返回而不让出事件循环。因此一个不断 write() 加 await drain() 的紧密循环仍可能长时间占用执行权;必要时可以显式使用 await asyncio.sleep(0) 让出调度机会。(docs.python.org)
这与 WebSocket 的发送队列是同一个思想层次的问题:
应用对象队列
│
▼
协议库发送队列
│
▼
操作系统 socket 写缓冲
│
▼
网络与客户端读取速度
每一层都可能暂存数据。调用了 send_json() 或 write(),不代表客户端已经处理了消息。生产系统需要区分:
- 应用层队列长度;
- 协议库或服务器层写缓冲;
- 网络 RTT;
- 客户端处理延迟;
- 最终确认或业务消费进度。
十、断开、取消与生命周期
10.1 断开不是异常中的异常
长连接关闭是正常状态转换:
Connected
│
├── 客户端主动关闭
├── 服务端主动关闭
├── 网络断开
└── 进程关闭
▼
Cleanup
WebSocket 中,接收操作通常通过 WebSocketDisconnect 暴露断开;发送操作也可能因连接已经关闭而抛出底层 OSError。ASGI 规范 2.4 起要求,向已关闭连接发送时应抛出服务器特定的 OSError 子类,但旧实现不一定保证这一点。(asgi.readthedocs.io)
因此清理逻辑应放在 finally:
try:
...
finally:
await hub.remove(client)
sender.cancel()
with contextlib.suppress(asyncio.CancelledError):
await sender
如果只在 except WebSocketDisconnect 中清理,发送任务抛出的异常、进程关闭时的取消,可能绕过清理逻辑。
10.2 为什么必须保存并取消任务
下面这种写法容易泄漏任务:
asyncio.create_task(websocket_sender(websocket, client))
如果没有保存任务对象:
- 无法在连接断开时取消;
- 无法等待它结束;
- 发送任务可能继续引用 WebSocket 和队列;
- 异常可能延迟到事件循环日志中才出现。
正确做法是保存任务:
sender = asyncio.create_task(...)
然后在 finally 中取消并等待。取消只是在任务上设置取消请求;任务必须运行到可取消的 await 才能真正结束。FastAPI 对流式响应生成器的说明也强调了这一点。(fastapi.tiangolo.com)
10.3 应用关闭时的后台任务
如果应用有后台发布任务、心跳任务或外部消息订阅任务,应把它们绑定到应用生命周期,而不是在模块导入时永久启动。
FastAPI 支持使用 lifespan 管理启动和关闭资源。ASGI 的 lifespan 规范定义了启动、启动完成、关闭和关闭完成等事件,并允许服务器向请求作用域提供共享状态。(asgi.readthedocs.io)
一个简化示例:
from contextlib import asynccontextmanager
from fastapi import FastAPI
async def consume_external_events() -> None:
while True:
await asyncio.sleep(1)
# 从 Redis、Kafka 或其他消息系统读取事件
@asynccontextmanager
async def lifespan(app: FastAPI):
task = asyncio.create_task(consume_external_events())
try:
yield
finally:
task.cancel()
with contextlib.suppress(asyncio.CancelledError):
await task
app = FastAPI(lifespan=lifespan)
关闭顺序应该避免在资源已经释放后仍有任务尝试发送:
停止接收新连接
│
▼
停止外部事件生产者
│
▼
通知或关闭现有客户端
│
▼
取消发送任务
│
▼
释放数据库、消息队列和其他资源
十一、常见错误与诊断方法
11.1 用无界列表保存待发送消息
错误模型:
pending_messages.append(event)
如果客户端网络变慢,列表会持续增长。诊断时应观察:
- 单连接队列长度;
- 队列最大值;
- 入队速度;
- 出队速度;
- 队列满次数;
- 客户端连接持续时间。
11.2 广播时直接串行发送
错误模型:
for websocket in clients:
await websocket.send_json(event)
失败表现:
- 某个客户端变慢后,所有客户端延迟一起上升;
- 广播调用耗时与最慢客户端相关;
- 事件循环中出现长时间等待发送。
改为每客户端队列后,需要继续观察发送任务的异常和队列积压,否则只是把问题从广播函数转移到了内存。
11.3 把 SSE 当成双向协议
错误模型:
const source = new EventSource("/events");
source.send("hello"); // 不存在
SSE 客户端到服务器的业务请求应通过普通 HTTP:
await fetch("/publish", {
method: "POST",
headers: {"content-type": "application/json"},
body: JSON.stringify({text: "hello"})
});
如果客户端也需要持续发送数据,WebSocket 通常更自然。
11.4 只依赖 TCP 连接状态判断在线
TCP 连接可能因为中间网络设备、移动网络切换或对端异常而处于长时间不可用状态。应用应使用协议级心跳或业务级确认,并为等待确认设置超时。
11.5 生产环境多进程下使用进程内广播
Hub 只存在于当前 Python 进程:
进程 1:客户端 A、B
进程 2:客户端 C、D
如果事件只发布给进程 1 的 Hub,C、D 不会收到。
多进程或多实例部署需要外部事件总线:
业务生产者
│
▼
Redis Pub/Sub、消息队列或其他事件系统
│
├──> 进程 1 Hub ──> A、B
└──> 进程 2 Hub ──> C、D
但 Pub/Sub 通常只保证“订阅期间转发”,不保证断线期间可恢复。需要可靠重放时,应使用带持久化和游标语义的事件日志,或在业务数据库中保存可重建状态。
十二、如何在 WebSocket 与 SSE 之间选择
选择标准不是“哪个性能更高”,而是通信方向和可靠性模型。
选择 WebSocket,通常是因为:
- 客户端需要持续向服务器发送消息;
- 需要低延迟的双向交互;
- 需要客户端确认、输入事件或实时协商;
- 连接协议和消息类型比较丰富。
选择 SSE,通常是因为:
- 主要是服务器向浏览器推送;
- 客户端使用
EventSource更简单; - 希望沿用 HTTP 代理、认证和流式响应路径;
- 事件是通知、日志、任务进度或状态更新。
无论选择哪一种,以下关系都不会改变:
事件生产速率 > 客户端消费速率
│
▼
必须采取阻塞、丢弃、合并、限速或断开策略
WebSocket 提供双向消息通道,SSE 提供 HTTP 事件流,心跳提供失活检测,广播提供一对多复制,而背压决定系统在客户端变慢时是否仍然可控。它们不是互相替代的功能,而是实时系统中的不同层次:
协议层:WebSocket / HTTP SSE
连接层:接受、发送、断开、心跳
调度层:asyncio 任务与取消
分发层:广播与订阅
资源层:有界队列与背压
可靠性层:事件 ID、重放、持久化
生命周期层:启动、关闭和故障清理
只有把这些层次分开,实时通信代码才不会停留在“能连上、能发消息”的演示阶段。
系列导航与关联阅读
- 系列入口:Python 完整学习路线:从语言模型、并发到 Web、数据、AI 与生产交付
- 上一篇:Python gRPC:Protobuf、流式调用、Deadline、拦截器和错误模型
- 下一篇:Python Web 认证与授权:Session、JWT、OAuth 2.0、RBAC 和 CSRF
- 延伸:Python asyncio 网络流:TCP、Reader、Writer、TLS 与背压
- 延伸:FastAPI 完整基础:路由、依赖注入、校验、异步和生命周期
官方资料
本文依据 Python 官方文档、相关 PEP 与生态项目官方文档重新梳理;正文、示例与工程清单由 WR BLOG 编写。

评论
0 条讨论