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 不缓冲
        },
    )

四件必须做的事:

  1. media_type="text/event-stream" —— 客户端据此识别;
  2. no-cache, no-transform —— no-transform 防止中间层做压缩/改写;
  3. X-Accel-Buffering: no —— nginx 专用指令(见第三节);
  4. finally 里清理 —— 客户端断开后必须释放订阅与资源,否则会泄漏。

二、检测客户端断开(否则你会一直推给空气)

if await request.is_disconnected():
    break

在 Starlette/FastAPI 里,is_disconnected() 通过读取底层 socket 判断。更可靠的做法是在每次 yield 之间检查,因为生成器只有在被消费时才会感知断开。

⚠️ 一个常见误区:以为 yield 之后代码会立刻继续。如果客户端不读,生成器会阻塞在 yield,finally 不会执行。所以:

# 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"     # 注释行,客户端不会派发事件

三个作用:

  1. 保活:让中间层(负载均衡、代理、NAT)不至于因空闲而断开;
  2. 探测断开:发不出去 = 客户端已走(TCP 会报错,触发 finally 清理);
  3. 让客户端知道"服务端还活着"(区分"没数据"和"连接死了")。

心跳间隔的取值:

⚠️ 心跳不要写进事实日志。如果心跳也被持久化,断线重连的客户端会回放出一堆无意义的心跳(与 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 里做同步阻塞操作 断开处理是否拖慢其它连接

自测题

  1. 一个 SSE 响应必须包含哪些响应头?no-transform 与 X-Accel-Buffering 各自解决什么?
  2. 为什么"本地能推、线上卡住"?按顺序列出四个可能的缓冲点。
  3. 心跳的三个作用是什么?间隔应该怎么取?
  4. 超时配置的分层原则是什么?举出三层并给出取值关系。
  5. 为什么推送路径上不能持有数据库会话?后果是什么?

进入 keel 阅读