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.10+
- LangChain: 0.1.0+
- LangGraph: 0.0.50+
- HTTP框架: 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”, streaming=True) # 注意这里streaming=True message = llm.invoke(f“将以下中文翻译成英文:{state[‘original_text’]}”) return {“translated_text”: message.content} def summarize_node(state: TranslationState): """总结节点""" llm = ChatOpenAI(model=“gpt-4o”, streaming=True) 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_stream:LLM开始流式生成(如果节点内LLM设置了streaming=True)。on_chat_model_stream会伴随多个on_llm_new_token事件,这才是真正的Token流!on_chain_end:节点执行结束。on_tool_end:工具调用结束(如果有)。
关键点:只有通过astream_events,并且节点内的LLM实例化时传入了streaming=True,你才能捕获到最细粒度的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_ascii=False)}\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时传入了streaming=True参数,否则不会触发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”, streaming=True) # 模拟关键词提取 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”, streaming=True) 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实例化时未设置 streaming=True。3. 使用了不支持流式的模型或API版本。 | 1. 检查astream_events的include_types或include_names参数,确保包含了“chat_model”或“llm”。2.务必在节点函数内实例化LLM时传入 streaming=True。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_ascii=False)}\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应用显得更加智能和响应迅速。