在这里插入图片描述

精读 LangChain 官方文档(五)Event Streaming 篇:把 Agent 运行过程升级为类型化事件流

本文基于 LangChain Python 官方文档整理:
Event Streaming:https://docs.langchain.com/oss/python/langchain/event-streaming
Markdown 版本:https://docs.langchain.com/oss/python/langchain/event-streaming.md
对应开源文档源码:
https://github.com/langchain-ai/docs/blob/main/src/oss/langchain/event-streaming.mdx

很多人学到 Streaming 时,会先记住一个画面:模型一边生成,一边把 token 推到前端。

这当然有用,但它还不是 Agent 产品里最难的部分。

真正进入工程现场后,我们关心的不只是“文字有没有出来”,还会关心:模型为什么调用了工具、工具参数是否已经成型、子智能体是谁在输出、当前 state 变成了什么、业务工具有没有自己的进度、日志系统能不能只订阅自己需要的那部分事件。

如果把所有东西都塞进一个普通 token 流里,前端会很难拆,后台也很难审计。最后的结果就是:聊天界面勉强能跑,但调试、监控、人工审核和多智能体协作都变得很痛苦。

LangChain 的 Event Streaming 文档要解决的正是这个问题:不再只把运行过程看成“字符串流”,而是把 Agent runtime 暴露成一组可以被选择、转换和组合的类型化事件投影。

这篇文档的核心主线可以概括成一句话:

Event Streaming = Agent Runtime 的类型化事件观察层。它把模型文本、工具调用、状态变化、子智能体输出和自定义业务事件,拆成可订阅、可组合、可审计的 projection。

也就是说,stream_events(..., version="v3") 不是 stream() 的简单替代,而是把 Agent 的运行过程提升成更清晰的工程接口:

  • stream.messages():看模型消息、文本、推理片段和工具调用。
  • stream.tool_calls():专门看工具调用输入与输出。
  • stream.values():看每一步 state 的完整快照。
  • stream.output():只关心最终输出。
  • stream.subagents() / stream.subgraphs():区分复杂运行结构里的来源。
  • stream.interleave(...):把多个投影重新合成一个有类型的事件流。
  • raw protocol 和 custom transformer:在默认投影不够时,自己定义观察方式。

下面这张图先把主线串起来:

事件流主线图

理解这条主线后,stream_eventsversion="v3"message.textmessage.tool_callsToolCallstream.values()StreamEventGenericStreamEvent 就不会是散落 API,而是同一层事件观察模型里的不同视角。



1. Event Streaming(事件流):它到底解决什么问题

它解决的问题:
Event Streaming 解决的是“Agent 运行时如何被结构化观察”,而不是只解决“模型输出如何逐字显示”。

传统流式输出更像这样:

用户输入 -> 模型生成 -> token token token -> 最终文本

事件流更像这样:

用户输入 -> Agent runtime -> typed projections -> UI / Logs / Audit / Debug

区别在于,事件流不会把所有运行信号压扁成字符串。它会把不同类型的运行信息拆出来,让应用可以按需订阅。

示例:

import os

from langchain.agents import create_agent
from langchain_openai import ChatOpenAI

model = ChatOpenAI(
    model="qwen3.7-max",
    api_key=os.environ["QWEN_API_KEY"],
    base_url=os.environ["QWEN_BASE_URL"],
)

agent = create_agent(
    model=model,
    tools=[],
    system_prompt="你是一名中文技术助手,需要清晰解释 Agent 运行过程。",
)

stream = agent.stream_events(
    {
        "messages": [
            {"role": "user", "content": "请解释 LangChain Event Streaming 的作用"}
        ]
    },
    version="v3",
)

for message in stream.messages():
    for delta in message.text:
        print(delta, end="")

