☰
从手写SSE帧到生产级流式问答:FastAPI + sse-starlette 实战踩坑记
2026/9/26 6:44:24 网站建设 项目流程

最近用 FastAPI 做了一个 RAG 流式问答系统,支持上传 PDF、多轮对话、打字机效果输出。本文记录从手写 SSE 协议帧,到引入 sse-starlette 的全过程,包括两个真实踩过的坑(依赖版本冲突、断连语义)和一个提前预研的问题(放到 Nginx 后面会怎样)。如果你也在做类似的东西,本文能帮你少走弯路。

一、为什么需要流式输出?

普通接口的做法:等大模型把答案全部生成完,再一次性返回 JSON。问题是用户要白屏等 5~10 秒,体验很差。

流式输出的目标是 "打字机效果"—— 大模型吐一个字,前端就显示一个字。技术方案对比:

方案

方向

复杂度

适合场景

普通 JSON 响应

请求 / 响应

低

短回答

SSE(Server-Sent Events)

服务器 → 客户端单向

中等

AI 流式输出(本文用这个)

WebSocket

双向

高

实时聊天、协作编辑

SSE 本质上就是一个普通 HTTP 长连接,服务器可以持续往里面写数据,浏览器自动逐条接收。

二、SSE 帧到底长什么样?

SSE 的数据格式非常简单,每一帧长这样:

data: {"type":"delta","content":"中国"}
注意三个关键点:
  1. data:是协议规定的前缀,不能改;
  2. 后面是一个完整的 JSON 字符串,我们自定义的 type、content 都打包在这个 JSON 里;
  3. 最后必须有一个空行(\n\n),而且这个空行在 JSON 引号外面—— 它是帧与帧之间的分隔符。空行位置写错,浏览器永远收不到事件。

三、第一版:手写 SSE 帧

最开始不引入任何第三方库,自己拼帧。FastAPI 里核心代码长这样:

import json from fastapi import FastAPI from fastapi.responses import StreamingResponse from pydantic import BaseModel app = FastAPI() class AskRequest(BaseModel): question: str def _sse_event(data: dict) -> str: """把事件 dict 序列化为 SSE 帧""" return f"data: {json.dumps(data, ensure_ascii=False)}\n\n" @app.post("/ask/stream") async def ask_stream(req: AskRequest): question = req.question.strip() if not question: raise HTTPException(status_code=400, detail="question 不能为空") def generate(): full_answer = "" for event in rag.ask_stream(question): if event["type"] == "delta": full_answer += event["content"] yield _sse_event({"type": "delta", "content": event["content"]}) elif event["type"] == "sources": yield _sse_event({"type": "sources", "sources": event["sources"]}) yield _sse_event({"type": "done"}) return StreamingResponse(generate(), media_type="text/event-stream")

这版能跑,但手写协议有三个坑:

坑 1:\n\n漏一个字符,前端就收不到事件。SSE 规定空行才是事件结束,少一个 \n,浏览器会一直等下一帧。

坑 2:中文必须写ensure_ascii=False。不写的话,"中国" 会被序列化成 \u4e2d\u56fd,帧体积变大,调试时也看不懂。

坑 3:StreamingResponse要手动设media_type="text/event-stream"。不设的话浏览器不知道这是 SSE 流,会当作普通 JSON 响应等全部下载完。

手写帧的好处是能看清协议本质,但生产环境不该自己维护这些细节 —— 容易错,还少了心跳、断线重连这些能力。

四、引入 sse-starlette,结果依赖冲突

标准做法是用 sse-starlette 这个库,它封装了 EventSourceResponse 和 ServerSentEvent:

pip install sse-starlette

然后改造代码:

from sse_starlette.sse import EventSourceResponse, ServerSentEvent @app.post("/ask/stream") async def ask_stream(req: AskRequest): question = req.question.strip() if not question: raise HTTPException(status_code=400, detail="question 不能为空") def generate(): full_answer = "" for event in rag.ask_stream(question): if event["type"] == "delta": full_answer += event["content"] yield ServerSentEvent(data={"type": "delta", "content": event["content"]}) elif event["type"] == "sources": yield ServerSentEvent(data={"type": "sources", "sources": event["sources"]}) yield ServerSentEvent(data={"type": "done"}) return EventSourceResponse(generate())

改动点:

  • _sse_event 辅助函数整个删掉,ServerSentEvent(data=...) 自动帮你做序列化和拼帧;
  • StreamingResponse(..., media_type=...) 换成 EventSourceResponse(generate()),MIME 类型自动设置。

踩坑:版本冲突

新装完启动直接报错:

fastapi 0.115.0 requires starlette<0.39.0,>=0.37.2, but you have starlette 1.7.0 which is incompatible.

原因:最新版 sse-starlette 3.x 要求 starlette ≥ 0.49,把 starlette 拉到了 1.7.0;而我的 FastAPI 0.115.0 锁死 starlette < 0.39。两个要求完全不重叠。

解决方法是降版本:

pip install "sse-starlette==1.8.2"

pip 会自动把 starlette 降回 0.38.6,和 FastAPI 匹配。这个版本的 API 和 3.x 用法完全一样,代码不用改。

五、断连后发生了什么?

有个问题我一开始想当然了:用户回答到一半关掉浏览器,数据库会留下什么?

