Appearance
事件流(Event Streaming)
普通的
stream()返回的是状态/消息流--你看到的是"状态怎么变的"。事件流返回的是执行事件的类型化投影--你看到的是"运行时发生了什么":哪个 run 开始了、哪个节点产出了值、哪个子图在执行。这让它成为可观测性、调试和实时 UI 推送的首选 API。
先修知识
- Streaming 流式处理 - 先理解
stream_mode的 7 种模式
必须深刻理解,不能跳过:
stream_mode和stream_events是两套不同的 API,解决不同的问题。stream_mode回答"状态的哪些切片要推送给我";stream_events回答"执行的哪些事件要推送给我"。前者是数据视角,后者是运行时视角。混用它们是新手最常犯的错误。
前端类比:先建立直觉
| 前端概念 | LangGraph 对应 | 说明 |
|---|---|---|
| Redux store subscribe | stream(stream_mode=...) | 订阅状态变化 |
| Redux middleware action log | stream_events(version="v3") | 订阅执行事件 |
| WebSocket 多频道 | channel 投影 (run.values 等) | 按频道过滤事件 |
| Chrome DevTools Performance | run.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 格式已过时,本页不展开。
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 |
|---|---|---|---|
| values | run.values | 完整状态快照 | stream_mode="values" |
| messages | run.messages | LLM token 级输出 | stream_mode="messages" |
| lifecycle | run.lifecycle | 执行生命周期事件(run 开始/结束、节点进入/退出) | 无直接对应 |
| subgraphs | run.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()和 v3stream_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 分流处理 |