这里:

  • agent.stream_events(...):启动事件流,它返回的不是单一字符串,而是一个可以继续投影的流对象。
  • version="v3":选择新版事件流协议。官方文档强调,v3 是使用类型化投影接口的关键版本。
  • stream.messages():从底层事件中投影出“消息相关事件”,方便消费模型文本、推理片段和工具调用。
  • message.text:表示当前 message 中的文本增量。它可以迭代,也可以在结束后转换成完整文本。

业务场景:
做企业客服 Agent 时,前端可以订阅 messages() 显示文本,后台可以订阅工具调用做审计,调试台可以订阅 state 快照。不同系统不必抢同一条混杂的 token 流。

最简记法:

Streaming 让你看到输出过程,Event Streaming 让你按类型理解运行过程。


2. Projection(投影):从原始事件里只取自己需要的视角

它解决的问题:
Projection 解决的是“同一段 Agent 运行过程,不同消费者应该看到不同形态的数据”。

一个 Agent 运行时会产生很多底层事件:模型开始、模型 token、工具调用、工具返回、状态更新、子图运行、自定义事件等。前端、日志系统、调试器不应该都处理同一堆原始事件。

所以官方文档提供了多种投影方法,把底层事件转换成更好用的对象:

投影方法 中文理解 常见用途
stream.messages() 消息投影 聊天 UI、token 展示、工具调用片段
stream.tool_calls() 工具调用投影 工具审计、参数展示、调用结果追踪
stream.values() 状态值投影 调试 Agent state、观察每步状态
stream.output() 最终输出投影 只拿最后结果
stream.events() 原始事件协议 自定义解析、底层诊断
stream.interleave(...) 多投影交错 同时消费多种事件

投影视角选择图

示例:

stream = agent.stream_events(
    {"messages": [{"role": "user", "content": "请帮我分析这笔订单是否需要人工介入"}]},
    version="v3",
)

for event in stream.interleave(
    "messages",
    "tool_calls",
    "values",
):
    print(event)

这里:

  • projection:中文可以理解为“观察视角”。它不改变 Agent 本身,只改变你如何消费运行事件。
  • stream.interleave(...):把多个投影重新合成一个事件流,并保留事件类型信息。
  • "messages""tool_calls""values":不是模型提示词,而是你希望从运行时拿到的事件类别。

业务场景:
如果做一个 Agent 调试控制台,左侧聊天区订阅 messages,中间工具面板订阅 tool_calls,右侧状态面板订阅 values。同一段运行,只是被投影成不同视图。

最简记法:

Projection 不是重新运行 Agent,而是给同一次运行换一个观察窗口。


3. messages()(消息投影):不只是文本 token

它解决的问题:
messages() 解决的是“如何用消息对象的方式消费模型输出”,而不是只消费裸字符串。

在事件流里,模型消息可以包含多种内容:普通文本、推理片段、工具调用请求、工具调用参数分片、消息结束后的完整对象等。messages() 会把这些内容包装成更稳定的投影对象。

示例:

stream = agent.stream_events(
    {"messages": [{"role": "user", "content": "请用三句话说明事件流的价值"}]},
    version="v3",
)

for message in stream.messages():
    for delta in message.text:
        print(delta, end="")

    final_text = str(message.text)
    print("\n完整文本:", final_text)

这里:

  • message:一次消息投影对象,不只是一个字符串。
  • message.text:文本内容的 projection,可以边迭代边显示,也可以最后用 str(...) 拿完整文本。
  • delta:本次流式返回的文本增量。
  • final_text:消息结束后聚合出来的完整文本。

业务场景:
聊天前端可以边迭代 message.text 边渲染,保存会话记录时再用完整文本入库。这样前端体验和数据库落库不会互相干扰。

最简记法:

messages() 给你的不是散 token,而是带结构的消息对象。


4. reasoning(推理片段):把思考过程作为可选信号处理

它解决的问题:
message.reasoning 解决的是“如何消费模型暴露出来的 reasoning 内容”,并且避免把它和普通文本混在一起。

