流式输出与订阅
第 11 章我们用普通 API 实现了终端逐字输出。LangGraph 把"过程"也做成了流的:graph.stream() 既能按节点产出状态更新,也能把模型回复按 token 逐个吐出来。前者适合观察智能体每一步在做什么,后者适合做"打字机"体验或对接服务端推送。本章沿用最小对话图做演示。
准备工作:一个最小对话图
import os
from typing import Annotated, TypedDict
from langchain_core.messages import HumanMessage
from langchain_openai import ChatOpenAI
from langgraph.graph import StateGraph, START, END
from langgraph.graph.message import add_messages
model = ChatOpenAI(
model=os.getenv("LLM_MODEL", "deepseek-chat"),
api_key=os.getenv("LLM_API_KEY"),
base_url=os.getenv("LLM_BASE_URL"),
)
class State(TypedDict):
messages: Annotated[list, add_messages]
def agent(state):
reply = model.invoke(state["messages"])
return {"messages": [reply]}
builder = StateGraph(State)
builder.add_node("agent", agent)
builder.add_edge(START, "agent")
builder.add_edge("agent", END)
graph = builder.compile()
invoke 是"等全部跑完拿最终结果";stream 是"边跑边看"。两者输入相同,都可以配 checkpointer 与 thread_id。
stream_mode="updates":按节点输出
每个产出是一份字典 {节点名: 该节点返回的部分更新},能看清执行了哪些节点:
for update in graph.stream(
{"messages": [HumanMessage("一句话介绍大模型")]},
stream_mode="updates",
):
for node, value in update.items():
# 本例只有 agent 节点,value 形如 {'messages': [AIMessage(...)]}
print("运行节点:", node)
# 输出:运行节点: agent(出现几次就代表图走了几个节点/几轮)
多节点图(如第 18 章的 ReAct)会依次打出 agent、tools、agent……直到结束,非常适合观察循环轮次与分支走向。
stream_mode="messages":逐 token 输出
messages 模式产出 (消息块, 元信息) 二元组,消息块可用 .content 累积成完整回复:
full = ""
for chunk, meta in graph.stream(
{"messages": [HumanMessage("一句话介绍大模型")]},
stream_mode="messages",
):
piece = chunk.content
if piece:
print(piece, end="", flush=True) # 逐字打印,形成打字机效果
full += piece
# 输出:大模型是通过海量文本训练出来的……(文字逐个出现)
print("\n[完整回复]", full)
meta 里带 langgraph_node、langgraph_checkpoint_ns 等字段,可用于判断消息来自哪个节点、是否需要跳过工具消息。
组合两种模式
想同时拿到"节点轨迹"和"逐字内容",传一个模式列表即可:
for kind, data in graph.stream(
{"messages": [HumanMessage("一句话介绍大模型")]},
stream_mode=["updates", "messages"],
):
if kind == "updates":
print("[节点更新]", list(data.keys()))
else:
chunk = data[0]
if chunk.content:
print(chunk.content, end="", flush=True)
# 输出:同时出现 [节点更新] ['agent'] 与逐字打印的回复内容(两者先后由框架调度决定)
订阅与推送
- 异步服务端用 graph.astream(async 版本),在协程里逐块处理即可。
- Web 服务把每个 chunk 包成 SSE(Server-Sent Events)推给浏览器,即可实现流式对话界面;完整 FastAPI + SSE 实现将在第 28 章给出。思路一句话:一个生成器里 yield f"data: {文本块}\n\n",浏览器用 EventSource 消费。
三种模式怎么选
| 模式 | 产出 | 典型用途 |
|---|---|---|
| updates | 每执行一个节点给一份 {节点: 状态更新} | 观察执行轨迹、调试循环 |
| messages | 每个 token 一个消息块 | 打字机效果、SSE 推送 |
| updates + messages | 两类事件按类型混排 | 界面同时显示步骤与文字 |
常见问题
- 为什么 messages 模式会吐出空内容块?工具调用、消息元数据等 chunk 不一定带文本,用 if piece 跳过即可。
- stream 与 invoke 的最终结果一致吗?一致,stream 只是提前暴露中间过程;想核对终态可收集全部内容,或用第 20 章的 get_state 查询。
- 想区分文字来自哪个节点?读消息块元信息 meta 里的 langgraph_node 字段。
- 输出乱序或丢字?服务端推送建议按 SSE 事件逐块下发、客户端按序追加,避免并发连接各自拼接。
小结:graph.stream 按 stream_mode 输出不同粒度——"updates" 逐节点看执行轨迹,"messages" 逐 token 拿模型回复;两者可组合,服务端再配 astream + SSE 就能把打字机效果推给用户。