本文目录
大模型 导购 或日志 tail:用户盯着空白页等 30 秒,体验像卡死。若每生成一段就推一块 JSON,浏览器可以边收边渲染——这就是 流式 HTTP:响应体不是一次性算完再 return,而是持续写出。SSE(Server-Sent Events)是 流式 里常用的一种 事件 格式,单向、基于 HTTP、前端 EventSource 能直接消费。FastAPI 用 StreamingResponse 承载字节 流;理解 事件 帧格式、背压 与断开,比背 API 名更重要。
流式 不是银弹:它改善的是「何时看到第一字节」,不改变 业务 正确性责任。该短 事务 提交 的写操作,不应因为输出 流式 就拖到 generator 里 commit。
普通响应 vs 流式响应
普通 return dict | StreamingResponse | |
|---|---|---|
| 内存 | 整份 body 在内存 | 可逐块生成 |
| TTFB | 业务算完才发 | 首块算完即发 |
| 中间件 | 完整经过 | 首包后 body 不再走部分中间件 |
| 客户端 | 一次读完 | 长连接读直到结束 |
ASGI 层:StreamingResponse 把 async generator / iterator 的产出 repeatedly send 给客户端。
从时间线看:客户端发请求 → 服务端返回 200 与 Header → body 尚未完整 → 连接保持 → 多个 chunk 顺序到达 → 连接关闭。TTFB(首字节时间)往往比总时长更能改善体感;流式 优化的是「什么时候开始看到内容」,不是峰值 QPS。
StreamingResponse 最小例
from collections.abc import AsyncIterator
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
app = FastAPI()
async def count_stream(limit: int) -> AsyncIterator[bytes]:
for i in range(limit):
yield f"chunk {i}\n".encode()
await asyncio.sleep(0.1)
@app.get("/stream/count")
async def stream_count() -> StreamingResponse:
return StreamingResponse(count_stream(10), media_type="text/plain")同步 generator 也可:StreamingResponse(sync_gen(), media_type=...),FastAPI 会放到线程池迭代——CPU 重的话优先 async generator + await 让出事件循环。
media_type 告诉客户端如何解析 bytes;SSE 必须是 text/event-stream,普通 NDJSON 流 常用 application/x-ndjson 或 text/plain,与前端约定一致即可。
SSE mental model
SSE 的 MIME 是 text/event-stream。每条 事件 是文本帧,以空行(\n\n)分隔:
data: 第一段文字
data: 第二段文字
event: done
data: [END]字段常见:
| 字段 | 含义 |
|---|---|
data: | payload,可多行 |
event: | 事件 类型名,前端 addEventListener |
id: | 断线重连 Last-Event-ID |
: comment | 注释,用于 keep-alive 心跳 |
生成 SSE 帧的小函数:
def sse_event(data: str, event: str | None = None, event_id: str | None = None) -> bytes:
lines: list[str] = []
if event_id is not None:
lines.append(f"id: {event_id}")
if event is not None:
lines.append(f"event: {event}")
for line in data.splitlines() or [""]:
lines.append(f"data: {line}")
lines.append("")
lines.append("")
return "\n".join(lines).encode("utf-8")路由:
async def guide_sse(user_id: int, product_id: int) -> AsyncIterator[bytes]:
async for token in guide_client.stream_tokens(user_id, product_id):
yield sse_event(token)
yield sse_event("[DONE]", event="done")
@app.get("/products/{product_id}/guide/stream")
async def stream_guide(product_id: int, user: UserDep) -> StreamingResponse:
return StreamingResponse(
guide_sse(user.id, product_id),
media_type="text/event-stream",
headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"},
)X-Accel-Buffering: no 提示 Nginx 别缓冲整段 body,否则 流式 变假 流式。
前端用 EventSource 订阅时,浏览器自动重连;服务端 id: 字段帮助断点续传。若 事件 带 JSON,把 JSON 字符串放在 data: 一行里,前端 JSON.parse(e.data)。
NDJSON 与 SSE 选型
| SSE | NDJSON 流 | |
|---|---|---|
| 格式 | data: 帧 + 空行 | 每行一个 JSON |
| 浏览器 | EventSource 原生 | 需 fetch + reader |
| 事件 类型 | event: 字段 | 自行约定行内 type |
| 代理缓冲 | 需关 buffering | 同样需注意 |
导购逐 token 打字机效果两种都能做;要浏览器零依赖选 SSE;要更灵活 Header 认证 常选 fetch 流。
依赖注入与异常:流式特殊点
路由函数 return StreamingResponse(...) 之前,Depends 链、参数校验、认证 仍正常执行;之后 generator 里抛错,全局异常处理器未必能改成 JSON 错误体——客户端可能已收到 200 和部分 事件。因此:
- 认证、权限、参数错误留在 generator 外解决;
- generator 内错误用 SSE
event: error+data: {...}约定,或记录 日志 后优雅yield结束帧。
课件里一点值得记住:流式 body 迭代完成后,依赖里带 yield 的 teardown 才跑完——若 generator 永不结束,Session 可能长期占用(见 事务 篇连接占用)。
因此 流式 路由里避免在 generator 内持有 DB 事务;读数据在 generator 之前 完成,或 generator 内只用已 materialize 的内存结构。长时间 流式 输出与「请求级 Session 占连接」是冲突的,边界 要在设计时划清。
背压与客户端断开
背压直觉:生产(generator yield)快、消费(客户端读)慢时,TCP 发送缓冲区填满,yield 侧 await 写 socket 会变慢——流式 自然限速,不必 busy loop。
客户端关 tab 或 EventSource.close():服务端应检测 Request.is_disconnected()(Starlette)并 break generator,停止调 LLM、释放 DB:
async def guide_sse(request: Request, ...) -> AsyncIterator[bytes]:
try:
async for token in guide_client.stream_tokens(...):
if await request.is_disconnected():
break
yield sse_event(token)
finally:
await guide_client.cancel()不处理断开 = 白算 token、白占连接。
也可在 finally 里打 日志 stream_aborted,带 request_id,便于统计「用户等不及关页面」的比例,评估是否真的需要 流式。
EventSource 与 token
浏览器原生 EventSource 不能自定义 Header,Bearer token 常放 query(?access_token=)或 Cookie。需要 Header 时用 fetch + ReadableStream 或第三方库。认证 设计 接口 时要与前端能力对齐,这不是 StreamingResponse 能单独解决的。
何时不要 流式
| 场景 | 更合适的做法 |
|---|---|
| 小 JSON(< 几十 KB) | 普通 return |
| 需要整体 事务 成功才展示 | 算完再返回 |
| 文件下载已知大小 | FileResponse / Content-Length |
| 双向实时 | WebSocket,非 SSE |
流式 换的是首字节时间与 UX,不是更高吞吐;Admin 导出 CSV 百万行往往用异步任务 + 轮询下载链接,而非 HTTP 流 挂一小时。
与 WebSocket 的边界
WebSocket 全双工,适合协同编辑、房间聊天;SSE 单向、走 HTTP,穿透企业代理通常更省心。导购只要「服务端推 token」用 SSE 足够;要客户端中途改 prompt 再推,才考虑 WebSocket。
与中间件、CORS
CORS 预检对 GET SSE 通常一次通过;若跨域带 credentials,需正确 Access-Control-Allow-Origin。某些中间件假设「响应 body 可重复读」——对 StreamingResponse 不成立;自定义中间件勿读光 body。
测试
TestClient 对 StreamingResponse 支持 with client.stream("GET", url) as resp: 迭代 resp.iter_lines()。断言首 chunk 时间可测,但集成测试里 sleep 要短。单元测试优先测 sse_event 格式与 generator 逻辑(mock 外部 流)。
Mock 外部 LLM 流 时,用 async generator yield 固定字符串,不必真等网络。断言最后一帧 event: done 是否存在,比 snapshot 整段 body 更稳。
常见误解
误解 1:流式 = WebSocket。SSE 仍是 HTTP 单向;WebSocket 双向、协议不同。
误解 2:每个 API 都该 SSE。多数 CRUD 无收益,还增加断开与缓存复杂度。
误解 3:generator 里 db.commit() 每 chunk 一次。事务 边界 仍应按 业务 单元,流式 只描述输出节奏。
生产 checklist:流式 上线前
| 检查项 | 原因 |
|---|---|
| 反向代理关闭 body 缓冲 | 否则客户端长时间无 chunk |
| generator 内检测 disconnect | 避免白算与连接泄漏 |
| 错误协议与前端约定 | 200 中途失败不能只靠断连 |
| 认证 方式与 EventSource 限制一致 | 避免生产才 discover 不能带 Header |
| 超时:代理、Uvicorn、上游 LLM 对齐 | 最短的一环会先断 |
| 日志 记录 stream 开始/结束/abort | 追踪 与计费 |
流式 接口的 SLA 往往比普通 JSON 难承诺:网络抖动、客户端提前关闭、上游 token 生成速度波动都会影响体感。文档里写清「非 事务 性、可能中断、建议前端重试策略」,比 silent fail 更专业。
心智模型图
用文字概括 SSE 数据流:
[Service 产出 token] → [sse_event 编码 bytes] → [StreamingResponse]
→ [Uvicorn send 循环] → [TCP] → [浏览器 EventSource onmessage]背压 发生在 Uvicorn send 与 TCP 之间:客户端读得慢,send await 变长,上游 generator 的 yield 也会慢下来——整条链自动限速。若你在 generator 里无 await 地狂 yield,仍可能占满内存缓冲区;CPU 型同步 generator 更要控制 chunk 大小。
中间件与 流式 body 的交互
Starlette 中间件在 call_next 返回后,若响应是 StreamingResponse,部分中间件只处理 headers,不再缓冲 body。自定义中间件若读取 response.body,对 流式 会失败或耗尽内存。Access 日志 中间件应在 call_next 前后计时:start 在进路由前,end 在 body 迭代完成时(或 ASGI 层 on_complete),否则 流式 接口的 duration_ms 会只反映首包时间,误导 日志 追踪。
用 fetch 消费 SSE(可带 Header)
原生 EventSource 不便带 Bearer 时,前端可用 fetch 读 流:
# 服务端仍是 StreamingResponse + text/event-stream
# 客户端伪代码思路:response.body.getReader() 逐块 decode,
# 按 \n\n 切帧,解析 data: 行 — 与 EventSource 同一套 SSE 格式这样 认证 与 流式 可以兼得;代价是前端多写解析逻辑。团队应在 接口 文档标明推荐接入方式,避免移动端用 EventSource 却带不上 token。
何时选普通 JSON 而非 流式
若导购 服务 整段响应小于几 KB、延迟可接受,普通 return GuideOut 更简单,测试与 事务 边界 也更清晰。流式 留给「生成时间明显长于用户耐心」的路径;系列里 导购 可先同步 接口 上线,再增 /guide/stream 作为 UX 增强,两者共用 Service 的 Port,Router 层分岔即可。
完整服务端 流式 骨架
把 认证、disconnect、SSE 帧编码收在一个可读例子里(省略 Port 内部细节):
async def stream_guide_body(
request: Request,
user_id: int,
product_id: int,
guide: ShoppingGuidePort,
) -> AsyncIterator[bytes]:
yield sse_event("", event="start")
try:
async for token in guide.stream_tokens(user_id, product_id):
if await request.is_disconnected():
logger.info("stream_aborted", extra=log_extra(user_id=user_id))
break
yield sse_event(token)
else:
yield sse_event("[DONE]", event="done")
except GuideError as exc:
yield sse_event(json.dumps({"detail": str(exc)}), event="error")路由在 generator 外完成 认证 与参数校验;generator 内只负责 事件 产出与断开检测。错误用 SSE event: error 告知前端,而不是 silent 断连——这是 流式 接口 与 JSON 接口 错误处理的重要差异。
WebSocket 一览(不展开)
双向通道、独立帧协议、常需心跳与 subprotocol 协商。若产品既要推送又要频繁改会话上下文,WebSocket 更合适;仅服务端推 token 序列则 SSE + HTTP 认证 通常更省工程。流式 选型先问「要不要客户端中途发很多消息」,再问「要不要 HTTP 兼容」。
代理与 HTTP/2
部分反向代理对 流式 HTTP/1.1 chunk 支持成熟;HTTP/2 下 流式 仍可行但缓冲策略因产品而异。上线前在「与生产同构的代理」后压测 SSE,观察首 chunk 延迟与 disconnect 行为,比在本地直连 Uvicorn 更能发现假 流式 问题。
与 导购 系列的衔接
同步 导购 接口 与 流式 导购 可共用 ShoppingGuidePort:同步方法 suggest,流式 方法 stream_tokens。Router 一个返回 GuideOut,一个返回 StreamingResponse;Service 层复用 用户 历史与 Catalog 校验逻辑,只在输出形态处分岔。
keep-alive 与注释行
长时间无 chunk 时,SSE 可发 : ping\n\n 注释保活,避免代理空闲断连。generator 每 N 秒 yield 注释帧即可;前端 EventSource 忽略 comment。背压 场景下 comment 不增加 业务 负载。
错误帧约定
event: error 的 data 建议用 JSON,字段与统一异常体一致,前端可复用同一套错误 UI。
TestClient 读 流
集成测试用 client.stream("GET", url) 迭代 lines,断言至少收到 event: start 与 event: done,以及中间若干 data: 行。测 disconnect 行为可 mock Request.is_disconnected 返回 True,断言 generator 提前结束且上游 cancel 被调用——流式 测试与 JSON 测试同样可自动化。
内存与 chunk 大小
极大 token 流若每字 yield 一次,框架开销偏大;可缓冲到词或句再 sse_event。在延迟与吞吐之间折中,通常按 UI 刷新粒度(词级)即可。
StreamingResponse 与 SSE 组合,是 FastAPI 里最常见的「边算边推」实现。事件 帧、text/event-stream、断开检测与代理缓冲四件事做好,流式 接口 才在生产环境真正「流」起来;否则只是代码里用了 generator,用户体验仍像一次性 JSON。
与 JSON 接口 共用同一套 认证、日志、异常码约定,流式 才是整体 架构 的一部分,而不是孤立 hack。下一篇部署用 Docker 打包时,流式 长连接超时配置也要写进 启动 参数与代理,与本地直连行为对齐。
压测 流式 时同时看 p95 TTFB 与总时长:前者反映 流式 价值,后者反映完整生成成本。若 TTFB 改善不明显,可能瓶颈在上游模型而非 FastAPI。
客户端解析 SSE 时要处理 \n\n 分帧与多行 data:;服务端 sse_event helper 统一编码,可减少前后端因换行符理解不一致导致的乱码或截断。
流式 响应头 Cache-Control: no-cache 与 Connection: keep-alive 常与 SSE 一起出现;CDN 若缓存 GET,需排除 流式 路径,否则客户端收到陈旧 事件 流。
长连接占用的 worker 线程/async 任务应在 disconnect 时释放;与 事务 篇一样,流式 也要关心资源何时归还。
Proxy 读超时若小于上游生成总时长,可能在 流式 未完成时切断连接;运维侧超时配置应与产品预期时长同量级,并写入运行手册,避免首上线才暴露代理超时。
SSE 适合服务端单向推送 事件;掌握 StreamingResponse、帧格式与断开处理,就能在 FastAPI 里稳妥落地 流式 体验,并与 JSON 接口 共用 认证 与 日志 约定。
上线前在 staging 用 curl -N 或浏览器 EventSource 试读 流,确认首 chunk 在可接受秒内到达;再在代理后重复,排除缓冲导致的「假 流式」。StreamingResponse 与 SSE 是 FastAPI 交付渐进式 UX 的标准组合,亦能与 导购 等 业务 接口 自然衔接,共用 Port 与 认证 链,减少重复实现,维护 流式 与同步双路由更轻松,扩展成本低。
小结
StreamingResponse 承载 async/sync generator 的字节 流;SSE 用 text/event-stream 与 data:/event: 帧推送 事件。注意断开检测、nginx 缓冲、generator 内错误表达方式,以及 Depends teardown 与连接占用。流式 适合「边算边展示」;小响应、强 事务 一致性、双向通信则选别的形态。把首包延迟、断开与错误帧写进接口约定,上线后排障会顺畅很多。