不同模型提供商对 reasoning token 的支持和格式并不完全相同。事件流的价值在于:应用侧可以把 reasoning 当成独立 projection 来处理,而不是在普通文本里硬拆。

示例:

stream = agent.stream_events(
    {"messages": [{"role": "user", "content": "请给出一个订单风控判断思路"}]},
    version="v3",
)

for message in stream.messages():
    for reasoning_delta in message.reasoning:
        print("推理片段:", reasoning_delta)

    for text_delta in message.text:
        print("回答文本:", text_delta)

这里:

  • message.reasoning:推理片段投影,用来处理模型提供的 reasoning 内容。
  • message.text:面向用户展示的普通回答文本。
  • reasoning_delta:推理内容增量,是否展示给用户要看产品策略和模型供应商能力。

业务场景:
内部代码审查助手可以把 reasoning 片段放进调试面板,帮助工程师理解模型判断路径;面向普通客户的客服机器人则可能只展示最终回答,把 reasoning 留给后台质检。

最简记法:

reasoning 是一种可选观察信号,不应该和最终回答混成一团。


5. tool_calls(工具调用):把工具参数和结果独立观察

它解决的问题:
tool_calls() 解决的是“如何在运行中观察工具调用”,尤其是工具名、参数、结果和调用时机。

Agent 产品里,工具调用往往比最终文本更关键。因为工具调用会连接订单系统、数据库、CRM、搜索引擎或内部审批流。只看最终回答,无法判断它有没有查对系统、有没有传错参数、有没有漏掉关键条件。

工具调用事件图

示例:

import os

from langchain.agents import create_agent
from langchain_openai import ChatOpenAI

# 查询订单当前状态,用于演示工具调用事件。
def query_order_status(order_id: str) -> str:
    """根据订单编号返回模拟订单状态。"""
    return f"订单 {order_id} 当前状态:已付款,待发货。"


model = ChatOpenAI(
    model="qwen3.7-max",
    api_key=os.environ["QWEN_API_KEY"],
    base_url=os.environ["QWEN_BASE_URL"],
)

agent = create_agent(
    model=model,
    tools=[query_order_status],
    system_prompt="你是一名中文订单助手,需要在必要时调用工具查询订单。",
)

stream = agent.stream_events(
    {"messages": [{"role": "user", "content": "请查询订单 A10086 的状态"}]},
    version="v3",
)

for tool_call in stream.tool_calls():
    print("工具名称:", tool_call.name)
    print("工具参数:", tool_call.args)
    print("工具结果:", tool_call.result)

这里:

  • query_order_status:一个普通 Python 函数,被注册为 Agent 工具。
  • order_id:工具参数,表示要查询的订单编号。
  • stream.tool_calls():工具调用投影,专门消费工具调用相关事件。
  • tool_call.name:工具名称。
  • tool_call.args:模型生成的工具参数。
  • tool_call.result:工具执行后的返回结果。

业务场景:
售后系统可以把 tool_calls() 的结果写入审计日志:用户问了什么、Agent 调了哪个工具、参数是什么、返回了什么。这样出现错误时可以定位是模型判断错了,还是业务工具返回错了。

最简记法:

tool_calls() 是 Agent 的工具审计通道。


6. subagents()(子智能体):区分复杂系统里的输出来源

它解决的问题:
subagents() 解决的是“多 Agent 或命名子任务里,当前输出到底来自谁”。

当系统只有一个 Agent 时,消息来源比较简单。但一旦进入多智能体、专家协作、子图执行或 supervisor 模式,同一条运行流里可能有多个来源。如果不区分来源,前端和日志都会很乱。

子智能体来源图

示例:

stream = agent.stream_events(
    {"messages": [{"role": "user", "content": "请分析客户问题,并给出售后处理建议"}]},
    version="v3",
)

for subagent_event in stream.subagents():
    print("来源名称:", subagent_event.name)
    print("事件内容:", subagent_event)

