给 RAG Agent 加流式输出:SSE 从设计到落地

给一个基于 FastAPI + LangChain/LangGraph 的 RAG Agent 加上"打字机"式流式输出。本文按自顶向下的顺序记录:从方案选型、协议设计,到后端/前端实现,最后是一段曲折的踩坑——dev 代理的 gzip 压缩会把流式悄悄变成一次性响应

先说结论:

流式对话用 SSE(Server-Sent Events)就够了——方案一(SSE 专用接口)成本可控、标准通用、能满足"打字机效果"的核心诉求。关键实现要点有三个:后端用 agent.astream(stream_mode="messages") 逐 token 产出、前端用原生 fetch + ReadableStream 解析(不能用 EventSource,因为要 POST + Authorization)、以及必须让响应头带 Cache-Control: no-cache, no-transform,否则开发代理的 gzip 会把流缓冲成一大块。


一、为什么需要流式

LLM 生成一段回复通常要 3~10 秒,Agent 场景还会先调工具检索,首字延迟更明显。如果没有流式,用户点击发送后只能盯着转圈,体验很差;而流式下首个 token 往往几百毫秒就到了,后续内容逐字出现,感知延迟大幅下降。

在动手前先明确约束:

  • 现有接口 POST /chat 一次性返回 JSON,前端 await 后整段渲染;
  • 后端已有 create_react_agent(llm, tools) 的 ReAct Agent(含 RAG 检索工具和 MCP 工具);
  • 前端是 React + umi(axios 封装),需要增量渲染能力。

二、方案选型:SSE / WebSocket / NDJSON

维度SSE(方案一)WebSocket(方案二)NDJSON(方案三)
实现成本
前端改动中(换 fetch)大(WS 客户端)小(axios 参数)
标准性 / 通用性
双向交互(停止 / 确认)
代理 / 负载均衡兼容需要特殊支持好(但怕压缩)
适用场景通用聊天流式,首选需要 Agent 交互控制快速验证 / 最小改动

选择 SSE 的决定性理由:

  1. 单向流正是聊天场景需要的——服务端推,客户端收,不需要上行通道;
  2. HTTP 协议、任何代理/CDN 都能透传,生产部署成本最低;
  3. 保留原有 POST /chat 不动,新端点并行存在,向后兼容;
  4. 把事件对象设计成与传输协议解耦({"type", "content", ...}),未来要升级 WebSocket 或加"停止生成"功能时,只换传输层即可。

一个容易踩的坑:EventSource(浏览器原生 SSE API)只支持 GET,而我们的接口需要 POST + Authorization 请求头,所以前端必须用 fetch + ReadableStream 手动解析 SSE 帧。

三、后端设计:事件协议

流式接口把原来的一次性 JSON 响应拆成一串事件,每个事件一行:

data: {"type": "text", "content": "地理信息系统(GIS"}       # 增量文本,前端追加渲染

data: {"type": "tool", "name": "info_retriever"}             # 工具调用,可展示"正在检索…"

data: {"type": "done", "response": "...", "session_id": "...", "message_count": 6}  # 结束

data: {"type": "error", "code": 5010, "message": "聊天调用超时"}  # 出错

SSE 帧格式是 data: <JSON>\n\n,一次连接推多个事件,最后自然结束。

四、后端实现

4.1 服务层:async generator 产出事件

核心是 LangGraph 的 agent.astream(..., stream_mode="messages"),它逐 token 产出 (message_chunk, metadata)。这里有个非常优雅的特性:工具调用阶段产出的 AIMessageChunk 内容为空(只有 tool_calls),自然被 elif content: 过滤掉;ToolMessage 单独透出为 tool 事件;最终答案逐 token 到达。这样"文本流"和"工具事件流"自动分离。

async def chat_stream(self, prompt, query, chat_model_name, ...):
    tools, messages = self._build_tools_and_messages(
        prompt=prompt, query=query, use_memory=use_memory,
        history=history, db_name=db_name, mcp_tools=mcp_tools,
        user_id=user_id,
    )
    llm = self._create_llm(chat_model_name)

    if tools:
        agent = create_react_agent(llm, tools)
        async for msg_chunk, _meta in agent.astream(
            {"messages": messages},
            config={"callbacks": [handler]},   # 工具审计回调不变
            stream_mode="messages",
        ):
            msg_type = getattr(msg_chunk, "type", "")
            content = getattr(msg_chunk, "content", "") or ""
            if msg_type == "tool":             # ToolMessage → 工具事件
                yield {"type": "tool",
                       "name": getattr(msg_chunk, "name", None) or "tool"}
            elif content:                      # AIMessageChunk → 增量文本
                yield {"type": "text", "content": content}
    else:
        async for chunk in llm.astream(messages):   # 纯对话直接流式
            content = getattr(chunk, "content", "") or ""
            if content:
                yield {"type": "text", "content": content}

