深入学 LangChain 官方文档(十三)Streaming 与 Event Streaming

发布时间:2026/7/21 14:53:21
深入学 LangChain 官方文档(十三)Streaming 与 Event Streaming 深入学 LangChain 官方文档十三Streaming 与 Event Streaming本篇对应的官方文档Streaming支撑stream/astream、updates、messages、custom与多模式订阅。Event streaming支撑stream_events(..., versionv3)、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可以传列表。versionv2返回的每个StreamPart至少包含type、ns和data消费端按type分支再从data读取载荷。fromlangchain.agentsimportcreate_agentfromlangchain_openaiimportChatOpenAI modelChatOpenAI(modelqwen3.7-plus,api_keyYOUR_API_KEY,base_urlYOUR_OPENAI_COMPATIBLE_ENDPOINT,)# 作用查询指定订单的当前售后状态示例省略真实数据库访问。defget_order_status(order_id:str)-str:返回订单状态生产环境应从业务系统读取。returnf订单{order_id}已签收正在核对退款资格。agentcreate_agent(modelmodel,tools[get_order_status])forchunkinagent.stream({messages:[{role:user,content:查询订单 A-2048 的退款资格}]},stream_mode[messages,updates],versionv2,):ifchunk[type]messages:token,metadatachunk[data]print(metadata[langgraph_node],token.content_blocks)elifchunk[type]updates:print(state update:,chunk[data])输入消息进入同一个 Agent run。messages分支实时收到模型产生的消息块updates分支在模型或工具步骤结束后收到状态更新。这里的关键不是两个print而是消费端已经承担了路由职责每增加一种 mode就要增加相应的解析、状态合并和错误处理。三、Event Streaming把一次运行投影成多个消费面当聊天区、工具面板、状态调试器和审批组件都要同时工作继续在一个循环里堆if chunk[type]会越来越难维护。LangChain 当前面向新应用和前端场景推荐stream_events(..., versionv3)。它返回的不是普通 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(...)把多个投影合并到一个循环同时保留投影名称。streamagent.stream_events({messages:[{role:user,content:查询订单 A-2048 的退款资格}]},versionv3,)forprojection,iteminstream.interleave(messages,tool_calls,values):ifprojectionmessages:fordeltainitem.text:print(delta,end,flushTrue)elifprojectiontool_calls:print(tool:,item.tool_name,item.input)fordeltainitem.output_deltas:print(tool delta:,delta)print(tool result:,item.output,item.error)elifprojectionvalues:print(state snapshot:,item)final_statestream.outputprint(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(..., versionv3)将同一次运行组织为 messages、tool calls、values、subgraphs 和 output 等 typed projections。工程上最重要的不是 API 名称而是保持生命周期边界。参数尚在生成时不能执行副作用工具State 已有快照时不能假定任务结束interrupt 等待人工时不能当作失败重试前端断线时不能让界面状态脱离持久线程。当这些边界都被写进事件适配、持久化、权限、幂等和监控合同Agent 才真正从“最后给一句答案”的脚本变成用户能够实时观察、可靠操作和安全恢复的应用。下一篇我们会沿着这份运行时数据合同继续向外走进入《MCP 模型上下文协议首讲》。Streaming 解决的是 Agent 运行过程中“发生了什么、怎样及时交给界面”MCP 要解决的则是 Agent 面对多个独立服务时怎样用统一协议发现并调用外部能力同时守住连接、会话和权限边界。