我的第一反应是 "大模型继续生成完,然后存库"。因为实际生活里的大模型是这样,但是我们目前做出来的demo还不是。

真实流程是:

  1. 用户关浏览器 → TCP 连接断开;
  2. Starlette 检测到客户端断开,关闭生成器(抛 GeneratorExit);
  3. for event in rag.ask_stream(...) 循环被打断;
  4. 循环后面的 db.save_message(...)根本不会执行;
  5. full_answer 拼到一半就随生成器销毁。

也就是说:断连后历史表里不会留下半截回答。这其实是好事 —— 历史表里不会出现 "用户问了一半、回答了两个字" 的脏数据。

正常流程:请求 → 读历史 → 流式生成 → 存用户问题 → 存完整回答 → done 断连流程:请求 → 读历史 → 流式生成一半 → 连接断 → 循环中断 → 什么都不存

如果产品要求 "哪怕断连也要存完整回答",就需要把生成任务放到后台协程里跑,让它独立于客户端连接。这个改动大概 40 行代码,demo 阶段可以zanshi需要做。

六、提前想一下:放到 Nginx 后面会遇到什么坑?

本地跑通后我在想:SSE 这种长时间保持连接的接口,放到反向代理后面会不会有问题?查了一下 Nginx 默认配置,果然有个 proxy_read_timeout:

proxy_read_timeout 60s;

它的意思是:Nginx 等后端数据超过 60 秒还没收到,就主动断开连接。那问题就来了:

  1. 大模型在 "思考"(检索资料、推理),这段时间没有 token 输出;
  2. SSE 连接上 60 秒没有任何数据流过;
  3. Nginx 判定超时,断开连接;
  4. 前端表现为 "卡住然后失败"。

这也是为什么 EventSourceResponse 自带心跳:每隔几秒自动发一个 SSE 注释帧 : ping\n\n,即使大模型没产出新 token,连接上也一直有数据流过,Nginx 就不会误判超时。

这也是为什么手写 StreamingResponse 虽然能跑,但生产环境更推荐 sse-starlette—— 心跳这种细节库已经帮你处理好了,自己写容易漏。

EventSourceResponse 自带心跳:每隔几秒自动发一个 SSE 注释帧 : ping\n\n,即使大模型没产出新 token,连接上也一直有数据流过,Nginx 就不会误判超时。

如果坚持用 StreamingResponse,需要自己在生成器里加心跳任务,定期往连接里写注释帧。

七、上传 PDF 的三道防线

顺便记录上传接口的安全设计,这部分和流式无关,但做 RAG 都要用到:

ALLOWED_EXTENSIONS = {".pdf"} MAX_UPLOAD_SIZE = 50 * 1024 * 1024 # 50MB @app.post("/upload") async def upload_pdf(file: UploadFile = File(...)): original_name = file.filename or "upload.pdf" ext = pathlib.Path(original_name).suffix.lower() if ext not in ALLOWED_EXTENSIONS: raise HTTPException(status_code=400, detail="仅支持 PDF 文件") content = await file.read() if len(content) > MAX_UPLOAD_SIZE: raise HTTPException(status_code=413, detail="文件超过 50MB 限制") # 防路径穿越:存储名用 uuid,原文件名只做展示 stored_name = f"{uuid.uuid4().hex}{ext}" pdf_path = UPLOAD_DIR / stored_name pdf_path.write_bytes(content) ...

三道防线:

  1. 扩展名白名单:只允许 PDF;
  2. 大小限制:50MB,防止超大文件撑爆内存;
  3. UUID 重命名:用户文件名可能是 ../../etc/passwd,直接拼路径会写穿目录。用随机名存储,原文件名只作为元数据。

八、最终项目结构

KubeRAG/ ├── app/ │ ├── main.py # FastAPI 路由层 │ ├── rag.py # PDF 切片、向量化、Chroma 检索 │ ├── db.py # SQLite 会话历史 │ └── agent.py # 工具调用 Agent ├── data/ │ ├── uploads/ # PDF 原文 │ └── chat_history.db ├── static/ │ └── index.html # 前端页面 └── requirements.txt

完整代码已上传 GitHub:BowliceDXY/KubeRAG (github.com)

九、面试可能会追问的问题

写这篇文章的过程中,我自己整理了几个面试官大概率会问的问题,供参考:

  1. SSE 和 WebSocket 怎么选?—— 单向推送选 SSE,双向交互选 WebSocket;
  2. \n\n为什么在 JSON 外面?—— 它是帧分隔符,不是数据内容;
  3. 用户中途关页面,历史会存半截吗?—— 不会,生成器被关闭,save_message 不执行;
  4. Nginx 超时断连怎么解决?—— 心跳 ping,保持连接活跃;
  5. sse-starlette 和手写 StreamingResponse 区别?—— 封装了序列化、MIME、心跳,少写协议细节。

总结

做流式输出本身不难,难的是处理那些教程文章里不太会写的边缘情况:依赖版本冲突、断连后发生什么、代理超时怎么办。

写这篇文章的初衷不是当教程,而是把自己这几天学习过程中踩过的坑完整记录下来。如果其中某一段刚好帮到正在做同样事情的同学,那就真的太好了。

我也是边学边做,文章里的理解不一定全对。如果你发现哪里有问题、或者有更优雅的实现方式,欢迎评论区指出来,咱们一起讨论。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询