顺带一个重构收益:把"构建工具 + 消息列表"抽成 _build_tools_and_messages(),普通 /chat 和流式 /chat/stream 共享,避免两套逻辑漂移。

4.2 路由层:StreamingResponse + 完整生命周期

流式不能牺牲原有保障,所以生命周期和一次性接口完全对齐:

  1. 先落库:用户消息、ChatRun 审计记录在流开始前写入(超时/断线也能恢复用户输入);
  2. 流转:把 generator 包进 StreamingResponse(media_type="text/event-stream")
  3. 流结束落库:AI 回复、会话标题、记忆、finish_chat_run("succeeded")
  4. 超时:用 asyncio.wait_for(iterator.__anext__(), timeout=...) 对每个 chunk 设超时,超时发 error 事件并记 timed_out
  5. 断线/取消except asyncio.CancelledError / finally 兜底记 cancelled,防止审计记录永远卡在 running。
def _sse_payload(event: dict) -> str:
    return f"data: {json.dumps(event, ensure_ascii=False)}\n\n"

@router.post("/chat/stream")
async def chat_stream(request: ChatRequest, ...) -> StreamingResponse:
    # ... 会话校验、先落库用户消息、create_chat_run ...
    return StreamingResponse(
        event_generator(),            # 内部 while 循环 + asyncio.wait_for 逐块消费
        media_type="text/event-stream",
        headers={
            "Cache-Control": "no-cache, no-transform",   # 关键!见踩坑部分
            "X-Accel-Buffering": "no",
            "Connection": "keep-alive",
        },
    )

五、前端实现

5.1 用 fetch 解析 SSE

为什么不用 EventSource:它只支持 GET,且无法自定义 Authorization 头。用 fetch + ReadableStream

export async function chatStream(body, onEvent, signal?) {
  const headers: Record<string, string> = { 'Content-Type': 'application/json' };
  const token = localStorage.getItem('token');
  if (token) headers['Authorization'] = `Bearer ${token}`;

  const resp = await fetch('/llm/v1/chat/stream', {
    method: 'POST', headers, body: JSON.stringify(body), signal,
  });
  if (!resp.ok) throw new Error(`HTTP ${resp.status}`);
  if (!resp.body) throw new Error('响应体不可用');

  const reader = resp.body.getReader();
  const decoder = new TextDecoder('utf-8');
  let buffer = '';

  while (true) {
    const { done, value } = await reader.read();
    if (done) break;
    buffer += decoder.decode(value, { stream: true });

    let sepIdx: number;
    while ((sepIdx = buffer.indexOf('\n\n')) >= 0) {   // 按空行切分 SSE 帧
      const frame = buffer.slice(0, sepIdx);
      buffer = buffer.slice(sepIdx + 2);
      for (const line of frame.split('\n')) {
        if (!line.startsWith('data: ')) continue;
        const payload = line.slice(6).trim();
        if (!payload || payload === '[DONE]') continue;
        onEvent(JSON.parse(payload));   // 交给 React 增量渲染
      }
    }
  }
}

注意:SSE 帧可能跨网络分块到达,所以必须用 buffer 拼接后按 \n\n 切分,不能假设一次 read() 就是一个完整事件。

5.2 增量渲染

状态更新只改最后一条 AI 消息:首个 text 事件创建消息(标记 streaming),后续事件原地追加内容。配合一个闪烁光标,就是完整的打字机效果。

let streamedText = '';
let streamMsgId: string | null = null;

await chatStream(requestData, (ev) => {
  if (ev.type === 'text' && ev.content) {
    streamedText += ev.content;
    if (streamMsgId === null) {
      streamMsgId = String(Date.now() + 1);
      setMessages(prev => [...prev,
        { id: streamMsgId!, type: 'ai', content: streamedText, streaming: true }]);
    } else {
      setMessages(prev => prev.map(m =>
        m.id === streamMsgId ? { ...m, content: streamedText } : m));
    }
  }
  // ev.type === 'tool' 可展示"正在检索…";ev.type === 'error' 提示错误
});

