Skip to content

事件流(Event Streaming)

普通的 stream() 返回的是状态/消息流--你看到的是"状态怎么变的"。事件流返回的是执行事件的类型化投影--你看到的是"运行时发生了什么":哪个 run 开始了、哪个节点产出了值、哪个子图在执行。这让它成为可观测性、调试和实时 UI 推送的首选 API。

先修知识

必须深刻理解,不能跳过stream_modestream_events 是两套不同的 API,解决不同的问题。stream_mode 回答"状态的哪些切片要推送给我";stream_events 回答"执行的哪些事件要推送给我"。前者是数据视角,后者是运行时视角。混用它们是新手最常犯的错误。

前端类比:先建立直觉

前端概念LangGraph 对应说明
Redux store subscribestream(stream_mode=...)订阅状态变化
Redux middleware action logstream_events(version="v3")订阅执行事件
WebSocket 多频道channel 投影 (run.values 等)按频道过滤事件
Chrome DevTools Performancerun.lifecycle 投影执行生命周期追踪

LangGraph 原生语义:事件流通过 stream_events(version="v3")(同步)或 astream_events(version="v3")(异步)获取。v3 API(langgraph >= 1.2,beta)返回类型化投影(typed projections),你可以按 channel 过滤,只接收你关心的事件类型。相比 v2 的 stream_mode,事件流提供了更丰富的执行上下文信息,特别适合调试和 UI 推送场景。

版本状态

stream_events(version="v3")v1.2 引入的 beta API,官方推荐新应用使用,但 API 可能在后续版本微调。生产使用前请核实最新文档。v1 格式已过时,本页不展开。

🔗 Event Streaming 官方文档


1. 事件流 vs 普通 stream()

核心区别

维度stream(stream_mode=...)stream_events(version="v3")
视角状态切片(数据)执行事件(运行时)
返回按 mode 格式化的 chunk类型化投影事件
过滤stream_mode按 channel / event 类型
子图subgraphs=True 展开内置 run.subgraphs 投影
最佳场景前端展示状态/token调试、可观测性、UI 推送
版本v1(默认)/ v2(推荐)v3(beta,推荐新应用)

什么时候用哪个

python
# 场景 1:前端打字机效果 -> 用 stream_mode="messages"
for chunk in graph.stream(input, stream_mode="messages"):
    msg_chunk, metadata = chunk
    print(msg_chunk.content, end="")

# 场景 2:调试"为什么节点没执行" -> 用 stream_events
async for event in graph.astream_events(input, version="v3"):
    if event["channel"] == "run.lifecycle":
        print(f"生命周期: {event['data']}")

2. v3 事件结构

每个事件是一个类型化的字典,包含 channel、namespace 和 data:

python
# 事件结构(概念性示意)
{
    "channel": "run.values",   # 事件所属的投影频道
    "ns": "agent",              # 命名空间(节点/子图路径)
    "data": {                   # 事件数据,随 channel 类型变化
        ...
    }
}

四种投影频道(projections)

投影channel 名数据内容类比 stream_mode
valuesrun.values完整状态快照stream_mode="values"
messagesrun.messagesLLM token 级输出stream_mode="messages"
lifecyclerun.lifecycle执行生命周期事件(run 开始/结束、节点进入/退出)无直接对应
subgraphsrun.subgraphs子图执行事件subgraphs=True

关键区别run.lifecycle 是事件流独有的投影,普通 stream_mode 没有等价物。它让你追踪"哪个 run 开始了、哪个节点正在执行、哪个结束了",这是调试和可观测性的核心需求。


3. 基本用法

同步:stream_events

python
from langgraph.graph import StateGraph, START, END, MessagesState
from langgraph.checkpoint.memory import InMemorySaver

def echo(state: MessagesState):
    return {"messages": [{"role": "assistant", "content": "hello"}]}

builder = StateGraph(MessagesState)
builder.add_node("echo", echo)
builder.add_edge(START, "echo")
builder.add_edge("echo", END)
graph = builder.compile(checkpointer=InMemorySaver())

