Streaming 流式事件
LangGraph 提供丰富的流式(Streaming)支持,允许在 Agent 执行过程中实时获取 LLM Token、节点输出、状态变更等事件。这对构建响应式的聊天界面至关重要。
为什么需要 Streaming
- 即时反馈:用户看到打字效果,体验更好
- 实时调试:观察每个节点的执行过程
- 增量处理:在完整输出到达前即可开始处理
Streaming 模式
LangGraph 支持多种流式模式:
| 模式 | 产出 | 适用场景 |
|---|---|---|
stream_mode="values" | 每步的完整 state | 调试、状态跟踪 |
stream_mode="updates" | 每步的 state 增量 | 追踪节点输出 |
stream_mode="messages" | 流式消息 Token | 聊天界面逐字显示 |
stream_mode="events" | 所有事件的完整流 | 高级自定义 |
stream_mode="custom" | 自定义事件 | 特定业务场景 |
stream_mode="debug" | 调试事件 | 开发调试 |
values 模式:完整状态流
每次 state 变更时,输出当前的完整 state:
python
from langgraph.graph import StateGraph, START, END
builder = StateGraph(State)
builder.add_node("step1", step1)
builder.add_node("step2", step2)
builder.add_edge(START, "step1")
builder.add_edge("step1", "step2")
builder.add_edge("step2", END)
graph = builder.compile()
for event in graph.stream({"input": "hello"}, stream_mode="values"):
print(event)
# 输出:
# {'input': 'hello'} # 输入时的 state
# {'input': 'hello', 'result': 'step1 done'} # step1 后的 state
# {'input': 'hello', 'result': 'step2 done'} # step2 后的 stateupdates 模式:增量流
只输出每个节点产生的 state 变更,不输出完整的 state:
python
for event in graph.stream({"input": "hello"}, stream_mode="updates"):
print(event)
# 输出:
# {'step1': {'result': 'step1 done'}} # 只有 step1 的变化
# {'step2': {'result': 'step2 done'}} # 只有 step2 的变化messages 模式:Token 级流
逐 token 输出 LLM 的响应,适合聊天 UI:
python
from langchain.chat_models import init_chat_model
model = init_chat_model("gpt-4")
def llm_node(state: MessagesState):
response = model.invoke(state["messages"])
return {"messages": [response]}
# 使用 stream_mode="messages"
async for event, metadata in graph.astream(
{"messages": [{"role": "user", "content": "讲个故事"}]},
stream_mode="messages"
):
if event.type == "ai":
print(event.content, end="", flush=True)
# 逐字输出 LLM 回复messages 模式的事件类型
python
async for event, metadata in graph.astream(inputs, stream_mode="messages"):
# event 的类型
print(type(event)) # AIMessageChunk, HumanMessage, ToolMessage
# metadata 中包含
print(metadata.keys())
# langgraph_node, langgraph_triggers, langgraph_task_idx, ...结合多种模式
LangGraph 允许同时订阅多种流模式:
python
for mode, event in graph.stream(
inputs,
stream_mode=["values", "messages"]
):
if mode == "values":
print(f"[STATE]: {event}")
elif mode == "messages":
print(f"[TOKEN]: {event.content}")异步 Streaming
python
async for event in graph.astream(inputs, stream_mode="updates"):
for node_name, update in event.items():
print(f"Node {node_name}: {update}")自定义 Streaming 事件
节点可以推送自定义事件:
python
from langgraph.types import emit_custom_event
def node_with_custom_events(state: State):
# 推送进度信息
emit_custom_event("progress", {"percent": 30, "message": "正在处理..."})
result = process(state["input"])
emit_custom_event("progress", {"percent": 100, "message": "完成"})
return {"result": result}
# 消费自定义事件
for event in graph.stream(inputs, stream_mode="custom"):
print(event) # {"percent": 30, "message": "正在处理..."}生产实践
- 开发用 values,生产用 messages(轮播 Token)
- 处理中断信号,让用户能停止 streaming
- 注意内存:长时间 streaming 不要保存全部 Token
- 异步 streaming(
astream)适合高并发场景