六、踩坑记录:gzip 把流式变成一次性响应

这是本次最值得记录的一环。

现象

后端直连、甚至通过 dev 代理用 curl 测试,事件都是逐块到达的(实测 31 个事件分布在 10~16 秒内)。但浏览器里就是等全部回复结束后一次性渲染,流式完全没有效果。

定位

在浏览器 Network 面板看 chat/stream 响应头,发现两个可疑项:

Content-Encoding: gzip
X-Powered-By: Express

X-Powered-By: Express 说明请求经过了 Umi dev 代理(Express 中间件)。Express 的 compression 中间件默认对接受 gzip 的客户端启用压缩——浏览器请求头带 Accept-Encoding: gzip, deflate, br,而 curl 不带,所以 curl 正常、浏览器异常,就是这个差异造成的假象。

为什么 gzip 会破坏流式?compression 中间件把响应包进 zlib 流,zlib 默认攒够一个块(约 16KB)才输出。对一段几百字的回复,数据会在流结束时一次性 flush,浏览器拿到的就是一整块压缩数据,解压后一次性渲染。

修复:Cache-Control: no-cache, no-transform

Umi 3 的 dev CLI 把 compress: true 写死,devServer.compress = false 覆盖不掉。最稳妥的方案是让后端对这个响应显式声明"不可转换"

headers={
    "Cache-Control": "no-cache, no-transform",
    ...
}

no-transformRFC 7234 §5.2.2.4 定义的标准指令。Express compression 中间件的 shouldTransform() 正是检查它:

var cacheControlNoTransformRegExp = /(?:^|,)\s*?no-transform\s*?(?:,|$)/

function shouldTransform(req, res) {
  var cacheControl = res.getHeader('Cache-Control')
  // Don't compress for Cache-Control: no-transform
  return !cacheControl || !cacheControlNoTransformRegExp.test(cacheControl)
}

命中后中间件直接跳过压缩,流就能原样逐块透传给浏览器。实测对比:同一代理,请求 /umi.js(普通资源)带 Accept-Encoding: gzip 返回 Content-Encoding: gzip;而 SSE 响应带 no-transform 后不再被压缩。

顺带说明:X-Accel-Buffering: no 是 nginx 的概念,对 Express 代理无效;真正让压缩中间件退场的是 no-transform

七、验证方法

不要靠"感觉",用带时间戳的实测脚本验证流是否真的逐块到达:

import time, json, urllib.request

req = urllib.request.Request(url, data=body, method="POST", headers=auth_headers)
t0 = time.time()
with urllib.request.urlopen(req, timeout=120) as r:
    for raw in r:                 # 逐行读取,天然带到达时间信息
        line = raw.decode().strip()
        if line.startswith("data: "):
            ev = json.loads(line[6:])
            print(f"+{(time.time()-prev)*1000:6.0f}ms {ev['type']}")
            prev = time.time()

判读标准:事件间隔呈几十到几百毫秒的分布、且总耗时接近模型生成时长(而不是第一帧就 0ms、最后一帧才全部出现)→ 流式生效。

八、小结与后续优化

这套实现的三个关键决策,值得带走:

  1. 协议选 SSE:单向聊天场景下,比 WebSocket 少 90% 的复杂度,且代理友好;
  2. stream_mode="messages" 让文本流与工具事件自动分离,Agent 的"思考-调用-作答"过程对前端透明化;
  3. no-transform 是流式接口的"免死金牌"——任何中间件只要尊重标准,就不会吃掉你的流。

后续可选优化方向:

  • 取消生成:加一个 AbortController + 后端取消信号(这需要双向通道,届时再考虑 WebSocket 或独立的 cancel 端点);
  • 工具事件可视化:前端把 tool 事件渲染成"正在检索知识库…“的步骤条;
  • Agent 状态持久化:用 LangGraph Checkpointer(SQLite/Postgres)替代内存 dict,跨请求保存 Agent 状态;
  • 重连语义:SSE 标准支持 Last-Event-ID,可配合请求 ID 实现断线续传。