这里:

  • stream.subagents():用于观察命名 Agent 运行事件的投影。
  • subagent_event.name:可以帮助你识别事件来自哪个 Agent 或哪段命名运行。
  • subagent_event:具体事件对象,不同场景下可能包含消息、状态或运行元信息。

业务场景:
一个客服 supervisor 可以协调“订单 Agent”“退款 Agent”“物流 Agent”。前端展示时,需要知道当前回复来自哪个专家;日志审计时,也需要知道是哪一个子 Agent 做出了判断。

最简记法:

subagents() 给复杂 Agent 系统加来源标签。


7. values()(状态快照):观察 Agent 每一步 state

它解决的问题:
values() 解决的是“如何看到 Agent state 在运行过程中的变化”。

Agent 不只是输入和输出,它还会维护运行状态,例如 messages、工具结果、中间变量、自定义状态字段等。只看 token 很难调试这些状态变化。

示例:

stream = agent.stream_events(
    {"messages": [{"role": "user", "content": "请帮我整理本周客户投诉重点"}]},
    version="v3",
)

for state in stream.values():
    print("当前状态快照:", state)

这里:

  • stream.values():把运行过程中的 state 快照投影出来。
  • state:当前步骤的完整状态值。具体字段取决于 Agent 或 LangGraph state schema。
  • messages:常见状态字段之一,通常保存对话消息列表。

业务场景:
做数据分析 Agent 时,state 里可能保存“已读取文件”“已生成 SQL”“已得到查询结果”“已生成图表”等中间结果。调试台订阅 values(),就能知道是哪一步把 state 写错了。

最简记法:

values() 看的是运行状态,不是输出文本。


8. output()(最终输出):只拿最后结果,不关心过程

它解决的问题:
output() 解决的是“有些调用方只需要最终结果,不想处理整个事件流”。

事件流提供了丰富过程信息,但不是每个业务模块都需要这些信息。比如后台批处理任务可能只关心最终报告,定时任务可能只关心最终 JSON,测试用例可能只想断言最后输出。

示例:

stream = agent.stream_events(
    {"messages": [{"role": "user", "content": "请生成一段客户回访摘要"}]},
    version="v3",
)

final_output = stream.output()
print(final_output)

这里:

  • stream.output():最终输出投影。
  • final_output:一次运行结束后的最终结果。
  • 过程事件仍然存在,只是这个消费者选择不处理它们。

业务场景:
同一个 Agent 可以同时服务两类调用方:聊天 UI 订阅 messages() 做流式体验;后台任务只用 output() 拿最终结果入库。

最简记法:

output() 是事件流里的最终答案入口。


9. interleave(交错组合):把多个投影合成一个事件时间线

它解决的问题:
interleave 解决的是“我想同时看多种事件,但又不想丢失它们的发生顺序”。

如果分别循环 messages()tool_calls()values(),你会得到三个视角,但它们之间的时间关系可能不好表达。interleave 的意义是把多个投影按照发生顺序组织起来,让 UI 或日志系统拿到一条统一时间线。

多投影交错图

示例:

stream = agent.stream_events(
    {"messages": [{"role": "user", "content": "请查询订单并生成处理建议"}]},
    version="v3",
)

for event in stream.interleave("messages", "tool_calls", "values"):
    if event.type == "messages":
        print("消息事件:", event.data)
    elif event.type == "tool_calls":
        print("工具事件:", event.data)
    elif event.type == "values":
        print("状态事件:", event.data)

这里:

  • event.type:当前交错事件的类型。
  • event.data:当前事件的数据载荷。
  • interleave("messages", "tool_calls", "values"):同时订阅消息、工具和状态,并保留事件顺序。

业务场景:
Agent 控制台可以按时间线展示:“模型开始回答 -> 生成工具参数 -> 工具返回结果 -> state 更新 -> 模型继续回答”。这比把三个列表分开看更接近真实运行过程。

最简记法:

