Appearance
流式 API 与 SSE
LLM 生成 token 需要几秒到几十秒,如果让用户盯着空白转圈,体验会非常糟。流式输出让用户看到「逐字蹦出来」的效果,是智能体应用的标配。本篇讲清楚:SSE 协议是什么、怎么用 FastAPI 把 LangGraph 的 astream 转成前端能消费的事件流、以及多 stream_mode 怎么聚合。
一、为什么是 SSE 而不是 WebSocket
| 协议 | 方向 | 复杂度 | 适用 |
|---|---|---|---|
| 普通 HTTP | 请求-响应 | 低 | 一次性结果 |
| SSE | 服务端→客户端单向 | 低 | LLM token 流、状态更新 |
| WebSocket | 双向 | 高 | 多人协作、实时游戏 |
LLM 流式输出本质是「服务端持续推、客户端只收不发」,SSE 是最契合的协议:
- 基于 HTTP,穿透防火墙/代理无压力
- 浏览器原生
EventSourceAPI,无需额外库 - 自动重连机制内置
- 比 WebSocket 简单一个数量级
给 Java 同学
SSE 在 Spring 里对应 SseEmitter,在 FastAPI 里对应 StreamingResponse,原理一致:保持长连接,按 text/event-stream 格式分块吐数据。
二、SSE 协议速览
SSE 响应头:
text
Content-Type: text/event-stream
Cache-Control: no-cache
Connection: keep-alive消息体格式——每条事件由若干行组成,事件之间用两个换行分隔:
text
event: token
data: {"content":"你"}
event: token
data: {"content":"好"}
event: done
data: [DONE]event:自定义事件名(前端可针对性监听)data:数据行(一般放 JSON 字符串)id:可选,用于断点续传retry:可选,重连间隔毫秒数
三、LangGraph 的流式输出回顾
LangGraph 的 astream 支持多种 stream_mode:
| stream_mode | 吐出内容 | 典型用途 |
|---|---|---|
values | 每步完整 state | 看状态演进 |
updates | 每步 state 增量 | 哪个节点改了啥 |
messages | LLM 的 token + 元数据 | 前端逐字显示 |
custom | 节点内自定义事件 | 进度条、自定义通知 |
最常用的是 messages,它会逐 token 吐出 LLM 输出。详见 流式输出 Streaming。
python
# 回顾:原生 astream messages 模式
async for msg, meta in graph.astream(
{"messages": [HumanMessage("讲个笑话")]},
config=config,
stream_mode="messages",
):
if msg.content: # 跳过空 content(如工具调用元数据)
print(msg.content, end="", flush=True)四、完整 FastAPI SSE 端点
1. 核心思路
mermaid
flowchart LR
Client[前端 EventSource] -->|POST/GET| API[FastAPI 端点]
API -->|astream messages| Graph[图]
Graph -->|逐 token| API
API -->|格式化 SSE 块| Client把 astream 的每个 chunk 转成一行 data: {...}\n\n,通过 StreamingResponse 持续 flush 给客户端。
2. 完整代码
python
import os
import uuid
import json
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
from pydantic import BaseModel
from langchain_core.messages import HumanMessage
from langgraph.checkpoint.memory import MemorySaver
import graph as graph_mod
app = FastAPI()
# 开发期用内存 checkpointer;生产换 PostgresSaver
graph = graph_mod.builder.compile(checkpointer=MemorySaver())
class StreamRequest(BaseModel):
message: str
thread_id: str | None = None
async def event_generator(message: str, thread_id: str):
"""把图流式输出转成 SSE 事件流"""
config = {"configurable": {"thread_id": thread_id}}
try:
async for msg, meta in graph.astream(
{"messages": [HumanMessage(content=message)]},
config=config,
stream_mode="messages",
):
# msg 是 AIMessageChunk,content 是本次新增 token
if msg.content:
payload = json.dumps({"content": msg.content}, ensure_ascii=False)
# SSE 格式:event 行 + data 行 + 空行
yield f"event: token\ndata: {payload}\n\n"
# 工具调用:通过 meta 判断节点
if meta and meta.get("langgraph_node") == "tools":
payload = json.dumps({"tool": "called"}, ensure_ascii=False)
yield f"event: tool\ndata: {payload}\n\n"
# 结束标记
yield "event: done\ndata: [DONE]\n\n"
except Exception as e:
# 错误也走 SSE,前端能收到
err = json.dumps({"error": str(e)}, ensure_ascii=False)
yield f"event: error\ndata: {err}\n\n"
@app.post("/chat/stream")
async def chat_stream(req: StreamRequest):
"""SSE 流式对话端点"""
thread_id = req.thread_id or str(uuid.uuid4())
generator = event_generator(req.message, thread_id)
return StreamingResponse(
generator,
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache",
"Connection": "keep-alive",
"X-Accel-Buffering": "no", # 关键:禁 nginx 缓冲
"X-Thread-Id": thread_id,
},
)启动:
bash
uvicorn app:app --host 0.0.0.0 --port 80003. 用 curl 验证
bash
curl -N -X POST http://localhost:8000/chat/stream \
-H "Content-Type: application/json" \
-d '{"message":"讲个三句话的笑话"}'-N 禁用缓冲,能实时看到 token 一段段蹦出来。
五、前端 EventSource 消费示例
javascript
// 用原生 EventSource(仅支持 GET,POST 需用 fetch + ReadableStream)
// 这里演示 POST 方案的 fetch 流式读法
async function streamChat(message, threadId) {
const resp = await fetch("http://localhost:8000/chat/stream", {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({ message, thread_id: threadId }),
});
const reader = resp.body.getReader();
const decoder = new TextDecoder();
let buffer = "";
let newThreadId = resp.headers.get("X-Thread-Id");
while (true) {
const { done, value } = await reader.read();
if (done) break;
buffer += decoder.decode(value, { stream: true });
// 按双换行切事件块
const parts = buffer.split("\n\n");
buffer = parts.pop(); // 最后一段可能不完整,留到下次
for (const part of parts) {
const lines = part.split("\n");
let event = "message", data = "";
for (const line of lines) {
if (line.startsWith("event:")) event = line.slice(6).trim();
if (line.startsWith("data:")) data = line.slice(5).trim();
}
if (event === "token") {
const { content } = JSON.parse(data);
// 把 content 追加到页面
appendToChat(content);
} else if (event === "done") {
console.log("结束");
} else if (event === "error") {
console.error("出错", data);
}
}
}
return newThreadId;
}EventSource 的局限
浏览器原生 EventSource 只支持 GET。若必须 POST(带 body),用 fetch + ReadableStream 手动解析,如上。也有 @microsoft/fetch-event-source 库封装好了 POST 场景。
六、多 stream_mode 聚合
生产中常要「token 流 + 节点状态」一起吐。astream 支持传多个 mode:
python
async for event in graph.astream(
{"messages": [HumanMessage(content=message)]},
config=config,
stream_mode=["messages", "updates"],
):
mode, chunk = event # 元组:第一个是 mode 名
if mode == "messages":
msg, meta = chunk
if msg.content:
yield f"event: token\ndata: {json.dumps({'content': msg.content}, ensure_ascii=False)}\n\n"
elif mode == "updates":
# chunk 是 {节点名: {state字段: 新值}}
node = next(iter(chunk))
yield f"event: node\ndata: {json.dumps({'node': node}, ensure_ascii=False)}\n\n"前端按 event 字段分发渲染:token 追加到气泡,node 更新流程图高亮。
七、事件格式设计建议
约定一套稳定事件协议,前后端省事:
| event | data 含义 | 时机 |
|---|---|---|
start | {"thread_id":"xxx"} | 流开始 |
token | {"content":"你"} | 每个 LLM token |
tool | {"name":"add","args":{...}} | 工具调用前 |
tool_result | {"name":"add","result":8} | 工具返回 |
node | {"node":"chatbot"} | 节点进入 |
error | {"error":"..."} | 异常 |
done | [DONE] | 流结束 |
八、常见踩坑
1. 流被 nginx 缓冲,前端一次性收到 症状:本地正常,上了 nginx 后变成「等几秒突然全出来」。解决:nginx 配置加 proxy_buffering off;,并在响应头加 X-Accel-Buffering: no。
nginx
location /chat/stream {
proxy_pass http://backend;
proxy_buffering off; # 关键
proxy_cache off;
proxy_set_header Connection '';
proxy_http_version 1.1;
chunked_transfer_encoding on;
}2. 客户端断开后服务端还在跑StreamingResponse 检测到客户端断开会停止迭代 generator,但 LLM 调用可能已在进行。可在 generator 里捕获 asyncio.CancelledError 做清理。生产可加 async with timeout 兜底。
3. 背压(backpressure) LLM 生成快、前端消费慢,token 堆积在内存。SSE 基于 TCP 流控天然有背压,但 generator 里别用 asyncio.Queue 无界缓冲。直接 yield 让框架处理即可。
4. Windows 下 curl 看不到流 Windows 的 curl 可能是 curl.exe,加 -N 不生效。改用 Invoke-WebRequest 或装真 curl,或直接用浏览器开发者工具的 Network 面板看 EventStream。
5. 多 worker 下流式中断 uvicorn --workers N 是多进程,每个请求落在一个进程,流式全程在同一进程,不受影响。但若用了 nginx 的 round-robin 且配了 proxy_buffering on,可能在 worker 间切换导致流断。保持 proxy_buffering off。
6. 中文乱码 确保 json.dumps(..., ensure_ascii=False),响应头 charset=utf-8:media_type="text/event-stream; charset=utf-8"。
九、小结
- SSE 是 LLM 流式输出最合适的协议:单向、轻量、浏览器原生支持
- FastAPI 用
StreamingResponse+ 异步 generator 把astream(stream_mode="messages")转成event/data块 - 上线务必关 nginx 缓冲(
proxy_buffering off),并约定稳定事件协议 - 多
stream_mode聚合能同时满足「逐字显示」和「流程图高亮」
下一篇讲怎么给服务加监控、日志、链路追踪,让生产环境可观测。