K 的一隅

Python FastAPI 后端实战

一边算一边推:流式响应与 SSE

StreamingResponse 与 SSE 事件格式;生成器逐块写出;背压与客户端断开;何时不必流式。

11 分钟阅读 更新于 2026-07-28
本文目录

大模型 导购 或日志 tail:用户盯着空白页等 30 秒,体验像卡死。若每生成一段就推一块 JSON,浏览器可以边收边渲染——这就是 流式 HTTP:响应体不是一次性算完再 return,而是持续写出SSE(Server-Sent Events)是 流式 里常用的一种 事件 格式,单向、基于 HTTP、前端 EventSource 能直接消费。FastAPI 用 StreamingResponse 承载字节 ;理解 事件 帧格式、背压 与断开,比背 API 名更重要。

流式 不是银弹:它改善的是「何时看到第一字节」,不改变 业务 正确性责任。该短 事务 提交 的写操作,不应因为输出 流式 就拖到 generator 里 commit。

普通响应 vs 流式响应

普通 return dictStreamingResponse
内存整份 body 在内存可逐块生成
TTFB业务算完才发首块算完即发
中间件完整经过首包后 body 不再走部分中间件
客户端一次读完长连接读直到结束

ASGI 层:StreamingResponse 把 async generator / iterator 的产出 repeatedly send 给客户端。

从时间线看:客户端发请求 → 服务端返回 200 与 Header → body 尚未完整 → 连接保持 → 多个 chunk 顺序到达 → 连接关闭。TTFB(首字节时间)往往比总时长更能改善体感;流式 优化的是「什么时候开始看到内容」,不是峰值 QPS。

StreamingResponse 最小例

python
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-ndjsontext/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 帧的小函数:

python
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")

路由:

python
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 选型

SSENDJSON
格式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:

python
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。

测试

TestClientStreamingResponse 支持 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 读

python
# 服务端仍是 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 内部细节):

python
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: errordata 建议用 JSON,字段与统一异常体一致,前端可复用同一套错误 UI。

TestClient 读

集成测试用 client.stream("GET", url) 迭代 lines,断言至少收到 event: startevent: done,以及中间若干 data: 行。测 disconnect 行为可 mock Request.is_disconnected 返回 True,断言 generator 提前结束且上游 cancel 被调用——流式 测试与 JSON 测试同样可自动化。

内存与 chunk 大小

极大 token 流若每字 yield 一次,框架开销偏大;可缓冲到词或句再 sse_event。在延迟与吞吐之间折中,通常按 UI 刷新粒度(词级)即可。

StreamingResponseSSE 组合,是 FastAPI 里最常见的「边算边推」实现。事件 帧、text/event-stream、断开检测与代理缓冲四件事做好,流式 接口 才在生产环境真正「流」起来;否则只是代码里用了 generator,用户体验仍像一次性 JSON。

与 JSON 接口 共用同一套 认证日志、异常码约定,流式 才是整体 架构 的一部分,而不是孤立 hack。下一篇部署用 Docker 打包时,流式 长连接超时配置也要写进 启动 参数与代理,与本地直连行为对齐。

压测 流式 时同时看 p95 TTFB 与总时长:前者反映 流式 价值,后者反映完整生成成本。若 TTFB 改善不明显,可能瓶颈在上游模型而非 FastAPI。

客户端解析 SSE 时要处理 \n\n 分帧与多行 data:;服务端 sse_event helper 统一编码,可减少前后端因换行符理解不一致导致的乱码或截断。

流式 响应头 Cache-Control: no-cacheConnection: keep-alive 常与 SSE 一起出现;CDN 若缓存 GET,需排除 流式 路径,否则客户端收到陈旧 事件 流。

长连接占用的 worker 线程/async 任务应在 disconnect 时释放;与 事务 篇一样,流式 也要关心资源何时归还。

Proxy 读超时若小于上游生成总时长,可能在 流式 未完成时切断连接;运维侧超时配置应与产品预期时长同量级,并写入运行手册,避免首上线才暴露代理超时。

SSE 适合服务端单向推送 事件;掌握 StreamingResponse、帧格式与断开处理,就能在 FastAPI 里稳妥落地 流式 体验,并与 JSON 接口 共用 认证日志 约定。

上线前在 staging 用 curl -N 或浏览器 EventSource 试读 ,确认首 chunk 在可接受秒内到达;再在代理后重复,排除缓冲导致的「假 流式」。StreamingResponseSSE 是 FastAPI 交付渐进式 UX 的标准组合,亦能与 导购业务 接口 自然衔接,共用 Port 与 认证 链,减少重复实现,维护 流式 与同步双路由更轻松,扩展成本低。

小结

StreamingResponse 承载 async/sync generator 的字节 SSEtext/event-streamdata:/event: 帧推送 事件。注意断开检测、nginx 缓冲、generator 内错误表达方式,以及 Depends teardown 与连接占用。流式 适合「边算边展示」;小响应、强 事务 一致性、双向通信则选别的形态。把首包延迟、断开与错误帧写进接口约定,上线后排障会顺畅很多。