KEEL · 龙骨 · A CURRICULUM FOR THE AI ERA
03 · 服务端实现:从"能跑"到"能上线" — keel 龙骨
这一章回答:一个 SSE 接口要满足哪些条件,才算生产可用。
这一章回答:一个 SSE 接口要满足哪些条件,才算生产可用。
一个能在本地跑通的 SSE 接口,通常还缺三样东西:缓冲被关掉、心跳在发、超时配置正确。缺任何一样,都会表现为某种"偶发卡住"。
一、最小可用实现(FastAPI)
import asyncio, json
from fastapi import APIRouter, Request
from fastapi.responses import StreamingResponse
router = APIRouter()
async def event_gen(request: Request, run_id: str, start_id: str | None):
try:
async for evt in source_events(run_id, after=start_id): # 你的事件源
if await request.is_disconnected():
break
yield f"id: {evt.id}\nevent: {evt.type}\ndata: {json.dumps(evt.payload, ensure_ascii=False)}\n\n"
yield "event: done\ndata: {}\n\n"
finally:
await cleanup(run_id) # 必须清理:取消订阅、释放资源
@router.get("/runs/{run_id}/stream")
async def stream(request: Request, run_id: str,
last_event_id: str | None = Header(default=None, alias="Last-Event-ID")):
return StreamingResponse(
event_gen(request, run_id, last_event_id),
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache, no-transform",
"Connection": "keep-alive",
"X-Accel-Buffering": "no", # 关键:让 nginx 不缓冲
},
)
四件必须做的事:
media_type="text/event-stream"—— 客户端据此识别;no-cache, no-transform——no-transform防止中间层做压缩/改写;X-Accel-Buffering: no—— nginx 专用指令(见第三节);finally里清理 —— 客户端断开后必须释放订阅与资源,否则会泄漏。
二、检测客户端断开(否则你会一直推给空气)
if await request.is_disconnected():
break
在 Starlette/FastAPI 里,is_disconnected() 通过读取底层 socket 判断。更可靠的做法是在每次 yield 之间检查,因为生成器只有在被消费时才会感知断开。
⚠️ 一个常见误区:以为 yield 之后代码会立刻继续。如果客户端不读,生成器会阻塞在 yield,finally 不会执行。所以:
- 用
asyncio.wait_for给关键等待加超时; - 或让事件源本身有超时机制(如 Redis
XREAD BLOCK 5000,每 5 秒返回一次,给 disconnect 检查留出机会)。
# Redis XREAD 的 BLOCK 参数天然提供了检查窗口
entries = await redis.xread({stream_key: last_id}, count=100, block=5000)
if not entries:
if await request.is_disconnected():
break
yield ": heartbeat\n\n" # 空闲时发心跳
continue
三、关闭一切缓冲(最重要的一节)
这是"本地正常、线上卡住"的头号原因。缓冲可能出现在四个地方:
| 位置 | 关闭方式 |
|---|---|
| 应用框架 | FastAPI StreamingResponse 本身不缓冲;但如果你先把内容攒成字符串再返回,就自己引入了缓冲 |
| nginx | proxy_buffering off; 或响应头 X-Accel-Buffering: no |
| 压缩 | gzip 会攒够一个块才输出 → gzip off;(或对 text/event-stream 类型关闭) |
| CDN / 云网关 | 各厂商不同,需要显式关闭"响应缓冲 / 智能压缩" |
location /api/runs/ {
proxy_pass http://127.0.0.1:8000;
proxy_http_version 1.1;
proxy_set_header Connection '';
proxy_buffering off; # 关键
proxy_cache off; # 不要缓存事件流
gzip off; # 压缩会攒块,破坏实时性
chunked_transfer_encoding on; # 保持分块传输
proxy_read_timeout 3600s; # 长连接必须调大(默认 60s 会掐断)
proxy_send_timeout 3600s;
proxy_connect_timeout 5s; # 连接建立仍然要短
send_timeout 3600s;
}
为什么 proxy_buffering off 如此关键:nginx 默认会把上游响应攒满一个缓冲区(几 KB)才发给客户端。SSE 的事件往往很小很稀疏,于是客户端要等很久才收到第一字节——表现为"连接成功了但一直没数据",而服务端日志显示明明在推。
验证方法:
curl -N http://localhost/api/runs/1/stream | ts # 加时间戳看是否逐条到达
# 或直接看首字节时间
curl -N -o /dev/null -w 'TTFB %{time_starttransfer}s\n' http://.../stream
首字节时间应该是毫秒级。如果是几秒甚至几十秒,就是被缓冲了。
四、心跳:不是为了活跃,是为了探测
yield ": ping\n\n" # 注释行,客户端不会派发事件
三个作用:
- 保活:让中间层(负载均衡、代理、NAT)不至于因空闲而断开;
- 探测断开:发不出去 = 客户端已走(TCP 会报错,触发
finally清理); - 让客户端知道"服务端还活着"(区分"没数据"和"连接死了")。
心跳间隔的取值:
- 常见代理/负载均衡的空闲超时是 60 秒(有些是 30 秒、有些是 4 分钟);
- 心跳间隔应设为 小于最小空闲超时的一半,通常 15~30 秒;
- 太频繁(1 秒)会浪费带宽与 CPU,尤其在万级连接时。
⚠️ 心跳不要写进事实日志。如果心跳也被持久化,断线重连的客户端会回放出一堆无意义的心跳(与 Harness 源码解剖室板块里讲「多宿主托管」的那门课第 03 章中"心跳只活在这条连接上"的设计同源)。
五、超时配置的完整清单
一条长连接涉及多层超时,任何一层比心跳短都会导致断开:
| 层 | 参数 | 建议 |
|---|---|---|
| 应用(uvicorn/gunicorn) | --timeout-keep-alive、h11 超时 |
大于心跳间隔;或设为 0(不限制)+ 由上层控制 |
| nginx | proxy_read_timeout / send_timeout |
长连接场景设 3600s |
| 负载均衡 / 云网关 | 空闲超时(idle timeout) | 显式调大(很多默认 60s) |
| CDN | 边缘连接超时 | 确认支持长连接(有些 CDN 不适合代理 SSE) |
| 客户端 | EventSource 自动重连 |
服务端下发 retry: 控制节奏 |
配置原则:外层超时 > 内层超时 > 心跳间隔,且留出至少 2 倍余量。
六、不要在推送过程中持有数据库资源
# ❌ 危险:连接期间一直持有会话/事务
async with async_session() as s:
async for evt in events:
yield format(evt) # 这个会话在整条连接期间都不释放
# ✅ 正确:事件来自独立的事件源(队列/Redis Stream/内存队列),不在推送路径上开事务
async for evt in event_source.subscribe(run_id):
yield format(evt)
万级连接时,前者会直接耗尽连接池(呼应《索引与并发实战》第 05 章的并发度计算)。推送路径上应该只有"读事件源 + 格式化",不碰数据库。
七、Node/Express 对照
app.get('/runs/:id/stream', async (req, res) => {
res.writeHead(200, {
'Content-Type': 'text/event-stream',
'Cache-Control': 'no-cache, no-transform',
'Connection': 'keep-alive',
'X-Accel-Buffering': 'no',
});
res.write('\n'); // 有些代理需要首个字节才认为响应已开始
const lastId = req.headers['last-event-id'];
const sub = subscribe(req.params.id, lastId);
const ping = setInterval(() => res.write(': ping\n\n'), 20000);
req.on('close', () => { clearInterval(ping); sub.close(); }); // 必须清理
for await (const evt of sub) {
res.write(`id: ${evt.id}\nevent: ${evt.type}\ndata: ${JSON.stringify(evt.payload)}\n\n`);
}
res.end();
});
要点一致:首字节、心跳、断开清理、关闭缓冲。
动手:可观察结果
| 产出 | 判断标准 |
|---|---|
| 一个生产可用的 SSE 接口 | 含正确响应头、断开检测、心跳、finally 清理 |
| nginx 配置 | proxy_buffering off / gzip off / 超时三项齐全,并验证 TTFB < 500ms |
| 断开清理验证 | 客户端强制断开后,服务端订阅数/资源在 10 秒内回落(无泄漏) |
| 心跳与超时配合验证 | 心跳 20s,代理空闲超时 60s,连接保持 10 分钟不断 |
| 推送路径无数据库验证 | 万级模拟连接下数据库连接数不随连接数增长 |
完成标志:你的 SSE 接口能在"经过 nginx + 压缩 + 负载均衡"的真实链路上连续推送 10 分钟不中断,且客户端断开后服务端无资源泄漏。
故障注入
| 注入方式 | 观察 |
|---|---|
打开 proxy_buffering(默认) |
首字节时间是否变成数秒/永不到达 |
| 开启 gzip | 事件是否攒批到达(实时性被破坏) |
把 proxy_read_timeout 保持 60s |
第 61 秒连接是否被掐断 |
| 不写心跳且连接空闲 | 负载均衡是否在空闲超时后断开 |
| 客户端强制断开后服务端不清理 | 订阅数/文件描述符是否持续增长(泄漏) |
| 在推送循环里开数据库事务 | 并发连接上升时连接池是否被耗尽 |
finally 里做同步阻塞操作 |
断开处理是否拖慢其它连接 |
自测题
- 一个 SSE 响应必须包含哪些响应头?
no-transform与X-Accel-Buffering各自解决什么? - 为什么"本地能推、线上卡住"?按顺序列出四个可能的缓冲点。
- 心跳的三个作用是什么?间隔应该怎么取?
- 超时配置的分层原则是什么?举出三层并给出取值关系。
- 为什么推送路径上不能持有数据库会话?后果是什么?