外观
事件流是大多数 LangGraph 应用代码推荐采用的进程内流式输出模型。它返回一个运行流对象,可以同时以多种方式消费。
快速入门
py
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.outputts
const stream = await graph.streamEvents(
{ messages: [{ role: "user", content: "What is 42 * 17?" }] },
{ version: "v3" }
);
for await (const message of stream.messages) {
for await (const token of message.text) {
process.stdout.write(token);
}
}
const finalState = await stream.output;要针对部署在 Agent Server 后方的图进行流式输出,请参阅 LangSmith Streaming API。
各部分如何协同工作
流式输出栈主要有两层:
- 流式输出从 Pregel 引擎发出原始图执行事件。
- 事件流对这些事件进行规范化处理,将其传递给流转换器,并暴露类型化投影。
事件路由器是两层之间的桥梁。它接收规范化后的 Pregel 事件,并将每个事件传递给已注册的流转换器。内置转换器会创建标准投影,例如 stream.messages、stream.values、stream.subgraphs 和 stream.output。自定义转换器可以在 stream.extensions 下添加特定于应用的投影。
事件流提供什么
运行流在单个底层事件流之上暴露类型化投影:
| Projection | Use |
|---|---|
stream | 遍历每个协议事件。 |
stream.messages | 流式输出对话模型消息和 token 增量。 |
stream.values | 遍历状态快照并等待最终值。 |
stream.output | 等待最终输出。 |
stream.subgraphs | 发现并观察嵌套图执行。 |
stream.interrupts | 检查人在回路中断的载荷。 |
stream.interrupted | 检查运行是否因等待人工输入而暂停。 |
stream.extensions | 消费自定义流转换器投影。 |
多个消费者可以并发读取这些投影。读取 stream.messages 不会消耗 stream.values、stream.subgraphs 或 stream.output 所需的事件。
事件流位于流式输出之上的一层,后者通过 stream_mode 模式(如 updates、values、messages、custom、checkpoints、tasks 和 debug)暴露原始图执行事件。当你需要对这些模式进行底层访问时,使用流式输出;当应用代码受益于类型化投影时,使用事件流。
流式输出消息
使用 stream.messages 获取对话模型输出:
py
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)ts
const stream = await graph.streamEvents(input, { version: "v3" });
for await (const message of stream.messages) {
const text = await message.text;
const usage = await message.usage;
console.log(text);
console.log(usage);
}在同步代码中,message.text 是可迭代的。遍历它以获得逐 token 的输出,或调用 str(message.text) 获取完整文本。
message.reasoning 暴露推理增量,message.tool_calls 暴露工具调用参数数据块。如果你需要严格按照到达顺序获取文本、推理和工具调用数据块,请遍历消息流的原始事件,而不是分别遍历每个投影。
message.text 既是异步可迭代对象,也是类 Promise 的值。遍历它以获得逐 token 的输出,或 await 它以获取完整文本。
流式输出子图
使用 stream.subgraphs 观察嵌套图工作,而无需解析命名空间字符串:
py
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)ts
const stream = await graph.streamEvents(input, { version: "v3" });
for await (const subgraph of stream.subgraphs) {
console.log(subgraph.name, subgraph.path);
for await (const message of subgraph.messages) {
console.log(await message.text);
}
}subgraph.graph_name 是编译后的图或智能体的 name。从工具派发的命名智能体(例如,通过 Deep Agents 的 task 工具调用的 create_agent(name=...))会以该名称在此呈现,而开启该作用域的 lifecycle 事件带有一个 cause,可回溯到派发它的工具调用。更多信息请参阅生命周期。
有关特定产品的流,请参阅 Deep Agents 事件流了解子智能体流,参阅 LangChain 智能体流了解工具调用和中间件事件。
流式输出状态
使用 stream.values 流式输出每一步之后的完整状态快照:
py
stream = graph.stream_events(input, version="v3")
for snapshot in stream.values:
print(snapshot)
final_state = stream.outputts
const stream = await graph.streamEvents(input, { version: "v3" });
for await (const snapshot of stream.values) {
console.log(snapshot);
}
const finalState = await stream.output;流式输出多个投影
要在异步代码中并发消费,请将 astream_events 与 asyncio.gather 一起使用:
py
import asyncio
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(...) 以严格的到达顺序消费多个投影:
py
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}")在 JavaScript 中需要多个投影时,请使用并发消费者:
ts
await Promise.all([
(async () => {
for await (const message of stream.messages) {
console.log(await message.text);
}
})(),
(async () => {
for await (const subgraph of stream.subgraphs) {
console.log(subgraph.path);
}
})(),
]);中断后恢复
当图因等待人工输入而暂停时,检查 stream.interrupted 和 stream.interrupts,然后使用 Command 再次调用 stream_events(..., version="v3") 以恢复执行。
恢复需要图使用检查点器编译,并且配置中带有线程 ID——请参阅持久化。
py
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.outputts
import { Command } from "@langchain/langgraph";
let stream = await graph.streamEvents(input, { version: "v3" });
for await (const message of stream.messages) {
console.log(await message.text);
}
if (stream.interrupted) {
console.log(stream.interrupts);
}
stream = await graph.streamEvents(
new Command({ resume: { decisions: [{ type: "approve" }] } }),
{ version: "v3" }
);
const finalState = await stream.output;流式输出所有协议事件
当你需要原始协议事件流时,直接使用运行对象本身:
py
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"])ts
const stream = await graph.streamEvents(
{ messages: [{ role: "user", content: "What is 42 * 17?" }] },
{ version: "v3" }
);
for await (const event of stream) {
const namespace = event.params.namespace;
console.log(namespace, event.method, event.params.data);
}每个事件都是一个 ProtocolEvent 信封,包裹着特定于通道的载荷。转换器的 process(event) 接收的也是这种相同的结构。
py
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`ts
interface ProtocolEvent {
readonly seq: number; // strictly increasing within a run; use for ordering
readonly method: string; // channel name: "messages", "values", "updates", "custom", "tools", "lifecycle", ...
readonly params: {
readonly namespace: string[]; // path of "<name>:<runtime_id>" segments from the root graph; [] is the root
readonly timestamp: number; // wall-clock milliseconds; can drift, don't rely on for ordering
readonly node?: string; // graph node that emitted this event, when applicable
readonly data: unknown; // channel-specific payload; shape depends on `method`
};
}namespace 是从根图到发出该事件的作用域的路径。根为空数组 []。每次子执行都会添加一个 "name:runtime_id" 段,因此子图内部的嵌套工具调用看起来像 ["researcher:6f4d", "tools:91ac"]。: 之前的名称是稳定的图或节点名称;后缀是每次调用的运行时 ID。当你只关心特定子树时,请自行按命名空间过滤原始事件——stream.subgraphs 已经为嵌套图执行做了这件事。
通道与事件生命周期
原始事件在通道上流动。通道名称以事件的 method 形式出现;每个通道都会发出特定形状的事件。
| Channel | Purpose |
|---|---|
values | 完整的图状态快照。 |
updates | 每个节点的状态增量。 |
messages | 以内容块为中心的对话模型输出。 |
tools | 工具调用开始、流式输出、结束和错误事件。 |
lifecycle | 运行、子图和子智能体状态变化。 |
checkpoints | 用于分支和时间旅行的轻量级检查点信封。 |
input | 人在回路输入请求和响应。 |
tasks | Pregel 任务创建和结果事件。 |
custom | 来自图代码的用户定义载荷。 |
custom:<name> | 应用定义的流转换器输出。 |
类型化投影(stream.messages、stream.values 等)由这些通道构建。当你直接遍历运行对象时,通道名称会以 method 字段的形式出现在原始事件上。
消息
messages 通道将输出建模为内容块。数据的 event 字段为以下之一:
message-startcontent-block-startcontent-block-deltacontent-block-finishmessage-finish
内容块具有明确的边界:一个块开始、发出零个或多个增量,然后结束,之后同一消息中的下一个块才开始。这使 token 流式传输、推理块、工具调用块和多模态内容变得明确,而无需提供商特定的格式。message-finish 可能包含 token 用量;无法恢复的模型调用失败会以消息错误事件的形式到达。
要直接消费原始内容块事件,而不是使用 stream.messages 投影:
py
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)ts
for await (const event of stream) {
if (event.method !== "messages") continue;
const data = event.params.data;
if (data.event !== "content-block-delta") continue;
const block = data.delta ?? {};
if (block.type === "text-delta") {
process.stdout.write(block.text ?? "");
} else if (block.type === "reasoning-delta") {
process.stdout.write(`[thinking]${block.reasoning ?? ""}`);
}
}工具
tools 通道暴露工具执行。数据的 event 字段为以下之一:
tool-startedtool-output-deltatool-finishedtool-error
工具事件通过工具调用 ID 关联,因此工具执行可以关联回 messages 通道上发起它的工具调用内容块。
生命周期
lifecycle 通道跟踪根运行、子图和子智能体状态。数据的 event 字段为以下之一:
startedrunningcompletedfailedinterrupted
除 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 引擎的流式输出。运行时将这些数据块规范化为协议事件,然后流处理器将每个事件路由到一组流转换器。
流处理器是单个流的中央调度器。对于每个协议事件,它都会:
- 按顺序调用每个已注册转换器的
process(event)钩子。 - 将命名
StreamChannel的推送重新接入协议事件流。 - 将事件存储到运行流中,除非某个转换器抑制了它。
- 在运行结束时对每个转换器调用
finalize()或fail()。
转换器是观察性的。它们不会回调图运行时。相反,它们消费事件,并将派生值推送到 StreamChannel、Promise 或其他投影对象中。
转换器结构
转换器实现 StreamTransformer 接口:
py
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:
...ts
interface StreamTransformer<TProjection = unknown> {
init(): TProjection;
process(event: ProtocolEvent): boolean;
finalize?(): void | PromiseLike<void>;
fail?(err: unknown): void;
}init()创建投影对象。用户转换器的投影出现在stream.extensions下。process()观察每个协议事件。ProtocolEvent的结构请参阅流式输出所有协议事件。仅当你故意要抑制原始事件时才返回false。finalize()在流成功结束后关闭或解析非通道投影。fail()将错误传播到非通道投影。
声明所需的流模式
required_stream_modes 控制底层图在流式过程中发出哪些 Pregel 流模式。运行时取所有已注册转换器的 required_stream_modes 的并集,并将该并集作为 stream_mode 参数传给图的 .stream() 调用。没有任何转换器请求的模式永远不会被发出——正是声明 ("custom",) 才使 custom 事件能够流过整个运行。
py
class CustomTransformer(StreamTransformer):
required_stream_modes = ("custom",)
def process(self, event: ProtocolEvent) -> bool:
if event["method"] == "custom":
...
return Trueprocess() 接收图发出的每个事件,并负责按 event["method"] 进行过滤。声明会开启上游的发出;它不会收窄 process() 所能看到的内容。有效值是 Pregel 流模式:"messages"、"tools"、"custom"、"values"、"updates"、"checkpoints"、"tasks"、"debug"。每个转换器都必须声明其操作的所有模式——未声明的模式不会被图发出,也永远不会到达 process()。
StreamChannel
StreamChannel 是转换器用于流式输出值的投影原语。它总是在 stream.extensions.<name> 上暴露一个可迭代流。构造函数参数决定每次 push() 是否也会以 custom:<name> 事件的形式流入运行的主事件流——也就是说,遍历原始协议事件时,该投影的值是否会出现。
| Need | Use |
|---|---|
| 仅旁路通道投影 | new StreamChannel<T>() |
| 同时将每次推送流入主事件流 | new StreamChannel<T>(name) |
| Need | Use |
|---|---|
| 仅旁路通道投影 | StreamChannel() |
| 同时将每次推送流入主事件流 | StreamChannel(name) |
命名通道的载荷必须是可序列化的,因为每个推送的值也会成为主流中的 custom:<name> 协议事件。请将 Promise、异步可迭代对象、类实例和其他进程内句柄放在未命名通道中。
流处理器拥有通道的生命周期。一旦 init() 返回通道,处理器就会在运行结束时为你关闭或使其失败。转换器只推送值。
示例:命名通道
向 StreamChannel 传入字符串名称,即可通过 stream.extensions 暴露流式投影,并将每个推送的值以 custom:<name> 协议事件的形式转发到运行的主事件流:
py
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 Truets
import { StreamChannel } from "@langchain/langgraph";
const toolActivityTransformer = () => {
const activity = new StreamChannel<{
name: string;
status: "started" | "finished" | "error";
}>("toolActivity");
return {
init: () => ({ toolActivity: activity }),
process(event) {
if (event.method === "tools") {
const data = event.params.data as { tool_name?: string; event?: string };
if (data.tool_name && data.event) {
activity.push({
name: data.tool_name,
status: data.event === "tool-error" ? "error" : "started",
});
}
}
return true;
},
};
};示例:未命名通道
没有名称时,通道只是旁路通道投影——可以在 stream.extensions 上访问,但对遍历原始事件的消费者不可见。对于持有无法序列化到主事件流上的进程内句柄(Promise、异步可迭代对象、类实例)的投影来说,这是正确的选择。
下面的示例将未命名通道与 get_stream_writer 搭配使用,后者让图节点能够发出 custom 通道事件,转换器随后将这些事件汇集到投影中:
py
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)ts
import { StreamChannel } from "@langchain/langgraph";
const customTransformer = () => {
const custom = new StreamChannel<unknown>();
return {
init: () => ({ custom }),
process(event) {
if (event.method === "custom") {
custom.push(event.params.data);
}
return true;
},
};
};示例:最终值投影
当投影不应流入主事件流时,请使用未命名流、Promise 或其他进程内对象:
py
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()ts
const statsTransformer = () => {
let totalTokens = 0;
let resolveTotal!: (value: number) => void;
const totalTokensPromise = new Promise<number>((resolve) => {
resolveTotal = resolve;
});
return {
init: () => ({ totalTokens: totalTokensPromise }),
process(event) {
if (event.method === "messages") {
const data = event.params.data as { usage?: { output_tokens?: number } };
totalTokens += data.usage?.output_tokens ?? 0;
}
return true;
},
finalize: () => resolveTotal(totalTokens),
};
};在调用时或编译时注册
在调用时传入转换器,便于本地实验:
py
stream = graph.stream_events(
input,
version="v3",
transformers=[StatsTransformer, ToolActivityTransformer],
)ts
const stream = await graph.streamEvents(input, {
version: "v3",
transformers: [statsTransformer, toolActivityTransformer],
});当该图的每次运行都应产生该投影时,将转换器编译进图中:
py
graph = builder.compile(
transformers=[StatsTransformer, ToolActivityTransformer],
)ts
const graph = builder.compile({
transformers: [statsTransformer, toolActivityTransformer],
});内置:ToolCallTransformer
LangGraph 内置了 ToolCallTransformer。注册它即可在普通的 StateGraph 上暴露 stream.tool_calls:
py
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 智能体事件流涵盖 ReAct 风格的智能体消息、工具调用和中间件更新。
- Deep Agents 事件流涵盖子智能体、嵌套消息和子智能体工具调用。
- LangChain 前端模式和 LangGraph 前端模式展示了构建在流式状态之上的 UI 用例。
- LangSmith Streaming API涵盖针对部署在 Agent Server 后方的图的流式输出。
线级事件和命令格式在 Agent Protocol 仓库中定义,可作为 PyPI 上的 langchain-protocol 和 npm 上的 @langchain/protocol 使用。