Skip to content

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 后的 state

updates 模式:增量流

只输出每个节点产生的 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": "正在处理..."}

生产实践

  1. 开发用 values生产用 messages(轮播 Token)
  2. 处理中断信号,让用户能停止 streaming
  3. 注意内存:长时间 streaming 不要保存全部 Token
  4. 异步 streaming(astream)适合高并发场景

参考

本站为非官方中文学习站点,不代表 LangChain 官方。部分内容参考官方文档并重新整理为中文学习笔记。