Skip to content

流式响应 Streaming

前置阅读:模型调用 Models · 智能体 Agent

本节解决什么问题

大语言模型的推理往往需要数秒甚至数十秒。非流式调用让用户面对空白等待,而流式响应通过增量返回部分结果,从根本上改善体验:

维度非流式流式
首字节延迟等待全部生成完毕首个 Token 即返回
用户感知"系统卡死了?""AI 正在打字"
工具反馈无中间状态实时进度更新
资源利用必须等完成可随时取消

流式显示即使不减少实际响应时间,也能将用户感知延迟降低 50% 以上——我们习惯看到对方逐字回答,而不是沉默许久后突然出现一大段话。

必须深刻理解,不能跳过:stream_mode 与 stream_events 是两套不同的流式 API

LangChain(create_agent 返回的是 LangGraph 图)的流式响应有两条路径:

  1. agent.stream(input, stream_mode=...) —— 按"模式"返回数据。stream_mode 是一个参数,取值 values | updates | messages | custom | checkpoints | tasks | debug(共七种,常用四种)。每种模式决定返回数据的粒度和结构。stream() 迭代器吐出的是该模式对应的数据块

  2. agent.stream_events(input, version="v3") / astream_events(input, version="v3")(langchain>=1.3)—— 事件流 API,返回类型化的事件投影,覆盖模型 token 流、工具调用开始 / 结束、链生命周期等全流程事件。它不是 stream_mode 的一个取值——"events" 不是合法的 stream_mode

简单记忆:stream() 是"选一个数据视角",stream_events() 是"订阅全生命周期事件"。新应用推荐 stream_events(version="v3"),因为它返回结构化事件,不需要手动解析元组或区分 node。

前端类比

流式响应对前端开发者并不陌生:

  • SSE (Server-Sent Events):服务器单向推送数据流,LangChain 的 HTTP 流式通常通过 SSE 实现
  • WebSocket:双向实时通信,适合客户端需要中途发送取消指令的场景
  • React Server Components Streaming:Next.js RSC 逐步将 UI 片段送达客户端,与 LLM Token 流式输出理念一致——"准备好一部分就先发一部分"
  • ReadableStream:Web Streams API 的流式原语,for await...of 消费 ReadableStream 和 agent.stream() 体验完全一样

原生语义

前端的 SSE / ReadableStream 是传输层流式,数据一旦到达就是完整的 JSON 或文本。LangChain 的流式是应用层流式:stream_mode="messages" 返回的是模型推理过程中的 Token 片段(可能是不完整的词),stream_mode="updates" 返回的是节点执行后的状态差量。理解这一差异才能正确处理数据拼接和 UI 更新。

🔗 LangChain Event Streaming 官方文档 · 🔗 LangGraph 流式概念

流式数据流全景

stream() 的四种常用模式

stream_mode粒度返回内容场景
"updates"步骤级节点执行后的状态差量追踪 Agent 决策、调试
"messages"Token 级每个 Token 增量 + 元数据元组聊天打字效果
"custom"自定义工具内发出的任意数据长任务进度条
"values"步骤级每步执行后的完整状态需要完整快照而非差量

注意

stream_mode 共有七种取值:values | updates | messages | custom | checkpoints | tasks | debug"events" 不是 stream_mode——事件流通过 stream_events() 方法获取,详见下文。

updates - 步骤级别更新

每个执行步骤完成后返回一次,包含该步骤的完整输出:

python
import os
from langchain.chat_models import init_chat_model
from langchain.agents import create_agent

def get_weather(city: str) -> str:
    """获取天气"""
    return f"{city}:晴,25°C"

model = init_chat_model(os.environ["LLM_MODEL"])
agent = create_agent(model, tools=[get_weather])

for chunk in agent.stream(
    {"messages": [{"role": "user", "content": "北京天气如何?"}]},
    stream_mode="updates",
):
    for node_name, node_output in chunk.items():
        last_msg = node_output["messages"][-1]
        if hasattr(last_msg, "tool_calls") and last_msg.tool_calls:
            for tc in last_msg.tool_calls:
                print(f"[{node_name}] 调用工具: {tc['name']}({tc['args']})")
        elif hasattr(last_msg, "content"):
            print(f"[{node_name}] {last_msg.content[:80]}")

输出:

[agent] 调用工具: get_weather({'city': '北京'})
[tools] 北京:晴,25°C
[agent] 北京今天天气很好!晴朗,气温 25°C,非常适合外出。

messages - Token 级别流式

最细粒度——每个 Token 生成后立即返回,实现"打字机效果":

