FastAPI 怎么做流式响应和 SSE?大模型流式输出怎么实现?
简化版
FastAPI 的流式响应靠 StreamingResponse——传入一个(异步)生成器,服务端边产出边发送,客户端不用等全部数据准备好。它有三类典型用途:① 大文件/大数据集导出(几百万行 CSV 不用先在内存里拼成一个大字符串);② SSE(Server-Sent Events)——服务端单向推送,格式固定为 data: xxx\n\n,浏览器用 EventSource 接收并自动重连;③ 大模型流式输出——把 LLM 逐 token 产出的内容实时转发给前端,这是现在最常见的场景。SSE 需要三个响应头:Content-Type: text/event-stream、Cache-Control: no-cache、以及 X-Accel-Buffering: no(关掉 Nginx 缓冲,否则流式会变成一次性返回)。最容易踩的三个坑:① 中间件破坏流式——BaseHTTPMiddleware(即 @app.middleware("http"))会缓冲整个响应体,让 SSE 和流式下载失效,要改用纯 ASGI 中间件;② 生成器里的阻塞会卡住整个事件循环(同步 for 循环里做重计算、用 requests 调上游),必须全程异步或丢线程池;③ 客户端断开时要能感知并停止——用户关掉页面后如果服务端还在跑(尤其是还在烧 LLM token),就是实打实的浪费,靠 await request.is_disconnected() 或捕获 asyncio.CancelledError 处理。SSE vs WebSocket 的选择:只需要服务端单向推送就用 SSE(走普通 HTTP、自动重连、能被代理和 CDN 正常处理、实现简单);需要双向实时通信才用 WebSocket。核心记忆:StreamingResponse + 异步生成器;SSE 格式是 data: ...\n\n;关 Nginx 缓冲和 gzip;别用 BaseHTTPMiddleware;要处理客户端断开。
详细版
三种流式场景对比:
| 大文件/导出 | SSE | WebSocket | |
|---|---|---|---|
| 方向 | 服务端 → 客户端 | 服务端 → 客户端 | 双向 |
| 协议 | 普通 HTTP | 普通 HTTP | 独立协议(升级) |
| 自动重连 | — | ✅ 浏览器内置 | ❌ 要自己实现 |
| 代理友好 | ✅ | ✅ | 需要额外配置 |
| 实现复杂度 | 低 | 低 | 中 |
| 适用 | 导出、下载 | 推送、LLM 流式 | 聊天、协作、游戏 |
from fastapi import FastAPI, Request
from fastapi.responses import StreamingResponse
import asyncio, json
app = FastAPI()
# ① ★基本流式:CSV 导出★
@app.get("/export.csv")
async def export_csv(db: DbDep):
async def generate():
yield "id,name,email\n"
result = await db.stream(select(User)) # ★★流式查询★★
async for user in result.scalars():
yield f"{user.id},{user.name},{user.email}\n"
return StreamingResponse(
generate(),
media_type="text/csv",
headers={"Content-Disposition": 'attachment; filename="users.csv"'},
)
# ② ★★SSE:服务器推送事件★★
@app.get("/events")
async def events(request: Request):
async def event_generator():
while True:
if await request.is_disconnected(): # ★★检测断开★★
break
msg = await get_next_message()
yield f"data: {json.dumps(msg)}\n\n" # ★★格式:两个换行结束★★
return StreamingResponse(
event_generator(),
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache",
"Connection": "keep-alive",
★"X-Accel-Buffering": "no"★, # ★★关 Nginx 缓冲★★
},
)
# ★完整的 SSE 消息格式★
# event: update\n ← 事件名(可选,前端 addEventListener 用)
# id: 42\n ← 事件 ID(断线重连时浏览器会带 Last-Event-ID)
# retry: 3000\n ← 重连间隔(毫秒)
# data: {"x": 1}\n ← 数据(可多行)
# \n ← ★★空行表示一条消息结束★★
# ③ ★★大模型流式输出(最常见场景)★★
@app.post("/chat")
async def chat(req: ChatRequest, request: Request):
async def token_stream():
try:
stream = await openai_client.chat.completions.create(
model="gpt-4", messages=req.messages, stream=True)
async for chunk in stream:
if await request.is_disconnected(): # ★★用户关页面就停★★
break
delta = chunk.choices[0].delta.content
if delta:
yield f"data: {json.dumps({'delta': delta})}\n\n"
yield "data: [DONE]\n\n" # ★★约定的结束标记★★
except asyncio.CancelledError:
logger.info("客户端断开,停止生成")
raise
except Exception as e:
logger.exception("生成失败")
yield f"data: {json.dumps({'error': str(e)})}\n\n" # ★★错误也要发出去★★
return StreamingResponse(token_stream(),
media_type="text/event-stream",
headers={"X-Accel-Buffering": "no"})
# ④ ★心跳(防止代理超时断开)★
async def with_heartbeat(gen, interval=15):
queue = asyncio.Queue()
async def producer():
async for item in gen:
await queue.put(item)
await queue.put(None)
task = asyncio.create_task(producer())
try:
while True:
try:
item = await asyncio.wait_for(queue.get(), timeout=interval)
if item is None: break
yield item
except asyncio.TimeoutError:
yield ": heartbeat\n\n" # ★★注释行,客户端忽略★★
finally:
task.cancel()
# ⑤ ★前端接收★
# const es = new EventSource("/events");
# es.onmessage = (e) => { if (e.data === "[DONE]") es.close();
# else render(JSON.parse(e.data)); };
# es.onerror = () => { /* ★浏览器会自动重连★ */ };
# ★注意:EventSource 只支持 GET 且不能自定义请求头★
# → ★要用 POST 或带 Authorization 就得用 fetch + ReadableStream★
⚠️ 三个必须记住的点:①
BaseHTTPMiddleware(@app.middleware("http"))会破坏流式响应。它的实现方式是把下游返回的响应完整收集后再转发,所以 SSE 会变成「等生成器全部跑完才一次性返回」——表现为「代码看起来对,但前端一直等到最后才收到全部内容」。这是流式功能最经典的踩坑,而且很难联想到是中间件的问题。解法是改写成纯 ASGI 中间件(实现async def __call__(self, scope, receive, send),用send_wrapper包装 send),或者至少确认应用里没有注册任何BaseHTTPMiddleware。② 链路上的每一层缓冲都会让流式失效:Nginx 的proxy_buffering(默认开启) 要用X-Accel-Buffering: no响应头关掉;gzip 压缩会攒够一块才输出,所以流式接口要排除在 GZipMiddleware 之外;CDN 通常也会缓冲,SSE 端点应该绕过 CDN。排查顺序是先用curl -N直连应用端口,如果那样是流式的,问题就在代理层。③ 必须处理客户端断开。用户关闭页面或切走后,如果服务端还在继续生成(尤其是还在调用 LLM 烧 token、或者还在跑一个大查询),就是实打实的资源和金钱浪费。两种方式:主动检查await request.is_disconnected()(在循环里判断),或者捕获asyncio.CancelledError(Starlette 在检测到断开时会取消任务)——后者还应该在finally里释放资源。
完整版教学
一、StreamingResponse 的机制
★ 普通响应 vs 流式响应:
┌──────────────────────────────────────────────────────────┐
│ ★普通响应★ │
│ 视图 return 数据 → ★序列化成完整的 body★ │
│ → 一次性 send(含 Content-Length) │
│ ★客户端等待全部数据准备好★ │
├──────────────────────────────────────────────────────────┤
│ ★StreamingResponse★ │
│ 视图 return 生成器 → ★ASGI 逐块 send★ │
│ → ★Transfer-Encoding: chunked★(没有 Content-Length) │
│ ★客户端边收边处理★ │
└──────────────────────────────────────────────────────────┘
★ ★ASGI 层面发生了什么★:
await send({"type": "http.response.start", "status": 200,
"headers": [...]}) # ★先发头★
async for chunk in generator:
await send({"type": "http.response.body",
"body": chunk, ★"more_body": True★}) # ★★逐块★★
await send({"type": "http.response.body", "body": b"",
★"more_body": False★}) # ★结束★
★ ★more_body 就是流式的核心★
★ ★同步生成器 vs 异步生成器★:
# ★异步生成器(★推荐★)★
async def gen():
async for row in db.stream(...):
yield f"{row}\n"
# ★同步生成器(Starlette 会用 iterate_in_threadpool 包装)★
def gen():
for row in cursor: # ★阻塞操作在线程池里★
yield f"{row}\n"
★ ✓ 同步生成器不会阻塞事件循环(★Starlette 帮你丢线程池★)
★ ✗ 但★受 40 线程池限制★
★ ★★生成器里的阻塞是致命的(异步生成器)★★:
async def gen():
for i in range(1000):
data = requests.get(url) # ★★阻塞事件循环★★
yield data
→ ★整个进程的所有请求都被卡住★
✓ 用 httpx.AsyncClient 或 run_in_threadpool
★ ★Content-Length 的消失★:
流式响应★没有 Content-Length★(因为长度未知)
→ 用 ★Transfer-Encoding: chunked★
→ ★影响:前端无法显示下载进度百分比★
✓ 如果长度已知,可以手动设:
StreamingResponse(gen(), headers={"Content-Length": str(size)})
★ ★三个相关的 Response 类★:
┌────────────────────┬──────────────────────────────────┐
│ ★StreamingResponse★ │ ★任意生成器★ │
│ ★FileResponse★ │ ★发送文件(内部就是流式 + sendfile)★│
│ Response │ 普通完整响应 │
└────────────────────┴──────────────────────────────────┘
# 大文件下载优先用 FileResponse
return FileResponse(path, filename="x.pdf",
media_type="application/pdf")
★ ★更好:受保护文件用 X-Accel-Redirect 交给 Nginx★
流式响应在 ASGI 层面的核心是 more_body——先 send 一个 http.response.start 发响应头,然后逐块 send http.response.body 并标记 more_body: True,最后发一个空 body 且 more_body: False 结束。因为长度未知,流式响应没有 Content-Length、走 Transfer-Encoding: chunked——副作用是前端无法显示下载进度百分比(长度已知时可以手动设)。同步生成器 Starlette 会用 iterate_in_threadpool 包装(不会阻塞事件循环,但受 40 线程池限制);而异步生成器里的阻塞调用是致命的——requests.get() 会卡住整个进程。大文件下载优先用 FileResponse(内部就是流式),受保护的文件更应该用 X-Accel-Redirect 交给 Nginx。
二、SSE 协议细节
★ ★SSE 的报文格式(★纯文本,很简单★)★:
data: hello\n
\n ← ★★空行 = 一条消息结束★★
★完整字段:★
event: update\n ← 事件名(前端 addEventListener("update", ...))
id: 42\n ← ★事件 ID★
retry: 3000\n ← 重连间隔(毫秒)
data: line1\n ← ★data 可以多行★
data: line2\n
\n
★注释行(★用作心跳★):★
: this is a comment\n\n ← ★客户端会忽略,但能保持连接活跃★
★ ★三个必须的响应头★:
Content-Type: ★text/event-stream★
Cache-Control: ★no-cache★
★X-Accel-Buffering: no★ ← ★★Nginx 专用,不加就不是流式★★
(Connection: keep-alive 通常由服务器自动处理)
★ ★★断线重连(SSE 的杀手锏)★★:
① 浏览器的 EventSource ★检测到连接断开会自动重连★
② 重连时★自动带上 Last-Event-ID 请求头★(值是最后收到的 id)
③ 服务端可以据此★续传★:
@app.get("/events")
async def events(request: Request):
last_id = request.headers.get("last-event-id")
async def gen():
async for msg in get_messages_after(last_id):
yield f"id: {msg.id}\ndata: {msg.body}\n\n"
return StreamingResponse(gen(), media_type="text/event-stream")
★ ★这是 WebSocket 需要自己实现而 SSE 白送的能力★
★ ★EventSource 的限制(★决定了什么时候不能用它★)★:
✗ ★只支持 GET★(不能 POST 大的请求体)
✗ ★不能自定义请求头★(★所以不能放 Authorization!★)
✗ ★HTTP/1.1 下浏览器对同一域名有 6 个连接的限制★
→ ★多个标签页打开就会互相挤占★(HTTP/2 下无此问题)
✓ 绕过方案:★用 fetch + ReadableStream 自己解析★
const resp = await fetch("/chat", {
method: "POST",
headers: {"Authorization": `Bearer ${token}`},
body: JSON.stringify(payload),
});
const reader = resp.body.getReader();
const decoder = new TextDecoder();
let buffer = "";
while (true) {
const {done, value} = await reader.read();
if (done) break;
buffer += decoder.decode(value, {stream: true});
const lines = buffer.split("\n\n");
buffer = lines.pop(); // ★★保留不完整的部分★★
for (const line of lines) { /* 解析 data: */ }
}
★ ★LLM 聊天基本都用这个方案★(要 POST + 要带 token)
★ ★心跳的必要性★:
★问题:★代理/负载均衡通常有 60 秒的空闲超时★
→ ★长时间没数据的 SSE 连接会被切断★
✓ ★每 15~30 秒发一个注释行★:
yield ": ping\n\n"
✓ 或发一个 event: ping 事件
★ ★错误处理★:
★问题:★HTTP 状态码在响应头发出后就固定了★
→ ★流到一半出错,没法改成 500★
✓ 只能★在流里发一条错误消息★:
yield f"data: {json.dumps({'error': 'xxx'})}\n\n"
✓ ★前端要能识别这种"业务内错误"★
★ ★所以:能在开始流之前校验的,一定要在之前做★
SSE 的报文格式很简单:data: xxx 加一个空行结束一条消息,还有 event(事件名)、id(事件 ID)、retry(重连间隔)三个可选字段,以及用作心跳的注释行 : ping。三个必须的响应头是 text/event-stream、no-cache、X-Accel-Buffering: no。SSE 的杀手锏是断线重连——浏览器自动重连并带上 Last-Event-ID 请求头,服务端据此续传,这是 WebSocket 需要自己实现而 SSE 白送的能力。但 EventSource 有三个硬限制:只支持 GET、不能自定义请求头(所以放不了 Authorization)、HTTP/1.1 下有 6 连接限制——所以 LLM 聊天基本都改用 fetch + ReadableStream 自己解析。还有两点:代理通常有 60 秒空闲超时,要每 15~30 秒发心跳;响应头一旦发出状态码就固定了,流到一半出错只能在流里发错误消息——所以能在开始流之前校验的一定要提前做。
三、大模型流式输出
★ ★完整链路★:
前端 fetch(POST /chat) → FastAPI → LLM API(stream=True)
↓ 逐 token
前端 ReadableStream ← SSE ← 转发 ←──┘
★ ★服务端实现要点★:
@app.post("/chat")
async def chat(req: ChatIn, request: Request, user: CurrentUser):
# ★★① 流开始前做完所有校验★★
if not await check_quota(user):
raise HTTPException(429, "额度不足") # ★这时还能返回正常的错误码★
async def stream():
full_text = []
try:
async with llm_client.stream(req.messages) as resp:
async for token in resp:
# ★★② 检测断开★★
if await request.is_disconnected():
logger.info("client disconnected, aborting")
break
full_text.append(token)
# ★★③ 每个 token 一条 SSE★★
yield f"data: {json.dumps({'delta': token})}\n\n"
yield "data: [DONE]\n\n"
except asyncio.CancelledError:
logger.info("stream cancelled")
raise # ★★必须重新抛出★★
except Exception as e:
logger.exception("llm error")
yield f"data: {json.dumps({'error': '生成失败'})}\n\n"
finally:
# ★★④ 无论如何都要记账和保存★★
await save_message(user, "".join(full_text))
await deduct_quota(user, len(full_text))
return StreamingResponse(stream(), media_type="text/event-stream",
headers={"X-Accel-Buffering": "no"})
★ ★四个关键设计★:
① ★校验前置★:一旦开始流就无法返回 4xx/5xx
② ★断开检测★:用户关页面 → ★立刻停止烧 token★
③ ★结束标记★:约定 [DONE],前端据此关闭连接
④ ★finally 里落库★:★哪怕中断也要保存已生成的部分 + 扣费★
★ ★为什么必须处理断开(★真金白银★)★:
用户提问后关掉页面
✗ 不处理:★LLM 继续生成完整回答★ → ★token 照样扣费★
✓ 处理:★检测到断开立刻 break★ → 省下剩余的钱
★ ★大模型场景下这是实打实的成本★
★ ★asyncio.CancelledError 的正确处理★:
★ Starlette 检测到客户端断开时会 ★cancel 这个 task★
★ 在生成器里表现为 ★抛出 CancelledError★
✓ except asyncio.CancelledError:
cleanup()
★raise★ # ★★必须重新抛出,不能吞掉★★
★ ✗ 吞掉会导致任务无法正常取消
★ ★首 token 延迟(TTFT)优化★:
★ 用户感知的"快"主要看 ★第一个 token 多久出来★
✓ ★流开始前的操作要快★(鉴权、检索、拼 prompt)
✓ ★先发一个空事件表示"已开始"★:
yield "data: {}\n\n" # ★立刻让前端进入"生成中"状态★
✓ ★RAG 场景:先把检索到的引用发出去,再流式生成正文★
★ ★并发与限流★:
★ 每个流式请求都★占用一个连接和一个任务★
✓ ★用信号量限制同时进行的生成数★:
sem = asyncio.Semaphore(50)
async def stream():
async with sem: ...
✓ ★超过就返回 429★,而不是让所有人都变慢
★ ★前端的完整处理★:
const resp = await fetch("/chat", {method:"POST", headers:{...},
body: JSON.stringify(msg),
signal: abortController.signal});
// ★★signal 让用户能主动取消(触发服务端的断开检测)★★
const reader = resp.body.getReader();
...
// 组件卸载时:abortController.abort();
大模型流式输出有四个关键设计:① 校验前置(一旦开始流就无法返回 4xx/5xx,额度检查、参数校验都要在之前做);② 断开检测(用户关页面就立刻停止——大模型场景下这是实打实的成本,不处理就是继续烧钱);③ 约定结束标记 [DONE];④ finally 里落库和扣费(哪怕中断也要保存已生成的部分)。asyncio.CancelledError 必须重新抛出,不能吞掉。首 token 延迟(TTFT)是用户感知「快慢」的关键——流开始前的操作要快,可以先发一个空事件让前端立刻进入「生成中」状态,RAG 场景可以先把检索到的引用发出去再流式生成正文。并发上要用信号量限制同时进行的生成数、超过就返回 429。前端要用 AbortController 的 signal 让用户能主动取消(这样服务端才能检测到断开)。
四、中间件与代理的坑
★ ★★坑一:BaseHTTPMiddleware 破坏流式★★
@app.middleware("http") # ★★基于 BaseHTTPMiddleware★★
async def add_header(request, call_next):
resp = await call_next(request)
resp.headers["X-Time"] = "..."
return resp
→ ★★SSE 变成"全部生成完才一次性返回"★★
★ 原因:BaseHTTPMiddleware 内部用 ★内存流桥接 ASGI★,
会等下游把响应体全部产出
✓ ★纯 ASGI 中间件(不影响流式)★:
class TimingMiddleware:
def __init__(self, app): self.app = app
async def __call__(self, scope, receive, send):
if scope["type"] != "http":
return await self.app(scope, receive, send)
start = time.perf_counter()
async def send_wrapper(message):
if message["type"] == "http.response.start":
dur = (time.perf_counter() - start) * 1000
message["headers"].append(
(b"x-process-time", f"{dur:.1f}".encode()))
await send(message) # ★★逐条转发,不缓冲★★
await self.app(scope, receive, send_wrapper)
app.add_middleware(TimingMiddleware)
★ ★坑二:GZipMiddleware★
app.add_middleware(GZipMiddleware, minimum_size=1000)
→ ★压缩需要攒够数据才能输出★ → ★破坏流式★
✓ 排除流式路径:
① 不用 GZipMiddleware,交给 Nginx(★但 Nginx 也要为 SSE 关 gzip★)
② 或自定义中间件跳过 text/event-stream
★ ★坑三:Nginx 缓冲★
# ★方案一:响应头(推荐,应用侧可控)★
headers={"X-Accel-Buffering": "no"}
# ★方案二:Nginx 配置★
location /events {
proxy_pass http://app;
★proxy_buffering off;★
★proxy_cache off;★
★proxy_read_timeout 3600s;★ # ★★长连接要调大★★
★gzip off;★
proxy_http_version 1.1;
proxy_set_header Connection ""; # ★★清掉 Connection: close★★
}
★ ★坑四:超时配置★
┌────────────────────┬──────────────────────────────┐
│ ★Nginx★ │ proxy_read_timeout(默认 60s)│
│ ★负载均衡/云 LB★ │ 空闲超时(★通常 60s,要调大★)│
│ ★uvicorn★ │ --timeout-keep-alive │
│ ★客户端★ │ fetch 的超时 │
└────────────────────┴──────────────────────────────┘
★ ★任何一层超时都会切断连接★ → 心跳 + 调大超时
★ ★坑五:TestClient 测流式★
✗ r = client.get("/events")
r.text # ★会一直读到结束(可能永远不结束)★
✓ with client.stream("GET", "/events") as r:
for line in r.iter_lines():
if line.startswith("data:"): ...
if 收够了: break
★ ★坑六:多 worker 下的广播★
★问题:★SSE 连接分散在不同 worker 进程★
→ ★进程 A 里的事件无法推给连接在进程 B 的客户端★
✓ ★用 Redis Pub/Sub 做跨进程广播★:
async def event_generator():
pubsub = redis.pubsub()
await pubsub.subscribe("events")
async for msg in pubsub.listen():
yield f"data: {msg['data']}\n\n"
★ ★这是"聊天室/通知"类功能的标准架构★
六个坑里前三个都和「缓冲」有关:BaseHTTPMiddleware 内部用内存流桥接 ASGI,会等下游把响应体全部产出——解法是写纯 ASGI 中间件用 send_wrapper 逐条转发;GZipMiddleware 需要攒够数据才能压缩输出,要排除流式路径;Nginx 的 proxy_buffering 默认开启,用 X-Accel-Buffering: no 响应头或在 location 里 proxy_buffering off。超时是第四个坑——Nginx、云 LB、uvicorn、客户端任何一层超时都会切断连接,要靠心跳加调大超时。测试流式要用 client.stream() 而不是直接读 r.text。最后一个是架构级的坑:多 worker 下 SSE 连接分散在不同进程,进程 A 的事件推不给连在进程 B 的客户端——要用 Redis Pub/Sub 做跨进程广播,这是聊天室和通知类功能的标准架构。
五、SSE vs WebSocket vs 轮询
★ ★三者对比★:
┌──────────┬──────────┬────────────┬────────────────┐
│ │ ★SSE★ │ ★WebSocket★ │ ★轮询★ │
├──────────┼──────────┼────────────┼────────────────┤
│ 方向 │ ★单向★ │ ★双向★ │ 单向(客户端拉)│
│ 协议 │ ★HTTP★ │ ws://(升级)│ HTTP │
│ 重连 │ ★自动★ │ 手动实现 │ 天然 │
│ 代理兼容 │ ★好★ │ ★需配置★ │ ★最好★ │
│ 认证 │ ★受限★ │ ★受限★ │ ★简单★ │
│ 二进制 │ ✗ │ ★✓★ │ ✓ │
│ 实现成本 │ ★低★ │ 中 │ ★最低★ │
│ 服务端资源│ 一连接一任务│ 同左 │ ★无长连接★ │
└──────────┴──────────┴────────────┴────────────────┘
★ ★决策树★:
需要客户端主动发消息吗?
├─ ★需要(聊天、协作编辑、游戏)★ → ★WebSocket★
└─ 不需要(只是服务端推)
├─ ★更新很频繁 / 要低延迟★ → ★SSE★
│ (LLM 流式输出、实时日志、进度推送、通知)
└─ ★更新不频繁(分钟级)★ → ★轮询就够了★
(★别为了"实时"上长连接,轮询简单可靠★)
★ ★SSE 相比 WebSocket 的优势(★常被低估★)★:
✓ ★就是普通的 HTTP★ → 现有的鉴权、日志、监控、限流全都能用
✓ ★浏览器自动重连 + Last-Event-ID 续传★
✓ ★代理和 CDN 不用特殊配置★(WebSocket 要配 Upgrade 头转发)
✓ ★调试简单★(curl 就能看)
✓ ★服务端实现是一个普通的路由★
★ ★WebSocket 才有的能力★:
✓ ★客户端 → 服务端的实时消息★
✓ ★二进制帧★
✓ ★更低的单条消息开销★(SSE 每条都有文本头)
★ ★轮询也别看不起★:
# 短轮询
setInterval(() => fetch("/status").then(...), 5000);
# ★长轮询(服务端 hold 住直到有数据或超时)★
@app.get("/poll")
async def poll(since: int):
try:
msg = await asyncio.wait_for(get_new_message(since), timeout=30)
return msg
except asyncio.TimeoutError:
return {"messages": []} # ★★空返回,客户端立刻再来★★
★ ✓ ★实现最简单、兼容性最好、无状态★
★ ✗ 延迟高(短轮询)或占连接(长轮询)
★ ★资源消耗的现实考量★:
★ 一个 SSE/WebSocket 连接 = ★一个 asyncio 任务 + 一个 socket★
★ ASGI 下几千个空闲连接问题不大(内存约几 KB/连接)
★ ✗ ★但 WSGI(Flask/Django 同步模式)下一个连接占一个 worker★
→ ★这就是"Flask 做 SSE 撑不住"的原因★
★ ★FastAPI/ASGI 是做长连接的正确选择★
三者的选择用一个决策树就清楚:需要客户端主动发消息就用 WebSocket;只是服务端推、更新频繁或要低延迟就用 SSE;更新不频繁(分钟级)轮询就够了——别为了「实时」两个字就上长连接。SSE 相比 WebSocket 的优势常被低估:它就是普通的 HTTP,所以现有的鉴权、日志、监控、限流全都能直接用,浏览器自动重连加 Last-Event-ID 续传,代理和 CDN 不用特殊配置,curl 就能调试。长轮询也是个被低估的方案(服务端 hold 住直到有数据或超时),实现最简单、无状态、兼容性最好。最后一个架构认知:一个长连接在 ASGI 下只占一个 asyncio 任务(几 KB 内存),而在 WSGI 下要占一个 worker——这就是「Flask 做 SSE 撑不住」的原因,FastAPI/ASGI 才是做长连接的正确选择。
六、实践清单
★ SSE 端点的标准模板:
@app.get("/stream")
async def stream_endpoint(request: Request, user: CurrentUser):
# ★① 流开始前完成所有校验和鉴权★
await check_permission(user)
async def gen():
try:
yield ": connected\n\n" # ★立刻建立连接感★
async for item in source():
if await request.is_disconnected():
break # ★② 断开检测★
yield f"data: {json.dumps(item)}\n\n"
yield "data: [DONE]\n\n" # ★③ 结束标记★
except asyncio.CancelledError:
raise # ★④ 不要吞★
except Exception:
logger.exception("stream error")
yield f"data: {json.dumps({'error':'内部错误'})}\n\n"
finally:
await cleanup() # ★⑤ 释放资源★
return StreamingResponse(gen(),
media_type="text/event-stream",
headers={"Cache-Control": "no-cache",
★"X-Accel-Buffering": "no"★})
★ 检查清单:
□ ★没有使用 BaseHTTPMiddleware(或已排除流式路径)★
□ ★加了 X-Accel-Buffering: no★
□ ★Nginx 配了 proxy_buffering off + read_timeout★
□ ★流式路径排除 gzip★
□ ★有心跳(15~30 秒)★
□ ★处理了 is_disconnected / CancelledError★
□ ★CancelledError 重新抛出了★
□ ★finally 里释放资源、落库、扣费★
□ ★校验和鉴权在流开始之前★
□ ★约定了结束标记★
□ ★用信号量限制并发流数★
□ ★多 worker 广播用 Redis Pub/Sub★
□ ★生成器里没有阻塞调用★
□ ★测试用 client.stream()★
★ ★排查"流式不生效"的顺序★:
① ★curl -N http://localhost:8000/events★(★直连应用,绕开代理★)
- ★是流式 → 问题在 Nginx/CDN★
- ★不是流式 → 问题在应用(中间件/gzip)★
② ★注释掉所有中间件再试★
③ ★检查响应头有没有 X-Accel-Buffering★
④ ★看 Content-Encoding 是不是 gzip★
⑤ ★确认生成器里没有一次性 yield 全部内容★
★ 一句话总结:
★"StreamingResponse + 异步生成器实现流式;SSE 的格式是
'data: xxx\\n\\n',必须配 text/event-stream + no-cache +
X-Accel-Buffering: no;最大的坑是 BaseHTTPMiddleware 和各层
缓冲会让流式失效(用 curl -N 定位);LLM 场景一定要处理客户端
断开——否则就是继续烧钱。"★
标准模板里的五个要点:校验在流开始前、断开检测、结束标记、CancelledError 不要吞、finally 里清理。排查「流式不生效」的第一步永远是 curl -N 直连应用端口——是流式就说明问题在 Nginx/CDN,不是流式就是应用侧(中间件或 gzip)。检查清单里最容易漏的是心跳(代理 60 秒空闲超时会切断连接)和多 worker 广播要用 Redis Pub/Sub。
记忆钩子:「FastAPI 的流式响应靠 ★StreamingResponse + (异步)生成器★,ASGI 层面的核心是 ★more_body: True 逐块 send★,因此★没有 Content-Length、走 Transfer-Encoding: chunked★(副作用是前端显示不了下载进度百分比)。三类用途:★大文件/CSV 导出、SSE 服务端推送、大模型逐 token 流式输出★。★SSE 的格式极简:
data: xxx+ 一个空行结束一条消息★,可选字段有 event(事件名)、★id(浏览器重连时会自动带 Last-Event-ID 请求头,服务端据此续传——这是 WebSocket 要自己实现而 SSE 白送的能力)★、retry,以及★用作心跳的注释行: ping★。★三个必须的响应头:text/event-stream + no-cache + X-Accel-Buffering: no★。★三个最容易踩的坑★:★① BaseHTTPMiddleware(@app.middleware(‘http’))会缓冲整个响应体,让 SSE 变成『全部生成完才一次性返回』★——要改写成★纯 ASGI 中间件用 send_wrapper 逐条转发★;★② 链路上每一层缓冲都会破坏流式★——Nginx 的 proxy_buffering(默认开)、★gzip 压缩(要攒够数据才输出)★、CDN,★排查第一步永远是 curl -N 直连应用端口:是流式说明问题在代理层,不是流式就在应用侧★;★③ 必须处理客户端断开★——用 ★await request.is_disconnected()★ 或捕获 ★asyncio.CancelledError(且必须重新抛出不能吞)★,★大模型场景下不处理就是用户关了页面还在继续烧 token 的真金白银★。LLM 流式的四个关键设计:★① 校验前置(响应头一发出状态码就固定了,流到一半出错只能在流里发错误消息)★ ② 断开检测 ③ ★约定 [DONE] 结束标记★ ④ ★finally 里落库和扣费(哪怕中断也要保存已生成的部分)★;还要注意★首 token 延迟 TTFT 是用户感知快慢的关键★(可以先发一个空事件让前端立刻进入生成中状态)。★EventSource 的三个硬限制:只支持 GET、不能自定义请求头(所以放不了 Authorization)、HTTP/1.1 下 6 连接上限★ → ★LLM 聊天基本都改用 fetch + ReadableStream 自己解析(要 POST + 要带 token),前端用 AbortController 让用户能主动取消★。选型:★需要客户端主动发消息才用 WebSocket;只是服务端推就用 SSE(就是普通 HTTP,鉴权/日志/监控/限流全能复用,curl 就能调试);更新不频繁(分钟级)轮询就够★。两个架构点:★多 worker 下 SSE 连接分散在不同进程,跨进程广播要用 Redis Pub/Sub★;★一个长连接在 ASGI 下只占一个 asyncio 任务(几 KB),在 WSGI 下要占一个 worker——这就是 Flask 做 SSE 撑不住的原因★。」
七、常见误区与追问
- 误区:写了
StreamingResponse加yield,前端就能看到内容陆续到达。 中间任何一层缓冲都会让它失效,而且现象很迷惑——「等了很久,然后所有内容一次性全出来」,看起来像代码写错了。四个缓冲点按出现频率排:①BaseHTTPMiddleware(即@app.middleware("http"))——它内部用内存流桥接 ASGI 接口,会等下游把响应体完整产出后再转发,这是 FastAPI 项目里最常见也最难联想到的原因;② Nginx 的proxy_buffering(默认开启)——用X-Accel-Buffering: no响应头或proxy_buffering off;关掉;③ gzip 压缩——压缩算法需要攒够一块数据才能输出,所以流式路径要排除在GZipMiddleware之外,Nginx 侧也要gzip off;④ CDN——SSE 端点通常应该绕过 CDN。排查的正确起点是curl -N http://localhost:8000/events直连应用端口(-N表示不缓冲):如果这样是流式的,问题就在代理层;如果不是,就在应用侧。 - 误区:客户端断开了服务端自然就停了,不用特意处理。 不会自动停。ASGI 服务器检测到连接断开后会尝试取消对应的任务,但如果你的生成器正在
await一个长时间的操作(比如等 LLM 返回下一个 token),取消不会立刻生效;更糟的是如果代码里吞掉了CancelledError,任务会继续跑到底。在大模型场景下这是实打实的金钱损失——用户提问后关掉页面,服务端还在向 OpenAI 请求剩下的几千个 token,账单照付。两种处理方式:① 主动检查await request.is_disconnected()——在生成循环里每次迭代都判断一下,检测到就break;② 捕获asyncio.CancelledError——在except里做清理然后必须raise重新抛出(吞掉会破坏 asyncio 的取消机制)。此外一定要有finally块:保存已生成的部分内容、按实际消耗扣费、释放上游连接。前端也要配合——用AbortController的signal,在组件卸载时abort(),服务端才能及时感知。 - 误区:SSE 和 WebSocket 差不多,用 WebSocket 更强大所以选它。 WebSocket 确实更强大,但如果你只需要服务端单向推送,SSE 的工程成本低得多。SSE 的优势常被低估:① 它就是普通的 HTTP GET 请求——现有的鉴权中间件、访问日志、APM 追踪、限流、错误处理全都能直接复用,而 WebSocket 是另一套协议,这些基础设施往往要重做一遍;② 浏览器内置自动重连,而且重连时自动带上
Last-Event-ID请求头让服务端续传——WebSocket 的重连和消息补发得自己实现(还要处理指数退避、状态恢复);③ 代理和网关不需要特殊配置(WebSocket 需要转发Upgrade/Connection头,很多云 LB 要单独开启支持);④ 调试极其简单——curl -N就能看到内容,浏览器 Network 面板也能正常展示。只有当客户端也需要实时向服务端发消息时(聊天输入、协作编辑、游戏操作),WebSocket 才是必需的。LLM 流式输出恰好是「服务端单向推送」,所以主流实现(OpenAI、Anthropic 的 API)用的都是 SSE。 - 误区:用
EventSource接 LLM 流式接口最标准。EventSource有两个硬限制让它不适合 LLM 场景:① 只支持 GET——而聊天请求要带完整的对话历史(可能几十 KB),塞进 URL 不现实;② 不能自定义请求头——这意味着放不了Authorization: Bearer <token>,只能退而求其次用 cookie 认证或把 token 放在 query string 里(后者会进 access log 和浏览器历史,是安全隐患)。此外 HTTP/1.1 下浏览器对同一域名有 6 个并发连接的限制,用户开几个标签页就会互相挤占(HTTP/2 下无此问题)。所以实践中 LLM 聊天基本都用fetch+ReadableStream自己解析 SSE 格式:可以 POST、可以带任意请求头、可以用AbortController取消。代价是要自己处理分块边界——网络分块和 SSE 消息边界不对齐,必须用一个 buffer 累积、按\n\n切分、并保留最后一段不完整的内容等下一块。 - 误区:多开几个 worker 就能支撑更多 SSE 连接。 worker 数确实影响承载量,但多 worker 会引入一个架构问题:连接分散在不同进程,进程间无法直接通信。典型场景:用户 A 的 SSE 连接落在 worker 1,而「有新消息」这个事件是在处理用户 B 的请求时(落在 worker 3)产生的——worker 3 里的代码无法把消息推送给连在 worker 1 上的那个生成器。如果你用进程内的
asyncio.Queue或全局字典管理订阅者,表现就是「有时候能收到推送,有时候收不到」,取决于连接和事件恰好落在哪个进程。正确架构是用 Redis Pub/Sub(或 NATS、Kafka)做跨进程广播:每个 SSE 生成器订阅一个 channel,产生事件的地方 publish 到那个 channel,所有 worker 里的订阅者都能收到。这是聊天室、通知中心、实时看板类功能的标准做法。顺带一提,连接数的承载能力在 ASGI 下很好——一个空闲的 SSE 连接只占一个 asyncio 任务加一个 socket(几 KB 内存),几千个连接完全没问题。 - 追问:流式响应中途出错了怎么办? 没法改状态码,只能在流里发错误消息。因为 HTTP 的响应状态码和响应头在第一个字节发出时就已经确定并发送了——一旦
StreamingResponse开始yield,200 OK就已经到了客户端,之后无论发生什么都改不成 500。这带来两个设计要求:① 所有能提前做的校验必须在流开始之前完成——鉴权、参数校验、额度检查、上游服务的可用性探测,这时候还能正常raise HTTPException(429)返回结构化的错误响应;② 流开始后的错误只能作为「业务内错误」发出去——比如yield f"data: {json.dumps({'error': '生成失败'})}\n\n",并且前端必须能识别这种消息(约定一个error字段或专门的event: error事件)。还有一个细节:如果错误发生在yield之前的准备阶段(生成器函数体的第一行到第一个 yield 之间),Starlette 实际上还没开始发送响应,这时抛异常仍然能返回正常的错误状态码——但这个边界很微妙,不要依赖它。 - 追问:SSE 连接为什么会莫名其妙断开? 几乎都是某一层的空闲超时。链路上有多个超时配置,任何一个到期都会切断连接:① Nginx 的
proxy_read_timeout(默认 60 秒)——指的是「多久没从上游读到数据就断开」,SSE 如果长时间没有事件就会中招;② 云负载均衡的空闲超时(AWS ALB 默认 60 秒、阿里云 SLB 类似);③ uvicorn 的--timeout-keep-alive;④ 客户端自己的超时(fetch 没设一般不超时,但有些 HTTP 客户端库有默认值)。标准解法是心跳:每 15~30 秒发一个 SSE 注释行: ping\n\n(客户端会自动忽略,但足以让各层认为连接是活跃的),同时把 Nginx 的proxy_read_timeout调大到 3600s。另外要注意EventSource断开后浏览器会自动重连(默认约 3 秒,可以用retry:字段调整)——这本身是好事,但如果服务端每次重连都从头开始推送数据,用户会看到内容重复;用id:字段配合Last-Event-ID实现续传才是完整的做法。 - 追问:怎么测试流式接口? 关键是不能一次性读完响应。用
TestClient时要用client.stream()而不是普通的client.get():with client.stream("GET", "/events") as r: for line in r.iter_lines(): ...——普通调用会阻塞直到生成器结束,而 SSE 的生成器可能是个无限循环,测试就永远挂住了。所以测试无限流时必须在收到足够的消息后主动break。几个实用的测试点:① 断言 SSE 格式(每条消息以data:开头、以\n\n结束);② 断言响应头(content-type是text/event-stream、有X-Accel-Buffering);③ 断言结束标记(收到[DONE]);④ 测试断开处理——提前退出循环后,验证服务端的清理逻辑执行了(比如检查数据库里保存了部分内容、额度被正确扣除);⑤ 测试错误路径——mock 上游抛异常,断言流里发出了 error 消息而不是连接直接断掉。异步测试用httpx.AsyncClient时对应的 API 是async with ac.stream(...)和async for line in r.aiter_lines()。
八、加强记忆
FastAPI 的流式响应靠 StreamingResponse + (异步)生成器,ASGI 层面的核心是 more_body: True 逐块 send,因此没有 Content-Length、走 Transfer-Encoding: chunked(副作用是前端显示不了下载进度百分比)。三类用途:大文件/CSV 导出、SSE 服务端推送、大模型逐 token 流式输出。SSE 的格式极简:data: xxx 加一个空行结束一条消息,可选字段有 event(事件名)、id(浏览器重连时会自动带上 Last-Event-ID 请求头,服务端据此续传——这是 WebSocket 要自己实现而 SSE 白送的能力)、retry,以及用作心跳的注释行 : ping。三个必须的响应头:text/event-stream + no-cache + X-Accel-Buffering: no。三个最容易踩的坑:① BaseHTTPMiddleware(@app.middleware("http"))会缓冲整个响应体,让 SSE 变成「全部生成完才一次性返回」——要改写成纯 ASGI 中间件用 send_wrapper 逐条转发;② 链路上每一层缓冲都会破坏流式——Nginx 的 proxy_buffering(默认开启)、gzip 压缩(要攒够数据才输出)、CDN,排查第一步永远是 curl -N 直连应用端口:是流式说明问题在代理层,不是流式就在应用侧;③ 必须处理客户端断开——用 await request.is_disconnected() 或捕获 asyncio.CancelledError(且必须重新抛出、不能吞掉),大模型场景下不处理就是「用户关了页面还在继续烧 token」的真金白银损失。LLM 流式的四个关键设计:① 校验前置(响应头一发出状态码就固定了,流到一半出错只能在流里发错误消息)、② 断开检测、③ 约定 [DONE] 结束标记、④ finally 里落库和扣费(哪怕中断也要保存已生成的部分);还要注意首 token 延迟(TTFT)是用户感知快慢的关键,可以先发一个空事件让前端立刻进入「生成中」状态。EventSource 有三个硬限制:只支持 GET、不能自定义请求头(所以放不了 Authorization)、HTTP/1.1 下 6 连接上限——所以 LLM 聊天基本都改用 fetch + ReadableStream 自己解析(要 POST、要带 token),前端用 AbortController 让用户能主动取消。选型判断:需要客户端主动发消息才用 WebSocket;只是服务端推就用 SSE(它就是普通 HTTP,鉴权、日志、监控、限流全能复用,curl 就能调试);更新不频繁(分钟级)轮询就够了。两个架构要点:多 worker 下 SSE 连接分散在不同进程,跨进程广播要用 Redis Pub/Sub;一个长连接在 ASGI 下只占一个 asyncio 任务(几 KB 内存),在 WSGI 下要占一个 worker——这就是「Flask 做 SSE 撑不住」的原因。