LangGraph实时流式输出Agent思考过程:从原理到实战
1. 项目概述为什么我们需要“看见”Agent的思考如果你正在用LangGraph或者类似的Agent框架开发应用大概率遇到过这样的场景你向Agent提了一个稍微复杂点的问题比如“帮我分析一下上个月的销售数据并给出下个月的营销建议”。然后你就开始了漫长的等待。屏幕上可能只有一个旋转的加载图标或者干脆一片寂静。几十秒甚至几分钟后它终于吐出了一段完整的回答。这中间发生了什么Agent是卡住了还是在努力思考它先查了数据还是先规划了步骤如果结果不尽如人意你根本无从下手调试因为整个过程就像一个黑盒。这就是“5.9 式输出”要解决的问题。这个标题里的“5.9”更像是一个版本号或者一个内部梗它核心指向的是一种实时流式输出Agent内部状态和思考过程的能力。简单说就是让Agent的“大脑”活动可视化让它一边想你一边能看到。这不仅仅是酷炫对于开发者而言这是调试、优化和理解Agent行为的生命线。想象一下你能看到Agent决定调用哪个工具、调用的参数是什么、工具返回的结果、以及它基于结果如何进行下一步推理整个决策链一目了然。这比只看最终输出效率提升了不止一个数量级。从技术上看这涉及到几个核心概念LangGraph作为一个基于状态图的框架天然地将Agent的执行过程分解为清晰的节点Nodes和边Edges。Stream模式则是实现实时输出的关键技术它允许我们将执行过程中的中间状态State增量地推送给前端或日志系统。而invoke方法是触发整个图执行的入口。所以这个项目的本质就是打通从LangGraph的invoke调用到以stream方式获取并呈现每一步状态变化的完整链路。我自己的体会是没有这种可视化能力之前调试一个多步骤的Agent就像在蒙着眼睛修手表。加上实时输出后不仅开发效率飙升你还能更深刻地理解你设计的图是否按预期运转哪里出现了循环或卡顿工具调用是否准确。这对于构建可靠、可信的AI应用至关重要。2. 核心原理LangGraph的状态流与Streaming机制要搞懂如何实时输出思考过程我们必须先拆解LangGraph是怎么工作的。很多人会把LangGraph和LangChain搞混这里简单厘清LangChain更像是一个“工具箱”和“胶水”它提供了连接各种组件模型、工具、记忆的标准接口而LangGraph是一个“编排框架”它专注于用有向图来定义和控制这些组件之间的执行流程与状态流转。你可以用LangChain的组件来构建LangGraph的节点。2.1 LangGraph的三要素State, Node, EdgeLangGraph的核心抽象非常清晰就三样东西State一个字典Dict它是在整个图执行过程中传递和修改的共享数据。它定义了Agent的“记忆”和当前上下文。通常你会定义一个StateGraph并明确State的结构比如{“messages”: List[BaseMessage] “intermediate_steps”: List[Tuple[AgentAction AgentFinish]] ...}。Node节点一个可调用的函数。它接收当前的State执行一些操作比如调用LLM、使用工具、处理数据然后返回一个更新后的State。一个节点可以很简单也可以很复杂。Edge边决定执行流程。分为条件边Conditional Edge和普通边。条件边根据当前State的内容决定下一步执行哪个节点普通边则固定指向下一个节点。这赋予了图强大的分支和循环能力。当你调用compiled_graph.invoke(input)时LangGraph会根据你定义的图结构从入口节点开始依次执行节点沿着边流转直到到达某个终点。而思考过程就蕴含在每一次State的变迁中。2.2 Streaming的本质拦截状态变迁标准的invoke是一次性、同步的调用它返回最终的State。而我们要的实时输出需要的是stream模式。在LangGraph中compiled_graph.stream(input)返回的是一个异步生成器。它不会一次性返回结果而是会在图执行的每一个步骤step完成后立即yield出该步骤执行后的State。这里的“步骤”通常对应一个节点的执行。所以Streaming的底层原理就是在图引擎执行每个节点后立即将当前的State对象抛出来而不是等到全部执行完毕。这就像是在Agent思考的每一个“心跳”处安装了一个探头把它的脑电图实时地绘制出来。2.3stream与astream的异同你可能会看到两个方法stream和astream。它们的核心功能一致都是流式返回状态。stream同步方法。它内部处理了异步循环返回一个同步的生成器。在常规的脚本或Jupyter Notebook中使用更方便。astream异步方法。返回一个异步生成器。在FastAPI、异步Web应用等场景下使用能更好地融入异步生态避免阻塞事件循环。对于实时查看思考过程这个需求两者都能满足。选择哪一个取决于你的应用运行环境。在本文的示例中我们会主要使用stream因为它更直观。3. 实战构建一个具备思考过程输出的查询Agent光说不练假把式。我们一起来构建一个具体的Agent并实现它的思考过程流式输出。这个Agent的功能是回答关于当前天气和地点附近餐馆的混合查询。比如用户问“北京今天天气怎么样另外推荐一家附近的川菜馆。”这个需求涉及两个不同的工具调用天气查询和地点搜索完美展示了Agent的规划与执行过程。3.1 定义State与工具首先我们定义Agent的State。为了清晰展示我们让State包含对话消息、已执行的步骤以及一个可能存放最终答案的字段。from typing import TypedDict List Annotated Union from langgraph.graph.message import add_messages import operator # 定义State结构 class AgentState(TypedDict): # 对话消息历史使用LangGraph提供的注解实现自动累加 messages: Annotated[List add_messages] # 记录Agent已经决定要执行的动作思考过程的关键 agent_action: Union[dict None] # 记录工具执行的结果 tool_outputs: List[str] # 最终答案 final_answer: Union[str None] # 模拟工具1获取天气 def get_weather(location: str) - str: # 这里应该是调用真实天气API我们模拟返回 print(f“[工具调用] 正在查询 {location} 的天气...) return f{location}今天天气晴朗气温25度微风。 # 模拟工具2搜索附近餐馆 def search_restaurants(location: str cuisine: str) - str: print(f“[工具调用] 正在在 {location} 搜索 {cuisine} 餐馆...) return f在{location}找到三家不错的{cuisine}餐馆A店评分4.5 B店评分4.3 C店评分4.7。注意Annotated[List add_messages]是LangGraph的一个语法糖它自动帮你处理消息列表的追加非常方便。agent_action和tool_outputs是我们为了追踪思考过程特意添加的字段。3.2 构建LangGraph节点与边接下来我们构建图的核心部分节点函数和边。from langgraph.graph import StateGraph END from langchain_core.messages import HumanMessage AIMessage ToolMessage from langchain_openai import ChatOpenAI from langchain.tools import tool from langchain.agents import create_tool_calling_agent AgentExecutor from langchain_core.prompts import ChatPromptTemplate import json # 初始化LLM llm ChatOpenAI(model“gpt-4o-mini” temperature0) # 将函数包装成LangChain工具 tools [ tool(get_weather) # 工具名会自动生成为 get_weather tool(search_restaurants) # 工具名会自动生成为 search_restaurants ] # 创建Agent的Prompt prompt ChatPromptTemplate.from_messages([ (“system” “你是一个有帮助的助手可以查询天气和餐馆。请根据用户问题决定是否需要调用工具以及调用哪个工具。每次只思考一步。”) (“placeholder” “{messages}”) # 这里会被历史消息填充 ]) # 创建Agent (基于LangChain的Tool Calling Agent) agent create_tool_calling_agent(llm tools prompt) # 节点1Agent决策节点 def agent_node(state: AgentState) - AgentState: Agent接收消息决定下一步行动回复或调用工具 print(f“\n--- Agent思考节点 ---”) print(f“当前消息历史: {state[‘messages’]}”) # 调用LangChain Agent来获取决策 agent_response agent.invoke({“messages”: state[“messages”]}) # agent_response 可能是一个 AIMessage包含 tool_calls 或者直接是内容 ai_message agent_response[“messages”][-1] # 取最新的一条AI消息 new_state {**state} new_state[“messages”] state[“messages”] [ai_message] # 关键记录Agent的决定到状态中供后续查看 if hasattr(ai_message ‘tool_calls’) and ai_message.tool_calls: action { “tool_name”: ai_message.tool_calls[0][‘name’] “tool_args”: ai_message.tool_calls[0][‘args’] } print(f“Agent决定调用工具: {action}”) new_state[“agent_action”] action else: print(f“Agent决定直接回复: {ai_message.content}”) new_state[“agent_action”] {“type”: “direct_reply”} new_state[“final_answer”] ai_message.content # 如果是最终回复记录下来 return new_state # 节点2工具执行节点 def tool_node(state: AgentState) - AgentState: 执行Agent决定的工具调用 print(f“\n--- 工具执行节点 ---”) action state[“agent_action”] if not action or action.get(“type”) “direct_reply”: # 如果没有待执行的动作直接返回原状态 return state tool_name action[“tool_name”] tool_args action[“tool_args”] # 根据工具名找到对应的函数并执行 tool_map {t.name: t for t in tools} if tool_name in tool_map: tool_to_call tool_map[tool_name] result tool_to_call.invoke(tool_args) print(f“工具 {tool_name} 执行结果: {result}”) else: result f“错误未找到工具 {tool_name}” print(result) # 将结果封装成ToolMessage并添加到消息历史 tool_message ToolMessage(contentstr(result) tool_call_id“fake_id”) new_state {**state} new_state[“messages”] state[“messages”] [tool_message] new_state[“tool_outputs”] state.get(“tool_outputs” []) [f“{tool_name}: {result}”] new_state[“agent_action”] None # 清空当前动作 return new_state # 构建图 workflow StateGraph(AgentState) # 添加节点 workflow.add_node(“agent” agent_node) workflow.add_node(“tool” tool_node) # 设置入口点 workflow.set_entry_point(“agent”) # 添加边和条件边 from langgraph.graph import START def route_after_agent(state: AgentState): 根据Agent的输出来决定下一步调用工具还是结束 last_message state[“messages”][-1] # 如果最新的AI消息包含工具调用就去执行工具 if hasattr(last_message ‘tool_calls’) and last_message.tool_calls: return “tool” # 否则结束流程 else: return END def route_after_tool(state: AgentState): 工具执行完后总是回到Agent进行下一步思考 return “agent” workflow.add_conditional_edges( “agent” route_after_agent {“tool”: “tool” END: END} ) workflow.add_edge(“tool” “agent”) # 编译图 compiled_graph workflow.compile()这个图形成了一个经典的Agent - (可能) Tool - Agent - ... - END的循环结构直到Agent不再调用工具直接给出最终答案。3.3 实现并解析Streaming输出现在最激动人心的部分来了实时流式输出。# 用户输入 user_input “北京今天天气怎么样另外推荐一家附近的川菜馆。” initial_state { “messages”: [HumanMessage(contentuser_input)] “agent_action”: None “tool_outputs”: [] “final_answer”: None } print(“ 开始流式执行Agent ) print(“用户问题:” user_input) print(“-” * 50) # 使用 stream 方法 for step output_state in enumerate(compiled_graph.stream(initial_state) 1): print(f“\n[步骤 {step}] 状态更新:”) # 打印出我们关心的部分 last_message output_state[“messages”][-1] if isinstance(last_message AIMessage): if hasattr(last_message ‘tool_calls’) and last_message.tool_calls: tc last_message.tool_calls[0] print(f“ Agent思考: 决定调用工具 {tc[‘name’]} 参数: {tc[‘args’]}”) else: print(f“ Agent回复: {last_message.content}”) elif isinstance(last_message ToolMessage): print(f“ ⚙️ 工具返回: {last_message.content}”) # 也可以打印整个状态的特定字段 if output_state.get(“agent_action”): print(f“ 记录的动作: {output_state[‘agent_action’]}”) if output_state.get(“tool_outputs”): print(f“ 累计工具输出: {output_state[‘tool_outputs’]}”) print(“-” * 50) print(“ 执行结束 ) print(“最终答案:” initial_state.get(“final_answer”) or output_state[“messages”][-1].content)运行这段代码你将在控制台看到类似下面的输出这就是Agent的“思考过程” 开始流式执行Agent 用户问题: 北京今天天气怎么样另外推荐一家附近的川菜馆。 -------------------------------------------------- --- Agent思考节点 --- 当前消息历史: [HumanMessage(content‘北京今天天气怎么样另外推荐一家附近的川菜馆。’ additional_kwargs{} response_metadata{})] [步骤 1] 状态更新: Agent思考: 决定调用工具 get_weather 参数: {‘location’: ‘北京’} 记录的动作: {‘tool_name’: ‘get_weather’ ‘tool_args’: {‘location’: ‘北京’}} --- 工具执行节点 --- [工具调用] 正在查询 北京 的天气... [步骤 2] 状态更新: ⚙️ 工具返回: 北京今天天气晴朗气温25度微风。 累计工具输出: [‘get_weather: 北京今天天气晴朗气温25度微风。’] --- Agent思考节点 --- 当前消息历史: [HumanMessage(...) AIMessage(content‘’ additional_kwargs{‘tool_calls’: [...]} response_metadata{}) ToolMessage(content‘北京今天天气晴朗气温25度微风。’ tool_call_id‘fake_id’ additional_kwargs{} response_metadata{})] [步骤 3] 状态更新: Agent思考: 决定调用工具 search_restaurants 参数: {‘location’: ‘北京’ ‘cuisine’: ‘川菜’} 记录的动作: {‘tool_name’: ‘search_restaurants’ ‘tool_args’: {‘location’: ‘北京’ ‘cuisine’: ‘川菜’}} --- 工具执行节点 --- [工具调用] 正在在 北京 搜索 川菜 餐馆... [步骤 4] 状态更新: ⚙️ 工具返回: 在北京找到三家不错的川菜餐馆A店评分4.5 B店评分4.3 C店评分4.7。 累计工具输出: [‘get_weather: 北京今天天气晴朗气温25度微风。’ ‘search_restaurants: 在北京找到三家不错的川菜餐馆A店评分4.5 B店评分4.3 C店评分4.7。’] --- Agent思考节点 --- 当前消息历史: [...] # 消息历史持续增长 [步骤 5] 状态更新: Agent回复: 北京今天天气晴朗气温25度微风。另外为您找到三家在北京不错的川菜馆A店评分4.5 B店评分4.3 C店评分4.7。您可以根据评分选择。 记录的动作: {‘type’: ‘direct_reply’} -------------------------------------------------- 执行结束 最终答案: 北京今天天气晴朗气温25度微风。另外为您找到三家在北京不错的川菜馆A店评分4.5 B店评分4.3 C店评分4.7。您可以根据评分选择。看整个过程清晰可见我们看到了Agent的两次“思考-决策”过程步骤1和3以及对应的两次工具执行步骤2和4最后一步步骤5是汇总信息的最终回复。这比只看到一个最终答案要有价值得多。4. 高级技巧与生产级优化上面的示例展示了基本原理但在实际生产环境中你需要考虑更多。下面分享几个我踩过坑后总结的进阶技巧。4.1 结构化输出与前端展示控制台打印只是第一步。在生产中你需要将Streaming事件发送到前端如WebSocket或结构化日志系统如JSON Lines。关键在于对State进行序列化和过滤。import json def stream_agent_for_api(question: str): 一个生成器用于API流式响应 initial_state { “messages”: [HumanMessage(contentquestion)] “agent_action”: None “tool_outputs”: [] “final_answer”: None } for step state in enumerate(compiled_graph.stream(initial_state) 1): # 1. 提取关键信息而不是发送整个State可能包含敏感信息或过大 event {“step”: step “type”: “state_update”} last_msg state[“messages”][-1] if isinstance(last_msg AIMessage): if hasattr(last_msg ‘tool_calls’): event[“subtype”] “agent_decision” event[“data”] { “action”: “call_tool” “tool”: last_msg.tool_calls[0][‘name’] “args”: last_msg.tool_calls[0][‘args’] } else: event[“subtype”] “agent_reply” event[“data”] {“content”: last_msg.content} elif isinstance(last_msg ToolMessage): event[“subtype”] “tool_result” # 可以对长的工具结果进行截断或摘要 event[“data”] {“result”: last_msg.content[:200] “...” if len(last_msg.content) 200 else last_msg.content} # 2. 转换为JSON字符串通过SSE或WebSocket发送 yield f“data: {json.dumps(event ensure_asciiFalse)}\n\n” # 最终结束事件 yield f“data: {json.dumps({‘type’: ‘stream_end’})}\n\n”前端可以通过EventSource API来接收这些事件并实时更新UI展示一个动态的思考链。4.2 处理“Stream Disconnected”等网络错误在热词列表中频繁出现stream disconnected before completion错误。这通常发生在长耗时任务中客户端浏览器或网络中间件如Nginx因为超时设置而断开了连接。解决方案保持连接活跃在流式响应中定期发送心跳包如注释行或空事件。# 在长时间处理节点中模拟 def slow_tool_node(state): # ... 执行耗时操作 time.sleep(30) # 模拟长时间运行 # 在真实场景中你需要在循环中定期 yield 心跳 # 对于HTTP流框架如FastAPI的中间件可能帮你处理但需要配置超时时间。 return state配置服务器超时如果你的后端是FastAPI、Flask等务必调整相关超时设置。FastAPI: 使用Starlette的Response或StreamingResponse并确保你的ASGI服务器如Uvicorn的timeout_keep_alive等参数设置得足够大。Nginx: 调整proxy_read_timeoutproxy_send_timeout等指令将其设置为一个较大的值如300秒。客户端重连机制前端实现自动重连逻辑。当EventSource触发onerror时尝试重新连接并从断点恢复如果服务端支持的话。使用更可靠的传输协议对于极其复杂的Agent考虑使用WebSocket替代Server-Sent Events (SSE)因为WebSocket是双向、全双工的对连接状态的控制力更强。4.3 性能优化与状态管理当图变得复杂、State很大时Streaming每一个完整State可能会成为性能瓶颈。选择性Streaming不要yield整个State。像上面的例子一样只提取和发送前端需要展示的关键字段。使用检查点CheckpointLangGraph支持检查点功能可以将State持久化。这对于需要暂停/恢复、或长时间运行的Agent至关重要。在Streaming时你可以选择只发送检查点ID或状态摘要。异步优化确保你的节点函数特别是那些调用外部API的是异步的使用async def并使用astream这样可以避免阻塞提高整体吞吐量。4.4 调试与日志集成将Streaming输出与你的日志系统如Loguru structlog集成可以方便地追踪线上Agent的行为。import logging from langgraph.checkpoint import MemorySaver # 配置一个专门的logger agent_logger logging.getLogger(“agent_stream”) agent_logger.setLevel(logging.INFO) # 在stream循环中记录 for step state in enumerate(compiled_graph.stream(initial_state) 1): # ... 提取事件信息 ... event_info {“step”: step “action”: extracted_action} agent_logger.info(json.dumps(event_info)) # 结构化日志便于ELK等系统收集分析 # ... 发送到前端 ...这样你不仅能在前端看到实时思考过程还能在日志平台按会话ID检索完整的执行轨迹对于排查线上问题无比方便。5. 常见问题与排查实录在实际操作中我遇到了不少坑。这里列一个速查表希望能帮你节省时间。问题现象可能原因解决方案调用graph.stream()没有任何输出1. 图没有正确编译或入口点设置错误。2.stream返回的生成器没有被迭代例如忘记写for ... in循环。3. 节点函数修改了State但没有返回新的State。1. 检查workflow.compile()是否成功并用print(compiled_graph.get_graph().draw_mermaid())可视化图结构。2. 确保代码中有for step in compiled_graph.stream(...):。3. 每个节点函数必须返回一个字典这个字典会被合并到总State中。stream disconnected before completion错误1. 网络连接超时最常见。2. 代理服务器如Nginx或负载均衡器超时设置过短。3. 某个节点执行时间过长客户端主动断开。1. 前端增加心跳/重连机制。2. 调整服务器配置Nginx的proxy_read_timeout。3. 优化节点逻辑或将长任务拆分为多个步骤分次yield进度。Streaming输出顺序混乱或重复在异步(astream)环境下如果处理不当事件的顺序可能错乱。确保你的处理逻辑是线程安全/协程安全的。对于关键顺序可以在State或事件中加入递增的序列号sequence_id。State过大导致传输缓慢或内存溢出在循环中不断向State添加大量数据如完整的对话历史。1. 使用Annotated的缩减器或定期清理State中的历史数据。2. 实现外部记忆体State中只存放索引或摘要。3. Streaming时只发送增量变化而非全量State。工具调用结果无法影响后续决策工具节点的结果没有正确地以ToolMessage格式添加回messages中。确保工具节点返回的State里messages列表末尾添加的是ToolMessage对象。LangChain Agent默认会读取这个消息类型。无法进入条件边或一直在循环条件边add_conditional_edges的判断函数逻辑有误始终返回同一个节点名。在判断函数中打印日志仔细检查输入State的结构和内容是否符合预期。使用graph.get_graph().draw_mermaid()检查边是否正确连接。一个独家避坑技巧在开发初期不要急于对接华丽的前端。先用最朴素的脚本把compiled_graph.stream()的每一个产出都print出来甚至把完整的State用json.dumps(state indent2 defaultstr)打印出来。这能帮你最直观地理解State是如何一步步变化的很多逻辑错误在这个过程中就暴露无遗了。可视化思考过程的前提是你自己先要能“看见”全部数据。