# 同步事件流
for event in graph.stream_events(
    {"messages": [{"role": "user", "content": "hi"}]},
    version="v3",
):
    print(f"[{event['channel']}] ns={event.get('ns', '')} data={event['data']}")

异步:astream_events(推荐用于 Web 服务)

python
import asyncio

async def consume_events():
    async for event in graph.astream_events(
        {"messages": [{"role": "user", "content": "hi"}]},
        version="v3",
    ):
        if event["channel"] == "run.lifecycle":
            print(f"生命周期: {event['data']}")
        elif event["channel"] == "run.messages":
            print(f"Token: {event['data']}")

asyncio.run(consume_events())

4. 按 channel 过滤

如果你只关心某一类事件,用 channel 过滤减少不必要的数据传输:

python
# 只接收生命周期事件(调试用)
async for event in graph.astream_events(
    input,
    version="v3",
    channel="run.lifecycle",
):
    print(f"生命周期: {event['data']}")

# 只接收 messages(前端打字机效果)
async for event in graph.astream_events(
    input,
    version="v3",
    channel="run.messages",
):
    token = event["data"]
    if token:
        print(token, end="", flush=True)

多 channel 订阅

python
# 同时接收 values 和 lifecycle
async for event in graph.astream_events(
    input,
    version="v3",
    channel=["run.values", "run.lifecycle"],
):
    if event["channel"] == "run.values":
        print(f"状态快照: {event['data']}")
    elif event["channel"] == "run.lifecycle":
        print(f"执行事件: {event['data']}")

5. 生命周期事件详解

run.lifecycle 是事件流最有价值的投影,它记录了执行的完整生命周期:

python
async for event in graph.astream_events(input, version="v3", channel="run.lifecycle"):
    data = event["data"]
    # 典型生命周期事件序列:
    # 1. run 开始
    # 2. 节点 "retrieve" 开始
    # 3. 节点 "retrieve" 结束
    # 4. 节点 "generate" 开始
    # 5. 节点 "generate" 结束
    # 6. run 结束
    print(f"[{event['ns']}] {data}")

事件流数据流

前端类比:这就像 Chrome DevTools 的 Performance 面板--你看到的不是"页面长什么样"(那是 stream_mode="values"),而是"哪个函数什么时候开始、什么时候结束、花了多久"(那是 run.lifecycle)。


6. 子图事件

当图包含子图时,run.subgraphs 投影让你追踪子图内部执行:

python
from langgraph.graph import StateGraph, START, END, MessagesState

# 子图
def sub_process(state: MessagesState):
    return {"messages": [{"role": "assistant", "content": "子图完成"}]}

sub_builder = StateGraph(MessagesState)
sub_builder.add_node("sub_process", sub_process)
sub_builder.add_edge(START, "sub_process")
sub_builder.add_edge("sub_process", END)
sub_graph = sub_builder.compile()

# 主图
def main_node(state: MessagesState):
    return {"messages": [{"role": "assistant", "content": "主图节点"}]}

main_builder = StateGraph(MessagesState)
main_builder.add_node("main", main_node)
main_builder.add_node("sub", sub_graph)
main_builder.add_edge(START, "main")
main_builder.add_edge("main", "sub")
main_builder.add_edge("sub", END)
graph = main_builder.compile()

# 接收子图事件
async for event in graph.astream_events(input, version="v3"):
    ns = event.get("ns", "")
    if "sub" in ns:
        print(f"[子图事件] {event['channel']}: {event['data']}")
    else:
        print(f"[主图事件] {event['channel']}: {event['data']}")

7. 前端消费方式

通过 FastAPI SSE 推送

python
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
import json

app = FastAPI()

