Skip to content

WebSocket 与 SSE 实时推送实战 ​

实时推送有两种做法:WebSocket 双向、SSE 单向。Robyn 两者都原生支持,官方文档告诉你 API 怎么调,这里讲的是上线时会真正咬人的部分:连接怎么保活、广播在多进程下为什么只推了一半、断线怎么续传。

一、先选 SSE 还是 WebSocket ​

维度SSEWebSocket
方向服务端 → 客户端双向
传输普通 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();

重连间隔务必加随机抖动,否则服务端重启瞬间会被全部客户端同时重试打崩(惊群)。

六、上线检查表 ​

  1. Nginx/CDN 是否透传 Upgrade 头,长连接超时是否已调大
  2. 使用 --processes > 1 时广播是否已改为 Redis/RabbitMQ 等外部总线
  3. 客户端是否有退避重连,服务端是否处理了 WebSocketDisconnect
  4. 单连接内存/CPU 开销预估:连接数 × 消息频率是否与 --workers 匹配
  5. SSE 是否设置了 id 以支持断线续传

基于 MIT 许可发布 · 隐私政策