python
for event in agent.stream(
    {"messages": [{"role": "user", "content": "搜索活跃用户"}]},
    stream_mode="messages",
):
    msg_chunk, metadata = event  # 元组:(消息片段, 元数据)
    node = metadata.get("langgraph_node", "")

    if node == "agent" and msg_chunk.content:
        print(msg_chunk.content, end="", flush=True)  # 逐字输出

    if node == "agent" and msg_chunk.tool_call_chunks:
        for tc in msg_chunk.tool_call_chunks:
            if tc.get("name"):
                print(f"\n[调用: {tc['name']}]")

    if node == "tools" and msg_chunk.content:
        print(f"\n[结果: {msg_chunk.content}]")

关键数据结构:

python
msg_chunk.content            # str - 文本片段
msg_chunk.tool_call_chunks   # list - 工具调用增量
metadata["langgraph_node"]   # "agent" | "tools"
metadata["langgraph_step"]   # int - 步骤序号

custom - 自定义流式更新

工具函数内部通过 get_stream_writer() 发送任意数据,适合长时间操作的进度反馈:

python
import time
from langgraph.config import get_stream_writer

def analyze_dataset(name: str) -> str:
    """分析数据集,过程中报告进度"""
    writer = get_stream_writer()

    for pct in [0, 25, 50, 75, 100]:
        writer({"phase": "loading", "progress": pct})
        time.sleep(0.3)

    writer({"phase": "analyzing", "progress": 100, "message": "分析完成"})
    return f"{name} 分析结果:50,000 条记录,日活 12,350"

agent = create_agent(model, tools=[analyze_dataset])

for chunk in agent.stream(
    {"messages": [{"role": "user", "content": "分析用户数据"}]},
    stream_mode="custom",
):
    if isinstance(chunk, dict):
        print(f"[{chunk.get('phase')}] {chunk.get('progress', 0)}%")

get_stream_writer() 要点:必须在工具函数内部调用;写入任意可序列化对象;仅 stream_mode"custom" 时才被消费。

values - 完整状态快照

updates 返回差量不同,values 在每步执行后返回完整的当前状态,适合需要完整快照的场景:

python
for chunk in agent.stream(
    {"messages": [{"role": "user", "content": "北京天气如何?"}]},
    stream_mode="values",
):
    # chunk 是完整状态字典,不是差量
    messages = chunk.get("messages", [])
    if messages:
        last = messages[-1]
        print(f"[step] {type(last).__name__}: {getattr(last, 'content', '')[:60]}")

values vs updatesupdates 只返回当步变化的部分(如 {"agent": {"messages": [...]}}),values 返回累积后的完整状态。需要差量做 diff 时选 updates,需要完整快照渲染 UI 时选 values

事件流:stream_events(version="v3")

当前推荐写法

langchain>=1.3 起,新应用推荐使用 stream_events(version="v3")。它返回类型化的事件投影,覆盖模型 token 流、工具调用生命周期、链执行等全流程,不需要手动解析 (msg_chunk, metadata) 元组或判断 langgraph_node

stream_events 是一个独立方法,不是 stream_mode 的取值。它返回的是事件对象,每个事件描述执行生命周期中发生的一件事:

python
import asyncio

async def stream_with_events():
    """事件流:覆盖模型 token 流 + 工具调用生命周期"""
    async for event in agent.astream_events(
        {"messages": [{"role": "user", "content": "北京天气如何?"}]},
        version="v3",
    ):
        kind = event["event"]

        # 模型 Token 流 —— 逐字输出
        if kind == "on_chat_model_stream":
            chunk = event["data"]["chunk"]
            if chunk.content:
                print(chunk.content, end="", flush=True)

        # 工具调用开始
        elif kind == "on_tool_start":
            print(f"\n[调用工具] {event['name']}")

        # 工具调用结束
        elif kind == "on_tool_end":
            output = event["data"].get("output", "")
            print(f"\n[工具结果] {output}")

asyncio.run(stream_with_events())

同步版本使用 agent.stream_events(input, version="v3")

python
for event in agent.stream_events(
    {"messages": [{"role": "user", "content": "北京天气如何?"}]},
    version="v3",
):
    kind = event["event"]
    if kind == "on_chat_model_stream":
        chunk = event["data"]["chunk"]
        if chunk.content:
            print(chunk.content, end="", flush=True)
    elif kind == "on_tool_start":
        print(f"\n[调用工具] {event['name']}")
    elif kind == "on_tool_end":
        print(f"\n[工具结果] {event['data'].get('output', '')}")

常见事件类型:

