ARTICLE DETAIL

资讯详情

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

SSE流式输出与LangChain结构化解析实战

SSE流式输出与LangChain结构化解析实战 1. 从打字机效果说起为什么流式输出不是锦上添花第一次做 AI 对话产品的人十有八九会卡在同一个体验问题上用户点下发送按钮界面转圈转了五六秒然后啪地一下整段回答全冒出来。功能是通的但用起来就是别扭——像打电话时对方憋了半分钟才一口气把话说完。打字机效果解决的正是这个等待焦虑。它的本质不是动画而是把一次完整的响应拆成很多个小块边生成边推送到前端。用户看到第一个字的时间从 5 秒缩短到 300 毫秒主观感受完全不同。这背后依赖的核心技术就是SSEServer-Sent Events一种基于 HTTP 的单向流式推送协议。但流式输出带来的麻烦也是实打实的。一旦响应变成碎片化的前端拿到的就不再是一个干净的 JSON而是一堆需要拼接、需要解析、需要处理边界情况的文本片段。更麻烦的是如果你还想让模型输出结构化数据比如固定字段的 JSON流式解析的难度会陡增——因为一个 JSON 对象在被完整生成之前本身就是语法不合法的。这篇内容就是围绕这条链路展开的从 SSE 的底层原理讲起到 LangChain 的结构化输出怎么和流式结合再到前端怎么把碎片拼成可用的对象。适合正在做 AI 应用、被流式解析折磨过的开发者也适合刚接触 LangChain、想搞清楚流式到底怎么落地的入门者。我会尽量把每一步的为什么讲透而不是只丢一段能跑的代码。2. SSE 到底是怎么把数据流过来的2.1 一次 HTTP 请求为什么能持续吐数据很多人对 HTTP 的直觉是请求-响应-结束一次交互就完事了。SSE 打破的正是这个直觉它用一次 HTTP 请求让服务端保持连接不关闭持续往客户端写数据。关键在响应头。服务端返回时带上这几个字段Content-Type: text/event-stream Cache-Control: no-cache Connection: keep-alivetext/event-stream告诉浏览器这不是普通文本是事件流no-cache防止中间层缓存住流keep-alive维持长连接。只要连接不断服务端就可以一直往里写。数据格式也有讲究每条消息长这样data: 你好 data: 我是 data: 流式返回的内容注意每条data:后面跟一个换行消息之间用空行分隔。这个空行是硬性规定——它是 SSE 协议里一条消息结束的标志。我见过不少人自己手写流式接口时忘了这个空行结果前端死活收不到消息排查半天才发现是格式问题。2.2 SSE 和 WebSocket 的选择题经常有人问既然要实时推送为什么不用 WebSocket答案在于需求方向。SSE 是单向的——服务端推给客户端客户端不能通过这条通道回推。而 AI 对话场景里用户发消息走的是普通的 POST 请求模型返回走的是流式推送两者本来就是分开的。这种一问一答、答案流式返回的模式SSE 天然契合。WebSocket 是全双工的能力更强但代价是要维护一套独立的连接协议、心跳、重连逻辑。对于服务端单向推流这个场景属于杀鸡用牛刀。而且 SSE 有个被低估的优势它跑在标准 HTTP 上能直接复用现有的鉴权、网关、负载均衡设施运维成本低得多。维度SSEWebSocket通信方向服务端到客户端单向双向协议基础HTTP独立协议握手用 HTTP自动重连浏览器原生支持需自行实现鉴权复用 HTTP 头/Cookie需自定义握手逻辑适用场景流式输出、通知推送实时协作、游戏、聊天选型结论很清晰AI 流式输出优先用 SSE除非你有明确的客户端回推需求。2.3 那些让人抓狂的 SSE 报错实际跑起来SSE 最常见的两个报错几乎每个做流式的人都遇到过。第一个是stream disconnected before completion: idle timeout waiting for sse。这个报错的意思是连接在流还没结束时被断开了原因是空闲超时。它通常不是代码 bug而是中间某一层反向代理、网关、负载均衡设置了空闲超时时间比如 60 秒。如果模型思考时间过长60 秒内一个字节都没吐出来这层就会认为连接死了直接掐断。解决办法有几个方向一是让服务端定期发送心跳注释SSE 支持以:开头的注释行客户端会忽略但能保活比如每 15 秒发一个: keep-alive二是调大中间层的超时配置三是确保模型侧尽快吐出第一个 token。我个人的经验是心跳 调超时双管齐下最稳只做其中一个总会在某些网络环境下翻车。第二个常见问题是消息被合并或截断。SSE 的解析必须严格按空行分隔来切分如果前端用fetch拿到的是原始字节流需要自己按\n\n切分缓冲区。很多人直接对每个 chunk 做JSON.parse结果因为一个 JSON 被拆在两个 chunk 里而报错。正确做法是维护一个字符串缓冲区每次收到新数据就追加然后按分隔符切出完整消息剩下的留在缓冲区等下一块。提示SSE 的data:字段如果内容本身包含换行需要拆成多个data:行客户端解析时再拼回一个换行。这个细节在返回 Markdown 或代码块时特别容易踩坑。3. LangChain 结构化输出让模型吐出能直接用的 JSON3.1 为什么让模型返回 JSON没那么简单直接跟模型说请返回 JSON 格式大部分时候它能照做但总有一些时候会出岔子多包了一层 Markdown 代码块、字段名拼错、该是数字的地方给了字符串、甚至前面加一句好的以下是结果。这些在人类看来无伤大雅但对程序来说就是解析失败。LangChain 的结构化输出Structured Output解决的正是这个最后一公里问题。它的思路是用 schema 约束模型的输出并在解析层做校验和重试。你定义一个 Pydantic 模型或者 JSON SchemaLangChain 负责把它转成模型能理解的格式指令拿到结果后再用 schema 校验不合法就触发重试或报错。from pydantic import BaseModel, Field from langchain_openai import ChatOpenAI class PersonInfo(BaseModel): name: str Field(description人物姓名) age: int Field(description年龄整数) skills: list[str] Field(description技能列表) llm ChatOpenAI(modelgpt-4o-mini) structured_llm llm.with_structured_output(PersonInfo) result structured_llm.invoke(介绍一下张三30岁会 Python 和 Go) print(result) # PersonInfo(name张三, age30, skills[Python, Go])这段代码的价值在于result直接就是一个PersonInfo对象result.age是int类型可以直接参与计算不需要你手动json.loads再取字段。3.2 底层是怎么逼模型输出合法 JSON 的不同模型厂商的实现路径不一样理解这一点对排查问题很有帮助。一类是基于函数调用Function Calling / Tool Calling。模型本身支持调用工具的能力LangChain 把 schema 包装成一个工具定义传给模型模型返回的就不是自由文本而是结构化的工具调用参数。这种方式最可靠因为模型在训练时就针对这种格式做过优化。另一类是基于提示词约束 输出解析器。对于不支持函数调用的模型LangChain 会把 schema 转成一段格式说明塞进提示词然后用PydanticOutputParser之类的解析器去解析返回文本。这种方式对模型的指令遵循能力要求更高偶尔需要重试。还有一类是厂商原生的结构化输出模式比如某些模型提供的 JSON mode在解码层面就限制了输出必须是合法 JSON。这种方式最稳但依赖具体模型支持。我在选型时的经验是能用函数调用的就用函数调用稳定性明显高一个档次。如果模型不支持再退到提示词约束方案并且一定要配重试。3.3 结构化输出和流式的天然矛盾到这里问题来了结构化输出要求最终是一个完整的、合法的 JSON而流式输出的本质是碎片化的、中间态不合法的文本。这两者怎么共存矛盾的核心在于一个 JSON 对象{name: 张三, age: 30}在被逐字生成的过程中会经历{name: 张、{name: 张三, age:这些语法不完整的中间状态。如果你在流式过程中对每个片段做json.loads必然报错。所以流式 结构化输出需要一套增量解析机制。这也是下一节要重点拆解的内容。4. 流式场景下的 JSON 增量解析实战4.1 增量解析的核心思路能解析多少解析多少增量解析Incremental Parsing的目标是每收到一段新文本就尝试从当前累积的缓冲区里提取出已经完整的部分而不是等整个 JSON 结束。举个直观的例子。假设模型在流式生成一个数组[{name: 张三, age: 30}, {name: 李四, age: 25}]当缓冲区累积到[{name: 张三, age: 30}时虽然整个数组还没闭合但第一个对象已经完整了。增量解析器应该能识别出这一点把第一个对象吐出来然后继续等第二个。实现这个能力靠的是流式 JSON 解析库。Python 里比较常用的是ijsonJavaScript 里可以用stream-json或者jsonparse。它们的共同特点是不要求输入是完整 JSON而是以事件流的方式逐个吐出解析到的 token。import ijson # 假设 buffer 是不断累积的字符串 def parse_incremental(buffer): parser ijson.items(buffer, item) for obj in parser: yield obj # 每解析出一个完整对象就 yield不过ijson对边收边解析的支持需要配合生成器实际用起来要处理解析到一半遇到不完整输入的异常。更省心的做法是用专门为流式设计的库比如partial-json-parser这类它能容忍不完整的 JSON 并返回当前能解析出的部分。4.2 用 LangChain 的流式事件拿到结构化片段LangChain 在流式 结构化输出这块提供了事件流接口。以astream_events为例你可以监听模型输出的每个 chunkasync for event in structured_llm.astream_events(input_text, versionv2): if event[event] on_chat_model_stream: chunk event[data][chunk] # chunk.content 是当前片段 print(chunk.content, end, flushTrue)这里有个关键点当使用结构化输出时流式返回的 chunk 内容往往是 JSON 的片段而不是自然语言。你需要把这些片段累积起来再用增量解析器处理。我的做法是维护一个buffer字符串每次收到 chunk 就追加然后调用增量解析函数尝试提取完整对象。提取出来的对象立刻推给前端或做后续处理缓冲区里保留还没解析完的部分。buffer async for event in structured_llm.astream_events(input_text, versionv2): if event[event] on_chat_model_stream: buffer event[data][chunk].content # 尝试从 buffer 中提取完整对象 for obj in try_extract_complete_objects(buffer): yield obj4.3 前端怎么接住这些碎片前端这一侧如果用fetch消费 SSE核心是按\n\n切分缓冲区。下面这段逻辑我用了很多次基本能覆盖绝大多数情况const response await fetch(/api/chat, { method: POST, headers: { Content-Type: application/json }, body: JSON.stringify({ message }) }); const reader response.body.getReader(); const decoder new TextDecoder(); let buffer ; while (true) { const { done, value } await reader.read(); if (done) break; buffer decoder.decode(value, { stream: true }); const parts buffer.split(\n\n); buffer parts.pop(); // 最后一段可能不完整留到下次 for (const part of parts) { const line part.trim(); if (!line.startsWith(data:)) continue; const data line.slice(5).trim(); if (data [DONE]) return; try { const obj JSON.parse(data); handleChunk(obj); } catch (e) { console.warn(解析失败可能是分片不完整, data); } } }这里有几个容易踩的坑我逐个说。第一decoder.decode(value, { stream: true })里的stream: true不能省。因为一个多字节字符比如中文可能被拆在两个 chunk 的边界上不加这个参数会解码出乱码。第二buffer.split(\n\n)之后一定要pop()出最后一段。因为最后一段很可能是半条消息直接解析会失败。把它留在缓冲区等下一块数据来了再拼。第三[DONE]这种结束标记是很多流式接口的约定收到它就该停止读取。但要注意不同服务端的结束标记可能不一样有的用[DONE]有的用空 data得按实际接口来。4.4 打字机效果和结构化解析怎么同时做这是最容易被问到的组合场景既要打字机效果又要结构化数据。我的建议是分两条路走。一条路负责展示把原始文本片段直接追加到界面上形成打字机效果另一条路负责数据把片段喂给增量解析器提取出结构化对象用于业务逻辑。let displayText ; let jsonBuffer ; function handleChunk(chunk) { // 路径一展示用直接拼接 displayText chunk.text; renderTypewriter(displayText); // 路径二数据用增量解析 jsonBuffer chunk.text; const objects tryExtractCompleteObjects(jsonBuffer); objects.forEach(obj updateBusinessData(obj)); }这样做的好处是两条路互不干扰。展示层不需要关心 JSON 是否合法数据层不需要关心打字机动画。如果强行用一条路要么展示层被 JSON 语法绑死要么数据层拿不到实时更新。注意如果模型输出的是自然语言 JSON混合内容增量解析会复杂很多。这种情况建议在提示词里明确要求只输出 JSON不要任何额外文字从源头减少解析难度。5. 一套能落地的完整链路设计5.1 后端FastAPI LangChain 的流式接口后端这块FastAPI 配合StreamingResponse是主流组合。核心是把 LangChain 的异步流包装成 SSE 格式from fastapi import FastAPI from fastapi.responses import StreamingResponse import json app FastAPI() async def event_generator(input_text: str): buffer async for event in structured_llm.astream_events(input_text, versionv2): if event[event] on_chat_model_stream: chunk event[data][chunk].content buffer chunk # 同时推送原始片段和已解析的完整对象 payload {raw: chunk} complete try_extract_complete_objects(buffer) if complete: payload[parsed] complete yield fdata: {json.dumps(payload, ensure_asciiFalse)}\n\n yield data: [DONE]\n\n app.post(/api/chat) async def chat(req: dict): return StreamingResponse( event_generator(req[message]), media_typetext/event-stream )这里我把raw和parsed一起推给前端前端可以按需使用。ensure_asciiFalse很重要否则中文会被转义成\uXXXX虽然不影响解析但可读性差、体积也大。5.2 心跳与超时让长连接活下来前面提到的idle timeout问题在后端加心跳就能解决。做法是在生成器里插入定时任务每隔一段时间 yield 一个注释行import asyncio async def event_generator_with_heartbeat(input_text: str): queue asyncio.Queue() async def produce(): async for event in structured_llm.astream_events(input_text, versionv2): if event[event] on_chat_model_stream: await queue.put(event[data][chunk].content) await queue.put(None) # 结束标记 asyncio.create_task(produce()) while True: try: chunk await asyncio.wait_for(queue.get(), timeout15) if chunk is None: break yield fdata: {json.dumps({raw: chunk}, ensure_asciiFalse)}\n\n except asyncio.TimeoutError: yield : keep-alive\n\n # 心跳注释 yield data: [DONE]\n\nasyncio.wait_for设了 15 秒超时超时就发一个心跳注释。注释行以:开头客户端会忽略内容但连接保持活跃中间层的空闲计时器就被重置了。5.3 前端Vue 里怎么优雅地消费 SSEVue 项目里消费 SSE我推荐用fetchReadableStream而不是EventSource。原因是EventSource只支持 GET 请求没法带复杂的 POST body而 AI 对话通常需要传一大坨上下文。把前面的解析逻辑封装成一个 composableimport { ref } from vue; export function useSSEChat() { const displayText ref(); const parsedData ref([]); const loading ref(false); async function send(message) { loading.value true; displayText.value ; parsedData.value []; const response await fetch(/api/chat, { method: POST, headers: { Content-Type: application/json }, body: JSON.stringify({ message }) }); const reader response.body.getReader(); const decoder new TextDecoder(); let buffer ; while (true) { const { done, value } await reader.read(); if (done) break; buffer decoder.decode(value, { stream: true }); const parts buffer.split(\n\n); buffer parts.pop(); for (const part of parts) { const line part.trim(); if (!line.startsWith(data:)) continue; const data line.slice(5).trim(); if (data [DONE]) { loading.value false; return; } try { const payload JSON.parse(data); if (payload.raw) displayText.value payload.raw; if (payload.parsed) parsedData.value.push(...payload.parsed); } catch (e) { /* 忽略不完整分片 */ } } } loading.value false; } return { displayText, parsedData, loading, send }; }这个 composable 把展示文本和结构化数据分开管理组件里直接绑定displayText做打字机渲染绑定parsedData做业务展示职责清晰。5.4 一个完整的排查清单流式链路出问题时按这个顺序排查效率最高现象可能原因排查方向完全收不到数据响应头不对检查Content-Type: text/event-stream收到数据但前端不显示分隔符不对确认消息间有空行\n\n中文乱码解码方式错TextDecoder加stream: true中途断开空闲超时加心跳注释、调大中间层超时JSON 解析报错分片不完整用缓冲区累积按分隔符切分结构化字段缺失schema 不匹配检查 Pydantic 模型定义和提示词这张表基本覆盖了我遇到过的 90% 的流式问题。每次出问题先对号入座能省下大量瞎试的时间。6. 几个只有踩过才知道的细节6.1 别在流式里做重逻辑流式回调触发非常频繁一个长回答可能触发几百次。如果在每次回调里做数据库写入、复杂计算、DOM 全量重渲染性能会肉眼可见地崩。我的做法是流式回调只做最轻量的缓冲和拼接重逻辑放到流结束后统一处理或者用防抖节流控制频率。6.2 结构化输出的 schema 越简单越稳我试过用嵌套三四层的复杂 schema结果模型经常在深层字段上出错。后来把 schema 拆成多个扁平的小模型分步调用稳定性明显提升。能用扁平结构就别嵌套这是血泪教训。6.3 增量解析要设最大缓冲如果模型输出异常缓冲区可能无限增长。加一个上限比如超过 1MB 就强制清空或报错避免内存泄漏。这个细节平时不显眼但线上出问题时能救命。6.4 测试时一定要模拟慢速流本地开发时模型响应快很多边界问题暴露不出来。我习惯在测试环境人为加延迟比如每个 chunk 之间 sleep 200ms这样能复现超时、分片、乱序等问题。快的时候一切正常慢的时候才是真相。6.5 结束标记要统一约定前后端一定要约定好结束标记[DONE]也好空 data 也好必须一致。我见过因为前后端对结束标记理解不同导致前端一直等、连接不释放的问题。这种问题排查起来特别费劲因为两边单独看都没错。7. 关于这套方案我自己的使用体会这套 SSE LangChain 结构化输出 增量解析的组合我在几个项目里反复用过最大的感受是流式的复杂度不在流本身而在边界。分片的边界、JSON 的边界、超时的边界、编码的边界——每一个边界都是一类 bug 的温床。真正让这套方案稳定的不是某个高级库而是几个朴素的习惯缓冲区永远留一段不解析、心跳永远比超时短、schema 永远往简单了设计、测试永远模拟慢速环境。这些习惯听起来不酷但能让你少熬很多夜。如果你刚开始做流式我的建议是先用最简单的方案跑通——一个StreamingResponse一个fetch循环一个split(\n\n)。跑通之后再逐步加结构化、加心跳、加增量解析。一次性把所有东西堆上去出了问题你根本不知道是哪一层的事。分层验证逐层加固这条路我走过确实最省心。
返回列表