interleave 是把多个观察窗口重新排成一条时间线。


10. events() 与 raw protocol(原始协议):默认投影不够时再往下看

它解决的问题:
events() 解决的是“内置投影无法满足需求时,如何直接处理底层事件协议”。

大多数业务不需要直接处理 raw protocol。因为 messages()tool_calls()values() 等投影已经足够好用。但如果你要做框架级调试、事件转发、统一埋点、或自己实现新的投影,就需要看原始事件。

示例:

stream = agent.stream_events(
    {"messages": [{"role": "user", "content": "请解释原始事件协议"}]},
    version="v3",
)

for raw_event in stream.events():
    print("事件类型:", raw_event.type)
    print("事件载荷:", raw_event.data)

这里:

  • stream.events():原始事件投影,返回更底层的事件对象。
  • raw_event.type:事件类型。
  • raw_event.data:事件载荷。
  • raw protocol:可以理解为“所有高级投影背后的共同事件格式”。

业务场景:
如果你要把 LangChain Agent 的运行事件统一转发到公司内部事件总线,比如 Kafka、Pulsar 或日志平台,就可能需要 raw protocol。普通聊天 UI 通常不需要直接用它。

最简记法:

raw protocol 是底层出口,业务开发优先用高级投影。


11. Custom Transformers(自定义转换器):把事件变成自己的产品协议

它解决的问题:
Custom Transformers 解决的是“官方投影很好,但我的产品需要自己的事件格式”。

真实系统里,前端或后台可能已经有自己的事件协议。例如:

  • 聊天 UI 要 {type: "assistant_delta", text: "..."}
  • 工具面板要 {type: "tool_started", name: "..."}
  • 审计系统要 {type: "audit_record", trace_id: "..."}
  • 业务看板要 {type: "progress", step: "..."}

这时,不一定要让所有调用方理解 LangChain 的原始事件。更好的做法是写一个 transformer,把 LangChain 事件转换成你的产品协议。

自定义转换器图

示例:

from collections.abc import Iterator


# 把 LangChain 事件转换成前端更容易消费的中文进度事件。
def to_ui_events(raw_events: Iterator[object]) -> Iterator[dict]:
    """把原始事件转换为前端 UI 事件。"""
    for raw_event in raw_events:
        yield {
            "type": "agent_event",
            "event_type": getattr(raw_event, "type", "unknown"),
            "payload": getattr(raw_event, "data", None),
        }


stream = agent.stream_events(
    {"messages": [{"role": "user", "content": "请生成一份客户跟进建议"}]},
    version="v3",
)

for ui_event in to_ui_events(stream.events()):
    print(ui_event)

这里:

  • to_ui_events:示例自定义转换器,把底层事件映射成前端 UI 事件。
  • raw_events:输入事件迭代器。
  • event_type:保留底层事件类型,便于排查问题。
  • payload:保留事件数据载荷,实际生产中可以按需筛选。

业务场景:
公司已有一套 WebSocket 消息协议时,不需要把 LangChain 的事件对象直接暴露给前端。可以在后端统一转换,让前端只理解公司自己的 agent_eventtool_eventprogress_event

最简记法:

Transformer 是从 LangChain 事件协议到产品事件协议的翻译层。


12. Custom Updates 与 Middleware:把业务进度和治理逻辑也放进事件流

它解决的问题:
Custom UpdatesMiddleware 解决的是“模型和工具之外的业务过程,如何进入同一条可观察事件流”。

Agent 应用里经常有一些非模型事件:正在读取文件、正在查询数据库、正在跑审批规则、正在调用第三方接口、正在脱敏输出、正在等待人工审核。这些事件不一定来自模型 token,但同样影响用户体验和后台审计。

业务进度与治理图

示例:

from langgraph.config import get_stream_writer


