ARTICLE DETAIL

资讯详情

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

SSE流式输出与LangChain结构化输出实战:增量JSON解析与打字机效果

SSE流式输出与LangChain结构化输出实战:增量JSON解析与打字机效果 1. 为什么流式输出不是锦上添花而是刚需如果你做过大模型应用一定遇到过这种场景用户点下发送按钮界面卡住十几秒然后啪地一下蹦出一大段完整回答。用户在这十几秒里不知道程序是死是活体验极差。而换成流式输出之后文字像打字机一样一个字一个字往外冒用户立刻就能感知到它在工作。这个差别不是视觉上的小修饰而是产品可用性的分水岭。流式输出的技术底座是SSEServer-Sent Events。它本质上是一种基于 HTTP 的单向推送协议服务端可以持续往客户端发送文本片段而不需要客户端反复轮询。相比 WebSocket 的双向通信SSE 更轻、更简单天然适合服务端持续吐字、客户端只负责渲染这种大模型对话场景。但光有 SSE 还不够。大模型吐出来的是一段段自然语言文本而你的业务代码往往需要的是结构化的 JSON 对象——比如从用户提问里提取出姓名、电话、地址或者让模型返回一个可以直接入库的字段集合。这就引出了第二个核心问题结构化输出。LangChain 提供了with_structured_output这类能力让模型直接返回符合 Pydantic 模型定义的 JSON但当你把流式和结构化输出放在一起用的时候坑就来了——流式吐出来的 JSON 是残缺的、不完整的你不能每收到一个 chunk 就去JSON.parse那样必然报错。这篇内容就是围绕这条链路展开的从 SSE 的底层原理讲起到 LangChain 结构化输出的实现方式再到流式场景下 JSON 的增量解析方案最后落到打字机效果的前端渲染。适合已经用过大模型 API、想进一步把流式体验做扎实的开发者也适合刚接触 LangChain、想搞清楚流式 结构化到底怎么配合的入门者。我会把每一步的为什么这么做讲清楚而不是只丢一段代码让你抄。2. SSE 流式原理数据到底是怎么一段段过来的2.1 SSE 的报文格式与传输机制很多人以为流式输出是什么高深的技术其实拆开看非常简单。SSE 的响应体就是一段纯文本服务端按照固定格式往里面写数据客户端按行读取并解析。它的 Content-Type 是text/event-stream每一条消息由若干字段组成字段之间用换行分隔消息之间用空行分隔。一条典型的 SSE 消息长这样data: {content: 你} data: {content: 好} data: {content: 世界}每个data:后面跟的就是一条消息的内容。注意消息之间必须有一个空行这是 SSE 协议规定的分隔符。客户端读到空行就知道上一条消息结束了。除了data字段还有event自定义事件类型、id消息编号用于断线重连、retry重连间隔这几个可选字段。大模型流式输出场景里最常用的就是data偶尔用event来区分正常内容和结束信号。为什么 SSE 能做到持续推送关键在于 HTTP 的chunked transfer encoding分块传输编码。普通 HTTP 响应会带一个Content-Length告诉客户端总共多少字节收完就结束。而流式响应不带Content-Length改用Transfer-Encoding: chunked服务端可以一块一块地写写完一块客户端就能收到一块直到服务端主动关闭连接。这就是打字机效果的物理基础——不是前端在模拟逐字显示而是数据真的是一段段到达的。2.2 为什么大模型场景偏爱 SSE 而不是 WebSocket这里有个常见的疑问WebSocket 不是更强大吗为什么大模型应用普遍用 SSE我实际做过几个项目之后总结下来有几个现实原因。第一大模型对话是典型的单向流。用户发一次请求服务端持续返回内容客户端在这个过程中基本不需要往服务端推数据。WebSocket 的双向能力在这里是浪费的反而增加了连接管理的复杂度。第二SSE 天然走 HTTP基础设施友好。你不需要额外开端口、不需要处理协议升级Upgrade、不需要担心某些网络环境对 WebSocket 的拦截。负载均衡、网关、日志系统对普通 HTTP 的支持都是现成的SSE 直接复用这一套。第三SSE 自带断线重连语义。协议里的id和retry字段就是为断线重连设计的浏览器端的EventSource会自动处理重连。虽然实际项目里我们经常自己用fetchReadableStream来读流因为EventSource不支持 POST 请求和自定义 header但协议层面的设计思路是值得借鉴的。提示如果你用EventSource它只能发 GET 请求没法带请求体。而大模型对话通常需要 POST 传 prompt所以生产环境更常见的是用fetch拿到response.body这个ReadableStream然后自己按 SSE 格式解析。这一点后面会详细讲。2.3 一个容易踩的坑idle timeout 与流中断热词里有个很典型的报错stream disconnected before completion: idle timeout waiting for sse。这个错误的本质是连接空闲时间超过了某一层的超时阈值被强制断开了。SSE 连接是长连接如果模型思考时间较长比如推理模型在想的时候可能几十秒不吐字这段时间连接上没有任何数据流动中间的任何一层——Nginx、负载均衡、API 网关、甚至客户端自己的超时设置——都可能判定这个连接僵死了然后把它掐掉。表现就是前端收到一半内容突然断了或者直接报上面那个错。解决思路分几个层面。服务端层面可以定期发送心跳注释。SSE 协议里以冒号开头的行是注释客户端会忽略但它的作用是保持连接活跃: keep-alive每隔 15 到 30 秒发一条这样的注释就能让中间层知道这个连接还活着。Nginx 层面需要关掉缓冲并调大超时location /api/stream { proxy_pass http://backend; proxy_buffering off; proxy_cache off; proxy_read_timeout 300s; proxy_set_header Connection ; proxy_http_version 1.1; chunked_transfer_encoding off; }proxy_buffering off是关键否则 Nginx 会把服务端的输出攒起来一起发流式效果就没了。proxy_read_timeout要设得比模型最长思考时间还大。客户端层面fetch本身没有超时限制但如果你用了某些封装库要检查它有没有默认超时。我在一个项目里就遇到过本地测试一切正常部署到线上后流式输出变成憋一大段再一次性出来。排查了半天最后发现是网关默认开启了响应缓冲。这类问题不看日志、不逐层排查光看代码是找不到的。3. LangChain 结构化输出让模型吐出能直接用的 JSON3.1 结构化输出解决的是什么问题大模型默认返回的是自然语言。你问它帮我提取这段文本里的人名和电话它可能回你好的这段文本里的人是张三电话是138xxxx。这对人来说没问题但你的代码没法直接用——你总不能写正则去抠张三和那串数字吧模型换个措辞你的正则就废了。结构化输出要做的就是约束模型必须返回一个符合预定 schema 的 JSON。比如你定义一个 Pydantic 模型from pydantic import BaseModel, Field class PersonInfo(BaseModel): name: str Field(description人物姓名) phone: str Field(description联系电话) city: str Field(description所在城市)然后让模型输出你期望拿到的是{name: 张三, phone: 13800000000, city: 杭州}这样你的代码就能直接PersonInfo.model_validate_json(...)拿到一个类型安全的对象字段缺失或类型不对会直接报错而不是悄悄给你一个错误的字符串。3.2 LangChain 里实现结构化输出的几种路子LangChain 提供了with_structured_output这个方法用法很直接from langchain_openai import ChatOpenAI llm ChatOpenAI(modelgpt-4o-mini) structured_llm llm.with_structured_output(PersonInfo) result structured_llm.invoke(张三住在杭州电话是13800000000) print(result) # PersonInfo(name张三, phone13800000000, city杭州)它底层其实做了两件事一是把你的 Pydantic 模型转成 JSON Schema二是通过模型的function calling / tool calling能力或者JSON mode来约束输出格式。不同模型支持的方式不一样LangChain 会帮你选一个可用的。这里有个关键选择点用 tool calling 还是用 JSON mode。tool calling 的约束更强模型被强制按照工具参数 schema 来输出格式正确率更高JSON mode 只是告诉模型请输出 JSON但不保证字段完全符合你的 schema。所以如果你的模型支持 tool calling优先用它。LangChain 的with_structured_output默认会走 tool calling如果模型不支持再降级。还有一种更土但更可控的方式在 prompt 里明确要求输出 JSON并给出示例然后自己解析。这种方式灵活但稳定性依赖模型能力字段一多就容易漏。我一般只在模型不支持 tool calling 的时候才用这招。3.3 结构化输出和流式的天然矛盾现在问题来了。结构化输出要求模型返回完整的、合法的 JSON而流式输出是一段段吐字符的。这两者放在一起就产生了一个根本矛盾你收到的每一个 chunk 都是残缺的 JSON单独拿出来根本没法解析。举个例子模型流式返回的内容可能是这样的chunk1: {na chunk2: me: 张 chunk3: 三, phone: chunk4: 13800000000 chunk5: , city: 杭州}你把 chunk1 拿去JSON.parse直接报错。你把 chunk1 到 chunk4 拼起来还是缺个右括号照样报错。只有拼到 chunk5才是一个完整的 JSON。那怎么办有两个思路。思路一流式只用来做展示结构化解析等流结束再做。也就是前端先把原始文本流式显示出来等流结束后服务端把完整文本解析成 JSON 返回。这个方案简单但用户体验上用户看到的是一堆 JSON 原文在往外冒不好看。思路二增量解析。每收到一个 chunk 就拼到缓冲区然后尝试解析缓冲区里的内容能解析出多少算多少。这个方案体验好但实现复杂需要一个能处理不完整 JSON的解析器。这就是下一节要重点讲的内容。4. 流式 JSON 的增量解析怎么在残缺状态下提取数据4.1 为什么不能直接 JSON.parse先明确一个事实标准的JSON.parse是全量解析它要求输入是一个完整的、合法的 JSON 字符串。少一个括号、多一个逗号、字符串没闭合都会直接抛异常。而流式场景下你拿到的缓冲区在绝大多数时刻都是不完整的所以直接 parse 必然失败。有人会说那我 try-catch 一下失败了就等下一个 chunk 呗。这个思路方向是对的但有个问题你没法区分因为不完整而失败和因为格式真的错了而失败。如果模型返回的 JSON 本身就有语法错误你一直等下去也等不到合法结果最后就是无限等待。所以需要一个更聪明的解析策略。4.2 增量解析的两种实用方案方案一括号配对 部分解析核心思路是维护一个缓冲区每次收到新 chunk 就追加进去然后扫描缓冲区统计括号{}和[]的配对情况同时处理字符串内的转义。当发现某个对象已经闭合时就尝试解析这个闭合的部分。这个方案的关键在于状态机你需要知道当前是否在字符串内部因为字符串里的{不算括号是否在转义状态\不算字符串结束。手写一个这样的扫描器大概几十行代码但很容易在边界情况上出错比如 Unicode 转义、嵌套对象、数组里套对象。方案二用现成的流式 JSON 解析库Python 生态里有ijson这类库但它主要面向超大 JSON 文件的流式读取对边收边解析的支持不算特别顺手。更实用的做法是用partial-json-parser或者自己封装一个容错解析函数先尝试标准解析失败后尝试补全缺失的括号再解析。我实际项目里用的是第三种思路基于 Pydantic 的增量校验。具体做法是维护一个字符串缓冲区每次追加后先尝试用json.loads解析成功就直接返回失败则尝试补全——统计未闭合的括号数量在末尾补上对应数量的}或]再解析。如果补全后能解析成功说明当前缓冲区是一个前缀合法的 JSON可以提取出已经完整的字段。import json def try_parse_partial(buffer: str): # 先尝试直接解析 try: return json.loads(buffer) except json.JSONDecodeError: pass # 尝试补全括号 stack [] in_string False escape False for ch in buffer: if escape: escape False continue if ch \\: escape True continue if ch : in_string not in_string continue if in_string: continue if ch in {[: stack.append(ch) elif ch in }]: if stack: stack.pop() # 补全未闭合的括号 closing for ch in reversed(stack): closing } if ch { else ] try: return json.loads(buffer closing) except json.JSONDecodeError: return None这个函数返回None就表示当前缓冲区还不足以解析出任何有效内容继续等下一个 chunk。返回字典就表示已经能提取出部分字段了。注意补全括号解析出来的结果某些字段可能是半截的。比如{name: 张补全成{name: 张}你拿到的 name 是张但实际完整值可能是张三。所以增量解析的结果只能用于展示不能用于最终入库。最终数据一定要等流结束后用完整内容重新解析一次。4.3 结构化输出流式场景下的特殊处理如果你用的是 LangChain 的with_structured_output配合流式情况会稍微不同。LangChain 在流式模式下astream返回的 chunk 可能是AIMessageChunk里面的tool_call_chunks字段会携带增量 JSON 片段。你需要把这些片段按顺序拼接才能得到完整的工具调用参数。from langchain_openai import ChatOpenAI llm ChatOpenAI(modelgpt-4o-mini) structured_llm llm.with_structured_output(PersonInfo) buffer async for chunk in structured_llm.astream(张三住在杭州电话13800000000): # chunk 是增量对象需要累积 buffer str(chunk) partial try_parse_partial(buffer) if partial: print(当前已解析:, partial)这里有个细节不同版本的 LangChain 对结构化输出流式的支持程度不一样。有些版本会直接把增量 JSON 片段放在chunk里有些版本则要求你用astream_events来监听更细粒度的事件。我建议先打印一下 chunk 的实际结构看清楚它到底吐的是什么再决定怎么拼。盲目照抄文档很容易踩版本差异的坑。5. 打字机效果的前端实现从数据流到视觉呈现5.1 用 fetch 读取流并逐块渲染前端要做的第一件事是拿到流。前面说过EventSource不支持 POST所以用fetchasync function streamChat(prompt, onChunk) { const response await fetch(/api/chat, { method: POST, headers: { Content-Type: application/json }, body: JSON.stringify({ prompt }) }); const reader response.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 格式切分消息 const lines buffer.split(\n\n); buffer lines.pop(); // 最后一段可能不完整留到下次 for (const line of lines) { if (line.startsWith(data: )) { const data line.slice(6); if (data [DONE]) return; onChunk(JSON.parse(data)); } } } }这里有几个关键点。decoder.decode(value, { stream: true })的stream: true很重要它保证多字节字符比如中文在跨 chunk 边界时不会被截断成乱码。buffer.split(\n\n)按 SSE 的消息分隔符切分lines.pop()把最后一段不完整的消息留在缓冲区等下一个 chunk 来了再拼。这个留尾巴的处理是流式解析的通用套路少了它就会出现消息被截断的 bug。5.2 打字机效果的两种实现层次拿到 chunk 之后怎么把它变成打字机这里有两个层次。层次一chunk 级渲染。模型每吐一个 chunk你就把 chunk 的内容追加到界面上。这是最简单的做法但效果取决于模型吐字的粒度。有些模型一次吐好几个字看起来就是一跳一跳的不够顺滑。层次二字符级渲染。把收到的 chunk 先放进一个队列然后用一个定时器比如每 20 毫秒从队列里取一个字符渲染出来。这样无论模型吐字粒度多大视觉上都是均匀的逐字显示。这个方案体验最好但要注意队列积压的问题——如果模型吐字速度远快于渲染速度队列会越积越长最后用户看到的内容严重滞后。我的经验是做一个自适应速度队列短的时候慢点渲染显得从容队列长的时候加快渲染速度追上进度。具体可以用队列长度动态调整定时器间隔或者一次渲染多个字符。let queue []; let rendering false; function enqueue(text) { queue.push(...text); if (!rendering) renderLoop(); } function renderLoop() { rendering true; if (queue.length 0) { rendering false; return; } // 队列越长一次渲染越多字符 const batch Math.max(1, Math.floor(queue.length / 10)); const chars queue.splice(0, batch).join(); appendToDom(chars); const delay queue.length 50 ? 5 : 20; setTimeout(renderLoop, delay); }5.3 流式渲染中的 Markdown 与代码块处理大模型的输出经常包含 Markdown 格式比如代码块、列表、加粗。如果你直接把原始文本塞进 DOM用户看到的就是一堆**和 体验很差。所以需要实时渲染 Markdown。但实时渲染 Markdown 有个坑不完整的 Markdown 语法会导致渲染错乱。比如模型正在输出一个代码块刚吐了py还没吐结束的这时候渲染器可能把后面的所有内容都当成代码。解决办法是延迟渲染对代码块这类需要成对出现的语法检测到开始标记后先缓存等结束标记出现再渲染。或者用一个容错性好的 Markdown 渲染器它对未闭合的语法有自己的处理策略。另一个细节是代码高亮。流式场景下代码是逐渐出现的如果每来一个字符就重新高亮整段代码性能会很差。我的做法是代码块在流式过程中先用等宽字体纯文本显示等这个代码块闭合后再做一次高亮。这样既保证了流式过程的流畅又保证了最终效果的美观。6. 把整条链路串起来一个可复现的完整方案6.1 服务端FastAPI LangChain 流式接口服务端我用 FastAPI 来搭因为它对异步和流式响应的支持很自然。核心是一个返回StreamingResponse的接口from fastapi import FastAPI from fastapi.responses import StreamingResponse from langchain_openai import ChatOpenAI import json app FastAPI() llm ChatOpenAI(modelgpt-4o-mini, streamingTrue) async def event_generator(prompt: str): try: async for chunk in llm.astream(prompt): if chunk.content: payload json.dumps({content: chunk.content}, ensure_asciiFalse) yield fdata: {payload}\n\n yield data: [DONE]\n\n except Exception as e: error json.dumps({error: str(e)}, ensure_asciiFalse) yield fdata: {error}\n\n app.post(/api/chat) async def chat(body: dict): return StreamingResponse( event_generator(body[prompt]), media_typetext/event-stream, headers{ Cache-Control: no-cache, X-Accel-Buffering: no, } )X-Accel-Buffering: no这个 header 是给 Nginx 看的告诉它不要缓冲这个响应。ensure_asciiFalse保证中文不被转义成\uXXXX减少传输体积。[DONE]是一个约定俗成的结束标记前端收到它就停止读取。6.2 结构化输出与流式展示的分离设计如果你的场景既要流式展示又要结构化数据我建议把两件事分开做。具体来说流式接口只负责把模型的自然语言输出吐给前端做展示等流结束后前端把完整文本回传给一个单独的解析接口服务端用with_structured_output或者直接调模型做结构化提取返回干净的 JSON。这样做的好处是职责清晰。流式接口只管快结构化接口只管准。如果硬要把两者揉在一起你会陷入既要增量解析又要保证字段完整的泥潭代码复杂度飙升还容易出 bug。当然如果你的场景对实时性要求极高比如用户边看边要看到提取出的字段那就只能用增量解析方案。这时候记住前面说的原则增量解析的结果只用于展示最终数据以流结束后的完整解析为准。6.3 实测中遇到的几个典型问题问题一中文乱码。前面提过TextDecoder要加stream: true。另外服务端json.dumps要加ensure_asciiFalse否则中文会被转义虽然不影响正确性但传输体积会变大。问题二流式输出憋大招。表现是前端等很久然后一次性收到全部内容。这几乎都是中间层缓冲导致的。排查顺序先看 Nginx 有没有proxy_buffering off再看网关有没有开响应缓冲最后看服务端框架有没有自己的缓冲机制。FastAPI 的StreamingResponse默认是不缓冲的但如果你在生成器里做了批量处理也可能导致这个问题。问题三结构化输出字段缺失。用with_structured_output时偶尔会遇到模型返回的 JSON 缺字段。这通常是因为模型能力不够或者 schema 太复杂。解决办法一是换更强的模型二是把 schema 拆简单点三是给字段加详细的description帮助模型理解。Pydantic 的Field(description...)不是摆设它会被转成 JSON Schema 的一部分传给模型描述写得越清楚模型填得越准。问题四流中断后无法恢复。SSE 协议本身支持通过Last-Event-ID做断点续传但大模型场景下实现起来比较麻烦因为模型不会从第 N 个字继续生成。实际项目里更常见的做法是检测到流中断后让前端提示用户网络波动请重试然后重新发起一次完整请求。如果业务允许也可以把已经生成的部分内容保留在界面上让用户知道断在哪里。7. 一些不那么显然的经验之谈做这类流式 结构化的项目有几个点是我踩过坑之后才真正理解的文档里通常不会强调。第一流式的粒度不由你控制由模型和网络共同决定。你可能期望每个 chunk 是一个字但实际可能一次来好几个字也可能因为网络原因几个 chunk 合并到达。所以前端渲染一定要做缓冲队列不能假设来一个 chunk 就渲染一个 chunk。这个假设一旦成立遇到网络抖动就会出现内容跳跃。第二结构化输出的 schema 要扁平。我试过用嵌套很深的 Pydantic 模型对象里套对象再套数组结果模型经常在深层字段上出错。后来把 schema 改成扁平的字段名用下划线连接表示层级关系正确率明显提升。模型对深层嵌套结构的处理能力比我们想象的要弱。第三超时设置要分层考虑。客户端、网关、服务端、模型 API 各有各的超时。任何一层的超时小于模型的最长响应时间都会导致流中断。我的做法是模型 API 超时设最大比如 300 秒网关和服务端设 310 秒客户端不设超时或者设一个很大的值。这样超时永远发生在最外层方便定位问题。第四测试流式一定要用真实网络环境。本地localhost测试流式几乎不会遇到缓冲和超时问题因为中间没有任何代理层。一部署到线上各种问题就冒出来了。所以有条件的话测试环境要尽量模拟生产环境的网络拓扑至少要有 Nginx 这一层。第五增量解析的补全逻辑要处理数组。前面给的try_parse_partial函数处理了对象和数组的括号补全但实际场景里还有一种情况数组元素是对象且对象还没闭合。比如{items: [{name: a}, {name: b这时候补全逻辑要能正确识别出需要补}]}。这个逻辑用栈来处理是最稳妥的但要注意字符串内的括号和转义字符。最后分享一个调试技巧在开发流式功能时我会在服务端把每个 chunk 的内容和到达时间打上日志前端也把每个 chunk 的接收时间打上日志。这样一旦出现卡顿或跳跃对比两端的时间戳就能快速定位是服务端吐得慢还是网络传输慢还是前端渲染慢。这个习惯帮我省下了大量排查时间。
返回列表