使用 v2 流式事件连接前端

让执行过程可见

本章目标

正确处理节点更新、自定义进度、模型 token 和 SSE。

完整代码:examples/streaming_demo.py

for part in graph.stream(
    {"question": "如何安全退款?"},
    stream_mode=["custom", "updates"],
    version="v2",
):
    if part["type"] == "custom":
        print(part["data"])
    elif part["type"] == "updates":
        print(part["data"])

v2 的每个事件都有稳定结构:

流式事件各司其职

{
  "type": "updates | messages | custom | checkpoints | tasks",
  "ns": (),
  "data": ...
}

接入真实 LangChain 模型后,token 通过 messages 模式读取:

async for part in graph.astream(
    inputs,
    stream_mode="messages",
    version="v2",
):
    if part["type"] == "messages":
        message_chunk, metadata = part["data"]
        print(message_chunk.content, end="", flush=True)

不应检查 additional_kwargs["type"] == "token";模型输出本身已经以消息块形式产生。

生产 API 的 SSE 端点位于 src/support_agent/api.py

sequence = 0
async for part in graph.astream(
    input_state,
    config=config,
    stream_mode=["updates", "tasks"],
    version="v2",
):
    if await request.is_disconnected():
        break
    sequence += 1
    payload = json.dumps(part, ensure_ascii=False, default=str)
    yield (
        f"id: {stream_id}:{sequence}\n"
        f"event: {part['type']}\n"
        f"data: {payload}\n\n"
    )

前端断线不等于运行状态丢失。SSE 只传输事件;恢复能力来自 Checkpointer 和稳定的 thread_id

传输与恢复是两件事

把事件格式当成 API 契约

前端不应直接依赖 LangGraph 内部对象的所有字段。服务端应定义稳定外层字段,并允许 data 随事件类型变化:

id: tenant-a:T-300:7
event: updates
data: {"type":"updates","ns":[],"data":{"query_order":{...}}}

id 用于日志关联和客户端去重;它不等于 Checkpoint ID。事件内容可能包含模型输入、工具参数或业务数据,上线前必须按事件类型做字段白名单和脱敏,不能简单把任意对象 default=str 后全部发给浏览器。

浏览器接入与断线语义

原生 EventSource 只支持 GET,示例端点使用 POST 提交问题,因此浏览器应使用 fetch() 读取流,或者采用更适合生产的两步协议:

  1. POST /runs 创建运行并返回 run_id
  2. GET /runs/{run_id}/events 使用 EventSource 订阅。

POST 流的最小客户端如下:

const response = await fetch("/tickets/T-300/stream", {
  method: "POST",
  headers: {"Content-Type": "application/json", "X-Tenant-ID": "tenant-a"},
  body: JSON.stringify({question: "订单退款"}),
  signal: abortController.signal,
});

const reader = response.body.getReader();
const decoder = new TextDecoder();
while (true) {
  const {value, done} = await reader.read();
  if (done) break;
  consumeSseFrames(decoder.decode(value, {stream: true}));
}

真实客户端必须处理半个 UTF-8 字符、跨 chunk 的半个 SSE frame、重复事件和非 2xx 响应;不要假设一次 read() 就是一条完整事件。

代理、心跳与背压

src/support_agent/api.py 为 SSE 返回 Cache-Control: no-cacheX-Accel-Buffering: no,并在检测到客户端断开后停止继续写流。部署时还应检查:

  • Nginx、Ingress 和 CDN 没有缓冲 text/event-stream
  • 空闲超时大于最长静默节点耗时,或定期发送注释心跳 : ping\n\n
  • 慢客户端有队列上限;超过上限时断开并让客户端从状态接口恢复。
  • 客户端重连不会重新执行退款等副作用。
  • 运行创建和事件订阅都经过同一租户授权。

客户端断开是否取消图运行不能由 TCP 状态自动决定。对长任务,常见做法是让运行继续落入 Checkpointer,前端重新连接后查询当前状态;只有明确标记为可取消的只读运行才传播取消信号。

本章验收

  • 能说明 updatesmessagescustomtasks 的消费者分别是谁。
  • SSE 响应包含事件 ID、事件类型、禁缓存和禁代理缓冲头。
  • 前端解析器可以处理拆包、重复事件和主动取消。
  • 断线重连不会导致同一业务动作执行两次。