# 检索客户资料,并在工具内部发送业务进度。
def search_customer_profile(customer_id: str) -> str:
    """根据客户编号检索模拟客户资料。"""
    writer = get_stream_writer()
    writer({"type": "progress", "message": "正在读取客户基础信息"})
    writer({"type": "progress", "message": "正在匹配最近的服务记录"})
    return f"客户 {customer_id} 最近 30 天有 2 次售后咨询。"

这里:

  • get_stream_writer():让工具内部主动写出自定义进度事件。
  • writer(...):把业务事件送入运行流。
  • customer_id:客户编号,是业务工具入参。
  • progress:示例事件类型,实际项目中可以统一成自己的业务事件协议。

官方文档还提到,事件流可以和 middleware 结合,做运行时治理。例如在输出进入前端前做 PII 脱敏,或把工具事件转换成更适合 UI 展示的形式。

业务场景:
企业内部知识库 Agent 在回答前可能要经历“搜索文档”“过滤权限”“生成摘要”“脱敏输出”。这些步骤都可以进入事件流,让前端显示真实进度,后台记录真实链路。

最简记法:

Custom updates 把业务过程放进流里,middleware 把治理规则放进流里。


13. 工程落地:前端、后端和审计系统怎么分工

它解决的问题:
这一节解决的是“知道 API 之后,真正做产品时怎么选投影”。

可以把事件流按消费者分成几层:

消费者 推荐投影 主要目的
聊天前端 messages() 展示回答文本、推理片段、工具调用提示
工具面板 tool_calls() 展示工具名、参数、结果
调试控制台 values() / events() 查看状态变化和底层事件
多 Agent 界面 subagents() / subgraphs() 区分事件来源
后台任务 output() 获取最终结果
统一日志 interleave(...) / custom transformer 保留完整时间线并转成内部协议

工程落地图

一个比较稳的后端设计是:

LangChain Agent
  -> stream_events(version="v3")
  -> 内部 transformer
  -> WebSocket / SSE / 日志平台 / 审计数据库

这里:

  • WebSocket:适合双向交互或复杂前端状态。
  • SSE:适合服务端单向推送事件,聊天流式输出里很常见。
  • 审计数据库:保存工具调用、关键状态、异常和人工审核记录。
  • transformer:把 LangChain 事件转换成公司内部稳定事件协议,避免前端直接依赖底层对象细节。

业务场景:
一个合同审查 Agent 可以这样落地:前端订阅 messages() 展示审查意见;风险面板订阅 tool_calls() 看调用了哪些规则工具;后台订阅 values() 保存每个阶段的 state;审计系统通过 transformer 保存关键事件。

最简记法:

不要让所有系统共用一条裸流;按消费者选择 projection,再用 transformer 稳定产品协议。


总结:Event Streaming 把 Agent 从“会输出”升级成“可观察”

可以把这篇文档压缩成下面这张心智表:

模块 一句话理解
stream_events(..., version="v3") 启动新版类型化事件流
messages() 消费模型消息、文本、推理片段和工具调用
tool_calls() 专门观察工具调用参数与结果
values() 观察 Agent state 每一步变化
output() 只拿最终输出
subagents() / subgraphs() 区分复杂运行结构中的事件来源
interleave(...) 把多个投影合成一条有顺序的时间线
events() 查看底层原始事件协议
custom transformer 把 LangChain 事件转换成产品自己的事件格式
custom updates / middleware 把业务进度和治理逻辑纳入事件流

这一篇最重要的抽象变化是:

从 token streaming
到 runtime event streaming
再到 product event protocol

也就是说,读完 Event Streaming 后,再回头看 messagestoolsstatesubagentsmiddleware,它们就不再是散落 API,而是同一个 Agent 工程结构里的可观察信号。

如果说上一篇 Streaming 解决的是“运行中怎么反馈”,那么这一篇 Event Streaming 解决的是“反馈本身怎么结构化、类型化、产品化”。

最简记法:

Event Streaming = 把 Agent 运行过程拆成可订阅的类型化事件。
Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