1115 分钟

流式输出(SSE)—— 让用户"看到"系统在思考

掌握 SSE 协议原理、LangGraph astream → Service → API 三层流式架构、前端 fetch + ReadableStream 消费模式,以及节点级到 token 级流式的升级路径。

SSE流式输出StreamingLangGraphFastAPIPython
进度保存在本机浏览器;验收通过后再点更稳妥

第 11 课:流式输出(SSE)—— 让用户"看到"系统在思考

本节目标:掌握 SSE 协议原理、LangGraph astream → Service → API 三层流式架构、前端 fetch + ReadableStream 消费模式,以及从节点级流式升级到 token 级流式的完整路径。

学完本课后,推荐阅读 附11:工程化部署与推理优化 —— vLLM 三大创新、量化方案对比、推理框架选型决策树。

前 10 课里,用户提问后要等整个 RAG 流程跑完(检索 → 生成 → 验证 → 格式化),3-5 秒后才一次性返回答案。这在 ChatGPT 时代是不可接受的体验。


1. 问题:用户等不了

Code
普通模式(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)

Code
┌──────────┐    HTTP POST     ┌──────────┐
│  Client  │ ──────────────→  │  Server  │
│ (前端)    │                  │ (FastAPI)│
│          │  Content-Type:   │          │
│          │  text/event-stream│         │
│          │ ←──────────────── │          │
│          │  data: {...}\n\n │          │  ← 第 1 条事件
│          │ ←──────────────── │          │
│          │  data: [DONE]\n\n│          │  ← 结束信号
└──────────┘                  └──────────┘

SSE vs WebSocket vs HTTP Polling

特性SSEWebSocketHTTP Polling
方向服务器 → 客户端(单向)双向客户端 → 服务器
连接一个 HTTP 长连接升级到 WS 协议反复建立新连接
复杂度最简单推荐 复杂推荐 中等
适合场景LLM 流式输出、进度推送聊天室、实时协作轮询状态
重连浏览器自动重连需自己实现不需要

LLM 流式输出用 SSE 就够了——只需要服务端单向推送给客户端。OpenAI、Claude 用的都是 SSE。

SSE 消息格式

code
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 一次:

python
# 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 封装

python
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 输出)

python
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=Falsejson.dumps(event, ensure_ascii=False)中文不转义,调试可读
default=strjson.dumps(..., default=str)datetime 等自动转 str,不崩
\n\nf"data: ...\n\n"SSE 协议要求双换行分隔
[DONE]yield "data: [DONE]\n\n"告诉前端"流结束了"
text/event-streammedia_type="text/event-stream"浏览器识别为 SSE 流

4. 流式事件序列

一次完整调用推送的事件:

code
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

javascript
const evtSource = new EventSource("/api/v1/rag-graph/ask?stream=true");
// EventSource 只支持 GET,不支持 POST body

正确方式:fetch + ReadableStream

javascript
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(推荐)

javascript
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. 测试流式

python
@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 是中文项目的必须项

Code
# 默认 → 中文变 \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 题)

练习 1rag_graph_service.stream()stream_mode="updates"。如果改成 stream_mode="values",每个事件的格式会变成什么?对前端有什么影响?

练习 2:如果 generate_node 执行到一半抛异常(LLM 超时),流式连接会怎样?前端怎么感知?当前代码有没有处理?

练习 3/rag-graph/askstream 参数在同一 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 里:

python
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:"不完整的字节序列先暂存,等下一块到了再拼"。

Code
// 不带 stream: true → 中文被拆断 → 乱码 "�"
// 带 stream: true → 暂存不完整字节 → 拼完整再输出

英文是 1 字节(ASCII),几乎不会被拆断。中文 3 字节,被拆断概率是英文的 3 倍。不加必乱码。

练习 5(最重要):升级到 token 级流式

需要改 4 层:

文件改动
LLM 层llm_service.py新增 chat_messages_stream() → 用 stream=True 逐 token yield
节点层nodes.pygenerate_nodeStreamWriter 在节点执行中发 token 事件
Service 层rag_graph_service.pystream() 处理 custom event(token)和 updates(节点)两种模式
API 层rag_graph.py不需要改——json.dumps(event) 自动适配
前端JS/TS区分 type: "token"(逐字追加)和普通节点事件(更新状态)

核心难点在节点层——LangGraph 节点默认返回 dict,需要 StreamWriter 回调才能边执行边推送。这就是分层架构的价值:改动从底层向上穿透,但最上层 API 层不变。


三句话带走第 11 课

  1. SSE 是 LLM 流式输出的最佳选择:单向推送、HTTP 原生、格式极简(data: {...}\n\n)。FastAPI 的 StreamingResponse + async generator 10 行代码搞定。
  2. 节点级 vs token 级是体验和复杂度的权衡:节点级流式足以展示进度;token 级流式需要改 LLM SDK + LangGraph 节点设计,但用户体验拉满。
  3. 流式的魔鬼在细节ensure_ascii=False(中文)、default=str(防序列化崩)、[DONE] 放 finally(防异常丢信号)、前端 TextDecoder({ stream: true })(防 UTF-8 拆断乱码)。