事件触发时机用途
on_chat_model_stream模型每生成一个 Token打字机效果
on_chat_model_start模型调用开始显示"正在思考"
on_chat_model_end模型调用结束获取完整 AIMessage
on_tool_start工具执行开始显示工具调用信息
on_tool_end工具执行结束显示工具结果
on_chain_start / on_chain_endAgent 链开始 / 结束生命周期追踪

stream_events vs stream(stream_mode="messages"):两者都能获取模型 Token 流,但 stream_events 额外提供工具调用、链生命周期等结构化事件,且返回格式统一(都是事件对象),不需要判断 langgraph_node。对于需要同时处理 Token 流和工具事件的前端应用,stream_events 更简洁。

旧版

stream_events(version="v1") 已弃用。version="v2" 在 langchain 1.1-1.2 期间可用。新应用直接使用 version="v3"(langchain>=1.3)。

多模式组合

stream_mode 设为列表即可同时获取多种粒度的数据:

python
for mode, chunk in agent.stream(
    {"messages": [{"role": "user", "content": "生成运营报告"}]},
    stream_mode=["messages", "custom"],
):
    if mode == "messages":
        msg_chunk, metadata = chunk
        if metadata.get("langgraph_node") == "agent" and msg_chunk.content:
            print(msg_chunk.content, end="", flush=True)
    elif mode == "custom":
        print(f"\n[进度] {chunk}")

每次迭代返回 (mode, chunk) 元组——通过 mode 区分数据来源,chunk 结构取决于对应模式。

注意

组合模式下不同模式的数据会交织出现。务必通过 mode 字段区分处理。

前端集成模式

SSE + FastAPI(推荐)

python
# server.py
import os
import json
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
from langchain.chat_models import init_chat_model
from langchain.agents import create_agent

app = FastAPI()

model = init_chat_model(os.environ["LLM_MODEL"])
agent = create_agent(model, tools=[])

async def event_generator(query: str):
    for event in agent.stream(
        {"messages": [{"role": "user", "content": query}]},
        stream_mode="messages",
    ):
        msg_chunk, metadata = event
        node = metadata.get("langgraph_node", "")
        if node == "agent" and msg_chunk.content:
            yield f"data: {json.dumps({'type': 'token', 'content': msg_chunk.content}, ensure_ascii=False)}\n\n"
    yield "data: [DONE]\n\n"

@app.get("/api/chat/stream")
async def stream_chat(query: str):
    return StreamingResponse(
        event_generator(query),
        media_type="text/event-stream",
        headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"},
    )
typescript
// React 前端
function useChatStream(query: string) {
  const [text, setText] = useState('')

  const start = useCallback(() => {
    const es = new EventSource(`/api/chat/stream?query=${encodeURIComponent(query)}`)
    es.onmessage = (e) => {
      if (e.data === '[DONE]') return es.close()
      const d = JSON.parse(e.data)
      if (d.type === 'token') setText((p) => p + d.content)
    }
    es.onerror = () => es.close()
    return () => es.close()
  }, [query])

  return { text, start }
}

React useStream Hook(@langchain/sdk)

LangChain 官方封装,自动处理连接管理和消息状态。适用于部署在 Agent Server 或自托管 langgraph up 的场景:

typescript
import { useStream } from '@langchain/sdk/react'

function Chat() {
  const { messages, start, stop, isStreaming } = useStream({
    apiUrl: 'http://localhost:8000',
    assistantId: 'my-agent', // Agent Server 中的 assistant ID
    messagesKey: 'messages',
  })

  return (
    <div>
      {messages.map((m, i) => <div key={i}>{m.content}</div>)}
      {isStreaming && <span>AI 正在输入...</span>}
      <button onClick={stop}>停止</button>
    </div>
  )
}

优势:自动 SSE 重连、内置状态管理、stop() 取消、与 Agent Server 无缝集成。

原生语义

前端的 useStream Hook 类似 React Query 的 useQuery,但数据是持续追加而非一次性返回。messages 数组在流式过程中会不断增长,每个 Token 到达时触发重渲染。这与 SSE 的 onmessage 手动拼接 text 的区别在于:useStream 自动维护结构化的消息列表,包括工具调用、中断等复杂状态。

WebSocket 方案

需要双向通信(如客户端实时取消)时使用:

python
from fastapi import WebSocket

@app.websocket("/ws/chat")
async def ws_chat(ws: WebSocket):
    await ws.accept()
    try:
        while True:
            data = await ws.receive_json()
            if data.get("type") == "cancel":
                break
            for event in agent.stream(
                {"messages": [{"role": "user", "content": data["content"]}]},
                stream_mode="messages",
            ):
                msg_chunk, meta = event
                if meta.get("langgraph_node") == "agent" and msg_chunk.content:
                    await ws.send_json({"type": "token", "content": msg_chunk.content})
            await ws.send_json({"type": "done"})
    finally:
        await ws.close()