@app.post("/chat/events")
async def chat_events(query: str):
    async def event_generator():
        async for event in graph.astream_events(
            {"messages": [{"role": "user", "content": query}]},
            version="v3",
        ):
            # 将事件序列化为 SSE 格式
            sse_data = json.dumps({
                "channel": event["channel"],
                "ns": event.get("ns", ""),
                "data": event["data"],
            }, ensure_ascii=False, default=str)
            yield f"data: {sse_data}\n\n"
        yield "data: [DONE]\n\n"

    return StreamingResponse(
        event_generator(),
        media_type="text/event-stream",
    )

前端接收端(TypeScript 示例)

typescript
// 前端按 channel 分流处理
const source = new EventSource('/chat/events?query=hello');

source.onmessage = (e) => {
  if (e.data === '[DONE]') {
    source.close();
    return;
  }
  const event = JSON.parse(e.data);

  switch (event.channel) {
    case 'run.lifecycle':
      updateExecutionStatus(event.data);
      break;
    case 'run.messages':
      appendToken(event.data);
      break;
    case 'run.values':
      updateStateSnapshot(event.data);
      break;
  }
};

前端类比:这就像在一个 WebSocket 连接中通过 event.type 分流处理--一条连接,多种事件类型,按需消费。


8. 调试与可观测性场景

场景 1:追踪执行卡在哪里

python
async for event in graph.astream_events(input, version="v3", channel="run.lifecycle"):
    print(f"[{event.get('ns', 'root')}] {event['data']}")
    # 如果某个节点开始后很久没有结束事件,说明它卡住了

场景 2:并行子图执行追踪

python
async for event in graph.astream_events(input, version="v3"):
    ns = event.get("ns", "")
    if event["channel"] == "run.lifecycle":
        print(f"[{ns}] {event['data']}")
    # 输出示例:
    # [root] run started
    # [root:fan_out_1] node started
    # [root:fan_out_2] node started   <- 并行子图
    # [root:fan_out_1] node ended
    # [root:fan_out_2] node ended
    # [root] run ended

场景 3:与 LangSmith 配合

事件流的 run.lifecycle 投影与 LangSmith 的 trace 结构对应。你可以在本地用事件流调试,在线上用 LangSmith 追踪,两者的执行视图一致。


9. v2 与 v3 的关键差异

如果你已经在用 v2 的 stream(stream_mode=..., version="v2"),以下是迁移到 v3 事件流时需要理解的差异:

维度v2 stream(version="v2")v3 stream_events(version="v3")
返回类型StreamPart dict(type/ns/data)类型化投影事件(channel/ns/data)
数据视角状态切片执行事件
子图处理StreamPart.ns 字段标识命名空间run.subgraphs 投影 + ns
生命周期run.lifecycle 投影
推荐场景状态/token 展示调试 / 可观测性 / UI 推送

不是替代关系:v2 stream() 和 v3 stream_events() 是互补的。状态展示用 v2,执行追踪用 v3。v1 格式已过时,不推荐新代码使用。


10. 不适用场景与能力边界

  • 简单 token 流:如果只需要前端打字机效果,stream(stream_mode="messages") 更简单直接,不需要事件流的额外复杂度。
  • 需要 updates 粒度:事件流没有直接等价于 stream_mode="updates"(节点增量更新)的投影。如果你需要"每个节点输出什么"的粒度,还是用 stream(stream_mode="updates")
  • beta 风险:v3 事件流仍是 beta,API 可能在后续版本调整。如果你对稳定性要求极高,先用 v2 stream(),等 v3 正式发布后迁移。
  • 性能开销:事件流比单一 stream_mode 产生更多事件。如果只需要一种数据,用 channel= 过滤避免不必要的事件处理。

要点回顾

概念一句话
事件流 vs stream事件流是运行时视角(发生了什么),stream 是数据视角(状态怎么变了)
stream_events(version="v3")v3 beta API,返回类型化投影事件
四种投影run.values / run.messages / run.lifecycle / run.subgraphs
run.lifecycle事件流独有,追踪执行生命周期,调试利器
channel 过滤channel= 只接收关心的事件类型
前端消费通过 SSE 推送,前端按 channel 分流处理

先修与下一步

参考

学习文档整合站点