第 11 课:流式输出(SSE)—— 让用户"看到"系统在思考
本节目标:掌握 SSE 协议原理、LangGraph astream → Service → API 三层流式架构、前端 fetch + ReadableStream 消费模式,以及从节点级流式升级到 token 级流式的完整路径。
学完本课后,推荐阅读 附11:工程化部署与推理优化 —— vLLM 三大创新、量化方案对比、推理框架选型决策树。
前 10 课里,用户提问后要等整个 RAG 流程跑完(检索 → 生成 → 验证 → 格式化),3-5 秒后才一次性返回答案。这在 ChatGPT 时代是不可接受的体验。
1. 问题:用户等不了
普通模式(stream=false):
用户提问 ──────────── 3-5s ──────────── 一次性返回答案
等待...等待...等待... "RAG 是检索增强生成 [1]。"
流式模式(stream=true):
用户提问 ─ 0.1s → "rewrite 完成"
─ 0.3s → "retrieve: 7 hits"
─ 0.5s → "judge: 5 kept"
─ 2.0s → "generate: RAG 是检索增强生成 [1]。"
─ 0.1s → "format: done"
─ 0.0s → [DONE]
体验差别:普通模式下 3 秒空白 → 用户以为卡了;流式模式下 0.1 秒就有反馈 → 用户看到进度 → 愿意等。
2. SSE 协议:最简单的服务端推送
什么是 SSE (Server-Sent Events)
┌──────────┐ HTTP POST ┌──────────┐
│ Client │ ──────────────→ │ Server │
│ (前端) │ │ (FastAPI)│
│ │ Content-Type: │ │
│ │ text/event-stream│ │
│ │ ←──────────────── │ │
│ │ data: {...}\n\n │ │ ← 第 1 条事件
│ │ ←──────────────── │ │
│ │ data: [DONE]\n\n│ │ ← 结束信号
└──────────┘ └──────────┘
SSE vs WebSocket vs HTTP Polling
| 特性 | SSE | WebSocket | HTTP Polling |
|---|---|---|---|
| 方向 | 服务器 → 客户端(单向) | 双向 | 客户端 → 服务器 |
| 连接 | 一个 HTTP 长连接 | 升级到 WS 协议 | 反复建立新连接 |
| 复杂度 | 最简单 | 推荐 复杂 | 推荐 中等 |
| 适合场景 | LLM 流式输出、进度推送 | 聊天室、实时协作 | 轮询状态 |
| 重连 | 浏览器自动重连 | 需自己实现 | 不需要 |
LLM 流式输出用 SSE 就够了——只需要服务端单向推送给客户端。OpenAI、Claude 用的都是 SSE。
SSE 消息格式
data: {"node": "rewrite", "delta": {"rewritten_query": "RAG 是什么"}}\n\n
data: {"node": "retrieve", "delta": {"hits": [...]}}\n\n
data: {"node": "generate", "delta": {"answer": "RAG 是..."}}\n\n
data: [DONE]\n\n
规则:每条消息以 data: 开头,以 \n\n(两个换行)结尾。[DONE] 是约定俗成的结束标志(OpenAI 也用这个)。
3. 三层架构:LangGraph → Service → API
第 1 层:LangGraph astream
LangGraph 的 CompiledGraph 提供 astream() 方法——异步迭代器,每个节点执行完后 yield 一次:
# LangGraph 内部机制
async for event in graph.astream(init, stream_mode="updates"):
# event = {"rewrite": {"rewritten_query": "...", "meta": {...}}}
# event = {"retrieve": {"hits": [...], "meta": {...}}}
stream_mode="updates":每个节点只返回该节点修改的字段(delta),不返回完整 state。
第 2 层:Service 封装
async def stream(
self, question: str, *, conv_id: str | None = None,
top_k: int = 5, retrieval_mode: str = "hybrid",
enable_rewrite: bool = True,
) -> AsyncIterator[dict[str, Any]]:
"""逐节点推流,每个 event 形如 {"node": "rewrite", "delta": {...}}."""
graph = get_rag_graph()
init: RagState = {
"question": question, "conv_id": conv_id,
"top_k": top_k, "retrieval_mode": retrieval_mode,
"enable_rewrite": enable_rewrite, "meta": {},
}
async for event in graph.astream(init, stream_mode="updates"):
for node_name, delta in event.items():
yield {"node": node_name, "delta": delta}
Service 做了什么:构建初始 state → 调用 graph.astream() → 把 {node_name: delta} 拆成 {"node": ..., "delta": ...} 统一接口。
为什么多一层 Service?API 层不该直接碰 LangGraph;
stream()和ask()共享统一参数签名;测试可以直接调service.stream()不需要起 HTTP。
第 3 层:API 路由(SSE 输出)
async def ask(req: GraphAskRequest):
if not req.question.strip():
raise ValidationError("question 不能只是空白字符")
if not req.stream:
return await rag_graph_service.ask(...) # 普通模式:一次性 JSON
async def _gen():
async for event in rag_graph_service.stream(...):
yield f"data: {json.dumps(event, ensure_ascii=False, default=str)}\n\n"
yield "data: [DONE]\n\n"
return StreamingResponse(_gen(), media_type="text/event-stream")
同一接口,stream 参数决定行为。关键细节:
| 细节 | 代码 | 为什么 |
|---|---|---|
ensure_ascii=False | json.dumps(event, ensure_ascii=False) | 中文不转义,调试可读 |
default=str | json.dumps(..., default=str) | datetime 等自动转 str,不崩 |
\n\n | f"data: ...\n\n" | SSE 协议要求双换行分隔 |
[DONE] | yield "data: [DONE]\n\n" | 告诉前端"流结束了" |
text/event-stream | media_type="text/event-stream" | 浏览器识别为 SSE 流 |
4. 流式事件序列
一次完整调用推送的事件:
data: {"node":"rewrite","delta":{"rewritten_query":"...","meta":{...}}}
data: {"node":"retrieve","delta":{"hits":[...],"meta":{"retrieve_count":5}}}
data: {"node":"judge","delta":{"hits":[...],"meta":{"judged_kept":3}}}
data: {"node":"generate","delta":{"answer":"RAG 是... [1]。","citations":[...],"finish_reason":"answered"}}
data: {"node":"format","delta":{"meta":{"total_elapsed":2.1}}}
data: [DONE]
| 事件 | 前端可以做什么 |
|---|---|
rewrite | 显示"正在理解你的问题…" |
retrieve | 显示"找到 N 条相关文档" |
judge | 显示"筛选出 N 条高质量结果" |
generate | 显示答案文字(核心!) |
format | 显示总耗时 |
[DONE] | 关闭加载动画 |
节点级 vs token 级:OpenAI 流式是 token 级(逐字推送),我们的是节点级(每个图节点完成后推送一次)。节点级已给即时反馈,token 级是进阶优化——需要改 LLM SDK 调用 + LangGraph 节点设计。
5. 前端如何消费 SSE
EventSource 不行——只支持 GET
const evtSource = new EventSource("/api/v1/rag-graph/ask?stream=true");
// EventSource 只支持 GET,不支持 POST body
正确方式:fetch + ReadableStream
async function streamAsk(question) {
const resp = await fetch("/api/v1/rag-graph/ask", {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({ question, stream: true }),
});
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 parts = buffer.split("\n\n");
buffer = parts.pop(); // 最后一段可能不完整,留着
for (const part of parts) {
if (!part.startsWith("data: ")) continue;
const payload = part.slice(6);
if (payload === "[DONE]") return;
const event = JSON.parse(payload);
if (event.node === "retrieve") showStatus(`找到 ${event.delta.meta?.retrieve_count} 条文档`);
if (event.node === "generate") showAnswer(event.delta.answer);
}
}
}
或用 @microsoft/fetch-event-source(推荐)
import { fetchEventSource } from "@microsoft/fetch-event-source";
await fetchEventSource("/api/v1/rag-graph/ask", {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({ question: "什么是 RAG?", stream: true }),
onmessage(ev) {
if (ev.data === "[DONE]") return;
const event = JSON.parse(ev.data);
console.log(event.node, event.delta);
},
});
封装了 POST + SSE 消费逻辑,自动处理重连和错误。
6. 测试流式
@pytest.mark.unit
@pytest.mark.asyncio
async def test_graph_stream_yields_per_node_events(monkeypatch) -> None:
from app.graph import nodes as nodes_mod
from app.services.rag_graph_service import rag_graph_service
monkeypatch.setattr(nodes_mod.hybrid_retrieval_service, "retrieve",
lambda *a, **kw: [_hit("x", 0.9)])
async def _fake_chat(messages, **_kw):
return {"answer": "ok [1]", "model": "mock", "latency_ms": 1}
monkeypatch.setattr(nodes_mod.llm_service, "chat_messages", _fake_chat)
nodes_seen = []
async for event in rag_graph_service.stream("q"):
nodes_seen.append(event["node"])
expected = {"rewrite", "retrieve", "judge", "generate", "format"}
assert expected.issubset(set(nodes_seen))
测试策略:mock 外部依赖 → 直接调 service.stream()(不经 HTTP)→ 收集所有 event["node"] → 断言关键节点都出现。不需要测 SSE 格式(那是 FastAPI StreamingResponse 的责任)。
7. 5 条工程教训
教训 1:ensure_ascii=False 是中文项目的必须项
# 默认 → 中文变 \uXXXX
json.dumps({"answer": "RAG 是检索增强生成"})
# → '{"answer": "RAG \\u662f\\u68c0\\u7d22\\u589e\\u5f3a\\u751f\\u6210"}'
# 设 ensure_ascii=False
json.dumps({"answer": "RAG 是检索增强生成"}, ensure_ascii=False)
# → '{"answer": "RAG 是检索增强生成"}'
教训 2:default=str 防止序列化崩溃 — state 里有 datetime → json.dumps 直接崩 → 流式连接断开 → 前端收到不完整流。
教训 3:[DONE] 信号不能省 — 不发 [DONE] = 前端不知道结束 → 继续等待 → 超时 → EventSource 自动重连 → 重复请求。
教训 4:buffer 处理防止消息被切断 — TCP 可能把一条 SSE 消息拆成两个包。前端必须用 buffer 拼接,按 \n\n 分割后再解析。
教训 5:节点级 vs token 级流式的权衡 — 节点级简单、能看到进度;token 级复杂、逐字显示体验最好。混合做法:非 LLM 节点发节点级事件,LLM 节点内部 stream=True 发 token 级事件。
8. 小练习(5 题)
练习 1:rag_graph_service.stream() 用 stream_mode="updates"。如果改成 stream_mode="values",每个事件的格式会变成什么?对前端有什么影响?
练习 2:如果 generate_node 执行到一半抛异常(LLM 超时),流式连接会怎样?前端怎么感知?当前代码有没有处理?
练习 3:/rag-graph/ask 用 stream 参数在同一 endpoint 切换;/verified-qa 用独立 /ask/stream。你选哪种?为什么?
练习 4:前端 decoder.decode(value, { stream: true }) 的 { stream: true } 参数作用是什么?中文项目为什么尤其重要?
练习 5(最重要):要把节点级流式升级为 generate_node 内部的 token 级流式(LLM 用 stream=True),需要改哪些层的代码?列出具体文件和改动要点。
答案与解析
练习 1:stream_mode="values" vs "updates"
| stream_mode | 每个事件包含 | 格式 |
|---|---|---|
"updates" | 该节点修改的字段 | {"node_name": delta_dict} |
"values" | 完整 state(所有字段快照) | {...全部字段...} |
对前端影响:values 数据量大(随节点递增),没有 node name 标记 → 前端不知道是哪个节点产生的。SSE 应该用 "updates"——带宽友好、能区分节点。"values" 适合调试。
练习 2:异常时流式连接的行为
generate_node 已捕获 LLMError → 返回 finish_reason="llm_error" → 流正常结束。但未预期异常会穿透到 astream → 流中断 → 前端收到 done: true 但没有 [DONE]。
更健壮的做法:[DONE] 放 finally 里:
async def _gen():
try:
async for event in rag_graph_service.stream(...):
yield f"data: {json.dumps(event, ensure_ascii=False, default=str)}\n\n"
except Exception as exc:
yield f"data: {json.dumps({'error': str(exc)})}\n\n"
finally:
yield "data: [DONE]\n\n"
练习 3:同一 endpoint 还是独立?
推荐独立 endpoint(/ask/stream):OpenAPI 文档清晰(每个 endpoint 有明确返回类型)、类型安全、CDN 不会把 SSE 缓存成 JSON、可针对流式配置不同超时。参数重复问题用共享 Request model 解决——两个 endpoint 共用一个 AskRequest。
练习 4:{ stream: true } 参数
TCP 传输时一个中文字符(UTF-8 3 字节)可能被拆成两个包。{ stream: true } 告诉 TextDecoder:"不完整的字节序列先暂存,等下一块到了再拼"。
// 不带 stream: true → 中文被拆断 → 乱码 "�"
// 带 stream: true → 暂存不完整字节 → 拼完整再输出
英文是 1 字节(ASCII),几乎不会被拆断。中文 3 字节,被拆断概率是英文的 3 倍。不加必乱码。
练习 5(最重要):升级到 token 级流式
需要改 4 层:
| 层 | 文件 | 改动 |
|---|---|---|
| LLM 层 | llm_service.py | 新增 chat_messages_stream() → 用 stream=True 逐 token yield |
| 节点层 | nodes.py | generate_node 用 StreamWriter 在节点执行中发 token 事件 |
| Service 层 | rag_graph_service.py | stream() 处理 custom event(token)和 updates(节点)两种模式 |
| API 层 | rag_graph.py | 不需要改——json.dumps(event) 自动适配 |
| 前端 | JS/TS | 区分 type: "token"(逐字追加)和普通节点事件(更新状态) |
核心难点在节点层——LangGraph 节点默认返回 dict,需要 StreamWriter 回调才能边执行边推送。这就是分层架构的价值:改动从底层向上穿透,但最上层 API 层不变。
三句话带走第 11 课
- SSE 是 LLM 流式输出的最佳选择:单向推送、HTTP 原生、格式极简(
data: {...}\n\n)。FastAPI 的StreamingResponse+ async generator 10 行代码搞定。 - 节点级 vs token 级是体验和复杂度的权衡:节点级流式足以展示进度;token 级流式需要改 LLM SDK + LangGraph 节点设计,但用户体验拉满。
- 流式的魔鬼在细节:
ensure_ascii=False(中文)、default=str(防序列化崩)、[DONE]放 finally(防异常丢信号)、前端TextDecoder({ stream: true })(防 UTF-8 拆断乱码)。