Appearance
流式响应 Streaming
前置阅读:模型调用 Models · 智能体 Agent
本节解决什么问题
大语言模型的推理往往需要数秒甚至数十秒。非流式调用让用户面对空白等待,而流式响应通过增量返回部分结果,从根本上改善体验:
| 维度 | 非流式 | 流式 |
|---|---|---|
| 首字节延迟 | 等待全部生成完毕 | 首个 Token 即返回 |
| 用户感知 | "系统卡死了?" | "AI 正在打字" |
| 工具反馈 | 无中间状态 | 实时进度更新 |
| 资源利用 | 必须等完成 | 可随时取消 |
流式显示即使不减少实际响应时间,也能将用户感知延迟降低 50% 以上——我们习惯看到对方逐字回答,而不是沉默许久后突然出现一大段话。
必须深刻理解,不能跳过:stream_mode 与 stream_events 是两套不同的流式 API
LangChain(
create_agent返回的是 LangGraph 图)的流式响应有两条路径:
agent.stream(input, stream_mode=...)—— 按"模式"返回数据。stream_mode是一个参数,取值values | updates | messages | custom | checkpoints | tasks | debug(共七种,常用四种)。每种模式决定返回数据的粒度和结构。stream()迭代器吐出的是该模式对应的数据块。
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 updates:updates 只返回当步变化的部分(如 {"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_end | Agent 链开始 / 结束 | 生命周期追踪 |
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 resultSSE 场景下将错误传播到前端:
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_mode的checkpoints/tasks/debug模式主要用于调试,不适合前端消费- 某些 Provider 可能不支持 Token 级流式(返回完整 chunk 而非逐 Token)
下一步
- 智能体 Agent - Agent 如何调度工具与模型
- 模型调用 Models - 不同模型对流式的支持差异
- 生产部署 - 流式 Agent 的生产环境部署