ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

从SSE流式到结构化输出:LangChain OutputParser与ToolCall实战

从SSE流式到结构化输出:LangChain OutputParser与ToolCall实战 做 AI 产品接入的时候最拧巴的一件事就是用户要的是“流式输出”像 ChatGPT 那样一个字一个字往外蹦但你的业务要的是“结构化输出”一条干净的 JSON 或者一个填好字段的对象能直接落库、能直接驱动前端渲染。这两件事放在一起就是标题里那句话——从 SSE 流式到结构化输出。它不是一个选择题而是一条完整的链路SSE 是传输层LangChain 的三大 OutputParser 是在文本流上做后期加工而 ToolCall 则从生成阶段就把输出钉死在 schema 里。我之前接过的项目里十有七八卡在同一个地方模型吐出来的东西“差不多能用”但前端要字段、后端要校验、运营要报表谁都不能拿一段带 Markdown 和废话的文本去解析。所以这篇文章就是把这条链路整个捋一遍SSE 接口怎么封装、三大 OutputParser 分别在什么场景用、ToolCall 怎么替代“事后解析”、以及线上跑 SSE 必踩的那些坑。适合正在做 LLM 应用的后端、全栈和 Agent 方向的同学照着抄完能少走至少两周弯路。1. 先定基调SSE、OutputParser、ToolCall 各管哪一层1.1 SSE 只是“水管”别把它和“结构化”混为一谈SSE 全称 Server-Sent Events是浏览器和服务端之间一条单向的 HTTP 流。它靠的不是什么黑魔法就是服务端不关连接一行一行地往外面写data: 内容两个连续换行代表一条消息结束。协议本身简单到可以手写这也是它在 LLM 场景里重新火起来的原因大模型生成是“一边算一边出”的天然适合这种半开连接逐步推数据的方式而不是等整段算完再一次性返回。至于为什么不直接用 WebSocket实际工程里理由很朴素SSE 走的是普通 HTTPNginx、网关、鉴权中间件全部不用额外适配客户端断线了可以用Last-Event-ID续传而且 LLM 生成的语义本来就是单向的不需要双向通道。很多团队碰上“通知 sse”“vue python sse”这类问题本质都是没搞明白这层协议是“分帧的文本流”而不是普通接口返回。你先记住一句话SSE 只解决“怎么把数据一段段送过去”它不承诺送过去的东西长什么样。1.2 OutputParser 和 ToolCall 的分工OutputParser 和 ToolCall 都在解决同一个痛苦——模型输出不可控但它们动手的时机完全不同。OutputParser 是“事后加工”模型已经吐出了一堆 token可能是纯文本、可能是半截 JSON、可能裹着 json 代码块解析器负责把这堆东西整理成你要的对象。它像一个流水线末端的质检员不管前面多乱到你手里必须变成合格零件。ToolCall 是“事前约束”模型在生成阶段就知道自己只能按某个 schema 输出直接返回tool_calls数组参数天生就是合法 JSON。它更像用模具压零件出来就是规格件没有多少质检空间。这个区别决定了选型方向如果你的接口只是想“把模型回答转成对象”OutputParser 足够了如果后面要接工具、要跑 Agent、要保证一万次调用里九千次都能直接解析别犹豫上 ToolCall。2. 三大 OutputParser 对比与选型2.1 StrOutputParser流式透传的默认项很多人以为 OutputParser 都要跟 JSON 打交道其实最常用的是最朴素的StrOutputParser。它的作用一句话把模型输出转成字符串。在 LangChain 的链式写法里prompt | llm | StrOutputParser()是标准配置它干的事就是吞掉 AIMessage 的包装直接吐出文本内容。别小看这层“透传”流式场景里它是默认选择。配合链的astream()你拿到手的每一个 chunk 都是增量文本后端只管把它包装成 SSE 的 token 事件发出去前端直接累加渲染。这个模式对应的是“对话型 UI”要的是逐字效果不需要结构。注意StrOutputParser 完全不校验内容。如果后续要落库、要提取字段就得再叠加一层解析或者干脆换下面的方案。2.2 JsonOutputParser带着“半成品容忍度”的 JSONJsonOutputParser是 LangChain 里专门为 JSON 输出准备的轻量解析器它的核心价值在于对“流式过程”友好。JsonOutputParser.parse_partial_json()能解析未闭合的 JSON 片段也就是说在流式没结束时你就可以拿它做增量预览不会因为少一个}就炸掉。这一点在给前端做“实时结构预览”时非常管用。使用方式也很直接在 Prompt 里拼上parser.get_format_instructions()要求模型只输出合法 JSON然后链的尾部挂parser。它适合字段少、容忍度高、不想为 Pydantic 校验付额外成本的对象。但缺点同样明显——它不做类型校验字段缺失、类型不对它全当没看见你必须自己在业务层兜底。2.3 PydanticOutputParser结构化输出的“正式合同”需要强约束的时候直接上PydanticOutputParser。它和 Pydantic 模型配合先定义好输出结构比如from pydantic import BaseModel, Field from typing import List class SummaryOutput(BaseModel): title: str Field(description一句话标题) keywords: List[str] Field(description关键词列表3-5个) score: float Field(ge0, le100, description相关度评分)然后PydanticOutputParser(pydantic_objectSummaryOutput)Prompt 里同样拼format_instructions。解析时模型输出会走 Pydantic 校验类型不对、枚举越界、字段缺失都会抛ValidationError。这就是“合同”的意义宁可这次调用失败也不让脏数据流进下游。如果解析失败LangChain 里还有OutputFixingParser它会拿一个 LLM 去“看”解析失败的原因并尝试修正输出。注意这会多一次模型调用延迟和成本都上去了所以我在实际项目里只会对低延迟要求的内部接口用对用户侧接口宁可重试一次。三大解析器选型表如下解析器输出类型流式兼容校验强度适用场景StrOutputParserstr强逐 chunk 透传无对话流式展示JsonOutputParserdict强支持 partial JSON弱仅 JSON 合法性轻量结构化、实时预览PydanticOutputParserBaseModel一般需等完整输出强类型字段校验落库、报表、下游强依赖顺带一提LangChain 还有一个StructuredOutputParser基于ResponseSchema生成说明可以理解为 JsonOutputParser 的“带说明版本”。真正常用的就上面三个这也是圈子里默认的那个“三件套”。3. ToolCall不再“事后解析”直接“按格式生成”3.1 ToolCall 的原理和坑ToolCall 背后的机制是模型原生支持的 function calling。模型在预训练和指令微调阶段就被训练成“当遇到工具调用需求时输出一段结构化的参数”而不是像普通文本那样自由发挥。在 LangChain 里模型返回的AIMessage会带tool_calls字段里面是[{ name: 工具名, args: {...}, id: 调用ID }]这样的结构args本身就是合法 JSON。它的最大价值就是让“解析失败”这件事几乎消失。模型在生成阶段就按 schema 约束输出你不再需要担心 Markdown 代码块、夹带解释文字、字段名漂移这些问题。代价是你必须自己实现“工具调用→拿到结果→回填→继续生成”的循环。这个循环其实不复杂但它是 Agent 的逻辑内核很多框架做 Agent 就是在封装这一层。3.2 with_structured_output 与 bind_tools 两条路LangChain 里有两种落地方式。第一种是with_structured_output适合“纯结构化返回”。你定义一个 Pydantic 模型然后structured_llm llm.with_structured_output(SummaryOutput) res await structured_llm.ainvoke(总结这段内容...) # res 是 SummaryOutput 实例直接用 res.title、res.keywords一行代码拿到对象底层自动走了 function calling省去手工解析。注意它的返回值是最终对象不再有流式中间过程——如果你又要流式又要结构得绕道下面的事件流方案。第二种是bind_tools适合“工具执行 Agent 循环”。先定义一个工具from langchain_core.tools import tool tool def get_weather(city: str) - str: 查询指定城市的当前天气参数为城市名。 return f{city}晴28℃东南风2级 llm_with_tools llm.bind_tools([get_weather])然后发起一次调用判断ai_msg.tool_calls是否存在存在就去执行工具把结果以role: tool的消息回填再继续要模型的最终回答。这个循环就是 Agent 最小实现。3.3 流式场景下 ToolCall 与吐字的“受体冲突”这里有个绕不开的矛盾ToolCall 必须先等模型把整个工具参数输出完你才能拿到结构化参数去执行工具执行期间模型是“暂停”的前端没有任何文本可以展示。而用户已经习惯了 ChatGPT 那种逐字流式体验一旦界面卡住两秒马上觉得服务坏了。解决思路是“事件流协议”把不同阶段的产出包装成不同类型的 SSE 事件。我的习惯是统一成五种事件token表示增量文本、tool_start表示开始执行工具、tool_end表示工具执行完成并带回结果、error表示出错、done表示整轮结束。前端拿到事件后按类型渲染把“工具执行中”设计成一个状态卡片这就既保住了流式体验又接住了结构化数据。所谓“封装sse 流式接口调用逻辑”封装的就是这一层事件路由而不是单纯把 text/event-stream 透传出去。4. FastAPI 后端 SSE 接口实战4.1 接口协议与响应头我习惯用 FastAPI 起一个POST /api/chat/stream请求体是{ messages: [...], mode: str }mode决定走哪种解析链路。返回必须用StreamingResponse而且media_type设成text/event-stream。响应头里有三个必须交代的东西Cache-Control: no-cache防止中间层缓存X-Accel-Buffering: no通知 Nginx 别开缓冲Connection: keep-alive保持长连接。如果漏了第二条Nginx 会把你的流攒一大段再一次性吐给前端表现就是“前端半天没动静突然冒出整段文字”这时候根本不是应用逻辑的问题是代理缓冲在作怪。4.2 一个可运行的 SSE 生成器骨架直接给一个能跑的骨架重点看event_gen这个异步生成器from fastapi import FastAPI, Request from fastapi.responses import StreamingResponse from pydantic import BaseModel, Field from typing import List import json from langchain_openai import ChatOpenAI from langchain_core.output_parsers import StrOutputParser, JsonOutputParser, PydanticOutputParser from langchain_core.prompts import ChatPromptTemplate app FastAPI() llm ChatOpenAI(modelgpt-4o-mini, temperature0) class ChatRequest(BaseModel): messages: List[dict] mode: str str async def event_gen(req: ChatRequest, request: Request): event_id 0 def emit(event_type: str, data: dict) - str: nonlocal event_id event_id 1 payload {type: event_type, **data} return fid: {event_id}\nevent: {event_type}\ndata: {json.dumps(payload, ensure_asciiFalse)}\n\n try: if req.mode str: chain ChatPromptTemplate.from_template(用户问题{input}) | llm | StrOutputParser() async for chunk in chain.astream({input: req.messages[-1][content]}): if await request.is_disconnected(): break yield emit(token, {content: chunk}) else: parser JsonOutputParser() prompt ChatPromptTemplate.from_template( 提取用户消息中的关键信息严格按 JSON 输出。\n{format_instructions}\n用户消息{input} ) result await (prompt | llm | parser).ainvoke({ input: req.messages[-1][content], format_instructions: parser.get_format_instructions(), }) yield emit(result, {data: result}) except Exception as e: yield emit(error, {message: str(e)}) finally: yield emit(done, {}) app.post(/api/chat/stream) async def chat_stream(req: ChatRequest, request: Request): return StreamingResponse( event_gen(req, request), media_typetext/event-stream, headers{ Cache-Control: no-cache, X-Accel-Buffering: no, Connection: keep-alive, }, )两个细节要敲黑板。第一循环里每拿到一个 chunk 就先检查await request.is_disconnected()客户端跑了就立刻停否则 token 还在烧钱。第二finally里必须发done前端靠这个事件结束加载态没有它前端就会一直转圈直到超时。ToolCall 模式的生成器在事件路由基础上加一层工具循环from langchain_core.messages import HumanMessage tool def get_weather(city: str) - str: 查询指定城市的当前天气参数为城市名。 return f{city}晴28℃东南风2级 async def tool_event_gen(req: ChatRequest, request: Request): # 省略 emit 定义同上 messages [HumanMessage(contentreq.messages[-1][content])] llm_with_tools llm.bind_tools([get_weather]) ai_msg await llm_with_tools.ainvoke(messages) if ai_msg.tool_calls: yield emit(tool_start, {name: ai_msg.tool_calls[0][name]}) for tc in ai_msg.tool_calls: result get_weather.invoke(tc[args]) yield emit(tool_end, {name: tc[name], result: result}) messages.append(ai_msg) messages.append({role: tool, tool_call_id: tc[id], content: result}) async for chunk in llm.astream(messages): if await request.is_disconnected(): break if chunk.content: yield emit(token, {content: chunk.content})这块产出的就是热词里提到的“基于 fastapi langchain 做 AI Agent”的最小闭环模型决定调用工具、执行工具、把结果交还模型、最终流式生成回答。4.3 客户端如何解析Vue 里的 fetch ReadableStream很多同学第一次接 SSE 会想用浏览器原生EventSource但它的限制很致命只支持 GET不能带自定义请求头也没法传 body。真实业务里要鉴权、要传上下文所以必须用fetch拿流自己解析。我封装了一个通用函数export async function createSSEStream(url, payload, handlers {}) { const res await fetch(url, { method: POST, headers: { Content-Type: application/json }, body: JSON.stringify(payload), }); if (!res.ok) throw new Error(HTTP ${res.status}); 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 }); const frames buffer.split(\n\n); buffer frames.pop(); // 最后一段可能是半包留在缓冲区 for (const frame of frames) { const dataLines frame.split(\n).filter(l l.startsWith(data:)); if (!dataLines.length) continue; const raw dataLines.map(l l.slice(5).trim()).join(\n); try { const msg JSON.parse(raw); (handlers[msg.type] || handlers.onMessage)?.(msg); } catch (e) { console.warn(SSE parse error, raw, e); } } } }关键点是buffer的处理网络包不会刚好按 SSE 的\n\n分帧你的数据可能半路被切断所以必须把读到的数据先拼进 buffer再按双换行切帧剩下的半帧留到下一轮继续拼。在 Vue 里使用时按事件类型注册回调就行createSSEStream(/api/chat/stream, { messages, mode: tool }, { token(msg) { answer msg.content; }, tool_start(msg) { showTool(工具运行中${msg.name}); }, tool_end(msg) { showTool(工具完成${msg.result}); }, error(msg) { showError(msg.message); }, done() { finishLoading(); }, });这套封装配合后端的id字段还能轻松做断线重连——记下最后收到的id重连时带回去请求服务端补发比无脑整轮重试省得多。5. 常见问题与排查技巧实录5.1 必须收藏的报错速查表这些是我在线上和群里反复帮人看过的典型问题先列成表现象根因解决办法stream disconnected before completion: idle timeout waiting for sse网关或 Nginx 的空闲超时把连接掐断了提高proxy_read_timeout每 15 秒发一条: ping心跳解析 JSON 报错模型输出带 json 代码块和废话模型没遵守 format instructionsPrompt 拼format_instructions或换 ToolCall 强制 schemaPydanticValidationError枚举值不在范围内模型自己“发明”了选项字段用Literal限定加默认值兜底必要时改用with_structured_output前端 JSON.parse 偶尔报错在半包数据上做了解析服务端发done后再最终解析预览用JsonOutputParser.parse_partial_json前端没动静突然整段文字冒出来Nginx 缓冲未关X-Accel-Buffering: noproxy_buffering off用户关闭页面后端还在疯狂调用模型没监听断开事件生成器里检查request.is_disconnected()429 限流流式中途断开频率超过配额指数退避重试SSE 场景优先断线续传而非全量重发5.2 实测心得这些坑只有跑线上才会遇到第一心跳不能等模型“有空”。大模型在思考、在调工具那段时间连接上是真的没有任何数据在流动很多网关默认 60 秒没数据就断。我在生成器里加了一个定时器超过 15 秒就发一条: ping\n\n——这是 SSE 规范里的注释行客户端解析时会直接忽略但服务端和网关都知道连接还活着。这个招数解决了一大半“流到一半断掉”的线上事故。第二with_structured_output不是免费的。它默认走 function calling相当于模型先“想想怎么调用”再生成答案比裸输出体系多一点点延迟和 token 消耗。对延迟敏感的产品我倾向于只在“结果要落库”的场景用纯对话展示一律用 StrOutputParser 流式透传。第三Pydantic 字段别全都写成Optional。我见过一张表为了“容错”把字段全设成可选结果校验形同虚设脏数据照样进库。正确的做法是关键字段必填枚举用Literal限制数量字段给ge/le上下限给非关键字段设默认值。第四版本差异是最大的隐性杀手。LangChain 0.2 和 0.3 之间 API 有调整网上教程代码经常对不上。别追新锁定项目依赖版本代码按你实际跑通的版本来写。我的做法是先在 notebook 里跑通最小链再搬到 FastAPI 里接 SSE每层单独验证不然排查问题时会同时面对“是流断了、还是解析挂了、还是框架版本变了”三个变量。结尾的一点私货我自己走完这套链路之后的体会是别一上来就上 LangGraph 那种重框架先把“一次请求 流式传输 结构化产出”这条最小链路跑通再谈多轮工具编排和状态管理。你手上有了一个能发 SSE 事件、能按 mode 切换 StrOutputParser / JsonOutputParser / PydanticOutputParser、能在 ToolCall 与文本吐字之间路由的骨架后面接 Agent、接工作流都是往上加东西的事。最后再分享一个印象最深的小技巧给每条 SSE 事件都带上自增的id这个字段平时用不上但一旦线上出现“stream disconnected before completion”这类中断它就是断点续传的救命稻草——记录最后收到的 id重连时从那个位置继续比重新花一次 token 便宜太多了。
返回列表