事件流是大多数 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

各组件如何协同工作

流式处理栈有两个主要层次:
  1. 流式处理 从 Pregel 引擎发出原始的图执行事件。
  2. 事件流 对这些事件进行规范化,将它们通过流转换器运行,并暴露类型化投影。
Pregel 引擎
运行图的各个步骤
发出
原始 Pregel 事件
updatesvaluesmessagescustomcheckpointstasksdebug
发送到
事件路由器
将每个事件通过转换器管道进行路由
级联通过
流转换器
ValuesTransformer
MessagesTransformer
自定义转换器
产生
事件流
为应用程序代码投影的事件
事件路由器是两层之间的桥梁。它接收规范化后的 Pregel 事件,并将每个事件传递给已注册的流转换器。内置转换器会创建标准投影,如 stream.messagesstream.valuesstream.subgraphsstream.output。自定义转换器可以在 stream.extensions 下添加特定于应用程序的投影。

事件流提供的功能

运行流在底层事件流之上暴露了类型化投影:
投影用途
stream迭代每个协议事件。
stream.messages流式传输聊天模型消息和 token 增量。
stream.values迭代状态快照并等待最终值。
stream.output等待最终输出。
stream.subgraphs发现并观察嵌套图的执行。
stream.interrupts检查人机交互中断的有效载荷。
stream.interrupted检查运行是否因等待人工输入而暂停。
stream.extensions消费自定义流转换器投影。
多个消费者可以并发读取这些投影。读取 stream.messages 不会消耗 stream.valuesstream.subgraphsstream.output 所需的事件。 事件流位于 streaming 之上,后者通过 stream_mode 模式(如 updatesvaluesmessagescustomcheckpointstasksdebug)暴露原始的图执行事件。当你需要底层访问这些模式时,使用 streaming;当应用程序代码能从类型化投影中受益时,使用事件流。

流式传输消息

使用 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 是可迭代的。迭代它可以逐 token 输出,或者调用 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。一个从工具中调度的命名代理(例如,通过 Deep Agents 的 task 工具调用的 create_agent(name=...))会在此处按该名称呈现,并且打开该作用域的 lifecycle 事件会携带一个 cause,链接回调度的工具调用。更多信息请参阅 Lifecycle 对于特定产品的流式处理,请参阅 Deep Agents 流式处理(针对子代理流),以及 LangChain 代理流式处理(针对工具调用和中间件事件)。

流式传输状态

使用 stream.values 在每一步后流式传输完整的状态快照:
stream = graph.stream_events(input, version="v3")

for snapshot in stream.values:
    print(snapshot)

final_state = stream.output

流式传输多个投影

在异步代码中并发消费时,将 astream_eventsasyncio.gather 结合使用:
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(...) 以严格的到达顺序消费多个投影:
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.interruptedstream.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

流式传输所有协议事件

当你需要原始协议事件流时,直接使用运行对象本身:
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) 接收的正是相同的形状。
class ProtocolEvent(TypedDict):
    seq: int                    # 在一次运行中严格递增;用于排序
    method: str                 # 通道名称:"messages"、"values"、"updates"、"custom"、"tools"、"lifecycle"、...
    params: ProtocolEventParams


class ProtocolEventParams(TypedDict):
    namespace: list[str]        # 从根图开始的 "<name>:<runtime_id>" 段路径;[] 表示根
    timestamp: int              # 挂钟时间毫秒;可能漂移,不要依赖它来排序
    data: Any                   # 特定于通道的有效载荷;形状取决于 `method`
namespace 是从根图到发出事件的作用域的路径。根是空数组 []。每个子执行会增加一个 "name:runtime_id" 段,因此子图中嵌套的工具调用看起来像 ["researcher:6f4d", "tools:91ac"]: 之前的部分是稳定的图或节点名称;后缀是每次调用的运行时 ID。当你只关心特定子树时,可以自己按命名空间过滤原始事件——stream.subgraphs 已经为嵌套图执行完成了这项工作。

通道和事件生命周期

