深入学 LangChain 官方文档(十三)Streaming 与 Event Streaming
本篇对应的官方文档
- Streaming:支撑 stream / astream、updates、messages、custom 与多模式订阅。
- Event streaming:支撑 stream_events(…, version="v3")、typed projections、工具生命周期、状态快照与子图流。
- Human-in-the-loop:支撑高风险工具调用的暂停、人工决策、持久化与同一线程恢复。
本篇讲解范围
本篇讲清基础 Streaming 与 Event Streaming 的职责差异,并用售后退款 Agent 串起消息、工具、状态、子 Agent、HITL 和前端消费。模型供应商的流式协议细节、LangGraph 底层事件协议、完整 WebSocket 服务与可观测性平台留给后续专题。
客服问题中,用户问:“订单 A-2048 的耳机为什么不能退?”Agent 需要查订单、检索售后政策,金额过高时还要等待人工审批。整个过程也许只花十几秒,但如果界面一直停在一个旋转图标上,用户看见的仍是黑箱:系统究竟在生成答案、调用工具,还是已经失败?
很多项目把 Streaming 理解成“文字一个字一个字冒出来”。这种效果确实能缩短首字等待时间,却只覆盖模型输出的一小段。Agent 的运行过程还包含工具参数生成、工具执行、状态写回、子 Agent 调用、人工中断和最终状态。真正的实时体验,是把这些不同事实以稳定结构交给各自消费者。
一次性返回只暴露最终答案,等待期间的模型、工具、状态和审批都不可见;实时运行把同一条执行链拆成连续更新。用户得到进度,前端得到结构化状态,工程团队也能更早定位卡住的位置。
这条链路首先要解决的不是传输格式,而是确认哪些运行事实值得在任务完成前交付。
一、Streaming 交付的不是“快”,而是运行中的事实
流式输出不会让模型推理或工具查询凭空变快。它改变的是结果交付方式:完整运行尚未结束,调用方已经可以消费已经发生的部分。
这会直接改变产品行为。模型开始回答时,聊天区可以显示文本增量;模型决定查订单时,工具面板可以进入 pending;工具返回后,同一张卡片更新为 completed;命中高金额规则时,界面显示等待审批,而不是继续假装“思考中”。
因此,判断一个流式接口是否有用,不应只问“能不能吐 token”,而要问三件事:它暴露了哪些运行事实;每类事实由谁消费;消费者能否区分增量、完成、错误和暂停。
基础 Streaming 与 Event Streaming 都建立在 LangChain Agent 的运行之上,但面向的复杂度不同。前者适合快速订阅几个预定义通道;后者适合前端和复杂应用把消息、工具、状态等消费面独立处理。
二、三种基础 Streaming 模式
agent.stream() 是同步入口,agent.astream() 是异步入口。二者通过 stream_mode 选择输出类型。当前官方文档保留三种核心模式。
updates 在 Agent 每一步之后给出状态更新。退款场景里,模型节点产生工具调用、工具节点返回订单结果、模型节点生成最终回复,会形成连续的 step updates。它适合进度面板和调试,但不是逐 token 文本流。
messages 输出 (token, metadata),既能拿到文本增量,也会经过产生工具调用的模型消息。消费者应结合 content_blocks 或消息类型判断当前拿到的是文字、reasoning 还是 tool_call_chunk,不能假定每个 chunk 都是可直接显示的字符串。
custom 由工具或图节点主动发出业务进度,例如“已查询 10/100 条记录”。它适合表达框架无法自动推断的领域状态,但也意味着事件名称与载荷要由应用自己维护合同。
updates 观察步骤后的 State 变化,messages 观察模型消息增量,custom 承载工具主动上报的业务进度。三者来自同一次 Agent run,却服务不同界面组件,混成一段字符串会丢失状态语义。
多模式订阅时,stream_mode 可以传列表。version="v2" 返回的每个 StreamPart 至少包含 type、ns 和 data;消费端按 type 分支,再从 data 读取载荷。
from langchain.agents import create_agent
from langchain_openai import ChatOpenAI
model = ChatOpenAI(
model="qwen3.7-plus",
api_key="YOUR_API_KEY",
base_url="YOUR_OPENAI_COMPATIBLE_ENDPOINT",
)
# 作用:查询指定订单的当前售后状态,示例省略真实数据库访问。
def get_order_status(order_id: str) –> str:
"""返回订单状态,生产环境应从业务系统读取。"""
return f"订单 {order_id} 已签收,正在核对退款资格。"
agent = create_agent(model=model, tools=[get_order_status])
for chunk in agent.stream(
{
"messages": [
{"role": "user", "content": "查询订单 A-2048 的退款资格"}
]
},
stream_mode=["messages", "updates"],
version="v2",
):
if chunk["type"] == "messages":
token, metadata = chunk["data"]
print(metadata["langgraph_node"], token.content_blocks)
elif chunk["type"] == "updates":
print("state update:", chunk["data"])
输入消息进入同一个 Agent run。messages 分支实时收到模型产生的消息块,updates 分支在模型或工具步骤结束后收到状态更新。这里的关键不是两个 print,而是消费端已经承担了路由职责:每增加一种 mode,就要增加相应的解析、状态合并和错误处理。
三、Event Streaming:把一次运行投影成多个消费面
当聊天区、工具面板、状态调试器和审批组件都要同时工作,继续在一个循环里堆 if chunk["type"] 会越来越难维护。LangChain 当前面向新应用和前端场景推荐 stream_events(…, version="v3")。
它返回的不是普通 chunk 迭代器,而是一个 run object。底层仍是同一次 Agent 运行,但上层提供了 typed projections(类型化投影):stream.messages 只消费模型消息,stream.tool_calls 只消费工具执行生命周期,stream.values 只消费 State 快照,stream.output 读取最终 State。
“投影”不是复制多次运行。它更像同一事件源的多个观察窗口:聊天区不需要理解工具错误字段,工具卡不需要从 token 中猜参数是否完成,状态面板也不必扫描所有消息。
基础 stream 把多种 mode 放进同一条 chunk 流,由消费端分支解析;Event Streaming 在同一 run 上提供 messages、tool calls、values 等独立投影。复杂界面的收益来自消费边界稳定,而不是事件数量更多。
如果只做命令行演示,基础 Streaming 已经足够。若产品需要多个并行消费者、独立失败边界或清晰的 TypeScript/Python 类型,Event Streaming 通常更合适。选择依据是应用消费模型,不是“v3 一定比 v2 高级”。
四、消息里的工具参数,不等于工具已经执行
工具调用最容易制造危险误解。模型生成 tool call 时,参数 JSON 可能以多个小块到达:第一个 chunk 只有工具名,后续才逐步拼出 {"order_id": "A-2048"}。这些内容属于 message.tool_calls,表示模型正在形成调用请求。
只有参数完整并解析成功后,Agent 才能进入真实工具执行。工具开始、输入、输出增量、最终输出与错误属于 stream.tool_calls 的生命周期。二者看起来都叫 tool calls,实际边界完全不同。
message.tool_calls 位于模型输出阶段,可能只是尚未闭合的参数片段;stream.tool_calls 从工具真正开始执行后记录输入、输出增量、最终结果和错误。副作用操作只能在参数完成、校验和策略检查之后执行。
前端可以用参数增量做“正在准备查询”的视觉反馈,却不能据此提前发邮件、扣款或删文件。生产系统至少要等待 finalized tool call,再经过 schema 校验、权限策略和必要的 HITL 审批。流式透明不应突破执行安全边界。
五、stream.values 是快照,stream.output 才是最终状态
消息流回答“模型正在说什么”,State 流回答“Agent 运行到这里已经保存了什么”。stream.values 会随着节点推进产生状态快照;stream.output 在运行结束后给出最终 Agent State。
退款 Agent 查到订单后,快照可能已经包含 order_status,但还没有 refund_decision;审批中断发生时,State 可能保存待审工具请求,却没有最终回复。前端若把任意快照当成结束结果,就会过早关闭加载状态或显示未确认结论。
stream.values 是节点执行后的连续 State 快照,字段会逐步补齐;stream.output 只在本次 run 正常完成后代表最终状态。暂停与失败也可能留下可恢复快照,因此“已有数据”和“任务完成”必须分开判断。
状态字段还要有明确所有权。订单事实来自业务系统,审批状态来自工作流,模型消息属于对话 State。Streaming 只是把变化暴露出来,不会自动解决字段冲突、敏感信息脱敏或业务数据库权威性问题。
六、子 Agent 运行也要有清晰命名和边界
第 12 篇已经把主 Agent、专业 Agent 和前端状态连接起来。Event Streaming 进一步解决“嵌套运行怎样被观察”:命名后的 create_agent 子 Agent 可以出现在专用投影中,普通 StateGraph 子图则通过 subgraph 投影暴露。
界面因此可以显示“订单 Agent 正在查询”或“政策 Agent 已完成”,而不是把所有 token 都归到 supervisor 名下。但命名只解决识别,不解决权限。子 Agent 的私有上下文、内部 reasoning 和敏感工具参数仍不应原样展示给用户。
命名子 Agent 的消息、工具和状态可以保持独立命名空间,普通子图也保留嵌套层级。前端按名称呈现公开进度,服务端仍负责过滤私有上下文与敏感载荷,避免“过程透明”演变为数据泄露。
多智能体场景还会出现并行完成顺序不同的问题。事件到达顺序只能说明实时发生顺序,不能自动说明业务依赖。政策查询与订单查询可以并行,但退款决策若依赖两者,就必须等待明确的汇合条件,而不是收到第一个 completed 就生成最终答案。
七、用 typed projections 写出可维护的消费代码
售后界面至少需要三类更新:聊天区读取文本,工具卡读取执行生命周期,调试或进度区读取 State 快照。同步代码可以用 stream.interleave(…) 把多个投影合并到一个循环,同时保留投影名称。
stream = agent.stream_events(
{
"messages": [
{"role": "user", "content": "查询订单 A-2048 的退款资格"}
]
},
version="v3",
)
for projection, item in stream.interleave(
"messages", "tool_calls", "values"
):
if projection == "messages":
for delta in item.text:
print(delta, end="", flush=True)
elif projection == "tool_calls":
print("tool:", item.tool_name, item.input)
for delta in item.output_deltas:
print("tool delta:", delta)
print("tool result:", item.output, item.error)
elif projection == "values":
print("state snapshot:", item)
final_state = stream.output
print("final state:", final_state)
messages 项是 ChatModelStream,文本增量从 item.text 消费;tool_calls 项承载真实工具执行,错误从 item.error 读取;values 项是当前 State 快照。循环结束后再读取 stream.output,才能取得最终状态。
异步服务不必把所有投影重新塞回一个循环。astream_events 可以配合 asyncio.gather,让聊天、工具和状态消费者各自异步迭代。无论同步还是异步,核心都是让每个组件只理解自己的 projection,并为它定义独立的超时、取消与降级策略。
同一 run 产生 messages、tool calls 与 values 三个投影,interleave 只负责按到达顺序合并消费,不改变各自类型。聊天区、工具卡和状态面板可以共享 run ID,却分别维护渲染与错误边界。
代码里的 print 在真实系统中应替换为事件适配层。这个层负责把 LangChain 对象转换成稳定的前端 DTO,保留 run、thread、tool call 等关联 ID,并删去不能公开的原始参数。
八、HITL 让“暂停”成为正式运行状态
高金额退款不能因为模型已经生成完整工具参数就自动执行。Human-in-the-loop Middleware 会根据 interrupt_on 策略检查工具调用,需要人工介入时发出 interrupt,并依靠 checkpointer 保存图状态。
调用方必须提供 thread_id,这样暂停和恢复才能指向同一条会话线程。审核者可以 approve、edit、reject;面向“询问用户”的工具还可以 respond。拒绝有副作用的操作应使用 reject,不能把 respond 当成拒绝,否则系统会把回复当作成功工具结果。
高风险 tool call 先经过 policy,命中规则后以 interrupt 保存到同一 thread;前端展示待审动作,人工决定再恢复执行。刷新页面或稍后处理都依赖持久 State,单纯在浏览器弹窗不能替代后端暂停。
Streaming 在这里承担两种职责:先把 interrupt 及时交给界面,再在恢复后继续输出工具结果和后续消息。消费端必须把 paused 与 failed 区分开。paused 表示系统在等待合法输入,盲目自动重试反而可能制造重复审批或副作用。
九、生产环境真正难的是消费侧治理
演示代码能打印事件,不代表已经具备生产能力。浏览器断线、慢消费者、重复投递、进程重启和多个并行工具都会改变事件消费方式。
首先是背压。模型 token 可能远快于界面渲染,状态快照也可能很大。服务端应合并可合并的增量、限制队列、允许取消,并确保工具错误和 interrupt 等关键事件不会被普通文本淹没。
其次是恢复。前端需要 run ID 与 thread ID 区分一次运行和可持续会话;断线重连后,先从持久状态恢复已确认事实,再续接实时事件。只依赖浏览器内存,会让刷新后的界面与后端真实状态分叉。
再次是幂等与顺序。工具执行不能因为客户端重连而重复触发;消费者要用稳定事件 ID 或 tool call ID 去重。并行事件的到达顺序不等于业务提交顺序,最终写入仍应由工作流和业务事务控制。
Agent run 经事件适配层进入有界队列,再由聊天、工具、状态和审批组件消费;thread 持久化负责恢复,稳定 ID 负责去重与关联,脱敏层阻止敏感参数外泄。实时性必须与正确性、恢复性和安全性一起验收。
最后是可观测性。至少要能用同一组标识关联 thread、run、model call、tool call、interrupt 和前端组件。日志记录事件类型、耗时、状态迁移和错误摘要,同时对用户数据与工具参数脱敏。只有这样,“页面一直等待”才能被定位到模型未结束、工具超时、审批未处理还是事件消费堵塞。
十、怎样选择 Streaming 或 Event Streaming
如果只是命令行打印 token、展示简单步骤,或者已有代码围绕 stream_mode 构建,基础 Streaming 足够直接。updates、messages、custom 三种模式清晰覆盖进度、模型消息和业务自定义更新。
如果应用有多个独立消费者,需要分别处理消息、工具、状态、子 Agent 和最终输出,优先考虑 Event Streaming。typed projections 减少分支解析,允许每个组件围绕稳定对象建立自己的生命周期。
无论选择哪一个,都应守住三个边界:消息增量不是最终消息,工具参数增量不是已经执行的工具,状态快照不是最终输出。再加上第四个产品边界:前端只呈现后端已经确认的运行事实,不自行猜测 Agent 状态。
第 12 篇把多智能体后端与前端状态接了起来,第 13 篇进一步把这条连接拆成可消费的实时通道。到这里,Streaming 不再只是打字动画,而是一份运行时数据合同:它规定什么事实何时出现、由谁消费、怎样暂停、怎样恢复,以及失败后如何找到正确的位置。
总结:实时体验来自可消费的运行状态
基础 Streaming 用 updates、messages、custom 快速暴露进度、模型消息和业务更新;Event Streaming 用 stream_events(…, version="v3") 将同一次运行组织为 messages、tool calls、values、subgraphs 和 output 等 typed projections。
工程上最重要的不是 API 名称,而是保持生命周期边界。参数尚在生成时不能执行副作用工具,State 已有快照时不能假定任务结束,interrupt 等待人工时不能当作失败重试,前端断线时不能让界面状态脱离持久线程。
当这些边界都被写进事件适配、持久化、权限、幂等和监控合同,Agent 才真正从“最后给一句答案”的脚本,变成用户能够实时观察、可靠操作和安全恢复的应用。
下一篇我们会沿着这份运行时数据合同继续向外走,进入《MCP 模型上下文协议首讲》。Streaming 解决的是 Agent 运行过程中“发生了什么、怎样及时交给界面”;MCP 要解决的则是 Agent 面对多个独立服务时,怎样用统一协议发现并调用外部能力,同时守住连接、会话和权限边界。

![[LangChain RAG] 01 大模型为什么需要 RAG:四个问题与标准流程-171主机测评](https://www.171host.com/wp-content/uploads/2026/08/20260825035331-6a8d11bb97bca-220x150.png)

