ARTICLE DETAIL

资讯详情

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

SSE流式传输+LangChain结构化输出:实现打字机效果与JSON稳定解析

SSE流式传输+LangChain结构化输出:实现打字机效果与JSON稳定解析 做 LLM 应用最痛苦的事情从来不是模型选型而是用户在前端盯着一个空白的加载框一秒、两秒、十秒过去才看到整段回答一口气蹦出来。我在实际项目里几乎每天都要跟 SSE 和结构化输出打交道这个标题正好把底层传输、模型输出解析、前端渲染三条线串到了一起。把这条链路理清楚你就能做到大模型一边生成、页面一边打字机输出同时还能把模型吐出来的关键数据解析成稳定的 JSON而不是靠正则碰运气。这篇文章我会从 SSE 流式传输原理讲起再落到 LangChain 的流式接口和结构化输出方案最后给出一套能直接抄的前后端代码。内容包括事件流格式、心跳保活、FastAPI 异步生成器、前端 ReadableStream 渲染、部分 JSON 的累积与修复以及我踩过的各种坑。适合正在做 AI 对话、报告生成、数据分析 Agent 的开发者不管你是刚接触 LangChain 还是在补全流式体验都可以直接参考里面的方案。1. 整体设计与思路拆解1.1 为什么“流式”和“结构化输出”必须放在一起考虑很多人的第一版 API 是“请求进来等大模型完整回答解析 JSON返回给前端”。在 demo 阶段这没问题一旦进入真实产品问题就来了大模型生成一段 1000 字的回答需要十几秒如果让用户盯着空白页面等这么久流失率是灾难性的更麻烦的是如果你要提取的结构化字段藏在回答中间等全部生成完再解析任何一个字段格式不对整个响应就废了又得重新调用模型。所以流式和结构化不是两个独立需求而是同一个问题的两面用户既要看到内容“一点一点出来”的实时反馈又希望最终交付的是稳定、可校验的数据结构。标题里把“打字机效果”和“JSON 解析”并列本质就是要解决“增量可见”和“最终可靠”之间的矛盾。这决定了我们不只要学 SSE 的 API 调用还要设计一套输出协议让模型生成的内容从一开始就具备可分割、可累积、可解析的结构。1.2 整体链路从模型 Token 到前端字符我最终落地的架构可以概括为三个环节第一层是模型输出层。LangChain 通过astream方法把大模型的 Token 一个接一个吐出来这一层负责和各家模型厂商打交道屏蔽不同 API 的差异。第二层是服务端转换层。FastAPI 用异步生成器接收 Token按 SSE 协议封装成data:事件流推给前端。同时在这一层完成 JSON 的片段解析和清洗把“对话文本”和“结构化数据”区分成不同的事件类型下发。第三层是前端渲染层。浏览器通过fetch拿到ReadableStream逐行读取 SSE 帧把文本追加到页面上形成打字机效果遇到 JSON 片段时先放入缓冲直到累积成完整 JSON 再渲染成表格或卡片。这个链路的关键在于SSE 是单向管道天然适配“模型推送给前端”的场景LangChain 的流式输出又是标准异步迭代器两者衔接非常自然。你不需要 WebSocket也不需要自己维护连接状态一个 HTTP 连接就能完成所有事情。1.3 方案选型SSE 还是 WebSocket、原生 fetch 还是 EventSource先聊传输层选型。SSEServer-Sent Events基于纯 HTTP服务端通过Content-Type: text/event-stream持续发送数据。它和 WebSocket 的核心区别在于WebSocket 是全双工客户端也能随时推消息给服务端SSE 是单向只能服务端往客户端推。对 LLM 对话这种“客户端发一次请求服务端持续回复”的场景SSE 是最省事的方案。你不需要额外的协议握手、心跳库、断线重连逻辑浏览器原生支持自动重连。但要注意一个限制原生EventSource只支持 GET 请求没法自定义 Header也没法发 POST body。对于需要携带长的 prompt 或需要鉴权的接口这非常不灵活。我的做法是放弃 EventSource直接用fetch读取响应流。fetch和 EventSource 一样能拿到流式数据但不受请求方式限制还方便中断请求和手动重连。后面的代码我都基于这个方案。2. SSE 流式传输原理与落地方案2.1 SSE 协议格式data、event、id、retry 以及被忽略的注释行SSE 的协议格式非常简单但很多人栽在细节上。每个事件由若干行组成字段包括data:、event:、id:和retry:事件之间使用一个空行分隔。最常用的就是data:它的值是事件内容event:可以指定事件类型前端用addEventListener监听对应类型id:用于断线重连时告诉服务端“我上次收到哪条消息”retry:告诉浏览器重连间隔。有一个经常被忽略的用法以冒号开头的行是注释行不会触发前端事件。注释行在实战中非常重要因为很多网关、负载均衡器会判定“连接空闲太久”并主动断开 TCP 连接。你只要每隔 15 秒发一行注释比如: ping\n\n连接就被判定为活跃这就是最实用的心跳保活方案。下面是一个最小完整事件帧event: delta data: {content: 你} event: done data: [DONE]前端收到event: delta就渲染增量字符收到event: done就关闭加载动画。这个格式比所有文本塞在data:里再在前端 split 要规范得多建议从一开始就按事件类型区分“增量文本”“结构化片段”“完成信号”。2.2 FastAPI 实现 SSE 流式接口StreamingResponse 与异步生成器FastAPI 里实现 SSE 非常简单核心就是StreamingResponse接收一个异步生成器。生成器每次yield一段字符串这段字符串就是符合 SSE 格式的完整事件。下面是一个接入 LangChain 的示例from fastapi import FastAPI from fastapi.responses import StreamingResponse from langchain_openai import ChatOpenAI import asyncio, json app FastAPI() llm ChatOpenAI(modelgpt-4o-mini, temperature0.7) async def event_stream(prompt: str): # 先发一个连接建立事件前端收到后立刻清理 loading 状态 yield fevent: start\ndata: {json.dumps({status: ok})}\n\n try: async for token in llm.astream(prompt): # token 是字符串片段作为增量内容下发 yield fevent: delta\ndata: {json.dumps({content: token})}\n\n except Exception as e: yield fevent: error\ndata: {json.dumps({message: str(e)})}\n\n return yield event: done\ndata: [DONE]\n\n app.post(/chat) async def chat(payload: dict): prompt payload.get(prompt, ) return StreamingResponse( event_stream(prompt), media_typetext/event-stream; charsetutf-8, headers{ Cache-Control: no-cache, Connection: keep-alive, X-Accel-Buffering: no, }, )这里有两个细节必须注意。第一个是media_type一定要带charsetutf-8否则前端TextDecoder用 UTF-8 解码时遇到多字节中文字符拆包可能会产生替换字符。第二个是X-Accel-Buffering: no这个头是给 Nginx 看的告诉它不要把这部分响应缓冲到完整再转发。如果你没加这个头Nginx 默认会缓冲后端响应导致前端等了很久才看到一整块数据打字机效果直接消失。2.3 心跳保活与超时问题idle timeout 的成因和解决搜索引擎热词里有一条非常典型“stream disconnected before completion: idle timeout waiting for sse”。这句话从 AWS、Cloudflare、各种网关的日志里都能看到含义是连接在等待数据时空闲超时被服务端主动断开。大模型生成通常需要 5 到 30 秒这期间如果有某个环节比较慢比如工具调用卡住、内部推理时间过长SSE 连接就可能被误杀。解决思路有两个层面。服务端应用层写一个循环每隔 15 秒yield : keepalive\n\n。注释行不产生事件但能证明连接是活跃的。网关配置层如果你用了 Nginxproxy_read_timeout要调大比如 300 秒如果是云厂商的负载均衡也要检查对应的“空闲超时”参数。我见过很多人只调了应用层忘了调网关结果问题依旧。另外要警惕一个伪心跳有些人会发event: pingdata: {}。这虽然也能保活但前端会多出一堆ping事件需要额外逻辑忽略。用注释行才是最干净的方案。2.4 断线重连与幂等客户端如何优雅恢复SSE 原生支持浏览器自动重连但基于fetch的方案没有这个能力需要自己实现。重连逻辑要考虑几个点指数退避第一次断开后等 1 秒第二次 2 秒最多不超过 10 秒重连后提示用户“连接已恢复”而不是静默如果应用场景允许重复消费服务端用id:字段标记序列号客户端在重连时带上Last-Event-ID头避免从零开始。对 LLM 场景还有一个更实用的经验如果连接真的断了与其拼命重连不如直接把已经生成的文本保留然后让用户手动点击“重新生成剩余部分”。大模型生成过程是不可恢复的服务端不知道你之前生成到哪了硬做断点续传的复杂度非常高。保留已有内容、提供重新生成的选项才是现实中可行的方案。3. LangChain 结构化输出与流式输出的矛盾处理3.1 with_structured_output 的原理Pydantic Schema 与底层映射LangChain 的with_structured_output是很多人的首选因为它把“从文本里提取 JSON”这件事彻底封装了。你只需要定义一个 Pydantic 模型from pydantic import BaseModel, Field class Article(BaseModel): title: str Field(description文章标题) summary: str Field(description一句话摘要) keywords: list[str] Field(description关键词列表) tone: str Field(description语气风格) structured_llm llm.with_structured_output(Article) result structured_llm.invoke(帮我写一段关于 AI 的短文)这段代码看起来简单底层却做了几件事LangChain 会把Article模型的字段名、类型、描述转换成 JSON Schema然后根据模型能力选择不同的调用方式。最优先的方式是 tool calling也就是把 Article 模型伪装成一个“工具”要求模型把结果填进去如果模型不支持 function calling就退化为 JSON mode要求模型只输出合法 JSON。最后 LangChain 还会把模型返回的内容校验并转换成 Pydantic 模型实例字段缺失或类型错误会直接抛异常。这套机制在非流式场景非常好用但你在流式请求里调用它会发现一个尴尬的事实with_structured_output默认会等待模型完整生成完毕再一次性返回解析结果。这也就意味着如果你直接用structured_llm.astream()你拿到的大概率不是逐 Token 的增量文本而是最终才冒出来的整个结构体打字机效果直接失效。3.2 流式场景下 JSON 增量两种可靠方案要兼顾流式和结构化就不能只依赖那个封装好的方法得把它拆开。我的方案有两种按场景选择。第一种是“双流分离”方案。如果模型是聊天模型且支持 tool calling你可以用astream_events监听底层事件从on_llm_stream事件里拿到模型增量 token同时对工具调用的参数片段做累积。LangChain 在不同版本里提供的astream_events接口变化比较大但核心思路一样不通过with_structured_output而是手动拼装增量。async for event in llm.astream_events(prompt, versionv2): if event[event] on_chat_model_stream: token event[data][chunk].content if token: yield fevent: delta\ndata: {json.dumps({content: token})}\n\n第二种是“文本与 JSON 分段”方案。让模型按照约定先输出自然语言再输出结构化数据二者用特殊标记分隔。这个方案不依赖模型供应商的流式 tool call 支持兼容性最好。我会在后面第 5 章展开。3.3 对比直接解析文本 JSON 与输出解析器LangChain 还提供了JsonOutputParser很多教程会告诉你它支持流式。我的实际使用体验是它在理想条件下确实能边生成边解析但一旦模型的 JSON 不规整比如字段值里有换行、错误地把单引号当成双引号、末尾多逗号这个解析器就会陷入半坏状态轻微的问题能修复严重的直接崩。它本质上依赖 JSON 的累积解析而这个解析本身对格式错误很敏感。相比之下我更喜欢在后端自己维护一个累积缓冲模型 token 到了就先拼起来不急着解析等收到[DONE]或者解析成功信号再做一次健壮性修复和 Pydantic 校验。这样做的代价是“JSON 字段不会实时显示”收益是“最终结果可校验、可失败重试”。需要实时展示 JSON 结构的情况用第 5 章的分隔标记方案更合适。4. 前端打字机效果从流式数据到逐字渲染4.1 用 fetch 读取 ReadableStream关键代码前端核心是把响应体当成流来处理。原生 EventSource 没法发 POST body所以我用fetchresponse.body.getReader()async function startStream(apiPath, payload) { const res await fetch(apiPath, { method: POST, headers: { Content-Type: application/json }, body: JSON.stringify(payload), }); if (!res.ok || !res.body) { 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) { handleSSEFrame(frame); } } }这里有两个非常重要的细节。第一decoder.decode(value, { stream: true })必须带stream: true否则一个汉字拆成两个字节到达时第二个字节会被丢弃导致乱码。第二SSE 帧以空行结尾但一次网络读取可能包含多个帧也可能只包含半个帧所以必须用缓冲区累积按\n\n切分处理完的帧消费掉剩下的保留给下一次循环。这个“半帧缓冲”处理看似简单但我见过不少前端把data:文本一个字符一个字符地处理完全没考虑跨包问题十次有八次会丢内容。4.2 打字机效果的正确实现不是每个 token 都 setState很多人实现打字机效果时拿到一个 token 就setState一次。这在 React 下会引发大量渲染尤其在内容很长时性能急剧下降。正确做法是维护两个状态fullContent保存完整累计文本displayedContent保存实际渲染出来的文本。每次收到增量先更新fullContent然后通过requestAnimationFrame批量地把displayedContent往fullContent方向推进。也就是说模型生成速度很快时前端每帧渲染的字符数多视觉上就是“一段一段蹦出来”模型生成速度慢时前端每帧渲染的字符数少看起来就是“一个字一个字蹦出来”。这个节奏是自适应的不会出现生成很快但渲染跟不上导致卡死的情况。代码类似function updateDisplay() { if (displayedContent.length fullContent.length) { displayedContent fullContent.slice(0, displayedContent.length 3); renderContent(displayedContent); requestAnimationFrame(updateDisplay); } }把核心逻辑抽象成“累计文本”和“渲染文本”两层比单纯 1:1 渲染 token 要稳得多。4.3 文本增量与 JSON 片段的渲染差异前端要区分两类内容自然语言文本和 JSON 数据块。我的做法是服务端在 SSE 事件里显式区分event: text和event: artifact。文本事件直接追加到对话气泡里artifact 事件里的 JSON 片段先加入一个独立缓冲同时前端实时展示“正在生成结构化数据”的占位效果等到 JSON 完整后整块解析渲染成表格或卡片。这里千万不要试图对半截 JSON 做真正的结构化渲染比如解析一半就渲染成表格再往后不断把表格拆了重建。这会带来巨大的闪烁和性能损耗。我见过有人写了一个几百行的“IncrementalJSONTable”组件试图把半个 JSON 对象渲染成表格效果非常糟糕。正确的姿势是把“结构感知”放到最后一步前面所有阶段只做纯文本的累积。用户看到的是 JSON 文字逐渐增多最终在结束时瞬间变成表格这种体验反而最顺滑。5. 流式 JSON 解析全方案5.1 问题本质半截 JSON 不能直接 parse流式 JSON 解析和普通 JSON 解析最大的区别在于你手里的数据永远不完整。直接JSON.parse(partialString)几乎必然抛异常因为字符串没闭合、括号数量不匹配、对象只有一半。所以任何方案的第一步都是“能忍”先累积不急着解析只有当内容明显完整时才尝试解析。这个策略的难点在于怎么判断“完整”。最简单粗暴的判断方法是直接尝试JSON.parse成功就算完整。对大字符串来说每次 token 到达都 parse 一次性能很差对中等规模 JSON这个开销其实可以接受。更精细的方法是统计花括号、方括号的配对数量并在识别到末尾闭合后确认引号状态正常。但这套状态机要考虑 JSON 字符串内部的转义和大括号逻辑复杂度不低很容易写错。我的建议是先试用try-parse性能不够再升级到状态机不要一上来就追求完美。5.2 实用方案一累积后修复再解析适合可靠场景对可靠性优先的场景后端维护一个JSONAccumulator把所有 token 都追加进去。结束时先用json_repair修复常见错误末尾逗号、单引号、未转义换行等再用 Pydantic 模型校验。import json from json_repair import repair_json class JSONAccumulator: def __init__(self): self.buffer def append(self, chunk: str): self.buffer chunk def is_complete(self) - bool: try: json.loads(self.buffer) return True except json.JSONDecodeError: return False def parse(self): # 先尝试直接解析失败再修复 try: return json.loads(self.buffer) except json.JSONDecodeError: return json.loads(repair_json(self.buffer))这个方案的好处是逻辑简单、可靠缺点是“解析”动作延迟到流结束后才发生无法实时感知结构。适合那些最终目标就是拿一个完整可靠 JSON 的场景比如数据入库、报表生成。5.3 实用方案二分隔标记让文本和 JSON 天然分离适合 Agent 场景这是我在 Agent 项目里最常用的一套。核心是约定模型输出格式【txt】这是给用户看的一段解释性文本 【json】{title: xxx, summary: yyy, keywords: [a, b]}服务端收到 token 后维护两个缓冲text_buffer和json_buffer。当前处于哪个标记模式token 就追加到哪个缓冲。当遇到【json】后后续 token 全部进 JSON 缓冲输出完成后JSON 缓冲内容是完整的一块再交给JSONAccumulator解析。这样天然解决了一个长期痛点模型经常会在 JSON 前后夹杂“这是你要的结果”、“注意 JSON 结尾多了一个逗号”这类废话污染 JSON。用分隔标记之后废话进了文本通道不会干扰解析。你可以在 prompt 里明确告诉模型必须严格使用这两个标记实测下来大多数模型都能遵守个别模型不遵守时服务端加一个兜底如果没找到【json】标记就把整段内容尝试解析解析失败则走重试逻辑。5.4 schema 校验与安全不要让模型输出直接进数据库无论用哪种解析方案最后一道关口都是校验。Pydantic 模型在这里发挥的作用不只是类型转换更是安全网字段缺失会被发现类型错误会抛异常额外的未知字段默认会报错。大模型输出本身是不可信的永远不要直接把模型返回的 JSON 存进数据库或直接传给下游系统。必须先校验、再清洗最后才能入库。还有两个安全问题值得注意。第一不要用eval()或exec()解析 JSONJSON 的true/false/null在 Python 里会被eval成True/False/None还算好但 JSON 里一旦混入恶意表达式后果严重。第二模型可能被 prompt injection 引导在 JSON 字段里插入命令或脚本所有文本字段都应当在上层做好转义和脱敏。5.5 备选方案第三方 partial-JSON 解析器如果你不想自己维护累积逻辑可以用现成的库。Python 生态有json-repair、partial-json-parser专门处理这种不完整的 JSON 流。这些库能在字符串未闭合的情况下尽最大努力返回一个合理的部分结构。我的体验是partial-json-parser适合“边生成边展示”的需求因为你能拿到解析到一半的对象但它对格式错误的容忍度不如json-repair。实际项目中最好两个都引入。前端场景还见过jsonstream这类 JS 库但流式场景在前端保持半成品状态意义不大我更倾向于前端只累积不解析。6. 常见问题与排查技巧实录6.1 问题速查表下面是我在实际项目中整理出的高频问题按出现频率排序现象可能原因解决方案前端等很久才看到内容Nginx 缓冲了响应或服务端未 flush设置X-Accel-Buffering: noStreamingResponse 确保逐段 yieldSSE 连接中途断开idle timeout waiting for sse网关认为连接空闲超时断开每 15 秒发注释行心跳调大网关读超时中文乱码出现media_type缺少charsetutf-8或前端TextDecoder未开stream模式服务端和前端同时修正编码JSON 解析失败末尾缺个}流式结束信号没触发或 token 丢失使用分隔标记确认按\n\n切帧时缓冲区逻辑正确LangChainwith_structured_output在流式下不输出增量该接口默认等待完整输出改用astream_events或自行累积解析前端重复渲染、性能卡顿每 token 都 setState用 accumulated 与 displayed 双状态requestAnimationFrame批量推进6.2 深度排查为什么我的打字机效果时灵时不灵“时灵时不灵”是所有流式问题里最让人抓狂的。排查方法其实非常固定先看网络面板里 SSE 帧是否在持续到达再看服务端日志里yield是否被阻塞最后看前端缓冲切分是否正确。绝大多数案例都出在“服务端 yield 了但没 flush”或“前端把半帧数据丢掉了”这两个位置。FastAPI 的StreamingResponse在每次yield时需要刷新到网络层如果你在生成器里做了大量的 CPU 阻塞操作流式就会被卡住。另一个隐蔽问题是本地开发正常、上线后不正常这几乎都是因为中间网关动了响应缓冲。你可以在前端fetch之外的开发者工具里直接查看curl -N的输出如果curl -N能看到流式数据而浏览器不行问题就在网关或代理配置。6.3 独家经验用超时机制防止模型“无限流式”还有一个很多教程不会讲的坑模型偶尔会输出极其长的内容或者因为内部异常一直不结束。服务端必须给流式生成加一个最大时长限制。我在生成器里用一个asyncio.wait_for或手动记录起始时间超过 90 秒强制中断向前端发送event: error。这个超时时间要根据模型和任务场景来定。简单的问答 30 秒足够复杂分析或带工具调用的 Agent 可能需要 120 秒以上。超时后前端要把已经生成的文本保留并提示“回答超时可重试”。这种做法看似粗暴但能避免用户因为一个无响应的连接挂在那里你的服务端线程也被长期占用。6.4 别忘了断电重连和重复消费最后提醒一个非常隐蔽的坑SSE 客户端重连后可能会重复消费部分事件。如果你用了id:字段前端在Last-Event-ID恢复时服务端通常需要支持按 id 跳过已发送内容。但对 LLM 流式生成来说这个场景不值得做太复杂因为服务端并不知道客户端断线前到底收到了多少 token。简单处理即可重连后重新生成整个响应前端清空旧内容从零开始或者让用户手动选择重新生成。这类问题在真实产品里经常以“用户发现回答重复了一半”“回答内容串线”的形式出现。复现时先抓包看请求头里的Last-Event-ID再对比服务端收到的重连请求时间就能定位到是哪一层出了问题。回到最初的话题。SSE 和 LangChain 的单点知识都不难难点在于把它们串成一条完整链路后各种极端情况会同时涌过来。我个人在实际操作中的体会是先把输出协议设计好再动代码。你希望前端拿到的是“纯自然语言文本”还是“结构化的 JSON”直接在 SSE 事件类型里定下来模型 prompt、服务端解析、前端渲染三段代码都围绕这个协议写后面基本不会别扭。如果一开始各层各写各的解析逻辑后面改协议会牵扯前后端所有代码非常痛苦。最后再分享一个小技巧给所有 SSE 事件加一个request_id字段前端把它和用户会话关联起来。排查问题的时候你只要拿这个 ID 去服务端日志里过滤就能看到整个请求从进入到流式生成到结束的完整链路定位问题的时间能少一半。这是我在多个项目里用下来最值的投资。
返回列表