错误处理

流式响应中错误可能在任意位置发生,需要专门的处理策略:

python
import time

def stream_with_retry(agent, query: str, max_retries: int = 3):
    """带指数退避重试 + 非流式回退"""
    for attempt in range(max_retries):
        try:
            for chunk in agent.stream(
                {"messages": [{"role": "user", "content": query}]},
                stream_mode="messages",
            ):
                yield chunk
            return  # 成功完成
        except (ConnectionError, TimeoutError) as e:
            wait = 2 ** attempt
            print(f"第 {attempt + 1} 次失败,{wait}s 后重试: {e}")
            time.sleep(wait)

    # 重试耗尽,回退到非流式
    result = agent.invoke({"messages": [{"role": "user", "content": query}]})
    yield result

SSE 场景下将错误传播到前端:

python
async def safe_event_generator(query: str):
    try:
        for event in agent.stream(
            {"messages": [{"role": "user", "content": query}]},
            stream_mode="messages",
        ):
            msg_chunk, metadata = event
            if metadata.get("langgraph_node") == "agent" and msg_chunk.content:
                yield f"data: {json.dumps({'type': 'token', 'content': msg_chunk.content}, ensure_ascii=False)}\n\n"
    except Exception as e:
        yield f"data: {json.dumps({'type': 'error', 'message': str(e)}, ensure_ascii=False)}\n\n"
    finally:
        yield "data: [DONE]\n\n"

流式取消

Python 端:break 即取消

python
for chunk in agent.stream(
    {"messages": [{"role": "user", "content": "..."}]},
    stream_mode="messages",
):
    msg_chunk, metadata = chunk
    if metadata.get("langgraph_node") == "agent" and msg_chunk.content:
        print(msg_chunk.content, end="", flush=True)
    if should_cancel():
        break  # 跳出循环即停止消费

前端:AbortController

typescript
const controller = new AbortController()

// 启动流式
fetch(`/api/chat/stream?query=${query}`, { signal: controller.signal })
  .then(async (res) => {
    const reader = res.body!.getReader()
    while (true) {
      const { done, value } = await reader.read()
      if (done) break
      processChunk(new TextDecoder().decode(value))
    }
  })
  .catch((err) => {
    if (err.name === 'AbortError') console.log('已取消')
  })

// 取消
controller.abort()

常见问题

Q:流式和非流式的总耗时有区别吗?

没有。LLM 推理时间相同,流式只是把"等完再返回"改为"生成一部分就返回一部分"。

Q:stream() 和 stream_events() 怎么选?

  • 只需要 Token 打字效果 -> stream(stream_mode="messages") 足够
  • 需要同时处理 Token 流 + 工具调用事件 + 生命周期 -> stream_events(version="v3")(新应用推荐)
  • 需要工具内自定义进度数据 -> stream(stream_mode="custom") 或组合 ["messages", "custom"]

Q:三种常用模式怎么选?

  • 聊天打字效果 -> messages
  • 执行步骤可视化 -> updates
  • 工具进度条 -> custom
  • 完整状态快照 -> values
  • 综合需求 -> ["messages", "custom"]stream_events(version="v3")

Q:前端选 SSE 还是 WebSocket?

大多数场景选 SSE——实现简单、HTTP 兼容、自动重连。只有需要客户端主动发数据(如实时取消、追加上下文)时才用 WebSocket。

Q:get_stream_writer() 可以在工具外使用吗?

不可以。它依赖 LangGraph 运行时上下文,仅在工具函数执行期间可用。

真实项目适用场景

  • 聊天界面stream_mode="messages"stream_events(version="v3") 实现 Token 级打字效果
  • Agent 执行可视化stream_mode="updates" 展示每步决策(调了什么工具、返回了什么)
  • 长任务进度条stream_mode="custom" + get_stream_writer() 从工具内部推送进度
  • 完整状态同步stream_mode="values" 在每步获取完整状态用于 UI 渲染

不适用场景 / 能力边界

  • 流式不能减少 LLM 推理总耗时,只是提前返回部分结果
  • stream_events(version="v3") 需要 langchain>=1.3,旧版本用 version="v2"
  • stream_modecheckpoints / tasks / debug 模式主要用于调试,不适合前端消费
  • 某些 Provider 可能不支持 Token 级流式(返回完整 chunk 而非逐 Token)

下一步

参考资源

学习文档整合站点