Skip to content

流式 API 与 SSE

LLM 生成 token 需要几秒到几十秒,如果让用户盯着空白转圈,体验会非常糟。流式输出让用户看到「逐字蹦出来」的效果,是智能体应用的标配。本篇讲清楚:SSE 协议是什么、怎么用 FastAPI 把 LangGraph 的 astream 转成前端能消费的事件流、以及多 stream_mode 怎么聚合。

一、为什么是 SSE 而不是 WebSocket

协议方向复杂度适用
普通 HTTP请求-响应一次性结果
SSE服务端→客户端单向LLM token 流、状态更新
WebSocket双向多人协作、实时游戏

LLM 流式输出本质是「服务端持续推、客户端只收不发」,SSE 是最契合的协议:

  • 基于 HTTP,穿透防火墙/代理无压力
  • 浏览器原生 EventSource API,无需额外库
  • 自动重连机制内置
  • 比 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 增量哪个节点改了啥
messagesLLM 的 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 8000

3. 用 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 更新流程图高亮。

七、事件格式设计建议

约定一套稳定事件协议,前后端省事:

eventdata 含义时机
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-8media_type="text/event-stream; charset=utf-8"

九、小结

  • SSE 是 LLM 流式输出最合适的协议:单向、轻量、浏览器原生支持
  • FastAPI 用 StreamingResponse + 异步 generator 把 astream(stream_mode="messages") 转成 event/data
  • 上线务必关 nginx 缓冲(proxy_buffering off),并约定稳定事件协议
  • stream_mode 聚合能同时满足「逐字显示」和「流程图高亮」

下一篇讲怎么给服务加监控、日志、链路追踪,让生产环境可观测。