使用 LangChain 构建 Agent 之后,最简单的执行方式是 invoke():提交一次输入,等待 Agent 完整执行,然后拿到最终结果。
但放到真实应用中,很快会遇到一个问题:一次 Agent 执行可能持续几秒甚至更久,中间还会经过模型思考、工具调用、数据查询等步骤。如果整个过程都没有反馈,前端看起来就像“卡住了”。
流式输出解决的正是这个问题。它并不会让模型本身执行得更快,而是让应用能够在执行尚未结束时,持续拿到 Token、状态变化、工具调用和自定义进度,并及时展示给用户。
本文主要介绍 LangChain 中 stream()、astream() 的使用,以及 updates、messages、custom 三种流式模式如何对应实际 Agent 的运行过程。
一、流式输出到底改变了什么
先看一次普通调用:
result = agent.invoke({
"messages": [
{"role": "user", "content": "查询北京今天的天气,并给出跑步建议"}
]
})
print(result["messages"][-1].content)
程序会一直等到 Agent 完成执行,最后一次性返回结果。
但 Agent 内部其实经历了完整的执行过程:
用户请求
↓
模型判断需要查询天气
↓
生成工具调用
↓
执行天气工具
↓
模型读取工具结果
↓
生成最终回答
invoke() 并不会把这些中间过程持续交给调用方。
而使用流式执行之后,我们可以在 Agent 尚未结束时不断拿到数据:
模型开始生成工具参数
↓
工具调用已经确定
↓
天气查询完成
↓
最终答案开始生成
↓
北京
北京今天
北京今天晴……
因此,流式输出并不只是聊天窗口里的“逐字显示”。
对于 Agent 应用来说,至少有三类信息值得实时展示:
| 信息 | 典型用途 |
|---|---|
| 模型输出片段 | 聊天内容逐步显示 |
| Agent 状态更新 | 显示当前执行到了哪一步 |
| 自定义业务进度 | 显示检索、下载、分析等耗时任务进度 |
LangChain 分别通过 messages、updates 和 custom 提供这些数据。
二、stream() 和 astream()
先创建一个简单的天气 Agent:
from langchain.agents import create_agent
def get_weather(city: str) -> str:
"""查询指定城市的天气。"""
return f"{city}今天晴,25℃"
agent = create_agent(
model="openai:gpt-5.4",
tools=[get_weather],
)
同步程序可以使用 stream():
for chunk in agent.stream(
{
"messages": [
{"role": "user", "content": "北京今天天气怎么样?"}
]
},
stream_mode="updates",
version="v2",
):
print(chunk)
异步程序则可以使用 astream():
async for chunk in agent.astream(
{
"messages": [
{"role": "user", "content": "北京今天天气怎么样?"}
]
},
stream_mode="updates",
version="v2",
):
print(chunk)
两者背后的流式机制基本相同,主要区别在调用方式。
stream() 返回同步迭代器,更适合普通脚本或同步应用。
astream() 返回异步迭代器,更适合 FastAPI、异步 Web 服务,以及需要同时处理大量连接的应用。
使用 version="v2" 时,返回的数据采用统一格式,每个 chunk 通常包含:
type
ns
data
其中:
type表示当前事件属于哪一种流式模式;data保存这次事件的实际内容;ns用来表示事件所属的命名空间。
这种格式有一个很实际的好处:当应用同时订阅多种流式事件时,可以直接通过 type 判断应该怎样处理,而不必依赖不同的返回结构去猜数据类型。
三、stream_mode 的三种模式
真正理解 LangChain 流式输出,关键不是记住 stream() 方法,而是理解 stream_mode。
同一次 Agent 执行,可以从三个不同角度观察。
1. updates:观察 Agent 执行到了哪一步
updates 返回的是 Agent 每一步执行完成后的状态更新。
例如:
for chunk in agent.stream(
{
"messages": [
{"role": "user", "content": "北京今天天气怎么样?"}
]
},
stream_mode="updates",
version="v2",
):
if chunk["type"] != "updates":
continue
for step, data in chunk["data"].items():
print(step)
print(data["messages"][-1])
一次典型输出可能是:
model
AIMessage(... tool_calls=[{
'name': 'get_weather',
'args': {'city': '北京'}
}])
tools
ToolMessage(content='北京今天晴,25℃')
model
AIMessage(content='北京今天晴,气温约25℃。')
这里看到的不是 Token,而是 Agent 的执行步骤。
第一次 model 表示模型判断需要调用工具。
接着进入 tools,天气工具执行完成,并产生 ToolMessage。
然后模型再次运行,根据工具结果生成最终回答。
所以,如果产品希望展示:
正在分析问题……
正在查询天气……
天气数据已返回……
正在整理回答……
这类状态通常应该来自 updates,而不是模型 Token。
2. messages:获取模型实时输出
如果希望实现类似聊天应用中的逐步文字显示,可以使用 messages。
for chunk in agent.stream(
{
"messages": [
{"role": "user", "content": "北京今天天气怎么样?"}
]
},
stream_mode="messages",
version="v2",
):
if chunk["type"] != "messages":
continue
message_chunk, metadata = chunk["data"]
print(
message_chunk.text,
end="",
flush=True,
)
最终看到的完整内容可能是:
北京今天晴,气温约25℃,天气条件比较适合户外活动。
但程序实际接收到的并不是这一整句话,而是一系列连续到达的消息片段,例如:
北京
今天
晴
,
气温
约
25℃
……
因此我们平时说的“Token 流式输出”,在 LangChain 中更准确地说,是模型不断产生 AIMessageChunk,应用收到这些 chunk 后持续拼接。
最终完整的 AIMessage,可以看作这些消息片段逐步累积之后的结果。
3. custom:发送业务自己的进度
还有一类信息既不是模型生成的内容,也不是 Agent 节点状态。
例如,一个工具需要处理 100 个文件:
开始读取文件
已处理 20/100
已处理 60/100
正在生成索引
索引创建完成
这些信息属于业务代码自己的执行进度。
这种场景可以使用 custom。
三种模式可以简单记成:
messages → 模型正在输出什么
updates → Agent 执行到了哪一步
custom → 业务代码执行到了哪里
四、消息流和状态流不是一回事
Agent 流式输出有一个很容易混淆的地方:
它并不是简单地把最终的 AIMessage 切成很多小块。
实际运行时,至少存在两条不同的数据流。
第一条是模型生成过程:
AIMessageChunk
↓
AIMessageChunk
↓
AIMessageChunk
↓
完整 AIMessage
这部分主要由 messages 暴露。
第二条是 Agent 本身的执行状态:
Model
↓
AIMessage(tool_calls)
↓
Tool
↓
ToolMessage
↓
Model
↓
AIMessage(final)
这部分更适合通过 updates 观察。
两者在工具调用场景中特别容易被混在一起。
例如模型正在生成:
get_weather
{"city": "北京"}
此时你可能已经能从 messages 中看到 tool_call_chunk。
但这并不代表天气工具已经执行。
只有当 Agent Runtime 完成参数解析、真正调用工具,并生成对应的 ToolMessage 后,工具才算真正产生了结果。
因此要区分:
tool_call_chunk
和:
ToolMessage
前者表示:
模型正在生成一次工具调用。
后者表示:
工具已经执行并返回结果。
这两个阶段在 UI 上往往也应该采用不同的状态。
五、工具调用如何实时展示
工具调用本身也可以通过 messages 被观察。
例如:
for chunk in agent.stream(
{
"messages": [
{"role": "user", "content": "北京今天天气怎么样?"}
]
},
stream_mode="messages",
version="v2",
):
if chunk["type"] != "messages":
continue
token, metadata = chunk["data"]
print(metadata["langgraph_node"])
print(token.content_blocks)
可能看到类似输出:
model
[
{
'type': 'tool_call_chunk',
'name': 'get_weather',
'args': '',
'index': 0
}
]
model
[
{
'type': 'tool_call_chunk',
'args': '{"city":"北京"}',
'index': 0
}
]
这里有一个细节值得注意。
工具参数本身也是模型生成出来的,所以参数可能分成多个片段返回,甚至暂时不是完整 JSON。
因此应用不应该在刚收到一个 tool_call_chunk 时就自行执行工具。
工具参数的拼接、解析,以及后续真正的 Tool 调用,都应该交给 Agent Runtime 完成。
如果界面既希望展示模型生成过程,又希望展示 Agent 节点状态,可以同时订阅多个模式:
for chunk in agent.stream(
{
"messages": [
{"role": "user", "content": "北京今天天气怎么样?"}
]
},
stream_mode=["messages", "updates"],
version="v2",
):
if chunk["type"] == "messages":
token, metadata = chunk["data"]
print("消息:", token.content_blocks)
elif chunk["type"] == "updates":
print("状态:", chunk["data"])
一次执行过程可能表现为:
消息: tool_call_chunk ...
状态: model -> AIMessage(tool_calls)
状态: tools -> ToolMessage(...)
消息: 北京
消息: 今天
消息: 晴
状态: model -> AIMessage(...)
这种方式很适合 Agent UI。
聊天区域由 messages 驱动,负责显示模型输出。
执行状态区或 Tool Card 则通过 updates 更新。
两者各管一层,前端逻辑会清晰很多。
六、耗时工具如何上报进度
真实项目里的工具通常不会瞬间执行完成。
例如一次知识库检索,内部可能经历:
解析查询
↓
检索候选文档
↓
重排序
↓
读取正文
↓
返回结果
如果只监听 updates,通常只能等整个 Tool 节点执行结束后,才看到对应结果。
但用户真正关心的可能是:
正在搜索……
找到 50 篇候选文档……
正在重排序……
正在读取前 5 篇正文……
这时可以让工具主动发送自定义事件。
from langgraph.config import get_stream_writer
def search_documents(query: str) -> str:
"""搜索文档。"""
writer = get_stream_writer()
writer({"stage": "search", "message": "正在检索候选文档"})
writer({"stage": "rerank", "message": "正在重新排序结果"})
writer({"stage": "done", "message": "检索完成"})
return f"找到与 {query} 相关的文档"
调用时订阅 custom:
for chunk in agent.stream(
{
"messages": [
{"role": "user", "content": "搜索 LangChain 流式输出资料"}
]
},
stream_mode="custom",
version="v2",
):
if chunk["type"] == "custom":
print(chunk["data"])
输出可能是:
{'stage': 'search', 'message': '正在检索候选文档'}
{'stage': 'rerank', 'message': '正在重新排序结果'}
{'stage': 'done', 'message': '检索完成'}
get_stream_writer() 获取的是当前 LangGraph 执行上下文中的 Stream Writer。
因此,这种写法依赖 LangGraph 的执行环境。如果直接脱离 Agent Runtime 调用这个普通 Python 函数,就不能指望 Writer 仍然存在。
在实际项目中,自定义事件最好使用结构化数据,而不是简单发送一句字符串。
例如:
writer({
"stage": "download",
"progress": 0.6,
"message": "已处理 60 个文件"
})
前端拿到这些字段后,可以自行决定显示成:
- 进度条;
- 状态文字;
- Tool Card;
- 时间线;
- 日志区域。
这样业务层和 UI 层之间不会被一段固定文案绑死。
七、前端的“逐字输出”是怎么实现的
前端看到的逐字效果,本质上也不是浏览器定时把一句话拆成一个个字显示。
真正的数据链路更接近:
模型产生 Message Chunk
↓
后端收到 Chunk
↓
HTTP Stream / SSE / WebSocket
↓
前端收到新数据
↓
追加到当前 AI Message
↓
组件重新渲染
也就是说,页面之所以看起来在“打字”,是因为消息状态一直在增长。
例如前端第一次收到:
北京
下一次收到:
今天
应用把它拼接成:
北京今天
随后继续追加:
晴
最终得到:
北京今天晴……
如果使用 LangChain 的前端能力,也可以通过 useStream 管理这些状态。
React 示例:
import { useStream } from "@langchain/react";
export default function Chat() {
const stream = useStream({
apiUrl: "http://localhost:2024",
assistantId: "agent",
});
return (
<div>
{stream.messages.map((message) => (
<div key={message.id}>
{message.text}
</div>
))}
</div>
);
}
随着后端持续返回新的内容,当前 AI Message 会不断增长:
北京
北京今天
北京今天晴
北京今天晴,气温约25℃……
这就是大多数聊天应用中“流式回答”的基本实现方式。
八、流式执行中的异常处理
流式调用和普通调用有一个很现实的区别:
异常发生之前,部分内容可能已经发给用户了。
例如:
正在查询数据库……
已找到 120 条记录
正在分析……
连接中断
前面的内容已经展示出去。
这时不能再把整次调用当成一个普通函数:
成功 → 有结果
失败 → 什么都没有
更合适的方式,是把一次流式任务看成一个持续变化的执行状态。
最基本的消费代码可以这样处理:
try:
async for chunk in agent.astream(
input_data,
stream_mode=["messages", "updates", "custom"],
version="v2",
):
await send_to_client(chunk)
except Exception as exc:
await send_error_to_client(str(exc))
发生异常时,前端可能已经显示:
正在查询数据……
正在分析结果……
随后再收到:
执行失败,请稍后重试
实际开发中有几个细节需要特别处理。
第一,不要把“最后一个 Token 已经输出”当成整次 Agent 执行已经成功。
模型输出之后,后面还可能存在 Tool、Middleware 或其他执行节点。
第二,前端最好单独维护任务状态,例如:
running
completed
error
不要只根据“有没有文本”来判断 Agent 是否执行完成。
第三,数据库、HTTP API 等外部调用仍然需要自己的超时、重试和错误恢复机制。
流式输出只是让执行过程变得可见,并不会自动解决可靠性问题。
第四,内部事件不要原样全部暴露给前端。
生产环境通常只需要发送用户真正关心的状态。内部 Prompt、调试信息、敏感参数和工具原始数据,都不应该因为开启了流式输出就直接发送出去。
九、stream() 和 invoke() 怎么选
流式输出并不是 invoke() 的替代方案。
两者适合不同场景。
| 场景 | 更适合 |
|---|---|
| 后台批处理 | invoke() |
| 只需要最终结构化结果 | invoke() |
| 很快结束的模型调用 | invoke() |
| 聊天应用 | stream() / astream() |
| 长时间 Agent 任务 | stream() / astream() |
| 需要展示工具执行过程 | stream() / astream() |
| 需要展示业务进度 | stream() / astream() |
例如:
输入一段文本
↓
生成分类结果
↓
写入数据库
整个过程完全发生在后台,用户也不关心中间步骤。
这种任务使用 invoke() 往往更简单,没有必要为了“流式”再增加事件消费、连接管理和前端状态维护。
但如果一个 Agent 要运行十几秒,并且中间会调用多个 Tool、执行多个节点,那么只使用 invoke(),用户会一直看不到系统正在做什么。
这种场景即使不需要逐 Token 输出,也可以通过 updates 或 custom 提供过程反馈。
所以真正需要判断的不是:
这个任务能不能流式输出?
而是:
用户有没有必要看到它的执行过程?
十、总结
LangChain 中的流式输出,可以理解为对 Agent 执行过程的持续观察。
messages 关注模型生成的 Message Chunk,适合聊天文字、工具参数等实时展示;updates 关注 Agent 每个执行步骤产生的 State 更新,适合呈现 Model、Tool 等执行过程;custom 则允许业务代码主动发送进度。
stream() 和 astream() 负责把这些数据持续交给应用,前端再通过 HTTP Stream、SSE、WebSocket 或 useStream 等方式将消息不断累积和渲染。
实际开发中,不需要为了“流式”而把所有信息都流出来。聊天内容使用 messages,Agent 运行状态使用 updates,耗时工具内部进度使用 custom,只需要最终结果的后台任务继续使用 invoke(),通常就是比较清晰的划分方式。
社区讨论
参与讨论
有问题或想法?欢迎继续讨论。