SSE 实时推流:测试平台日志与状态推送的完整实践

测试平台的实时性需求很朴素:任务在跑,浏览器要现在看到日志和进度。技术选型时 WebSocket 是惯性答案,但我们的场景(单向推送 + HTTP 基建 + 内网环境)用 SSE 更简单更稳。这篇是 Flask + SSE 的完整实践:从事件设计到断线重连到性能边界。

一、为什么选 SSE 而不是 WebSocket

维度 SSE WebSocket
方向 单向(服务端→客户端)✅ 正好 双向(用不上)
断线重连 浏览器原生EventSource 自动带 Last-Event-ID 自己写
基建穿透 普通 HTTP,代理/网关/鉴权全复用 部分内网设备要开 ws 端口
流量开销 每帧带 HTTP 头(~200B) 2~14B
连接数 受 HTTP 限制(同域 6 连接) 无此限制

我们的场景:日志推送、任务状态、SSE 事件流全是单向;同域 6 连接上限通过事件合并到 1~2 个流解决;200B 帧头对日志场景无感。省掉的"重连 + 握手 + 基建"复杂度,比帧头钱值钱得多

二、事件设计:一个流,多频道

不要每个功能开一个 SSE 端点(吃连接数),合并成一个事件流 + 类型字段

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
@bp.route("/ui/events")
def events():
sub = live_log.subscribe() # 拿到专属队列

def stream():
last_ping = time.time()
try:
while True:
try:
ev = sub.get(timeout=15)
except queue.Empty:
yield ": ping\n\n" # 心跳:保活 + 检测死连接
last_ping = time.time()
continue
if ev is SENTINEL: # 服务端通知关闭
break
# SSE 帧格式:event 类型 + id(断线重连续传)+ data
yield f"id: {ev.seq}\nevent: {ev.type}\ndata: {ev.payload}\n\n"
finally:
live_log.unsubscribe(sub)

return Response(stream(), mimetype="text/event-stream",
headers={"Cache-Control": "no-cache",
"X-Accel-Buffering": "no"}) # 关 nginx 缓冲!

三个关键细节:

  1. id: seq 单调递增:断线重连时浏览器自动带 Last-Event-ID,服务端从该 seq 续传——环形缓冲保留最近 2000 帧即可覆盖"断网 30 秒"的场景;
  2. 15s 心跳:内网代理普遍 30~60s 空闲断连,心跳必须比它短;心跳同时是死连接检测器——yield 给死连接会抛 BrokenPipe,finally 里退订,不泄漏订阅队列
  3. X-Accel-Buffering: no:nginx 默认缓冲 SSE 响应,日志会"攒一批才出来"——这个头不加,所有实时性都白搭(我们第一个版本就栽在这,表现为"日志有 5 秒延迟")。

客户端:

1
2
3
4
const es = new EventSource("/api/ui/events");
es.addEventListener("log", (e) => appendLog(JSON.parse(e.data)));
es.addEventListener("task_state", (e) => updateBadge(JSON.parse(e.data)));
// 重连由浏览器原生处理,Last-Event-ID 自动带上,服务端续传

三、发布侧:进程内注册表 + 背压

Flask 单进程内,发布者是执行引擎的后台线程:

1
2
3
4
5
6
7
8
9
10
11
12
13
class LiveLogRegistry:
"""任务 ID → 订阅者队列列表;发布 O(订阅者数),慢订阅者不阻塞发布"""

def emit(self, task_id, type, payload):
seq = self._seq.next()
self._ring.append((task_id, seq, type, payload)) # 环形缓冲,断线续传
for q in self._subs.get(task_id, []):
if q.full():
# 慢订阅者:丢弃旧帧 + 发一条 "dropped:N" 告知客户端跳帧
q.get_nowait()
q.put((task_id, seq, type, payload))
else:
q.put((task_id, seq, type, payload))

背压策略:发布永不阻塞。执行引擎是生产核心,一个卡住的浏览器连接不能拖慢任务执行;代价是慢订阅者丢帧——丢帧通过 seq 缺口暴露(客户端发现 id 跳号,拉一次 /ui/events/backfill?from=seq 补全),实时性让位于完整性兜底

四、性能边界与实测

单台 Flask 机器(4C8G)实测:

并发 SSE 连接 日志吞吐(发布侧) 事件循环影响
50 500 条/s 无感
200 800 条/s 发布线程 CPU +3%
500 1000 条/s(帧合并后) 需要帧合并(见下)

帧合并是 500 连接档的关键:日志高频时不逐条 yield,而是 100ms 窗口内的多条日志合并成一帧(data: [ev1, ev2, ev3])——浏览器端本来也是批量渲染的,逐条推是浪费。合并窗口 100ms 对"看日志"场景的实时性无感知。

Gunicorn 注意:SSE 长连接会占死 worker。生产配置 --threadsgevent workerSSE 端点单独限流(内网工具场景,单用户 2 个连接足够)。

五、踩过的坑

  1. 浏览器自动重连风暴:服务端重启瞬间,200 个 EventSource 同时重连 + 补拉,日志接口被打满——重连加客户端随机抖动setTimeout(random 0~3s) 后再 new EventSource),补拉接口限流;
  2. yield 里的异常吞掉:generator 里 BrokenPipe 不捕获,订阅队列泄漏,表现为"越跑内存越高"——finally 退订是硬要求;
  3. 时区/换行:日志 payload 里的 \n 必须转义成 \\n——SSE 协议里裸 \n 是帧分隔符,多行日志直接破坏帧结构(这个 bug 排查了一下午)。

六、小结

SSE 实践的四个支柱:单流多频道(省连接数)、seq + 心跳 + 缓冲续传(断线无损)、发布不阻塞 + seq 缺口补全(背压与完整性兼顾)、帧合并(规模档性能)。WebSocket 不是更好的答案——它是更贵、更复杂、我们场景里不需要的答案。