事件流是大多数 LangGraph 应用程序代码推荐的进程内流式传输模型。它返回一个运行流对象,可以同时以多种方式消费。
快速入门
stream = graph.stream_events({
"messages": [{"role": "user", "content": "What is 42 * 17?"}],
}, version="v3")
for message in stream.messages:
for token in message.text:
print(token, end="", flush=True)
final_state = stream.output
要针对部署在 Agent Server 后面的图进行流式传输,请参阅 LangSmith 流式传输 API.
各部分如何协同工作
流式传输堆栈有两个主要层:
- **流式传输** 从 Pregel 引擎发出原始图执行事件。
- **事件流** 规范化这些事件,通过流转换器运行它们,并公开类型化投影。
Pregel 引擎 运行图步骤 发出 原始 Pregel 事件 <code>更新</code>, <code>值</code>, <code>消息</code>, <code>自定义</code>, <code>检查点</code>, <code>任务</code>, <code>调试</code> 发送到 事件路由器 通过转换器管道路由每个事件 通过以下方式级联 流转换器 ValuesTransformer MessagesTransformer ... 自定义转换器 生成 事件流 应用程序代码的投影事件
事件路由器是两 层之间的桥梁。它接收规范化的 Pregel 事件,并通过注册的流转换器传递每个事件。内置转换器创建标准投影,例如 stream.messages, stream.values, stream.subgraphs和 stream.output。自定义转换器可以在 stream.extensions.
事件流提供的内容
运行流公开了一个底层事件流上的多个类型化投影:
| 投影 | 用途 |
|---|---|
stream | 遍历每个协议事件。 |
stream.messages | 流式传输聊天模型消息和 token 增量。 |
stream.values | 遍历状态快照并等待最终值。 |
stream.output | 等待最终输出。 |
stream.subgraphs | 发现并观察嵌套图执行。 |
stream.interrupts | 检查人工介入中断的有效载荷。 |
stream.interrupted | 检查运行是否因人工输入而暂停。 |
stream.extensions | 消费自定义流转换器投影。 |
多个消费者可以并发读取这些投影。读取 stream.messages 不会消费 stream.values, stream.subgraphs, or stream.output.
事件流位于 流式处理,它通过以下方式暴露原始图执行事件 stream_mode 等模式,例如 updates, values, messages, custom, checkpoints, tasks和 debug。当您需要低级访问这些模式时,请使用流式处理;当应用程序代码受益于类型化投影时,请使用事件流式处理。
流式消息
使用 stream.messages 用于聊天模型输出:
stream = graph.stream_events(input, version="v3")
for message in stream.messages:
text = str(message.text)
usage = message.output.usage_metadata
print(text)
print(usage)
message.text 在同步代码中可迭代。迭代它以获取逐标记输出,或调用 str(message.text) 获取完整文本。
message.reasoning 暴露推理增量, message.tool_calls 暴露工具调用参数块。如果您需要按精确到达顺序获取文本、推理和工具调用块,请迭代消息流的原始事件,而不是分别迭代每个投影。
流式子图
使用 stream.subgraphs 来观察嵌套图工作,而无需解析命名空间字符串:
stream = graph.stream_events(input, version="v3")
for subgraph in stream.subgraphs:
print(subgraph.graph_name, subgraph.path)
for message in subgraph.messages:
print(message.text)
subgraph.graph_name 是 name 编译图或代理的。当从工具调度命名代理(例如, create_agent(name=...) 通过 Deep Agents 调用 task 工具)时,它会在该名称下出现在此处,并且 lifecycle 事件会打开作用域并携带 cause ,链接回调度工具调用。参见 生命周期 获取更多信息。
对于特定产品的流,请参见 Deep Agents 流式处理 获取子代理流, LangChain 代理流式处理 获取工具调用和中间件事件。
流式状态
使用 stream.values 在每个步骤后流式传输完整状态快照:
stream = graph.stream_events(input, version="v3")
for snapshot in stream.values:
print(snapshot)
final_state = stream.output
流式传输多个投影
对于异步代码中的并发使用,请使用 astream_events 配合 asyncio.gather:
stream = await graph.astream_events(input, version="v3")
async def consume_messages():
async for message in stream.messages:
print(f"[llm] node={message.node}")
async def consume_subgraphs():
async for subgraph in stream.subgraphs:
print(f"[subgraph] path={subgraph.path}")
await asyncio.gather(consume_messages(), consume_subgraphs())
对于同步代码,请使用 stream.interleave(...) 以严格到达顺序使用多个投影:
stream = graph.stream_events(input, version="v3")
for name, item in stream.interleave("values", "messages", "subgraphs"):
if name == "values":
print(f"[state] keys={list(item)}")
elif name == "messages":
print(f"[llm] node={item.node}")
elif name == "subgraphs":
print(f"[subgraph] path={item.path}")
中断后恢复
当图暂停等待人工输入时,请检查 stream.interrupted 和 stream.interrupts,然后通过调用 stream_events(..., version="v3") 再次使用 Command.
恢复需要一个使用检查点编译的图和一个带有线程 ID 的配置——参见 持久化.
from langgraph.types import Command
stream = graph.stream_events(input, version="v3")
for message in stream.messages:
print(message.text)
if stream.interrupted:
print(stream.interrupts)
stream = graph.stream_events(
Command(resume={"decisions": [{"type": "approve"}]}),
version="v3",
)
final_state = stream.output
流式传输所有协议事件
当您需要原始协议事件流时,请直接使用 run 对象:
stream = graph.stream_events({
"messages": [{"role": "user", "content": "What is 42 * 17?"}],
}, version="v3")
for event in stream:
namespace = event["params"]["namespace"]
print(namespace, event["method"], event["params"]["data"])
每个事件都是 ProtocolEvent 包装特定通道负载的信封。同样的结构也是转换器的 process(event) receives.
class ProtocolEvent(TypedDict):
seq: int # strictly increasing within a run; use for ordering
method: str # channel name: "messages", "values", "updates", "custom", "tools", "lifecycle", ...
params: ProtocolEventParams
class ProtocolEventParams(TypedDict):
namespace: list[str] # path of "<name>:<runtime_id>" segments from the root graph; [] is the root
timestamp: int # wall-clock milliseconds; can drift, don't rely on for ordering
data: Any # channel-specific payload; shape depends on `method`
是 namespace 是 从根图到发出事件的范围的路径。根是空数组 []。每个子执行添加一个 "name:runtime_id" 段,因此子图中嵌套的工具调用类似于 ["researcher:6f4d", "tools:91ac"]。 : 之前的名称是稳定的图或节点名称;后缀是每次调用的运行时 ID。当您只关心特定子树时,请自行按命名空间过滤原始事件—— stream.subgraphs 已经为嵌套图执行完成了此操作。
通道和事件生命周期
原始事件通过通道流动。通道名称作为事件的 method出现;每个通道发出一个特定的事件结构。
| 通道 | 用途 |
|---|---|
values | 完整图状态快照。 |
updates | 每个节点的状态增量。 |
messages | 以内容块为中心的聊天模型输出。 |
tools | 工具调用启动、流式输出、完成和错误事件。 |
lifecycle | 运行、子图和子代理状态变更。 |
checkpoints | 用于分支和时间旅行的轻量级检查点信封。 |
input | 人机交互输入请求和响应。 |
tasks | Pregel 任务创建和结果事件。 |
custom | 来自图代码的用户定义负载。 |
custom:<name> | 应用程序定义的流转换器输出。 |
类型化投影(stream.messages, stream.values等)从这些通道构建。当您直接迭代运行对象时,通道名称作为 method 字段出现在原始事件上。
消息
messages 通道将输出建模为内容块。数据的 event 字段是以下之一:
- -
message-start - -
content-block-start - -
content-block-delta - -
content-block-finish - -
message-finish
内容块具有明确的边界:一个块开始,发出零个或多个增量,然后在同一消息中的下一个块开始之前完成。这使得令牌流式传输、推理块、工具调用块和多模态内容显式化,而无需特定于提供商的格式。 message-finish 可能包含令牌使用量;不可恢复的模型调用失败作为消息错误事件到达。
要直接消费原始内容块事件而不是使用 stream.messages projection:
for event in stream:
if event["method"] != "messages":
continue
data = event["params"]["data"][0]
if not isinstance(data, dict):
continue
if data.get("event") != "content-block-delta":
continue
block = data.get("delta") or {}
if block.get("type") == "text-delta":
print(block.get("text", ""), end="", flush=True)
elif block.get("type") == "reasoning-delta":
print(f"[thinking]{block.get('reasoning', '')}", end="", flush=True)
工具
tools 通道公开工具执行。数据的 event 字段是以下之一:
- -
tool-started - -
tool-output-delta - -
tool-finished - -
tool-error
工具事件通过工具调用 ID 关联,因此工具执行可以与其在 messages channel.
生命周期
lifecycle 通道跟踪根运行、子图和子代理状态。数据的 event 字段是以下之一:
- -
started - -
running - -
completed - -
failed - -
interrupted
除了 event之外,生命周期数据可能包含一个可选的 graph_name, error和 cause ,描述子范围启动的原因(父工具调用、扇出发送、边转换)。
构建您自己的投影
流转换器是事件流中的投影层。它们观察协议事件,维护自己的状态,并暴露运行时的派生视图——如工具活动、token总数、进度事件、产物或另一个协议的消息。 StreamChannel 是转换器用来发布这些视图的投影原语。
内置投影(stream.messages, stream.values, stream.subgraphs, stream.output)和产品特定投影(LangChain的 stream.tool_calls,Deep Agents的 stream.subagents)本身就是使用此契约的转换器。用户转换器通过编译时或调用时注册堆叠在它们之上,它们的投影显示在 stream.extensions.
当现有投影不符合应用程序需要的形状时,请编写一个。
转换器的工作原理
事件流从LangGraph Pregel引擎的流式输出开始。运行时将这些块规范化为协议事件,然后流处理器通过流转换器堆栈路由每个事件。
flowchart TD
A[Pregel modes] --> B[Events]
B --> C[Built-in projections]
C --> D[User transformers]
D --> E[Run projections]
流处理器是一个流的中央调度器。对于每个协议事件,它会:
- 按顺序调用每个已注册转换器的
process(event)钩子。 - 连接命名
StreamChannel推回协议事件流。 - 将事件存储在运行流中,除非转换器抑制了它。
- 在运行结束时调用
finalize()orfail()每个转换器。
转换器是观察性的。它们不会回调图运行时。相反,它们消费事件并将派生值推入 StreamChannel、promise或其他投影对象。
转换器形状
转换器实现 StreamTransformer interface:
from langgraph.stream import ProtocolEvent, StreamTransformer
class MyTransformer(StreamTransformer):
def init(self) -> dict:
...
def process(self, event: ProtocolEvent) -> bool:
...
def finalize(self) -> None:
...
def fail(self, err: BaseException) -> None:
...
- -
init()创建投影对象。用户转换器投影显示在stream.extensions. - -
process()观察每个协议事件。参见 流式传输所有协议事件 获取ProtocolEvent形状。仅当你有意抑制原始事件时才返回false。 - -
finalize()在流成功结束后关闭或解析非通道投影。 - -
fail()将错误传播到非通道投影。
声明所需的流模式
required_stream_modes 控制底层图在流期间发出哪些Pregel流模式。运行时取每个已注册转换器的 required_stream_modes 并集,并将其作为 stream_mode 参数传递给图的 .stream() call. **没有转换器请求的模式永远不会发出** ——声明 ("custom",) 才是导致 custom 事件在运行中流动的原因。
class CustomTransformer(StreamTransformer):
required_stream_modes = ("custom",) # [!code highlight]
def process(self, event: ProtocolEvent) -> bool:
if event["method"] == "custom":
...
return True
process() 接收图发出的每个事件,并负责按 event["method"]进行过滤。声明会开启上游发射;不会缩小 process() 看到的内容。有效值为Pregel流模式: "messages", "tools", "custom", "values", "updates", "checkpoints", "tasks", "debug"。每个转换器必须声明它所操作的每个模式——省略的模式不会由图发出,也永远不会到达 process().
StreamChannel
StreamChannel 是转换器用于流式值的投影原语。它始终在 stream.extensions.<name>. 构造函数参数决定每个 push() 也作为 custom:<name> 事件——即,当迭代原始协议事件时,投影的值是否出现。
| 需求 | 使用 |
|---|---|
| 仅侧通道投影 | StreamChannel() |
| 同时将每次推送流入主事件流 | StreamChannel(name) |
命名通道的有效载荷必须可序列化,因为每个推送的值也会成为 custom:<name> 主流中的协议事件。请将 promises、异步迭代器、类实例和其他进程内句柄保留在未命名通道中。
流处理器负责通道的生命周期。一旦 init() 返回通道,处理器会在运行结束时为你关闭它或使其失败。转换器只能推送值。
示例:命名通道
传递一个字符串名称给 StreamChannel 通过 stream.extensions *和* 将每次推送的值转发到运行的 main event stream 作为 custom:<name> 协议事件:
from typing import TypedDict
from langgraph.stream import ProtocolEvent, StreamChannel, StreamTransformer
class ToolActivity(TypedDict):
name: str
status: str
class ToolActivityTransformer(StreamTransformer):
required_stream_modes = ("tools",)
def __init__(self, scope: tuple[str, ...] = ()) -> None:
super().__init__(scope)
self.activity = StreamChannel[ToolActivity]("tool_activity")
def init(self) -> dict:
return {"tool_activity": self.activity}
def process(self, event: ProtocolEvent) -> bool:
if event["method"] != "tools":
return True
data = event["params"]["data"]
if isinstance(data, dict) and data.get("tool_name") and data.get("event"):
status = "error" if data["event"] == "tool-error" else "started"
self.activity.push({"name": data["tool_name"], "status": status})
return True
示例:未命名通道
没有名称的通道是仅侧通道投影——可在 stream.extensions 上访问,但对迭代原始事件的消费者不可见。这对于包含进程内句柄(promises、异步迭代器、类实例)的投影是正确的选择,这些句柄无法序列化到主事件流上。
下面的示例将未命名通道与 get_stream_writer配对,让图节点发出 custom通道事件,然后由转换器将其排入投影:
from langgraph.config import get_stream_writer
from langgraph.stream import ProtocolEvent, StreamChannel, StreamTransformer
def node(state):
writer = get_stream_writer()
writer({"kind": "progress", "message": "retrieving context"})
return state
class CustomTransformer(StreamTransformer):
required_stream_modes = ("custom",)
def __init__(self, scope: tuple[str, ...] = ()) -> None:
super().__init__(scope)
self.log = StreamChannel()
def init(self) -> dict:
return {"custom": self.log}
def process(self, event: ProtocolEvent) -> bool:
if event["method"] == "custom":
self.log.push(event["params"]["data"])
return True
stream = graph.stream_events(input, version="v3", transformers=[CustomTransformer])
for item in stream.extensions["custom"]:
print(item)
示例:最终值投影
当投影不应流入主事件流时,请使用未命名流、promises 或其他进程内对象:
from langgraph.stream import ProtocolEvent, StreamChannel, StreamTransformer
class StatsTransformer(StreamTransformer):
required_stream_modes = ("messages",)
def __init__(self, scope: tuple[str, ...] = ()) -> None:
super().__init__(scope)
self.total_tokens = 0
self.total_tokens_log = StreamChannel[int]()
def init(self) -> dict:
return {"total_tokens": self.total_tokens_log}
def process(self, event: ProtocolEvent) -> bool:
data = event["params"]["data"]
if isinstance(data, dict):
usage = data.get("usage") or {}
self.total_tokens += usage.get("output_tokens") or 0
return True
def finalize(self) -> None:
self.total_tokens_log.push(self.total_tokens)
self.total_tokens_log.close()
在调用时或编译时注册
在调用时传递转换器以进行本地实验:
stream = graph.stream_events(
input,
version="v3",
transformers=[StatsTransformer, ToolActivityTransformer],
)
将转换器编译到图中,以便该图的每次运行都产生投影:
graph = builder.compile(
transformers=[StatsTransformer, ToolActivityTransformer],
)
Built-in: ToolCallTransformer
LangGraph 附带 ToolCallTransformer 作为内置功能。注册它以在普通 stream.tool_calls 上公开 StateGraph:
from langgraph.prebuilt import ToolCallTransformer
stream = graph.stream_events(input, version="v3", transformers=[ToolCallTransformer])
for tool_call in stream.tool_calls:
print(tool_call.tool_name, tool_call.input)
相关
LangGraph 定义了流式原语。关于将流式功能与 LangChain 或 Deep Agents 结合使用,请参阅相关产品文档:
- - LangChain agent 流式传输 涵盖 ReAct 风格 agent 消息、工具调用和中间件更新。
- - Deep Agents 流式传输 涵盖子 agent、嵌套消息和子 agent 工具调用。
- - LangChain 前端模式 和 LangGraph 前端模式 展示基于流式状态构建的 UI 用例。
- - LangSmith Streaming API 涵盖对部署在 Agent Server 后面的图进行流式传输。
线级事件和命令格式定义在 Agent Protocol 仓库中,可作为 langchain-protocol 在 PyPI 上和 @langchain/protocol 在 npm 上。