免费获取学习方案
ARTICLE DETAIL

资讯详情

深耕编程基础知识与建站技术分享的一线实战洞察。

SSE流式传输、OutputParser与ToolCall:一条链路打通LangChain智能体服务

SSE流式传输、OutputParser与ToolCall:一条链路打通LangChain智能体服务 最近在把一套基于 FastAPI LangChain LangGraph 的智能体服务整个从前到后打通前端用 Vue 通过 SSE 流式接收大模型输出后端用 OutputParser 把模型吐出来的文本整理成结构化数据同时接上 ToolCall 让 Agent 去查业务库、调第三方接口。整个过程走下来最深的感受是很多人把 SSE、OutputParser、ToolCall 当成三个孤立知识点但它们实际上是一条链路上的三个环节任何一个环节处理不好前端用户体验都会直接崩掉。这篇把我实操过程中的链路设计、代码细节和踩坑记录整理出来给正在做类似流式对话、智能体服务的同学一个可直接参考的落地方案。1. 先理清三层问题流式传输、结构化解析、工具调用1.1 为什么前端的“通知”必须用 SSESSEServer-Sent Events本质上是服务端单向推送的一套 HTTP 约定连接建立后服务端可以持续往同一个响应里写数据每条消息以data:开头以两个换行符\n\n结束。它跟 WebSocket 最大的区别就是简单——底层就是普通 HTTP不需要升级协议不需要单独维护心跳断线了浏览器还会自动重连EventSource自带重连机制。大模型生成回答天然是“挤牙膏”式的不是等全部生成完一次性返回而是一边推理一边吐 token。用户能接受多等几秒钟但接受不了十几秒的空白加载。所以现在的对话式 AI 应用基本都默认用 SSE 把增量输出推给前端让用户看到“正在输入”的过程。提示SSE 是单向的服务端到客户端。如果还需要客户端频繁向服务端发指令比如游戏操作那才需要 WebSocket。对话场景下用户的输入本来就是一个普通 POST 请求服务端把响应流式返回SSE 刚好匹配。1.2 解析器和工具调用是两件事不是一件事很多新手会把 OutputParser 和 ToolCall 混在一起其实它们的职责完全不同OutputParser解决“模型必须以我要求的格式返回”的问题。你问它“把这段话翻译成 JSON”它可能偶尔多解释几句、偶尔少个括号Parser 负责把模型输出改成可靠的结构化数据。ToolCall解决“模型要调用外部能力”的问题。比如用户问“北京明天天气”模型本身不掌握实时天气它需要按约定返回一个“我想调用 get_weather 函数参数是北京”的结构由我们的代码去真正执行函数再把结果喂回给模型。两者经常同时出现但不要混为一谈。调试的时候要先分清是结构不对还是工具没调起来不然会浪费很多时间。1.3 这条链路在真实业务中的完整形态我实际落地时需求比“做个聊天框”复杂得多。我把整套链路拆成了这几个环节前端 POST 一条用户消息给后端。后端把消息交给 LangChain 的 Agent内部用 LangGraph 管理状态和循环。Agent 判断是否需要调用工具需要则进入 ToolCall 循环把工具结果作为上下文继续推理不需要则直接进入生成回答阶段。模型输出通过 SSE 增量推给前端。前端拿到的是分片数据要在本地重组、增量渲染并且把“正在调用工具”这类过程状态也展示出来。这个链路里只要有一个环节不同步用户看到的现象就会很诡异。比如工具调用期间前端空白了十几秒比如 JSON 解析偶尔失败导致整轮对话报错再比如 SSE 连接被中间层掐掉。后面我逐个环节展开。2. 三大 OutputParser 的逐个拆解与选型LangChain 生态里 OutputParser 非常多日常用得最顺手的就是这三个PydanticOutputParser、StructuredOutputParser、JsonOutputParser。我分别说用法、原理和坑。2.1 PydanticOutputParser字段多、类型严的首选这个 Parser 的核心逻辑是你定义一个 Pydantic 模型LangChain 会把模型的 JSON Schema 转成一段格式说明塞进提示词里模型按这个格式返回 JSON最后 Parser 用模型类把 JSON 校验并实例化成对象。from pydantic import BaseModel, Field from langchain.output_parsers import PydanticOutputParser from langchain_core.prompts import ChatPromptTemplate class MovieReview(BaseModel): title: str Field(description电影名) rating: float Field(description评分0到10之间的数字) tags: list[str] Field(description标签列表例如 [剧情, 悬疑]) recommendation: str Field(description一句话推荐语) parser PydanticOutputParser(pydantic_objectMovieReview) prompt ChatPromptTemplate.from_template( 根据用户的描述输出电影评论信息。\n{format_instructions}\n用户描述{input} ) chain prompt | parser result await chain.ainvoke({ input: 我昨晚看了《流浪地球3》特效不错但剧情略显拖沓我给7.5分推荐喜欢科幻的人看。, format_instructions: parser.get_format_instructions(), }) print(result.title) # 直接拿到对象属性关键点在format_instructions它会把 JSON Schema 转换成像“请用如下 JSON 格式输出{...}”这样的指令。如果你用的是prompt | model | parser这种链式写法需要把format_instructions显式塞进 prompt 变量里。很多人漏了这一步模型乱输出Parser 直接报错。实际项目中我踩过的坑有两个第一Pydantic 版本差异。Pydantic v2 里.dict()改成了.model_dump()LangChain 新版本已经兼容但如果项目里同时混用 v1 和 v2字段校验规则会互相干扰。建议统一pydantic2.5。第二字段描述必须写清楚。模型的格式遵循来自你的 Field description如果你只写rating: float不写范围模型可能给个 8.5也可能给个 85解析时虽然不报错但下游业务就可能出现诡异数据。注意PydanticOutputParser 适合结构复杂、字段多、需要类型校验的场景。它对解析失败的兜底能力一般一旦模型输出完全不符合 JSON会直接抛OutputParserException。要自己接重试或修复逻辑后面的“问题排查”章节会讲。2.2 StructuredOutputParser简单键值对时最省事StructuredOutputParser 不需要定义 Pydantic 模型只用 ResponseSchema 列字段名和描述就能生成格式说明解析后得到一个普通字典。from langchain.output_parsers import StructuredOutputParser, ResponseSchema response_schemas [ ResponseSchema(nameteam, description团队名称), ResponseSchema(namemembers, typelist, description成员姓名列表), ResponseSchema(nameplan, description一句话行动计划), ] parser StructuredOutputParser.from_response_schemas(response_schemas) prompt ChatPromptTemplate.from_template( 根据需求拆分项目团队\n{format_instructions}\n需求{input} ) chain prompt | parser data chain.invoke({ input: 我们要做一个智能客服需要产品、开发、测试, format_instructions: parser.get_format_instructions(), }) print(type(data)) # class dict它的优点是上手快不需要额外维护 Pydantic 类适合那种“我就是想把模型输出转成几个固定 key 的字典”的轻量需求。缺点也很明显嵌套结构支持差。你要是想输出“成员数组里再套一个技能对象”StructuredOutputParser 就有点吃力了而且生成的格式指令比较啰嗦对一些指令理解弱的模型反而容易出错。我个人的习惯是字段不超过 5 个、且没有深层嵌套时用它一旦结构复杂就换 PydanticOutputParser。2.3 JsonOutputParser跟 JSON Mode 配合最稳JsonOutputParser 是三者里代码最简洁的用途也最纯粹——就是让模型输出一个合法 JSON 对象。from langchain.output_parsers import JsonOutputParser from langchain_core.prompts import ChatPromptTemplate parser JsonOutputParser() prompt ChatPromptTemplate.from_template( 把所有要点转成 JSON格式如下\n{format_instructions}\n输入{input} ) chain prompt | parser json_data chain.invoke({ input: 预算5000块以内推荐3个适合程序员聚会的地方, format_instructions: parser.get_format_instructions(), }) print(json_data)但在 2024 年以后我更推荐直接配合with_structured_output(JsonOutputParser)或with_structured_output(SomePydanticModel)来用因为新版本里 LangChain 会尽可能在能支持 JSON Mode 的模型上强制开启 JSON 模式这样模型几乎不会输出多余废话解析成功率能提高一个档次。from langchain_openai import ChatOpenAI model ChatOpenAI(modelgpt-4o-mini, temperature0) structured_model model.with_structured_output(JsonOutputParser()) result structured_model.invoke(把我今天不想上班翻译成英语给我JSON) # {translated: I dont want to go to work today.}JsonOutputParser 还有个独门武器它可以解析“流式中不完整的 JSON”。配合 LangChain 的parse_partial_json函数可以把流式输出的半截 JSON 尽量恢复成字典。这是后面 SSE 流式场景下很有用的点。2.4 三种 Parser 怎么选Parser适用场景优点缺点流式友好度PydanticOutputParser复杂嵌套结构、需要类型校验基于 Schema 强校验、对象访问方便需维护模型类、提示词要额外塞指令一般需整体拼接后解析StructuredOutputParser简单扁平字典快速、无需定义类嵌套弱、指令冗长一般JsonOutputParser纯 JSON 输出配合 JSON Mode代码最简、兼容 with_structured_output结构校验弱较好支持 partial 解析提示如果模型本身支持 JSON ModeOpenAI 系列、Qwen 的enable_thinking关闭后加 JSON 约束等都支持优先model.with_structured_output(pydantic_model)而不是手动塞format_instructions。前者走的是模型的原生结构化输出通道稳定性和一致性远高于让模型“读提示词猜格式”。3. ToolCall 方案的完整落地3.1 Function Calling 的本质是什么ToolCall也叫 Function Calling / Tool Calling解决的是“让模型决定什么时候调用什么函数”的问题。模型本身不会执行函数它只是根据工具的 JSON Schema 描述输出一个足够规范的结构化意图{ name: get_weather, arguments: {\city\: \北京\} }我们的后端代码真正去解析这个结构、查天气、把结果拼成一条消息再喂回模型模型基于工具结果给出最终回答。整个过程对模型来说像是“伸手拿了张纸条”它并不理解纸条内容背后的计算只看得到纸条上的字。LangChain 提供tool装饰器能直接把一个普通 Python 函数转成模型可识别的工具定义函数的 docstring 自动变成工具描述。3.2 bind_tools 与调度循环我项目里最常用的是 OpenAI 兼容接口LangChain 里绑工具的写法from langchain_core.tools import tool from langchain_core.messages import SystemMessage, HumanMessage, ToolMessage from langchain_openai import ChatOpenAI import json tool def get_weather(city: str) - str: 获取指定城市的实时天气返回文字描述。 # 这里接真实天气 API return f{city} 当前晴25 度适合出门 tool def search_product(keyword: str) - str: 根据关键词搜索商品返回商品名称和价格列表。 return f搜索结果{keyword} 相关商品 3 件价格 99 元起 tools [get_weather, search_product] model ChatOpenAI(modelgpt-4o-mini, temperature0) model_with_tools model.bind_tools(tools) messages [ SystemMessage(content你是购物助手可以查询天气和搜索商品。), HumanMessage(content北京今天天气如何顺便帮我找一下鼠标), ]调用一次模型后重点看返回对象里的tool_callsresponse model_with_tools.invoke(messages) if response.tool_calls: for tool_call in response.tool_calls: print(tool_call[name], tool_call[args])拿到的结构大致是get_weather {city: 北京} search_product {keyword: 鼠标}接下来要做的就是真正执行工具并把结果回传给模型for tool_call in response.tool_calls: tool_name tool_call[name] tool_args tool_call[args] # 根据名称找到对应工具函数 tool_map {get_weather: get_weather, search_product: search_product} selected_tool tool_map[tool_name] tool_result selected_tool.invoke(tool_args) messages.append(ToolMessage(contenttool_result, tool_call_idtool_call[id])) # 再让模型基于工具结果生成最终回答 final_response model_with_tools.invoke(messages) print(final_response.content)这里有个工程细节必须把ToolMessage的tool_call_id跟模型返回的tool_call[id]对上否则模型无法关联工具结果整个上下文会错乱。这是新手最容易忽略的点。还需要设置工具调用轮数的上限否则模型可能陷入“调工具→看到结果→再调工具”的死循环。一般我会控制在 3~5 轮超出就强制终止并返回错误提示。3.3 工具调用与流式的节奏配合工具调用如果也走流式会出现一个很麻烦的问题模型可能在流式过程中先吐出一部分工具调用信息再吐正文。前端如果直接把所有流都当正文渲染用户会看到一大段乱码 JSON。我的方案是后端做一次“阶段化管理”在 Agent 判定需要调用工具时先通过 SSE 给前端推一个{type: tool_call, tool: get_weather, args: {...}}事件。前端收到这个事件后显示“正在查询天气...”之类的过程卡片。工具执行完、模型开始生成最终正文时后端再推{type: message, delta: ...}数据流。这样前端实际只负责聊天气泡和状态卡片不需要理解工具调用的内部结构。注意如果模型在流式过程中先输出tool_callsLangChain 的流式 chunk 里会带有tool_call_chunks字段。你需要额外处理增量聚合把零散的 name/arguments 片段拼起来。建议在 Agent 层面提前判断是否要工具调用是的话就“先工具后文本”避免在流式中间突然冒出工具调用。4. SSE 流式接口从后端到前端的完整封装4.1 FastAPI 端如何封装流式接口后端我用 FastAPI 的StreamingResponse来做 SSE。核心是生成器不断yield符合 SSE 格式的文本import json from fastapi import FastAPI from fastapi.responses import StreamingResponse from pydantic import BaseModel app FastAPI() class ChatRequest(BaseModel): message: str session_id: str | None None async def event_generator(message: str): # 这里调用 LangChain / LangGraph 的 Agent 流程产出不同类型的增量事件 yield fdata: {json.dumps({type: status, content: 思考中}, ensure_asciiFalse)}\n\n # 假设这是一个流式生成器 async for delta in some_agent_stream(message): event {type: message, delta: delta} yield fdata: {json.dumps(event, ensure_asciiFalse)}\n\n # 结束事件 yield fdata: {json.dumps({type: done}, ensure_asciiFalse)}\n\n app.post(/api/chat) async def chat(req: ChatRequest): return StreamingResponse( event_generator(req.message), media_typetext/event-stream, headers{ Cache-Control: no-cache, Connection: keep-alive, X-Accel-Buffering: no, # 重要禁用 Nginx 缓冲 }, )这里X-Accel-Buffering: no非常关键。如果前端和 FastAPI 之间隔了 Nginx 或 CDN默认中间层很有可能把整个 SSE 响应缓冲起来结果就是前端等了好久才一次性收到全部内容“流式”直接失效甚至超时断开。4.2 前端 Vue 如何接收并解析 SSE前端的坑也很明确EventSource只支持 GET 请求。但我们的对话接口是 POST要传消息体和 session_id所以不能用EventSource得用 Fetch API ReadableStream 手动解析。我在 Vue 里封装了一个简单的事件解析器// useChat.js export async function streamChat(message, onEvent) { const res await fetch(/api/chat, { method: POST, headers: { Content-Type: application/json }, body: JSON.stringify({ message }), }); if (!res.ok || !res.body) { throw new Error(请求失败); } const reader res.body.getReader(); const decoder new TextDecoder(utf-8); let buffer ; while (true) { const { done, value } await reader.read(); if (done) break; buffer decoder.decode(value, { stream: true }); // SSE 按 \n\n 分割事件 const parts buffer.split(\n\n); buffer parts.pop(); for (const part of parts) { const lines part.split(\n); for (const line of lines) { if (line.startsWith(data:)) { const dataStr line.slice(5).trim(); const event JSON.parse(dataStr); onEvent(event); } } } } }如果对可靠性要求再高一点可以在这个基础上加一个line.startsWith(event:)的处理以及用id:字段做断点续传。但大多数场景下data:一条事件贯穿到底就够用了。4.3 增量数据与状态机的映射前端拿到事件之后需要根据事件类型维护一个简单的状态机。我用四种事件类型事件类型含义前端处理statusAgent 正处于某阶段显示状态提示tool_call正在调用工具展示工具名和参数message模型正文增量追加到当前消息气泡done本轮结束结束 loading允许发送下一条这样做的好处是前端拿到的不再是一堆裸文本而是一套可驱动的画面状态。真正常用的 Agent 产品其实都是这种事件驱动的模式。5. 高频问题与排查实录5.1 “stream disconnected before completion: idle timeout waiting for sse” 的根因与解法这个报错我在调试时遇到太多次了。字面意思是“SSE 等待超时连接在完成前被断开了”。常见原因有以下几种第一服务端长时间没有任何数据发出。比如模型推理比较慢中间十几秒没有输出而前端、Nginx 或云网关的超时时间设置得比较短就把连接掐了。解决办法是在没有实际内容输出时定期发送心跳。第二Nginx 缓冲或代理超时。需要在 Nginx 配置里加上proxy_buffering off; proxy_read_timeout 300s; proxy_send_timeout 300s;第三FastAPI 服务被 uvicorn 或云负载均衡断开。如果前面套了阿里云/腾讯云 SLB 之类注意看它们的连接空闲超时设置一般要调到 300 秒以上。心跳的写法很简单在生成器里每隔一段时间发一个注释行import asyncio async def event_generator(): try: # 业务处理循环 while True: delta await some_async_generator() if delta is not None: yield fdata: {json.dumps({type: message, delta: delta})}\n\n except asyncio.CancelledError: # 客户端断开时清理 raise finally: # 兜底心跳减少空闲断连 yield f: ping\n\nSSE 规范里以冒号开头的行是注释客户端会忽略但连接仍然活跃。有的前端实现会把data:里没有任何数据也当心跳比如data: {type: heartbeat}两者都可以看兼容性选择。5.2 JSON 解析失败的兜底策略OutputParser 解析失败是最常见的。模型偶尔就是会给出不完整的 JSON、多输出一句废话、或者干脆把{写成了[。我提供三个兜底层级第一层利用模型的 JSON Mode。OpenAI、Qwen、DeepSeek 等模型都支持通过参数强制 JSON 输出。LangChain 里用with_structured_output()时它会尽量帮你开启这个模式所以优先用它。第二层对不完整 JSON 做修复。LangChain 自带的OutputFixingParser可以包装一个普通 Parser解析失败时自动把“原输出 报错信息”发给模型让模型修正后重新解析。from langchain.output_parsers import OutputFixingParser base_parser JsonOutputParser() fixing_parser OutputFixingParser.from_llm( parserbase_parser, llmChatOpenAI(modelgpt-4o-mini, temperature0) )第三层流式场景下的 partial 解析。如果你正在用 SSE 流式输出 JSON可以在前端或后端把半截 JSON 用parse_partial_json尝试解析解析成功就提前渲染结构失败就继续攒数据from langchain_core.output_parsers.json import parse_partial_json partial {title: 流浪地 parsed parse_partial_json(partial) # 可能返回 None也可能返回 {title: 流浪地}注意parse_partial_json不是万能的。对于未闭合的复杂嵌套结构它经常返回None。所以这条策略适合“能解析就提前展示解析不了就等待下一次数据”的增量思路不适合作为唯一的保底方案。5.3 ToolCall 状态与前端渲染不同步前端经常出现的一种情况是后端已经进入工具调用阶段但前端还停留在“正在思考”的加载动画上用户等得很焦躁。而且如果工具本身执行很慢甚至会出现“用户以为服务挂了”的情况。我的做法是后端把 Agent 的每个阶段都显式抛成事件。async def agent_events(message: str): yield {type: status, content: 理解需求中} # LangGraph 中在节点之间切换时可以插入状态事件 if need_tool: yield {type: tool_call, tool: tool_name, args: tool_args} async for chunk in final_stream: yield {type: message, delta: chunk}这样前端可以做出“先显示状态卡片→再显示工具卡片→最后出现正文”的完整过渡用户体验会有质的提升。用 LangGraph 的话还可以在节点里自定义 event 的发送或者用一个共享队列把节点执行状态推给 SSE 生成器。注意如果 Agent 使用了 LangGraph 的 checkpoint 持久化比如基于 Postgres 保存状态断线重连后可以从上次节点继续而不是从头再来。这个机制对长耗时工具链特别有用也是 LangGraph 相比普通链式调用最大的优势之一。6. 在实际项目里把这条链路跑通后的一点体会开头提到的整套方案我在真实项目里已经跑了一段时间。有几个细节想放在最后说。一是 架构选型上不要为了流式而流式。如果只是给内部工具做一个单轮问答PydanticOutputParser 和 ToolCall 可能反而是过度设计直接普通 HTTP 请求同步返回就行。但一旦涉及“用户感知等待时间超过 3 秒”或者“Agent 要多次调用工具”SSE 事件驱动就是刚需。二是 关于工具调用与流式的节奏我最终选择了“工具阶段不流式、正文阶段流式”的混合方案。工具执行过程是计算密集且不确定的流式输出中间状态意义不大而正文阶段是用户的耐心瓶颈必须流式。这个取舍让代码简单很多也让前端状态机变得清晰。三是 后续扩展空间很大。等工具数量多了之后可以用 LangChain 的 Agent 系列模块做工具路由也可以参考 LangGraph 的多 Agent 协作框架把“让 AI 下地干活”变成“让多个 AI 角色分工干活”。不管怎么扩展底层这套 SSE 数据协议、结构化解析、工具调用循环的骨架是不变的建议先把骨架打扎实再往上加复杂度。最后再分享一个小技巧调试 SSE 时不要只依赖浏览器 Network 面板用 curl 最直接。curl -N -X POST http://localhost:8000/api/chat \ -H Content-Type: application/json \ -d {message: 北京今天天气如何} \ --max-time 120-N参数表示禁用缓冲能实时看到服务端推送的每一行 SSE 数据。这一步能帮你快速区分问题到底出在后端、中间层还是前端解析逻辑省下大量排查时间。
返回列表