ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

LangChain流式输出与结构化返回实战:SSE与OutputParser协同指南

LangChain流式输出与结构化返回实战:SSE与OutputParser协同指南 1. 为什么流式输出和结构化返回总是打架做 GenAI 应用做了这几年我最大的感受是流式输出和结构化返回天然就是一对矛盾体。大模型默认吐出来的是自然语言你要它一段一段地流式返回给前端打字机效果又要它在最后给出一个能被程序直接校验、入库、调用工具的 JSON 结构这中间如果不做设计十有八九会翻车。先说一个最典型的场景你用 FastAPI 起了一个 Agent 服务前端用打字机效果展示大模型的回答同时后端还要把是否调用了某个工具工具参数是什么最终结果字段有哪些解析出来用于日志审计和业务判断。如果只简单地把模型输出整段塞给前端前端拿到的是一堆夹杂着废话的纯文本如果只等全部生成完再一次性返回流式体验就没了。SSEServer-Sent Events恰好是这两者之间的桥——它能保证文本一块一块地推给前端又能在每个事件块里携带结构化元数据。这篇文章我就完整复盘一下怎么用 LangChain 的三大 OutputParser 配合 ToolCall在 SSE 流式场景下既保体验、又保结构化。这篇文章适合谁看已经在用 LangChain 写 Agent、但被流式解析折磨过的后端工程师或者正准备把 LLM 能力封装成标准 API 服务、又不想丢掉打字机效果的前后端同学。我会把每一步的取舍和原理都讲清楚不是让你照抄代码而是让你下次遇到类似需求时能自己拍板选型。2. 先搞清楚 SSE 在这个场景里到底扮演什么角色2.1 SSE 和 WebSocket为什么聊天场景常选 SSE很多人一想到实时推送就默认 WebSocket但在大模型应用里SSE 往往是更务实的选择。SSE 是单向的服务器往客户端推数据客户端不需要也不应该频繁回传。这和 LLM 生成的场景天然匹配——用户提问之后剩下的就是模型一直说前端一直听。协议层面SSE 就是普通的 HTTP 响应Content-Type 设为text/event-stream然后在响应体里按固定格式写事件event: message data: {type: token, content: 你} event: message data: {type: token, content: 好} event: done data: {type: done, session_id: abc123}每个事件之间用空行隔开data字段是真正的载荷event是事件类型id可以用来做断点续传retry是客户端自动重连的时间间隔。就这么简单。WebSocket 呢它是全双工要握手、要维护长连接心跳、要处理断线重连逻辑在只需要服务器单向推送的场景里属于杀鸡用牛刀。更现实的问题很多企业内网的网关、Nginx 配置对 WebSocket 的升级请求支持不友好但对 SSE 这种普通 HTTP 长响应基本上零成本透明转发。我在实际项目里用 SSE 遇到过的最大坑反而很简单网关的超时时间设太短模型思考超过 60 秒连接直接被掐断前端就收到一个stream disconnected before completion: idle timeout waiting for sse。这个问题后面在排查章节我会详细说。2.2 从 LangChain 的流式机制到 SSE 事件拼装LangChain 从Runnable体系开始把流式能力统一成了stream/astream两个接口。chain.stream(query)会一块一块地yield输出。需要注意的是这个一块不一定是模型吐的一个 token而是 LangChain 每个Runnable步骤产出的一个完整单元。比如Retriever步骤吐出文档列表LLM步骤吐出 token 片段。在 Agent 场景下中间可能还夹着Tool的执行结果。所以设计 SSE 接口时我的做法是先约定一套内部事件协议而不是直接把 token 裸推出去。每个事件我用 JSON 串里面至少带两个字段type和payload。event: message data: {type: start, payload: {session_id: uuid}} event: message data: {type: token, payload: {content: 你好}} event: message data: {type: tool_call, payload: {name: search_news, args: {keyword: 人工智能}}} event: message data: {type: tool_result, payload: {result: ..., duration_ms: 1200}} event: message data: {type: structured, payload: {title: ..., summary: ...}} event: message data: {type: done, payload: {finish_reason: stop}}前端只需要根据type决定怎么渲染token追加到正文tool_call可以展示正在调用工具的动画structured存到表单里做后续业务。这样流式体验和结构化数据就各归其位了。后面我讲的三大 OutputParser本质都是在最后那个结构化事件这个环节里保证你拿到的payload是干净、可校验的 JSON。3. 三大 OutputParser把模型的话转成程序能用的结构OutputParser 在 LangChain 里的定位是从 LLM 的原始输出里抽取出程序需要的结构。注意它通常不是魔法它靠的是提示词约束 格式校验两步走。让模型按指定格式输出再用解析器校验、纠错。理解了这一点你就能明白为什么换个模型解析失败中文环境下 JSON 不标准这类问题会反复出现了。3.1 PydanticOutputParser给 JSON 上一份类型保票PydanticOutputParser是项目里最常用的一个。它结合了 Pydantic 的BaseModel把输出格式约束成明确的字段类型——字符串、整型、列表、嵌套对象类型不对直接校验失败。用法上核心三步定义模型类、创建解析器、把格式指令塞进提示词。from typing import List from pydantic import BaseModel, Field from langchain.output_parsers import PydanticOutputParser class ArticleSummary(BaseModel): title: str Field(description生成的文章标题不超过20字) keywords: List[str] Field(description3-5个关键词) summary: str Field(description100字以内的核心摘要) confidence: float Field(description模型对摘要质量的自信度0到1之间) parser PydanticOutputParser(pydantic_objectArticleSummary) prompt PromptTemplate( template请分析下面这段文本并严格按照格式要求输出。\n文本{text}\n{format_instructions}\n, input_variables[text], partial_variables{format_instructions: parser.get_format_instructions()}, )parser.get_format_instructions()生成的那段指令本质是把你定义的字段、类型、约束翻译成自然语言模板告诉模型必须输出一个 JSONkey 有哪些value 是什么类型。它还会补充一句不要输出其他内容之类的强调。实际效果上模型越强GPT-4 级别、Claude 3.5遵循度越高小模型经常把 JSON 包在 Markdown 代码块里或者多解释一句。解析器的parse方法还内置了纠错能力如果模型输出的文本可以被eval成 JSON但类型不对它会尝试用 LLM 自动修复。不过这个修复是有损的——它需要额外调一次模型速度和成本都要考虑。所以在流式场景里我通常不在最后阶段用自动修复而是做两段式先流式展示再后端静默校验校验失败才触发修复。3.2 StructuredOutputParser轻量到不需要定义类StructuredOutputParser适合那种不想建 Pydantic 模型只要几个简单字段的场景。它通过ResponseSchema列表来声明字段名、类型和描述使用起来比 Pydantic 版更轻。from langchain.output_parsers import StructuredOutputParser, ResponseSchema response_schemas [ ResponseSchema(nameanswer, description对问题的直接回答, typestring), ResponseSchema(namesource, description答案的参考来源如果没有则为null, typestring), ] parser StructuredOutputParser.from_response_schemas(response_schemas)从源码实现看StructuredOutputParser内部并没有把输出严格转成 Pydantic 对象而是返回一个字典。它对字段顺序、缺失字段的处理更宽松底层用的其实是类似正则 字典提取的简易逻辑。所以它适合内部接口、日志记录、字段不多的场景一旦你的下游真的要用强类型做入参校验还是 Pydantic 版更省心。这里有个经验如果返回结构里嵌套层级很深或者字段会因为模型输出风格波动我建议直接用 Pydantic 版不要在 Structured 版上强行造轮子。轻量方案省下的代码量会在排障时加倍还回去。3.3 JsonOutputParser最容易被低估的那个第三个其实是JsonOutputParser—— LangChain 里专门用来只要 JSON不要类定义的解析器。它和PydanticOutputParser最大的区别是不要求目标类型是 Pydantic 模型你用普通 dict 声明一个期望的 JSON 结构模板它就能照这个模板校验。from langchain.output_parsers import JsonOutputParser parser JsonOutputParser() prompt PromptTemplate( template输出JSON格式结果字段包括: title(string), items(array of string)。\n{format_instructions}\n, input_variables[], partial_variables{format_instructions: parser.get_format_instructions()}, )它做的事情是拿到模型文本提取并解析出 JSON 对象。它不会去做严格类型强转保持了 dict 的灵活性。实际项目中我经常把它作为流式过程中增量解析 JSON 的工具——不是等模型完整输出后一次性解析而是配合字节流每次拿到新的 token 片段就去试着解析能解析出部分字段就先缓存。这个思路在长回答场景里特别有用你可以在模型还在生成正文时就把标题、关键列表等轮廓字段提前推给前端。总结一下三者的选择逻辑要强类型校验、下游要严格入库选 Pydantic只要几个字段、随拿随用选 Structured既要 JSON 又不想绑定模型定义、或需要在流中做增量解析选 Json。没有绝对好坏只看约束强度。4. ToolCall 方案让模型把工具意图直接交出来4.1 为什么不用让模型自己拼工具调用文本早期 LangChain 的 Agent 实现里模型是用纯文本的方式假装调用工具——输出一行Action: search_news\nAction Input: 人工智能然后 AgentExecutor 去解析这段文本。这种方案在模型能力弱的时候还算勉强能用但问题很明显模型一旦在文本里多加一句解释、少写一个换行整个解析就崩了。而且文本格式因模型而异换模型就要调解析规则。现在主流方案是ToolCall函数调用。模型在生成时除了输出自然语言还可以输出一个结构化的工具调用意图——方法名、参数 JSON。OpenAI 的 function calling、Claude 的 tool use、国产模型不少也兼容这个协议。LangChain 里的做法是把工具定义绑定到模型上这一类模型能力称为bind_tools。from langchain_openai import ChatOpenAI from langchain_core.tools import tool tool def search_news(keyword: str, limit: int 5) - list: 搜索新闻资讯keyword为关键词limit为返回条数。 # 这里写真实检索逻辑 return [{title: 示例新闻, url: https://example.com}] llm ChatOpenAI(modelgpt-4o-mini, temperature0) llm_with_tools llm.bind_tools([search_news]) response llm_with_tools.invoke(帮我搜一下今天人工智能领域的新闻)关键在于response是一个AIMessage如果模型决定调用工具它的tool_calls属性里会带上结构化调用信息response.tool_calls # [{name: search_news, args: {keyword: 人工智能, limit: 5}, id: call_xxx}]这个id字段很重要尤其是做并发工具调用时它用来关联工具结果和对应的调用请求。4.2 ToolCall 事件在 SSE 里怎么推既然AIMessage.tool_calls是结构化的那么从流式事件角度它也能流式地分片到达——模型先生成工具名再一点一点生成参数 JSON。LangChain 的astream_events可以让你捕获on_chat_model_stream事件从而拿到 token 级的流。但我要提醒你工具参数这种 JSON前端完全没必要做打字机效果。你只要在tool_call开始事件里推一条正在调用工具等工具结果出来再推一条结构化结果就行参数 JSON 在中间过程可以直接攒在后端。这是我的实践结论——不要一上来就把所有 token 都推给前端做逐字渲染那样只会让前端渲染逻辑又复杂又容易出错。工具执行完你拿到的结果同样建议包装成结构化事件推给前端。同时把工具结果作为新的上下文消息再喂回给模型让它基于结果生成最终回答。这个模型→工具→模型的循环如果自己用for循环写很容易在异常分支和超时控制上出问题——这也是为什么存在 LangGraph 这类带状态编排的框架。不过对于单轮工具调用场景手动循环完全可控不需要上重型框架。4.3 OutputParser 和 ToolCall 怎么配合这就是这个方案的精髓了ToolCall 解决模型要调用什么工具、参数是什么的结构化提取OutputParser 解决模型最终要返回给业务的最终结论的结构化提取。两者是在一次请求的不同阶段各司其职。class FinalAnswer(BaseModel): reply: str Field(description面向用户的最终回答) used_tools: List[str] Field(description本次实际使用到的工具名称列表) data_source: List[str] Field(description参考信息的来源列表) final_parser PydanticOutputParser(pydantic_objectFinalAnswer)流程大致是用户提问 → 模型决定调用工具ToolCall 结构化→ 执行工具 → 把工具结果拼进上下文 → 模型生成最终回答OutputParser 结构化→ 通过 SSE 推送给前端。中间环节的结构化靠tool_calls最后的业务结构靠 OutputParser。两条线泾渭分明谁也不会干扰谁。5. FastAPI LangChain 完整落地一条 SSE 接口打通全流程5.1 服务端异步流式接口的分层设计我强烈建议把模型调用逻辑和HTTP 流式协议分开。模型调用逻辑是一个普通的异步生成器它只负责产出结构化事件字典HTTP 层只负责把事件字典按 SSE 协议编码。这样拆开你可以对模型逻辑单测不用每次起服务。from fastapi import FastAPI from fastapi.responses import StreamingResponse import json, asyncio from langchain_openai import ChatOpenAI from langchain_core.output_parsers import PydanticOutputParser from pydantic import BaseModel, Field from typing import List app FastAPI() llm ChatOpenAI(modelgpt-4o-mini, temperature0.3, streamingTrue) class FinalAnswer(BaseModel): reply: str Field(description面向用户的最终回答) keywords: List[str] Field(description3-5个关键词) parser PydanticOutputParser(pydantic_objectFinalAnswer) async def event_generator(prompt: str): # 1. 自动构建带格式约束的提示词 formatted_prompt ( 请回答用户问题并输出严格JSON。\n问题{q}\n{fmt}\n ).format(qprompt, fmtparser.get_format_instructions()) # 2. 第一个事件告知开始 yield { event: message, data: json.dumps({type: start, payload: {time: asyncio.time()}}, ensure_asciiFalse) } # 3. 流式输出 token 事件 collected async for chunk in llm.astream(formatted_prompt): collected chunk.content yield { event: message, data: json.dumps({type: token, payload: {content: chunk.content}}, ensure_asciiFalse) } # 控制推送节奏避免瞬间把积压的token全倒出去 await asyncio.sleep(0) # 4. 结构化解析并推给前端 try: parsed parser.parse(collected) yield { event: message, data: json.dumps({type: structured, payload: parsed.model_dump()}, ensure_asciiFalse) } except Exception as exc: yield { event: message, data: json.dumps({type: parse_error, payload: {error: str(exc)}}, ensure_asciiFalse) } # 5. 结束事件 yield { event: message, data: json.dumps({type: done, payload: {finish_reason: stop}}, ensure_asciiFalse) } app.post(/chat/stream) async def chat_stream(payload: dict): prompt payload.get(prompt, ) return StreamingResponse( event_generator(prompt), media_typetext/event-stream, headers{ Cache-Control: no-cache, Connection: keep-alive, X-Accel-Buffering: no, # 重要禁止Nginx缓冲 }, )X-Accel-Buffering: no这个 header是很多人在 Nginx 反代下 SSE 不流式的元凶。Nginx 默认会缓冲响应攒满 4KB 或者等连接结束才发给前端你明明在服务端yield了前端却半天没动静。加了这个 header 就是明确告诉 Nginx 别缓冲。5.2 前端事件分发与渲染解耦前端用fetch配合ReadableStream解析 SSE 就够了不一定非要引eventsource-parser这类库但引了确实省事。核心逻辑是读到一行data:解析 JSON按type走不同的渲染函数。async function streamChat(prompt) { const resp await fetch(/chat/stream, { method: POST, headers: { Content-Type: application/json }, body: JSON.stringify({ prompt }), }); const reader resp.body.getReader(); const decoder new TextDecoder(); let buffer ; while (true) { const { value, done } await reader.read(); if (done) break; buffer decoder.decode(value, { stream: true }); // SSE事件以空行分隔 let sepIndex; while ((sepIndex buffer.indexOf(\n\n)) ! -1) { const rawEvent buffer.slice(0, sepIndex); buffer buffer.slice(sepIndex 2); const dataLine rawEvent.split(\n).find(line line.startsWith(data:)); if (!dataLine) continue; const msg JSON.parse(dataLine.slice(5).trim()); handleEvent(msg); } } } function handleEvent(msg) { switch (msg.type) { case token: appendText(msg.payload.content); break; case structured: fillMetaPanel(msg.payload); break; case tool_call: showToolIndicator(msg.payload.name); break; case done: stopLoading(); break; } }注意千万别直接用浏览器的原生EventSource—— 它只支持 GET 请求而我们往往需要 POST 传递 prompt。原生EventSource也没有自定义 header 的能力鉴权都麻烦。用fetch流式读取是最通用的方案。5.3 一边流式一边结构化增量解析的实践前面提到的JsonOutputParser增量解析在实际项目中可以这样用每收到一段新 token就把collected追加后尝试parser.parse如果解析成功哪怕还不完整只要能出部分字段就把部分结果推给结构化预览事件如果解析失败因为 JSON 还没闭合忽略即可不算错误。这个模式我用来解决一个具体的痛点用户问帮我总结这份文档并给出三个要点模型正文还没写完我希望前端右侧栏已经先把要点标题渲染出来。虽然严格说最终结果要以最后完整解析为准但增量预览的体验提升非常明显。代价是每次追加 token 都会触发一次 JSON 解析token 非常长时会有少量 CPU 开销。我实测下来对普通问答长度的文本这个开销可以忽略。6. 常见问题与排查技巧实录6.1 SSE 流中途断开空闲超时是头号杀手提示stream disconnected before completion: idle timeout waiting for sse这类报错绝大多数不是代码问题是链路中的代理/网关配置问题。我遇到过最典型的三层排查顺序第一层本地测试。先用curl -N直接打服务接口观察事件是不是正常持续输出。curl -N能实时打印服务器推来的每个事件如果这一步正常问题就不在后端。第二层查反向代理。Nginx 的proxy_read_timeout默认 60 秒模型思考时间一旦超过代理直接断连。调大或用proxy_read_timeout 300s可以缓解。另外确认proxy_buffering off;或X-Accel-Buffering: no已生效。第三层查云厂商网关。很多云负载均衡器对长连接也有空闲超时限制比如 60 秒、120 秒。尽量用 WebSocket 或 SSE 都能走的长连接配置同时后端在流式传输过程中即使没有数据也要定期发一个: ping注释行作为心跳。SSE 规范里以冒号开头的行是注释客户端会忽略它但能刷新代理的空闲计时器。async def keepalive(): while True: yield : ping\n\n await asyncio.sleep(15)把这个生成器和主事件生成器用asyncio.gather合并就能在模型长时间思考时维持连接活性。6.2 OutputParser 拿到半截 JSON或 Markdown 代码块模型输出里常见的脏格式有两种一是把 JSON 藏在json代码块里二是前后夹带解释文字。PydanticOutputParser本身会尝试从文本里提取 JSON 块但并不可靠。我的兜底方案是写一个基础清洗函数在喂给解析器之前先做预处理。import re, json def extract_json_string(text: str) - str: text text.strip() # 去掉首尾的 markdown 代码块标记 code_block_pattern re.compile(r(?:json)?\s*(.*?)\s*, re.DOTALL) match code_block_pattern.search(text) if match: return match.group(1) # 尝试从第一个 { 到最后一个 } 截取 start text.find({) end text.rfind(}) if start ! -1 and end ! -1 and end start: return text[start:end1] return text然后统一走parse。注意如果清洗后还是解析失败再去触发 LLM 原文本修复。一定不要默认让每个失败都走修复不然成本和延迟都会失控。6.3 并发请求下事件错乱上下文变量与队列如果你的 FastAPI 服务同时处理多个 SSE 会话每个会话的生成器是独立的理论上不会串。但我踩过一个实际的坑在生成器内部用了模块级的全局变量缓存工具结果两个用户同时触发同一个工具调用时A 用户的结果可能被 B 用户覆盖。解决方案很简单每个会话的事件生成器必须是自包含的所有状态都放在生成器内部不要依赖模块级可变对象。需要跨函数传状态就用contextvars或者干脆把 session_id 作为 key 放进一个字典管理队列。我在项目里用的模式是每个会话一个asyncio.Queue生成器往队列放事件SSE 层从队列取事件编码输出。这个抽象能让你在后续扩展多 Agent 编排时游刃有余。7. 一些我踩过坑之后的固定习惯先说工具声明。LangChain 的tool装饰器会读取函数的 docstring 和类型注解来生成工具的 schema。docstring 里的描述、参数的类型提示、默认值都会成为传给模型的 tool schema 的一部分。所以我在写工具函数时会强制自己把每个参数的单位、取值范围、边界情况写进 docstring这不是为了写注释好看而是直接决定模型能不能正确填参。一个只写 keyword: 搜索关键词 的工具和一个写着 keyword: 搜索关键词最长20字符不要带引号 的工具在模型调参准确率上差很多。然后是解析器的temperature。做结构化输出时模型温度不建议设太高0 到 0.3 之间最稳。温度高了模型更容易发挥创造力去改格式、加注释这对我们的 JSON 解析是灾难。如果既要创意又要结构化我一般拆两条链一条低温度出结构化摘要一条高温度润色成自然语言回复最后再把两段结果拼进 SSE 事件里。最后再分享一个小技巧给事件协议加一个trace_id字段。每个 SSE 会话生成一个trace_id在start事件里发给前端同时在服务端日志里打出来。前端报 bug 时直接甩这个 ID 给你你能在日志里把整条链路还原出来省掉大量你刚才问的什么来着的沟通成本。我在生产环境靠这个字段排查过很多偶发问题尤其是模型偶发出错那种玄学问题有 trace_id 才能对上号。这套方案跑稳之后你会发现流式体验和结构化返回其实不是二选一关键是让它们各走各的通道文本走 token 事件结构化走独立事件中间用统一的 JSON 协议串联。理解了这层设计后续换模型、加工具、上 LangGraph 编排都不会再被到底是文本还是 JSON这个问题卡住。
返回列表