9.流式与实时交互
LangGraph 的流式与实时交互机制,本质上是为了解决一个核心问题:如何让 AI 应用的响应过程对用户可见。传统的请求-响应模式像在黑箱中操作,用户只能等待最终结果。而流式交互则像打开了一扇窗,让整个过程变得透明、可感知。
这种透明性在现代 AI 应用中至关重要。当 LLM 生成一段长文本时,用户希望看到文字逐字出现;当工具在后台执行耗时操作时,用户需要进度提示;当多步骤工作流运行时,开发者需要调试信息。LangGraph 的流式系统正是为满足这些需求而设计。
流模式类型配置
LangGraph 提供了六种流模式,每种模式对应不同的数据粒度。理解它们的区别,是合理使用流式功能的第一步。
六种流模式详解
values 模式会推送完整的图状态。每次节点执行后,整个状态对象都会被发送。这种模式适合需要完整上下文的场景,比如调试时查看所有变量。但缺点是数据量大,网络开销高。
updates 模式只推送状态的变化部分。它返回节点名称和该节点产生的更新,格式为 {'node_name': {'field': 'new_value'}}。这是大多数应用的首选模式,既保持了信息完整性,又避免了冗余数据传输。
messages-tuple(或 messages)模式专门用于流式传输 LLM 生成的 Token。它返回一个元组 (message_chunk, metadata),其中 message_chunk 包含 Token 内容,metadata 包含节点信息、标签等上下文。构建聊天界面时,这个模式几乎是必选项。
custom 模式允许从图的任意位置发送自定义数据。通过 get_stream_writer() 获取写入器,可以推送任意 JSON 可序列化的数据。这在工具执行进度、业务状态通知等场景中非常有用。
debug 模式会推送尽可能多的执行信息,包括节点名称、完整状态、执行时间等。它主要用于开发和调试阶段,生产环境应谨慎使用,避免性能损耗。
events 模式会流式传输所有事件,包括图状态变化、LLM 调用、工具执行等。这个模式主要用于迁移旧版 LCEL 应用,新应用通常不需要。
模式组合使用
实际应用中,往往需要同时获取多种类型的数据。LangGraph 支持传入模式列表来实现这一点:
# Python 示例:同时获取状态更新和 LLM Token
for mode, chunk in agent.stream(
{"messages": [{"role": "user", "content": "讲个笑话"}]},
stream_mode=["updates", "messages-tuple"]
):
if mode == "updates":
print(f"状态更新: {chunk}")
elif mode == "messages-tuple":
message_chunk, metadata = chunk
print(f"Token: {message_chunk.content}")
// JavaScript 示例:组合模式
for await (const [mode, chunk] of await agent.stream(
{ messages: [{ role: "user", content: "讲个笑话" }] },
{ streamMode: ["updates", "messages-tuple"] }
)) {
if (mode === "updates") {
console.log(`状态更新: ${JSON.stringify(chunk)}`);
} else if (mode === "messages-tuple") {
const [messageChunk, metadata] = chunk;
console.log(`Token: ${messageChunk.content}`);
}
}
代码中,stream_mode 参数接受列表形式,流式输出会返回 (mode, chunk) 元组。通过判断 mode 值,可以区分不同类型的数据。这种模式组合让前端能够同时渲染对话历史和实时生成的文本。
LLM Token 实时推送
Token 流式传输是聊天应用的核心功能。LangGraph 的实现有几个关键点值得注意。
基本用法
即使使用 .invoke() 而非 .stream() 调用 LLM,只要设置了 stream_mode="messages-tuple",Token 依然会被流式传输。这是因为 LangGraph 在运行时层面拦截了 LLM 调用,自动处理流式逻辑。
# Python 示例:Token 流式传输
from langchain.chat_models import init_chat_model
from langgraph.graph import StateGraph, START
llm = init_chat_model(model="openai:gpt-4o-mini")
def call_model(state):
# 使用 .invoke() 而非 .stream()
response = llm.invoke([{"role": "user", "content": f"讲个关于{state['topic']}的笑话"}])
return {"joke": response.content}
graph = StateGraph(State).add_node(call_model).add_edge(START, "call_model").compile()
# 但依然可以流式接收 Token
async for message_chunk, metadata in graph.astream(
{"topic": "猫"},
stream_mode="messages-tuple"
):
if message_chunk.content:
print(message_chunk.content, end="|", flush=True)
// JavaScript 示例
import { ChatOpenAI } from "@langchain/openai";
const llm = new ChatOpenAI({ model: "gpt-4o-mini" });
const callModel = async (state: any) => {
const response = await llm.invoke([
{ role: "user", content: `讲个关于${state.topic}的笑话` }
]);
return { joke: response.content };
};
// 流式接收 Token
for await (const [messageChunk, metadata] of await graph.stream(
{ topic: "猫" },
{ streamMode: "messages-tuple" }
)) {
if (messageChunk.content) {
process.stdout.write(messageChunk.content + "|");
}
}
这种设计简化了开发流程。无需在业务逻辑中考虑流式细节,只需在调用图时指定模式即可。.invoke() 和 .stream() 的区别仅在于调用方式,流式能力由运行时统一提供。
过滤特定 LLM 调用
复杂应用中可能包含多个 LLM 调用。通过 tags 参数可以区分它们:
# Python 示例:为不同 LLM 添加标签
joke_model = init_chat_model(model="openai:gpt-4o-mini", tags=['joke'])
poem_model = init_chat_model(model="openai:gpt-4o-mini", tags=['poem'])
async def call_model(state, config):
topic = state["topic"]
# 生成笑话
joke_response = await joke_model.ainvoke(
[{"role": "user", "content": f"写个关于{topic}的笑话"}],
config
)
# 生成诗歌
poem_response = await poem_model.ainvoke(
[{"role": "user", "content": f"写首关于{topic}的诗"}],
config
)
return {"joke": joke_response.content, "poem": poem_response.content}
# 只接收笑话模型的 Token
async for msg, metadata in graph.astream(
{"topic": "猫"},
stream_mode="messages-tuple"
):
if metadata["tags"] == ["joke"]:
print(msg.content, end="|", flush=True)
// JavaScript 示例
const jokeModel = new ChatOpenAI({ model: "gpt-4o-mini", tags: ["joke"] });
const poemModel = new ChatOpenAI({ model: "gpt-4o-mini", tags: ["poem"] });
// 只接收笑话模型的 Token
for await (const [msg, metadata] of await graph.stream(
{ topic: "猫" },
{ streamMode: "messages-tuple" }
)) {
if (metadata.tags?.includes("joke")) {
process.stdout.write(msg.content + "|");
}
}
tags 参数在初始化模型时设置,会包含在元数据中。通过过滤 metadata["tags"],可以精确控制显示哪个 LLM 的输出。这在多模型协作的场景中特别有用,比如一个模型生成内容,另一个模型评估质量。
按节点过滤
除了按标签过滤,还可以按节点名称过滤:
# Python 示例:只接收特定节点的 Token
for msg, metadata in graph.stream(
inputs,
stream_mode="messages-tuple"
):
# 检查节点名称
if metadata["langgraph_node"] == "write_poem":
print(msg.content, end="|", flush=True)
// JavaScript 示例
for await (const [msg, metadata] of await graph.stream(
inputs,
{ streamMode: "messages-tuple" }
)) {
if (metadata.langgraph_node === "writePoem") {
process.stdout.write(msg.content + "|");
}
}
metadata["langgraph_node"] 字段记录了产生该 Token 的节点名称。这在并行执行多个 LLM 调用时非常实用,可以分别渲染不同节点的输出。
工具执行进度流
工具执行可能耗时较长,比如查询数据库、调用外部 API。让用户感知到进度,能显著提升体验。
在工具中发送进度
Python 中通过 get_stream_writer() 获取写入器,JavaScript 中通过 config.writer 实现:
# Python 示例:工具进度流
from langgraph.config import get_stream_writer
from langchain_core.tools import tool
@tool
def query_database(query: str) -> str:
"""查询数据库"""
writer = get_stream_writer()
# 发送进度更新
writer({"type": "progress", "data": "已检索 0/100 条记录"})
# 模拟耗时操作
# ... 执行查询 ...
writer({"type": "progress", "data": "已检索 50/100 条记录"})
# ... 更多操作 ...
writer({"type": "progress", "data": "已检索 100/100 条记录"})
return "查询完成"
# 使用 custom 模式接收进度
for chunk in graph.stream(
{"query": "SELECT * FROM users"},
stream_mode="custom"
):
if chunk.get("type") == "progress":
print(f"进度: {chunk['data']}")
// JavaScript 示例
import { tool } from "@langchain/core/tools";
const queryDatabase = tool(
async (input: any, config: any) => {
// 发送进度更新
config.writer({ type: "progress", data: "已检索 0/100 条记录" });
// 模拟耗时操作
// ... 执行查询 ...
config.writer({ type: "progress", data: "已检索 50/100 条记录" });
// ... 更多操作 ...
config.writer({ type: "progress", data: "已检索 100/100 条记录" });
return "查询完成";
},
{
name: "query_database",
description: "查询数据库",
schema: z.object({ query: z.string() })
}
);
// 接收进度
for await (const chunk of await graph.stream(
{ query: "SELECT * FROM users" },
{ streamMode: "custom" }
)) {
if (chunk.type === "progress") {
console.log(`进度: ${chunk.data}`);
}
}
进度数据可以是任意结构。示例中使用 {"type": "progress", "data": "..."} 格式,前端可以根据 type 字段区分消息类型,渲染不同的 UI 组件。
注意事项
在 Python 异步代码中(Python < 3.11),get_stream_writer() 无法自动获取上下文。需要显式传递 writer 参数:
# Python 3.10 及以下版本需要这样做
from langgraph.types import StreamWriter
async def query_database(query: str, writer: StreamWriter) -> str:
writer({"type": "progress", "data": "开始查询"})
# ... 异步操作 ...
return "结果"
# 图定义中不需要特殊处理
graph = StateGraph(State).add_node(query_database).compile()
LangGraph 会自动检测函数签名中的 writer 参数并注入。这个限制在 Python 3.11+ 中已解决,因为 asyncio 支持了上下文传播。
自定义事件数据流
除了进度,还可以发送任意业务相关的事件数据。比如,在电商客服机器人中,可以推送订单状态更新;在数据分析工具中,可以发送中间计算结果。
节点内发送自定义数据
# Python 示例:节点内自定义事件
from langgraph.config import get_stream_writer
def process_order(state):
writer = get_stream_writer()
# 发送订单确认
writer({"event": "order_confirmed", "order_id": state["order_id"]})
# 模拟处理
# ... 处理逻辑 ...
# 发送发货通知
writer({"event": "order_shipped", "tracking_number": "12345"})
return {"status": "completed"}
# 接收自定义事件
for chunk in graph.stream(
{"order_id": "ORD-001"},
stream_mode="custom"
):
if chunk.get("event") == "order_confirmed":
print(f"订单已确认: {chunk['order_id']}")
elif chunk.get("event") == "order_shipped":
print(f"订单已发货,追踪号: {chunk['tracking_number']}")
// JavaScript 示例
import { LangGraphRunnableConfig } from "@langchain/langgraph";
const processOrder = async (state: any, config: LangGraphRunnableConfig) => {
config.writer({ event: "order_confirmed", order_id: state.order_id });
// ... 处理逻辑 ...
config.writer({ event: "order_shipped", tracking_number: "12345" });
return { status: "completed" };
};
// 接收自定义事件
for await (const chunk of await graph.stream(
{ order_id: "ORD-001" },
{ streamMode: "custom" }
)) {
if (chunk.event === "order_confirmed") {
console.log(`订单已确认: ${chunk.order_id}`);
} else if (chunk.event === "order_shipped") {
console.log(`订单已发货,追踪号: ${chunk.tracking_number}`);
}
}
自定义事件让图的状态机具备了事件驱动能力。前端可以订阅特定事件,触发相应的 UI 更新,比如显示通知、更新进度条、切换页面等。
与非 LangChain LLM 集成
如果使用的 LLM 没有 LangChain 集成,可以通过 custom 模式手动流式传输:
# Python 示例:集成自定义 LLM
from langgraph.config import get_stream_writer
from openai import AsyncOpenAI
client = AsyncOpenAI()
async def stream_custom_llm(state):
writer = get_stream_writer()
response = await client.chat.completions.create(
model="gpt-4o-mini",
messages=[{"role": "user", "content": state["query"]}],
stream=True
)
full_content = ""
async for chunk in response:
if chunk.choices[0].delta.content:
content = chunk.choices[0].delta.content
full_content += content
writer({"custom_llm_chunk": content}) # 手动发送 Token
return {"response": full_content}
# 接收自定义 LLM 的 Token
for chunk in graph.stream(
{"query": "讲个笑话"},
stream_mode="custom"
):
if "custom_llm_chunk" in chunk:
print(chunk["custom_llm_chunk"], end="|", flush=True)
这种方式提供了最大灵活性。无论 LLM 提供何种接口,只要支持流式响应,就可以集成到 LangGraph 的流式体系中。
React 应用流集成
前端集成是流式功能的最终落地点。LangGraph 提供了 useStream Hook,大幅简化了 React 应用的开发。
基础用法
// React 组件示例
"use client";
import { useStream } from "@langchain/langgraph-sdk/react";
import type { Message } from "@langchain/langgraph-sdk";
export default function ChatApp() {
const thread = useStream<{ messages: Message[] }>({
apiUrl: "http://localhost:2024",
assistantId: "agent",
messagesKey: "messages",
});
return (
<div>
{/* 渲染消息历史 */}
<div>
{thread.messages.map((message) => (
<div key={message.id}>
<strong>{message.type}:</strong> {message.content as string}
</div>
))}
</div>
{/* 输入表单 */}
<form
onSubmit={(e) => {
e.preventDefault();
const form = e.target as HTMLFormElement;
const message = new FormData(form).get("message") as string;
form.reset();
// 提交消息并自动处理流式响应
thread.submit({ messages: [{ type: "human", content: message }] });
}}
>
<input type="text" name="message" placeholder="输入消息..." />
{/* 根据加载状态显示不同按钮 */}
{thread.isLoading ? (
<button type="button" onClick={() => thread.stop()}>
停止
</button>
) : (
<button type="submit">发送</button>
)}
</form>
</div>
);
}
useStream Hook 封装了所有复杂逻辑:创建线程、管理状态、处理流式响应、更新消息列表。开发者只需关注 UI 渲染。
关键特性
自动消息管理:thread.messages 会自动更新。当 LLM 生成新 Token 时,Hook 会合并到最新消息中,触发组件重新渲染。
加载状态:thread.isLoading 指示是否有正在进行的流式请求。可以用来禁用输入框、显示加载动画。
中断支持:thread.stop() 可以取消当前流式请求。这在用户想停止生成时很有用。
分支管理:thread.branch 和 thread.setBranch() 支持对话分支。用户可以回到历史消息,从该点开启新对话分支。
错误处理:thread.error 包含任何发生的错误。可以在 UI 中显示错误提示。
事件回调
对于更细粒度的控制,可以使用回调函数:
// 高级事件处理
const thread = useStream<{ messages: Message[] }>({
apiUrl: "http://localhost:2024",
assistantId: "agent",
messagesKey: "messages",
// 自定义事件
onCustomEvent: (data) => {
if (data.type === "progress") {
console.log(`进度: ${data.data}`);
}
},
// 错误处理
onError: (error) => {
console.error("流式错误:", error);
alert("发生错误,请重试");
},
// 完成回调
onFinish: (state) => {
console.log("流式完成:", state);
},
// 元数据事件(包含 runId 和 threadId)
onMetadataEvent: (data) => {
console.log("运行元数据:", data);
}
});
这些回调让应用能够响应各种事件,实现进度条、错误提示、日志记录等功能。
断线重连
useStream 支持自动重连,这在页面刷新后恢复对话时特别有用:
// 自动重连配置
const thread = useStream<{ messages: Message[] }>({
apiUrl: "http://localhost:2024",
assistantId: "agent",
reconnectOnMount: true, // 启用自动重连
});
// 自定义存储
const thread = useStream<{ messages: Message[] }>({
apiUrl: "http://localhost:2024",
assistantId: "agent",
reconnectOnMount: () => window.localStorage, // 使用 localStorage
});
默认使用 sessionStorage 存储运行 ID。启用后,组件挂载时会自动尝试恢复未完成的流式请求,确保不丢失任何消息。
TypeScript 类型支持
对于 LangGraph.js 用户,可以复用状态类型定义:
import {
Annotation,
MessagesAnnotation,
type StateType,
type UpdateType,
} from "@langchain/langgraph/web";
const AgentState = Annotation.Root({
...MessagesAnnotation.spec,
context: Annotation<string>(),
});
const thread = useStream<
StateType<typeof AgentState.spec>,
{ UpdateType: UpdateType<typeof AgentState.spec> }
>({
apiUrl: "http://localhost:2024",
assistantId: "agent",
messagesKey: "messages",
});
这种方式保证了类型安全,避免手动定义接口的繁琐和错误。
总结
LangGraph 的流式系统提供了从底层到前端的全栈解决方案。六种流模式覆盖了从调试到生产的各种需求,Token 流式传输让 LLM 响应实时可见,工具进度流提升了长操作的体验,自定义事件为业务逻辑提供了事件驱动能力,React 集成则让前端开发变得简单高效。
流式交互不仅是技术优化,更是用户体验的升级。它让 AI 应用从"黑箱"变为"白箱",从"等待"变为"参与"。在下一章中,我们将探讨 LangGraph 的持久化与时间旅行功能,看看如何让对话状态可追溯、可回退,进一步提升应用的可靠性和可控性。