原始事件在通道上流动。通道名称作为事件的 method 出现;每个通道发出特定形状的事件。
通道用途
values完整的图状态快照。
updates每个节点的状态增量。
messages以内容块为中心的聊天模型输出。
tools工具调用开始、流式输出、完成和错误事件。
lifecycle运行、子图和子代理的状态变化。
checkpoints用于分支和时间旅行的轻量级检查点信封。
input人机交互输入请求和响应。
tasksPregel 任务创建和结果事件。
custom来自图代码的用户定义有效载荷。
custom:<name>应用程序定义的流转换器输出。
类型化投影(stream.messagesstream.values 等)正是基于这些通道构建的。当你直接迭代运行对象时,通道名称会作为原始事件上的 method 字段出现。

消息

messages 通道将输出建模为内容块。数据的 event 字段是以下之一:
  • message-start
  • content-block-start
  • content-block-delta
  • content-block-finish
  • message-finish
内容块有明确的边界:一个块开始,发出零个或多个增量,然后完成,之后同一消息中的下一个块才开始。这使得 token 流式传输、推理块、工具调用块和多模态内容显式化,无需依赖特定于提供商的格式。message-finish 可能包含 token 使用量;无法恢复的模型调用失败会作为消息错误事件到达。 要直接消费原始的内容块事件,而不使用 stream.messages 投影:
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 通道上的原始工具调用内容块关联起来。

生命周期

lifecycle 通道跟踪根运行、子图和子代理的状态。数据的 event 字段是以下之一:
  • started
  • running
  • completed
  • failed
  • interrupted
除了 event 之外,生命周期数据还可能包含可选的 graph_nameerrorcause,用于描述子作用域启动的原因(父工具调用、扇出发送、边转换)。

构建自己的投影

流转换器是事件流中的投影层。它们观察协议事件,维护自己的状态,并暴露运行的派生视图——例如工具活动、token 总数、进度事件、工件或用于另一个协议的消息。StreamChannel 是转换器用于发布这些视图的投影原语。 内置投影(stream.messagesstream.valuesstream.subgraphsstream.output)和特定于产品的投影(LangChain 的 stream.tool_calls,Deep Agents 的 stream.subagents)本身就是使用相同契约的转换器。用户转换器通过编译时或调用时注册叠加在其上,它们的投影会出现在 stream.extensions 下。 当现有投影与应用程序所需的形状不匹配时,就可以编写一个。

转换器的工作原理

事件流从 LangGraph Pregel 引擎的流式输出开始。运行时会将这些数据块规范化为协议事件,然后流处理器将每个事件通过一个流转换器栈进行路由。 流处理器是一次流中的中央调度器。对于每个协议事件,它会:
  1. 按顺序调用每个已注册转换器的 process(event) 钩子。
  2. 将命名的 StreamChannel 推送重新连接到协议事件流上。
  3. 将事件存储在运行流中,除非转换器抑制了它。
  4. 在运行结束时,对每个转换器调用 finalize()fail()
转换器是观察性的。它们不会回调到图运行时。相反,它们消费事件并将派生值推入 StreamChannel、Promise 或其他投影对象中。

转换器形状

转换器实现 StreamTransformer 接口:
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() 调用。没有转换器请求的模式永远不会被发出——声明 ("custom",) 是使 custom 事件在运行中流动的原因。
class CustomTransformer(StreamTransformer):
    required_stream_modes = ("custom",)

    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> 协议事件。将 Promise、异步可迭代对象、类实例和其他进程内句柄保留在未命名通道中。 流处理器拥有通道的生命周期。一旦 init() 返回一个通道,处理器会在运行结束时为你关闭或标记失败。转换器只负责推送值。

示例:命名通道

将字符串名称传递给 StreamChannel,可以通过 stream.extensions 暴露流式投影,并且 将每次推送的值作为 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 上访问,但对迭代原始事件的消费者不可见。对于持有无法序列化到主事件流的进程内句柄(Promise、异步可迭代对象、类实例)的投影,这是正确的选择。 下面的示例将一个未命名通道与 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)

示例:最终值投影

当投影不应流入主事件流时,使用未命名流、Promise 或其他进程内对象:
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],
)

内置:ToolCallTransformer

LangGraph 内置了 ToolCallTransformer。注册它可以在普通的 StateGraph 上暴露 stream.tool_calls
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 结合使用,请查看相关产品文档: 线级事件和命令格式在 Agent Protocol 仓库中定义,并可作为 langchain-protocol(PyPI)和 @langchain/protocol(npm)来使用。