WebSocket 与 SSE 实时推送实战
实时推送有两种做法:WebSocket 双向、SSE 单向。Robyn 两者都原生支持,官方文档告诉你 API 怎么调,这里讲的是上线时会真正咬人的部分:连接怎么保活、广播在多进程下为什么只推了一半、断线怎么续传。
一、先选 SSE 还是 WebSocket
| 维度 | SSE | WebSocket |
|---|---|---|
| 方向 | 服务端 → 客户端 | 双向 |
| 传输 | 普通 HTTP,自动走现有代理 | 需要 Upgrade,部分代理要额外配置 |
| 断线重连 | 浏览器原生支持 | 需自己实现 |
| 适合的场景 | 通知、仪表盘、进度、日志流 | 聊天、协同编辑、游戏 |
经验:先考虑 SSE。它不需要 Nginx 额外配 Upgrade 头,浏览器自带重连,能覆盖大多数"实时"需求。只有确实需要客户端主动推消息时才上 WebSocket。
二、SSE:把重做一次的活交给浏览器
python
import asyncio
import json
import time
from robyn import Robyn, SSEResponse, SSEMessage
app = Robyn(__file__)
@app.get("/events")
async def stream_events(request):
async def event_generator():
last_id = request.headers.get("last-event-id")
i = int(last_id) + 1 if last_id and last_id.isdigit() else 0
while True:
data = {"id": i, "ts": time.time(), "message": f"更新 {i}"}
# id 是关键:客户端重连时会带上 Last-Event-ID
yield SSEMessage(json.dumps(data), event="update", id=str(i))
i += 1
await asyncio.sleep(2)
return SSEResponse(event_generator())前端什么都不用写,浏览器会自己连回来,并把最后一个事件通过 Last-Event-ID 请求头带回来:
javascript
const es = new EventSource("/events");
es.addEventListener("update", (e) => console.log(JSON.parse(e.data)));要点:
SSEMessage(data, event=..., id=...)里的id一定要给,它是断线续传的唯一依据- 生成器可以是同步的也可以是异步的,里面有 IO 就用异步生成器(见 服务器发送事件)
- 连接数是有上限的:HTTP/1.1 下浏览器对同一域名并发连接约 6 个,SSE 会一直占着。多开标签页时考虑用 HTTP/2
三、WebSocket:连接、广播与优雅断开
python
from robyn import Robyn
from robyn import WebSocketDisconnect
app = Robyn(__file__)
@app.websocket("/chat")
async def chat(websocket):
try:
while True:
msg = await websocket.receive_text()
await websocket.broadcast(f"{websocket.id}: {msg}")
except WebSocketDisconnect:
print("client disconnected", websocket.id)
@chat.on_connect
def on_connect(websocket):
return f"欢迎,{websocket.id}!"
@chat.on_close
def on_close(websocket):
return "再见!"常用 API(完整表格见 WebSockets):
| 方法 | 作用 |
|---|---|
await ws.receive_text() / receive_json() | 阻塞等待下一条消息 |
await ws.send_text() / send_json() | 只发给当前连接 |
await ws.broadcast() | 广播给同一端点的所有连接 |
await ws.close() | 服务端主动断开 |
ws.id / ws.query_params | 连接 UUID 与 URL 查询参数 |
四、多进程下的广播陷阱(最容易翻车的地方)
broadcast() 只作用于当前进程内的连接注册表。一旦你用 --processes=4 启动,请求会被分发到不同进程:
bash
python app.py --processes=4 # 每个进程各自维护一份连接表结果就是:客户 A 连在进程 1,客户 B 连在进程 2,A 发消息 broadcast,B 收不到。这在本地单进程调试时完全看不出来,上线才发现。
解决办法是外置一层 Pub/Sub:
python
import asyncio
import redis.asyncio as redis
r = redis.from_url("redis://localhost:6379")
@app.websocket("/chat")
async def chat(websocket):
pubsub = r.pubsub()
await pubsub.subscribe("chat")
# 收 Redis → 推给本连接
asyncio.create_task(pump(websocket, pubsub))
# 收本连接 → 发到 Redis
try:
while True:
msg = await websocket.receive_text()
await r.publish("chat", msg)
except WebSocketDisconnect:
await pubsub.close()
async def pump(websocket, pubsub):
async for message in pubsub.listen():
if message["type"] == "message":
await websocket.send_text(message["data"].decode())架构上只需记住一句话:WebSocket 连接是进程本地的,跨进程广播必须走外部消息总线。
五、心跳与重连
Nginx 默认 60 秒会掐掉空闲连接,proxy_read_timeout 之外还需要在应用层发心跳:
nginx
location /ws {
proxy_pass http://robyn_app;
proxy_http_version 1.1;
proxy_set_header Upgrade $http_upgrade;
proxy_set_header Connection "upgrade";
proxy_read_timeout 3600s; # 长连接务必调大
}客户端侧:
javascript
function connect() {
const ws = new WebSocket("wss://example.com/chat");
ws.onmessage = (e) => render(e.data);
ws.onclose = () => setTimeout(connect, 1000 + Math.random() * 2000); // 抖动的退避重连
}
connect();重连间隔务必加随机抖动,否则服务端重启瞬间会被全部客户端同时重试打崩(惊群)。
六、上线检查表
- Nginx/CDN 是否透传
Upgrade头,长连接超时是否已调大 - 使用
--processes > 1时广播是否已改为 Redis/RabbitMQ 等外部总线 - 客户端是否有退避重连,服务端是否处理了
WebSocketDisconnect - 单连接内存/CPU 开销预估:连接数 × 消息频率是否与
--workers匹配 - SSE 是否设置了
id以支持断线续传