LangGraph-Stream流式输出与实时反馈
Agent 后端逻辑做得再好,如果没有实时反馈,用户感知到的就是"这个系统很慢、很卡"。用户发了一个问题,然后面对空白界面等了两分钟,没有任何进度提示,不知道程序是在跑还是卡住了——这是很多 Agent 应用的真实状态。
LangGraph Stream:流式输出与实时反馈
Agent 后端逻辑做得再好,如果没有实时反馈,用户感知到的就是"这个系统很慢、很卡"。用户发了一个问题,然后面对空白界面等了两分钟,没有任何进度提示,不知道程序是在跑还是卡住了——这是很多 Agent 应用的真实状态。
LangGraph 提供了完整的流式输出机制,可以在 Agent 执行过程中实时推送进度——哪个节点在执行、LLM 输出了什么字符、工具调用的结果是什么。
1.1 LangGraph 的四种流式模式
LangGraph 四种 stream_mode——values(全量)、updates(增量)、messages(token 流)、events(事件)
LangGraph 的 stream 方法支持通过 stream_mode 参数切换不同的输出粒度,针对不同场景有不同的选择。
1.1.1 模式一:`values`——每次输出完整 State
每个节点执行完之后,输出当前完整的 State 快照。
for event in app.stream(initial_state, stream_mode="values"):
print(event) # 输出的是整个 State 字典
特点:每次事件都是完整的 State,字段多的时候数据量很大。适合需要实时看到完整 State 变化的调试场景,不适合传给前端。
1.1.2 模式二:`updates`——每次只输出变化的部分
每个节点执行完之后,只输出这个节点改变了哪些字段。
for event in app.stream(initial_state, stream_mode="updates"):
# event 是 {"节点名称": {"字段名": "新值", ...}} 格式
for node_name, updates in event.items():
print(f"节点 [{node_name}] 更新了:{list(updates.keys())}")
特点:数据量小,能清楚看到每个节点产出了什么。这是最常用的模式,推荐作为默认选择。
1.1.3 模式三:`messages`——逐 token 输出 LLM 的响应
当节点内部调用了 LLM,可以把 LLM 生成的 token 流实时推出来,实现打字机效果。
for chunk, metadata in app.stream(initial_state, stream_mode="messages"):
if hasattr(chunk, "content") and chunk.content:
print(chunk.content, end="", flush=True) # 逐字打印
特点:最接近用户体验的模式,让用户感觉 AI 在"实时思考和回答"。适合对话类场景。
1.1.4 模式四:`astream_events`——最细粒度的事件流
异步方法,输出所有内部事件:节点开始、节点结束、LLM token、工具调用开始、工具调用结束等。
async for event in app.astream_events(initial_state, version="v2"):
event_type = event["event"]
if event_type == "on_chat_model_stream":
# LLM 正在流式输出
chunk = event["data"]["chunk"]
print(chunk.content, end="", flush=True)
elif event_type == "on_chain_start":
# 节点开始执行
print(f"\n>>> 节点 [{event['name']}] 开始执行")
特点:粒度最细,可以监控到每一个内部事件。代价是数据量大,处理逻辑也最复杂。适合需要精细监控的场景,比如在界面上显示工具调用的实时状态。
1.2 完整示例:四种模式的实际使用
下面是一个完整可运行的示例,包含一个多步骤 Agent 和四种 stream 模式的演示。
import asyncio
from typing import TypedDict, List, Annotated
import operator
from langgraph.graph import StateGraph, START, END
try:
from langchain_openai import ChatOpenAI
HAS_LLM = True
except ImportError:
HAS_LLM = False
# ============================================================
# 1. 定义 State
# ============================================================
class ResearchState(TypedDict):
question: str
search_results: List[str]
analysis: str
final_answer: str
# 使用 Annotated + operator.add 让列表字段支持追加而非覆盖
logs: Annotated[List[str], operator.add]
# ============================================================
# 2. 定义节点(不依赖真实 LLM,全部模拟)
# ============================================================
def search_node(state: ResearchState) -> dict:
"""搜索节点:模拟网络搜索"""
import time
time.sleep(0.5) # 模拟耗时操作
question = state["question"]
results = [
f"搜索结果1:关于「{question}」的基础知识介绍",
f"搜索结果2:「{question}」的实际应用案例",
f"搜索结果3:「{question}」的最新研究进展",
]
return {
"search_results": results,
"logs": [f"搜索完成,找到 {len(results)} 条结果"],
}
def analyze_node(state: ResearchState) -> dict:
"""分析节点:对搜索结果进行分析"""
import time
time.sleep(0.3)
result_count = len(state["search_results"])
analysis = (
f"基于 {result_count} 条搜索结果的分析:\n"
f"主要发现了三个维度的信息——基础概念、实际应用和前沿动态。"
)
return {
"analysis": analysis,
"logs": ["分析完成"],
}
def answer_node(state: ResearchState) -> dict:
"""回答节点:生成最终答案"""
import time
time.sleep(0.2)
answer = (
f"针对您的问题「{state['question']}」,综合分析如下:\n\n"
f"{state['analysis']}\n\n"
f"如需了解更多细节,建议参考搜索结果中的具体资料。"
)
return {
"final_answer": answer,
"logs": ["最终答案生成完成"],
}
# ============================================================
# 3. 构建图
# ============================================================
def build_agent():
graph = StateGraph(ResearchState)
graph.add_node("search", search_node)
graph.add_node("analyze", analyze_node)
graph.add_node("answer", answer_node)
graph.add_edge(START, "search")
graph.add_edge("search", "analyze")
graph.add_edge("analyze", "answer")
graph.add_edge("answer", END)
return graph.compile()
# ============================================================
# 4. 演示 stream_mode="updates"(推荐,最常用)
# ============================================================
def demo_stream_updates():
print("\n" + "=" * 60)
print("演示 stream_mode='updates':只输出每个节点的变化部分")
print("=" * 60)
app = build_agent()
initial_state = {
"question": "LangGraph 流式输出如何工作",
"search_results": [],
"analysis": "",
"final_answer": "",
"logs": [],
}
for event in app.stream(initial_state, stream_mode="updates"):
for node_name, updates in event.items():
# 过滤掉内部节点(以 __ 开头)
if node_name.startswith("__"):
continue
updated_fields = list(updates.keys())
print(f"\n[{node_name}] 更新字段:{updated_fields}")
# 显示关键字段的内容
if "logs" in updates:
for log in updates["logs"]:
print(f" 日志:{log}")
if "final_answer" in updates:
print(f" 最终答案(前50字):{updates['final_answer'][:50]}...")
print("\n[完成] 所有节点执行完毕")
# ============================================================
# 5. 演示 stream_mode="values"(适合调试)
# ============================================================
def demo_stream_values():
print("\n" + "=" * 60)
print("演示 stream_mode='values':每次输出完整 State")
print("=" * 60)
app = build_agent()
initial_state = {
"question": "LangGraph 流式输出如何工作",
"search_results": [],
"analysis": "",
"final_answer": "",
"logs": [],
}
for i, state_snapshot in enumerate(app.stream(initial_state, stream_mode="values")):
non_empty_fields = {
k: v for k, v in state_snapshot.items()
if v and k != "question"
}
print(f"\n快照 {i + 1},已有数据的字段:{list(non_empty_fields.keys())}")
# ============================================================
# 6. 演示打字机效果(需要真实 LLM)
# ============================================================
def demo_typewriter_effect():
"""
演示用 stream_mode="messages" 实现打字机效果
需要配置 OPENAI_API_KEY 环境变量
"""
if not HAS_LLM:
print("\n[跳过] 未安装 langchain_openai,跳过打字机效果演示")
print("安装方式:pip install langchain-openai")
return
print("\n" + "=" * 60)
print("演示 stream_mode='messages':打字机效果")
print("=" * 60)
from langchain_openai import ChatOpenAI
from langchain_core.messages import HumanMessage, AIMessage
from typing import Annotated
import operator
class ChatState(TypedDict):
messages: Annotated[List, operator.add]
def llm_node(state: ChatState) -> dict:
llm = ChatOpenAI(model="gpt-4o-mini", streaming=True)
response = llm.invoke(state["messages"])
return {"messages": [response]}
chat_graph = StateGraph(ChatState)
chat_graph.add_node("llm", llm_node)
chat_graph.add_edge(START, "llm")
chat_graph.add_edge("llm", END)
chat_app = chat_graph.compile()
initial = {
"messages": [HumanMessage(content="用三句话介绍 LangGraph")]
}
print("\nAI 回复(打字机效果):")
for chunk, metadata in chat_app.stream(initial, stream_mode="messages"):
if hasattr(chunk, "content") and chunk.content:
print(chunk.content, end="", flush=True)
print()
# ============================================================
# 7. 只关心特定节点的输出
# ============================================================
def demo_filter_nodes():
print("\n" + "=" * 60)
print("演示:只关心 answer 节点的输出,忽略其他节点")
print("=" * 60)
app = build_agent()
initial_state = {
"question": "如何过滤 LangGraph 的节点输出",
"search_results": [],
"analysis": "",
"final_answer": "",
"logs": [],
}
NODES_TO_WATCH = {"answer"} # 只关心这些节点
for event in app.stream(initial_state, stream_mode="updates"):
for node_name, updates in event.items():
if node_name not in NODES_TO_WATCH:
continue # 跳过不关心的节点
print(f"[{node_name}] 输出:")
for field, value in updates.items():
if field != "logs":
print(f" {field}: {str(value)[:100]}")
if __name__ == "__main__":
demo_stream_updates()
demo_stream_values()
demo_filter_nodes()
demo_typewriter_effect()
1.3 结合 FastAPI 提供流式 API
光在后端打印没有意义,实际场景需要把流式事件推给前端。FastAPI(Python 的高性能 Web 框架,专门用于快速构建 HTTP API 接口)的 StreamingResponse 配合 Server-Sent Events(SSE,服务器推送事件:一种服务器主动向浏览器实时推送数据的技术,前端无需轮询,适合进度展示)是最常见的方案。
import asyncio
import json
from typing import AsyncGenerator
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
from langgraph.graph import StateGraph, START, END
from typing import TypedDict, List, Annotated
import operator
app_fastapi = FastAPI()
class ResearchState(TypedDict):
question: str
search_results: List[str]
analysis: str
final_answer: str
logs: Annotated[List[str], operator.add]
def search_node(state: ResearchState) -> dict:
import time
time.sleep(0.5)
return {
"search_results": [f"搜索结果:关于「{state['question']}」"],
"logs": ["搜索完成"],
}
def analyze_node(state: ResearchState) -> dict:
import time
time.sleep(0.3)
return {
"analysis": f"分析完成:基于搜索结果的综合分析",
"logs": ["分析完成"],
}
def answer_node(state: ResearchState) -> dict:
return {
"final_answer": f"针对「{state['question']}」的最终回答:{state['analysis']}",
"logs": ["回答生成完成"],
}
def build_langgraph_agent():
graph = StateGraph(ResearchState)
graph.add_node("search", search_node)
graph.add_node("analyze", analyze_node)
graph.add_node("answer", answer_node)
graph.add_edge(START, "search")
graph.add_edge("search", "analyze")
graph.add_edge("analyze", "answer")
graph.add_edge("answer", END)
return graph.compile()
agent = build_langgraph_agent()
async def stream_agent_events(question: str) -> AsyncGenerator[str, None]:
"""
把 LangGraph 的 stream 事件转换成 SSE 格式
SSE 的格式要求:每条消息以 "data: " 开头,以 "\n\n" 结尾
"""
initial_state = {
"question": question,
"search_results": [],
"analysis": "",
"final_answer": "",
"logs": [],
}
# 在线程池里运行同步的 stream,避免阻塞事件循环
import concurrent.futures
loop = asyncio.get_event_loop()
def run_stream():
return list(agent.stream(initial_state, stream_mode="updates"))
with concurrent.futures.ThreadPoolExecutor() as pool:
events = await loop.run_in_executor(pool, run_stream)
for event in events:
for node_name, updates in event.items():
if node_name.startswith("__"):
continue
# 构造要发送给前端的事件数据
payload = {
"type": "node_update",
"node": node_name,
"fields": list(updates.keys()),
"logs": updates.get("logs", []),
}
# 如果是最终答案节点,把答案也带上
if "final_answer" in updates:
payload["final_answer"] = updates["final_answer"]
# SSE 格式
yield f"data: {json.dumps(payload, ensure_ascii=False)}\n\n"
# 发送结束信号
yield f"data: {json.dumps({'type': 'done'}, ensure_ascii=False)}\n\n"
@app_fastapi.get("/api/research/stream")
async def research_stream(question: str):
"""
流式 API 端点,返回 Server-Sent Events 格式的实时进度
前端消费示例:
const source = new EventSource('/api/research/stream?question=xxx');
source.onmessage = (e) => {
const data = JSON.parse(e.data);
if (data.type === 'node_update') {
console.log(`节点 ${data.node} 执行完成`);
} else if (data.type === 'done') {
source.close();
}
};
"""
return StreamingResponse(
stream_agent_events(question),
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache",
"X-Accel-Buffering": "no", # 禁用 Nginx 的响应缓冲
}
)
# 运行:uvicorn <filename>:app_fastapi --reload
前端收到的事件流格式:
data: {"type": "node_update", "node": "search", "fields": ["search_results", "logs"], "logs": ["搜索完成"]}
data: {"type": "node_update", "node": "analyze", "fields": ["analysis", "logs"], "logs": ["分析完成"]}
data: {"type": "node_update", "node": "answer", "fields": ["final_answer", "logs"], "final_answer": "...", "logs": ["回答生成完成"]}
data: {"type": "done"}
前端根据 node 字段更新进度条,根据 final_answer 展示最终结果,根据 type === 'done' 关闭连接。整个交互是实时的,用户看到的不是空白等待,而是"正在搜索 → 正在分析 → 正在生成答案"的逐步推进。
1.4 流式输出的数据流
1.5 不同场景用哪种模式
Agent 进度展示("正在执行第 X 步"):用 stream_mode="updates"。每个节点完成后推一条事件,前端更新进度。数据量小,逻辑简单。
对话类应用(用户看到 AI 打字回复):用 stream_mode="messages"。对接 LLM 的 token 流,实现逐字输出效果。注意需要 LLM 本身支持 streaming。
调试和监控(想看到所有细节):用 astream_events。可以看到工具调用的输入输出、LLM 的原始请求和响应、每个节点的执行时长。代价是处理逻辑较复杂。
开发阶段快速验证:用 stream_mode="values" 或直接用 invoke。不需要实时性的时候,invoke 代码最简单。
一个常见的误区是所有场景都用 astream_events,认为信息最全。但全不代表好——前端处理大量事件也有开销,而且大多数事件前端其实不需要。根据场景选择合适的粒度,才是正确的做法。
1.6 小结
LangGraph 的流式接口设计得很干净,几种模式各有侧重,选对了直接用。
难点不在技术,在于判断:后端推什么、前端展示什么。从用户视角想一遍——等这几分钟,用户需要看到什么才不会以为程序卡死了。
到这里,LangGraph 的核心机制:状态机、条件分支、人机协作、子图模块化、持久化、流式输出,都覆盖了。下一篇讲 Multi-Agent。