1. 项目概述为什么我们需要关注LangGraph的流式输出最近在折腾一个AI应用的后端服务其中一个核心场景是处理需要多步骤推理的复杂任务比如一个智能客服需要先理解用户意图、再查询知识库、最后生成回答。这种场景下传统的“请求-等待-响应”模式用户体验很差用户看着空白的界面干等十几秒心里肯定在骂娘。于是流式输出Streaming Output就成了刚需它能让AI“边想边说”用户能实时看到思考过程和部分结果体验流畅度直接拉满。在LangChain生态里LangGraph凭借其强大的有状态、多步骤工作流编排能力迅速成为了处理这类复杂链式或图式AI任务的首选框架。但是当我把一个在普通LangChain链上跑得挺顺的流式接口迁移到LangGraph上时问题就来了输出卡顿了或者流出的内容格式不对又或者子图Subgraph的流式响应根本出不来。这促使我专门进行了一次深入的“LangGraph流式输出特性测试”。这次测试的目标很明确不是简单地跑通一个Demo而是要摸清在真实生产环境中如何稳定、高效、可控地驾驭LangGraph的流式能力尤其是那些官方文档可能一笔带过但实际开发中一定会踩到的坑。简单来说如果你也在用或打算用LangGraph来构建需要实时反馈的AI应用比如对话机器人、自动化报告生成、交互式数据分析工具那么关于流式输出的这些细节——从基础的astream调用到复杂的子图流式传播再到与Spring Boot、FastAPI等Web框架的对接——就是你迟早要面对和解决的问题。接下来我就把这次测试中梳理出的思路、方案、代码和踩过的坑毫无保留地分享出来。2. 核心概念与工具选型解析在深入代码之前我们得先统一一下认知理解几个关键概念和为什么选它们。2.1 LangGraph 与流式输出Streaming的本质LangGraph可以看作是对LangChain的增强它引入了“图”Graph和“状态”State的概念。一个工作流被定义为由节点Node和边Edge组成的图每个节点执行特定功能如调用LLM、查询数据库而状态对象则在节点间传递保存着整个工作流的上下文信息。其核心魅力在于支持循环Loop和条件分支非常适合多轮对话、迭代式任务。流式输出在LLM语境下通常指的是服务器端一边生成Token文本块一边就通过网络发送给客户端而不是等全部生成完毕再一次性返回。对于LangGraph流式输出有了更丰富的内涵节点级流式单个节点尤其是调用LLM的节点可以流式输出其生成的内容。图级流式整个图的执行过程可以被流式化你可以实时看到执行跳转到哪个节点、每个节点的输入输出是什么、状态如何变化。这对于调试和用户展示“思考过程”至关重要。最终结果流式用户最关心的通常是最终答案的流式生成。LangGraph提供了astream、astream_log、astream_events等多个异步流式方法它们返回的都是异步迭代器Async Iterator这是我们实现流式响应的基础。2.2 关键工具与版本说明本次测试基于以下环境不同的版本可能在API细节上有差异Python: 3.10LangChain: 0.1.0LangGraph: 0.0.50HTTP框架: FastAPI (用于构建流式API端点它原生支持异步和流式响应比Django等更合适)LLM: 主要使用OpenAI GPT-4o的API也测试了通义千问、DeepSeek等兼容OpenAI格式的本地模型。这里有一个重要的选型考量为什么用asyncio和异步迭代器因为流式本质上是长时间运行的I/O密集型任务不断等待LLM生成下一个Token。同步阻塞的写法会独占服务器资源导致并发能力极差。而异步编程模型允许服务器在等待一个请求的LLM响应时去处理其他请求极大地提高了资源利用率和吞吐量。FastAPI asyncio LangGraph的astream系列方法是天作之合。注意如果你在旧版本的LangGraph中找不到某些方法如astream_events请务必升级。流式相关的API在近期版本中迭代很快新版本的功能和稳定性通常更好。3. 基础流式测试从astream到astream_events我们从最简单的图开始逐步增加复杂度。3.1 构建一个简单的链式图假设我们有一个“翻译-总结”工作流先将用户输入翻译成英文再总结英文内容的要点。from typing import TypedDict, Annotated from langgraph.graph import StateGraph, END from langchain_openai import ChatOpenAI import operator # 1. 定义状态 class TranslationState(TypedDict): original_text: str translated_text: str summary: str # 2. 定义节点函数 def translate_node(state: TranslationState): 翻译节点 llm ChatOpenAI(model“gpt-4o”, streamingTrue) # 注意这里streamingTrue message llm.invoke(f“将以下中文翻译成英文{state[‘original_text’]}”) return {“translated_text”: message.content} def summarize_node(state: TranslationState): 总结节点 llm ChatOpenAI(model“gpt-4o”, streamingTrue) message llm.invoke(f“总结以下英文文本的要点{state[‘translated_text’]}”) return {“summary”: message.content} # 3. 构建图 builder StateGraph(TranslationState) builder.add_node(“translate”, translate_node) builder.add_node(“summarize”, summarize_node) builder.set_entry_point(“translate”) builder.add_edge(“translate”, “summarize”) builder.add_edge(“summarize”, END) graph builder.compile()3.2 测试astream方法astream方法是最直接的它流式返回每个节点执行后整个状态的更新值。import asyncio async def test_astream(): initial_state {“original_text”: “LangGraph是一个用于构建多步骤AI工作流的强大框架。”} async for chunk in graph.astream(initial_state): print(f“流式块: {chunk}”) # 运行 asyncio.run(test_astream())输出可能类似于流式块: {‘translate’: {‘translated_text’: ‘LangGraph is a powerful framework for building multi-step AI workflows.’}} 流式块: {‘summarize’: {‘summary’: ‘- LangGraph is a framework.\n- It is used for building AI workflows.\n- These workflows can involve multiple steps.\n- It is described as powerful.’}}看到了什么astream返回的是每个节点执行完成后对状态对象的增量更新delta。第一个块来自translate节点更新了translated_text字段第二个块来自summarize节点更新了summary字段。它流式的是“节点执行的结果”而不是节点内部LLM调用产生的Token流。3.3 测试astream_events方法更强大astream_events是更强大的调试和展示工具它提供了执行过程中的事件流粒度更细。async def test_astream_events(): initial_state {“original_text”: “测试流式输出。”} async for event in graph.astream_events(initial_state, version“v1”): print(f“事件类型: {event[‘event’]}, 内容: {event}”) asyncio.run(test_astream_events())输出事件会更丰富包括on_chain_start图开始执行。on_chat_model_streamLLM开始流式生成如果节点内LLM设置了streamingTrue。on_chat_model_stream会伴随多个on_llm_new_token事件这才是真正的Token流on_chain_end节点执行结束。on_tool_end工具调用结束如果有。关键点只有通过astream_events并且节点内的LLM实例化时传入了streamingTrue你才能捕获到最细粒度的on_llm_new_token事件从而实现真正的“逐词输出”效果。astream流的是节点输出astream_events流的是执行过程事件。3.4 实战将流式接入FastAPI现在我们把上面的流式能力通过一个HTTP API暴露出来。from fastapi import FastAPI from fastapi.responses import StreamingResponse import asyncio from pydantic import BaseModel app FastAPI() class Request(BaseModel): text: str async def stream_graph_response(text: str): 生成器函数用于流式响应 initial_state {“original_text”: text} # 方法一使用 astream流式返回状态更新JSON格式 # async for chunk in graph.astream(initial_state): # # 将chunk转换为字符串并格式化为SSE (Server-Sent Events) 格式 # yield f“data: {json.dumps(chunk, ensure_asciiFalse)}\n\n” # 方法二使用 astream_events流式返回LLM生成的Token更适合前端显示 async for event in graph.astream_events(initial_state, version“v1”): if event[‘event’] ‘on_llm_new_token’: # 提取Token并发送 token event[‘data’][‘chunk’].content if token: # 这里可以封装成SSE格式也可以直接返回纯文本 yield f“data: {token}\n\n” # 你也可以选择发送节点开始/结束等事件给前端用于更新UI状态 elif event[‘event’] ‘on_chain_start’: yield f“event: node_start\ndata: {event[‘name’]}\n\n” elif event[‘event’] ‘on_chain_end’: yield f“event: node_end\ndata: {event[‘name’]}\n\n” app.post(“/stream-process”) async def process_stream(request: Request): return StreamingResponse( stream_graph_response(request.text), media_type“text/event-stream” # 使用SSE协议 )前端简单示例如何接收const eventSource new EventSource(‘/stream-process?text你的输入’); eventSource.onmessage (event) { const data event.data; console.log(‘收到数据:’, data); // 将data追加到页面上的某个元素中 document.getElementById(‘output’).innerHTML data; }; eventSource.onerror (error) { console.error(‘流式连接错误:’, error); eventSource.close(); };实操心得1选择正确的流式方法如果前端只需要最终结果的分段输出如先显示翻译再显示总结用astream更简单它返回结构化的状态更新。如果前端需要实现“打字机”效果实时显示LLM正在生成的内容必须使用astream_events并过滤on_llm_new_token事件。同时确保节点函数内实例化LLM时传入了streamingTrue参数否则不会触发Token级事件。astream_log则更适合后台调试它会流式输出非常详细的日志包括所有内部调用数据量很大一般不适合直接推给前端。4. 高级特性测试子图Subgraph的流式传播子图是LangGraph中实现模块化复用的关键特性。但子图的流式输出行为需要特别注意。4.1 创建与嵌套子图假设我们的“总结”节点本身也是一个复杂的子图包含“提取关键词”和“生成摘要”两个步骤。from langgraph.graph import StateGraph as SubGraphBuilder # 定义子图的状态 class SummarySubState(TypedDict): input_text: str keywords: list[str] final_summary: str # 构建子图 sub_builder SubGraphBuilder(SummarySubState) def extract_keywords(state: SummarySubState): llm ChatOpenAI(model“gpt-4o”, streamingTrue) # 模拟关键词提取 message llm.invoke(f“从以下文本提取3个关键词{state[‘input_text’]}”) # 假设LLM返回逗号分隔的关键词 keywords [k.strip() for k in message.content.split(‘,’)] return {“keywords”: keywords} def generate_summary(state: SummarySubState): llm ChatOpenAI(model“gpt-4o”, streamingTrue) kw_str ‘, ‘.join(state[‘keywords’]) prompt f“基于关键词({kw_str})为以下文本生成一段摘要{state[‘input_text’]}” message llm.invoke(prompt) return {“final_summary”: message.content} sub_builder.add_node(“extract”, extract_keywords) sub_builder.add_node(“summarize”, generate_summary) sub_builder.set_entry_point(“extract”) sub_builder.add_edge(“extract”, “summarize”) sub_builder.add_edge(“summarize”, END) # 编译子图 sub_graph sub_builder.compile() # 在主图中将子图作为一个节点 from langgraph.graph import START def summarize_with_subgraph(state: TranslationState): # 准备子图输入 sub_state {“input_text”: state[‘translated_text’]} # 运行子图并获取最终结果 final_result sub_graph.invoke(sub_state) return {“summary”: final_result[“final_summary”]} # 修改主图将原来的summarize_node替换为新的子图节点 builder StateGraph(TranslationState) builder.add_node(“translate”, translate_node) builder.add_node(“summarize_complex”, summarize_with_subgraph) # 使用子图节点 builder.set_entry_point(“translate”) builder.add_edge(“translate”, “summarize_complex”) builder.add_edge(“summarize_complex”, END) complex_graph builder.compile()4.2 子图流式输出的挑战与解决方案现在我们流式执行complex_graph。问题来了当你使用astream_events时默认情况下子图内部的事件如extract和summarize节点内部的LLM Token流可能不会被传播到主图的事件流中。你只能看到主图节点summarize_complex的开始和结束事件看不到其内部的细节。解决方案使用astream_events的include_names或include_types参数并确保递归包含。async def test_subgraph_stream(): initial_state {“original_text”: “这是一个测试文本用于验证子图内部的流式输出是否能被捕获。”} async for event in complex_graph.astream_events( initial_state, version“v1”, include_names[“translate”, “summarize_complex”, “extract”, “summarize”], # 明确包含子图节点名 # 或者使用 include_types[“chat_model”] 来包含所有LLM事件 ): if event[‘event’] ‘on_llm_new_token’: # 现在这个Token可能来自主图的translate节点也可能来自子图的extract或summarize节点 token event[‘data’][‘chunk’].content node_name event[‘name’] # 通过name字段区分来源 print(f“[{node_name}] 生成Token: {token}”)更优雅的方案封装子图的流式执行。如果子图逻辑复杂更好的做法是在子图节点函数内部也实现流式并以某种方式将流“冒泡”到主图。import json async def summarize_with_subgraph_streaming(state: TranslationState): 一个能内部流式执行的子图节点 sub_state {“input_text”: state[‘translated_text’]} # 我们不在这个节点直接调用invoke而是流式执行子图 async for event in sub_graph.astream_events(sub_state, version“v1”): # 这里是一个关键点我们需要将子图的事件“转发”出去。 # 但节点函数本身无法直接yield给主图的流。 # 一种常见模式是将子图的事件写入一个队列或者作为状态的一部分传递。 # 更实用的生产级方案是将子图也视为一个可流式调用的单元在主图的流式循环中处理。 pass # 为了简化我们先获取结果 final_result await sub_graph.ainvoke(sub_state) return {“summary”: final_result[“final_summary”]}实操心得2子图流式的设计模式对于复杂的嵌套流式一个清晰的设计模式是主图负责协调和最终输出主图使用astream_events驱动。子图作为可流式单元每个子图节点函数本身也设计为异步生成器接收输入yield内部产生的事件或Token。主图消费子图流在主图的流式循环中调用子图节点函数并遍历其生成的异步迭代器将子图的yield值包装后yield给主图的调用者如FastAPI的StreamingResponse。使用唯一ID关联在流式事件中携带一个唯一的run_id或parent_id方便前端区分不同层级的输出来源。这需要更精细的架构设计但能实现最深度的流式控制。对于大多数场景使用include_names参数来捕获子图内部LLM事件已经足够满足“展示Token流”的需求。5. 生产环境问题排查与性能调优在实际部署中流式输出会遇到各种预料之外的问题。5.1 常见问题速查表问题现象可能原因排查步骤与解决方案流式连接过早关闭1. 网络超时Nginx、LB、浏览器。2. 服务器端异常未捕获导致生成器中断。3. LLM API调用超时或频控。1. 检查代理服务器Nginx配置增加proxy_read_timeout,proxy_buffering off。2. 在流式生成器函数内用try...except包裹发生错误时yield一个错误事件而非崩溃。3. 实现LLM客户端的重试和退避机制使用tenacity库。前端收不到Token或接收不连续1. SSE格式不正确缺少\n\n分隔符。2. 前端EventSource解析错误。3. 服务器端缓存如gzip干扰。1. 严格保证SSE格式data: {content}\n\n。用json.dumps确保内容正确转义。2. 前端检查event.data而非event。使用onmessage和onerror回调。3. 在StreamingResponse中设置headers{‘Cache-Control’: ‘no-cache’}禁用中间件缓存。流式输出内容被截断或丢失1. LangGraph的astream_events配置问题某些事件类型被过滤。2. LLM实例化时未设置streamingTrue。3. 使用了不支持流式的模型或API版本。1. 检查astream_events的include_types或include_names参数确保包含了“chat_model”或“llm”。2.务必在节点函数内实例化LLM时传入streamingTrue。3. 确认模型支持流式如OpenAI的gpt-4支持某些开源模型配置可能不同。内存占用随时间增长1. 状态State对象在流式过程中不断累积中间数据未清理。2. 异步任务未正确取消导致资源泄漏。1. 设计状态结构时考虑将需要流式输出的内容与庞大的中间数据分离。使用pydantic模型并设置arbitrary_types_allowed来管理复杂类型。2. 在FastAPI中处理客户端断开连接时主动取消异步生成器任务。可以利用request.is_disconnected()或asyncio的Task取消机制。与Spring Security等权限框架集成时流式中断1. 权限拦截器或过滤器未正确处理StreamingResponse类型。2. CSRF、CORS策略阻止了长连接。1. 在权限框架中为流式端点配置特殊的拦截规则或将其路径排除在常规鉴权链之外但需有其他方式如Token验证。2. 确保CORS配置允许text/event-stream的Content-Type并正确设置Access-Control-Allow-Origin等头。对于CSRF流式端点可能需要禁用或使用Token验证。5.2 性能与可靠性调优要点连接管理与超时服务器端为流式端点设置合理的超时时间。太短会断开长任务太长会占用连接资源。可以根据任务类型动态设置。客户端实现自动重连逻辑。当EventSource触发onerror时可以等待几秒后重新连接并携带上一个接收到的消息ID如果服务端支持以继续。错误处理与优雅降级async def robust_stream_generator(text: str): try: async for event in graph.astream_events({“text”: text}, version“v1”): # ... 处理事件 yield formatted_data except asyncio.CancelledError: # 客户端断开连接正常清理 logging.info(“Streaming connection cancelled by client.”) raise except Exception as e: # 其他异常返回错误信息而不是让服务器崩溃 logging.error(f“Streaming error: {e}”) yield f“event: error\ndata: {json.dumps({‘msg’: ‘处理过程发生错误’})}\n\n”状态序列化优化 LangGraph的状态在节点间传递。如果状态中包含大型对象如图片、长文本频繁的序列化/反序列化会影响性能。考虑使用引用在状态中存储数据库ID或文件路径而非数据本身。使用pickle或cloudpickle处理复杂Python对象注意安全性和版本兼容性。对于超长工作流研究LangGraph的检查点Checkpoint功能将状态持久化避免内存压力。并发与限流 LangGraph本身是异步的但底层LLM API调用可能有速率限制。在生产中需要使用像asyncio.Semaphore或更高级的限流库如slowapi来控制并发请求数避免触发上游API的频控。6. 进阶技巧自定义流式内容与前端协同流式不仅仅是传Token我们可以传递更丰富的结构化信息。6.1 传递结构化事件除了on_llm_new_token我们可以定义自己的事件类型用于前端更新进度条、切换UI状态等。async def stream_with_custom_events(): initial_state {“query”: “请解释量子计算。”} async for event in graph.astream_events(initial_state, version“v1”): event_type event[‘event’] if event_type ‘on_chain_start’: # 通知前端某个节点开始了 yield { “type”: “node_start”, “node”: event[‘name’], “timestamp”: time.time() } elif event_type ‘on_llm_new_token’: yield { “type”: “token”, “token”: event[‘data’][‘chunk’].content, “node”: event[‘name’] } elif event_type ‘on_tool_start’: yield { “type”: “tool_call”, “tool”: event[‘name’], “input”: str(event[‘data’].get(‘input’)) } # ... 其他事件处理 # 流结束时发送完成事件 yield {“type”: “stream_end”, “status”: “completed”} # 在FastAPI中将这些字典转换为JSON字符串再通过SSE发送 async for custom_event in stream_with_custom_events(): yield f“data: {json.dumps(custom_event, ensure_asciiFalse)}\n\n”6.2 前端处理结构化流前端根据收到的事件类型更新不同的UI组件。eventSource.onmessage (e) { const event JSON.parse(e.data); switch(event.type) { case ‘node_start’: updateProgressBar(event.node); addLog(开始执行: ${event.node}); break; case ‘token’: appendToOutput(event.token); // 追加Token到答案区 break; case ‘tool_call’: addLog(调用工具: ${event.tool} 输入: ${event.input}); break; case ‘stream_end’: eventSource.close(); showCompletionMessage(); break; case ‘error’: showError(event.msg); eventSource.close(); break; } };6.3 流式控制暂停、继续与取消这是一个高级需求。LangGraph的CompiledStateGraph本身不直接提供暂停/继续的API但我们可以通过状态State和外部信号来实现一个简单的协作式控制。思路在状态中定义一个pause_requested或cancel_requested的布尔标志。暴露一个额外的API端点如POST /workflow/{run_id}/pause来修改这个标志需要将状态存储在有状态的后端如Redis。在每个节点的开始或结束处检查这个标志。如果pause_requested为真则让节点进入一个循环等待直到标志被清除。如果cancel_requested为真则抛出一个特定异常终止图的执行。流式响应端需要能处理这种“等待”状态可能发送一个“paused”事件给前端。这实现起来较为复杂需要仔细设计状态管理和任务生命周期。对于大多数应用如果只是需要取消更简单的做法是直接关闭前端的EventSource连接并在服务器端的流式生成器中捕获asyncio.CancelledError来清理资源。经过这一系列从基础到进阶的测试和探索LangGraph的流式输出特性虽然在某些细节上需要小心处理但其灵活性和强大功能足以支撑起生产级复杂AI应用的实时交互需求。关键在于理解不同流式方法astream,astream_events的粒度差异妥善处理子图嵌套以及做好生产环境的错误处理、性能监控和前后端协同。最后记住流式不仅仅是技术实现更是用户体验的一部分设计好流式的事件协议能让你的AI应用显得更加智能和响应迅速。