1. 流式输出为什么成了 AI 应用的标配
做过对话类产品的朋友应该都有体会,用户对"等待"这件事的容忍度低得离谱。早几年我们做接口,一个请求发出去转圈三五秒,用户也就忍了,因为那时候大家默认"AI 思考需要时间"。但现在不一样了,ChatGPT 把打字机效果带火之后,用户的心理预期被彻底拉高——你这边还在转圈,他那边已经开始怀疑是不是卡死了,甚至直接刷新页面走人。
这就是SSE(Server-Sent Events)流式输出要解决的核心问题。它不是什么新鲜技术,早在 HTML5 时代就有了,但真正被大规模用起来,还是因为大模型应用的爆发。简单说,SSE 让服务端可以像"挤牙膏"一样,把生成的内容一个字一个字地推给前端,前端收到一段就渲染一段,用户看到的就是文字在屏幕上"长出来"的效果,也就是我们常说的打字机效果。
但光有流式还不够。实际项目里你会发现两个绕不开的坎:第一,流式传输的底层原理和连接管理,稍不注意就会遇到各种断连、超时、乱码问题;第二,流式输出的是自然语言,而业务系统往往需要的是结构化数据——比如从一段对话里提取出订单号、金额、日期,或者让模型返回一个能直接塞进数据库的 JSON。这就引出了LangChain 结构化输出和JSON 解析的话题。
这篇内容我打算把这两条线串起来讲:从 SSE 的底层原理讲起,聊清楚流式接口怎么封装、怎么解析、怎么处理各种异常;再切到 LangChain 的结构化输出,讲清楚怎么让模型稳定吐出 JSON,以及流式和结构化这两个看似矛盾的需求怎么在同一个系统里共存。适合正在做 AI 应用、被流式接口折磨过、或者想让模型输出更"听话"的开发者参考。不管你是刚接触 LangChain 的新手,还是已经踩过一些坑的老手,应该都能从里面找到点有用的东西。
2. SSE 流式原理拆解:它到底是怎么把数据"挤"出来的
2.1 从 HTTP 的"一问一答"说起
要理解 SSE,得先理解普通 HTTP 请求的痛点。传统的 HTTP 是"请求-响应"模型,客户端发一个请求,服务端处理完,一次性把完整响应返回,然后连接关闭。这个模型对于"我查个数据"这种场景没问题,但对于"模型要生成 500 个字"这种场景就很尴尬——服务端明明可以边生成边返回,却非要等全部生成完才吐出来,中间这段时间连接是空闲的,用户是干等的。
有人会说,那我用轮询(Polling)行不行?客户端每隔一秒问一次"生成好了没"。这方案能用,但很蠢:一是延迟高,二是大量无效请求,三是服务端要维护"生成到哪了"的状态。长轮询(Long Polling)稍微好点,但本质上还是"一问一答",每次都要重新建立请求,开销不小。
SSE 的思路完全不同。它基于 HTTP,但把响应改成了持续不断的流。客户端发一次请求,服务端返回的Content-Type是text/event-stream,然后这个连接就不关闭了,服务端可以随时往里面写数据,写一条客户端就收一条,直到服务端主动结束或者连接断开。
注意:SSE 是单向的,只能服务端推给客户端。如果需要双向通信,那得用 WebSocket。但对于大模型这种"用户问一句、模型答一长串"的场景,单向推送完全够用,而且 SSE 比 WebSocket 轻量得多,基于纯 HTTP,不需要额外的协议升级,代理和网关的兼容性也更好。
2.2 SSE 的数据格式:别小看那几个换行
SSE 的报文格式看着简单,但细节不少,很多人第一次手写 SSE 响应就是栽在格式上。它的基本单位是"事件(event)",每个事件由若干字段组成,字段之间用换行分隔,事件之间用空行分隔。核心字段有这么几个:
data:消息内容,可以有多行,多行会被拼接event:事件类型,客户端可以据此区分不同种类的消息id:事件 ID,用于断线重连时告诉服务端"我从哪继续"retry:重连等待时间,单位毫秒
一个典型的 SSE 响应长这样:
data: {"content": "你"} data: {"content": "好"} data: {"content": "呀"}注意每个data后面跟一个换行,然后再跟一个空行表示这个事件结束。这个空行是必须的,少了它客户端会一直等,以为事件还没结束。我见过太多人调试 SSE 时前端一直收不到消息,最后发现就是少了个空行。
还有一个坑:data字段里的内容如果本身包含换行,需要拆成多个data:行,客户端会自动用换行符拼接。比如要发送line1\nline2,得写成:
data: line1 data: line22.3 浏览器端的 EventSource:好用但有局限
浏览器原生提供了EventSource对象来接收 SSE,用起来很简单:
const es = new EventSource('/api/stream'); es.onmessage = (e) => { console.log(e.data); }; es.onerror = (err) => { console.error('连接出错', err); };但EventSource有几个硬伤,实际项目里经常不够用:
第一,它只支持 GET 请求,没法带复杂的请求体。而对话场景往往需要把历史消息、参数配置一起发过去,用 GET 拼 URL 又丑又有长度限制。
第二,它不能自定义请求头,想加个鉴权 token 都费劲。
第三,它的重连机制是自动的,但有时候我们并不想它自动重连,或者想自定义重连逻辑,EventSource就不够灵活了。
所以现在主流的做法是用fetch配合ReadableStream手动处理流。这样既能用 POST,又能自定义 header,还能完全掌控解析和重连逻辑。代价就是得自己写解析代码,把流里的字节按 SSE 格式切出来。后面实操部分我会给一份完整的解析实现。
2.4 服务端怎么"推":以 Python 为例
服务端这边,不管什么语言,核心就一件事:把响应的Content-Type设成text/event-stream,然后保持连接,分块写入数据。以 FastAPI 为例,最简洁的写法是用StreamingResponse:
from fastapi import FastAPI from fastapi.responses import StreamingResponse import asyncio app = FastAPI() async def event_generator(): for char in "你好呀,这是流式输出": yield f"data: {char}\n\n" await asyncio.sleep(0.1) @app.get("/stream") async def stream(): return StreamingResponse( event_generator(), media_type="text/event-stream", headers={ "Cache-Control": "no-cache", "Connection": "keep-alive", "X-Accel-Buffering": "no", }, )这里有几个关键点值得说。media_type必须是text/event-stream,这是 SSE 的标识。Cache-Control: no-cache是防止中间层缓存流式响应,不然你可能收到的是攒了一大坨的旧数据。X-Accel-Buffering: no是给 Nginx 看的,告诉它别缓冲这个响应——这个头非常关键,很多人本地测试好好的,一上 Nginx 就变成"一次性返回",就是因为它默认会缓冲。
实操心得:如果你用的是 Nginx 做反向代理,除了加
X-Accel-Buffering: no,还要确认proxy_buffering off;和proxy_read_timeout设置合理。默认的 60 秒超时对于长对话来说太短了,模型思考久一点连接就被掐了,前端就会报那个经典的stream disconnected before completion: idle timeout waiting for sse。
3. 封装 SSE 流式接口:从字节流到可用消息
3.1 为什么必须自己封装解析逻辑
前面说了EventSource的局限,实际项目里我们基本都用fetch手动处理。但fetch返回的response.body是一个ReadableStream,你读到的是一块一块的Uint8Array字节,这些字节的切分完全不保证和 SSE 事件对齐。也就是说,一次read()可能读到半个事件,也可能读到三个半事件。你必须自己维护一个缓冲区,把字节解码成字符串,按\n\n切分,把完整的事件取出来处理,剩下的留在缓冲区等下一块数据。
这个过程听起来简单,但坑很多。比如中文是多字节的 UTF-8,一个汉字占 3 个字节,如果一次read()正好把某个汉字切成了两半,你直接TextDecoder解码就会得到乱码。正确的做法是用TextDecoder的stream: true模式,它会把不完整的字节序列暂存起来,等下一块数据来了再一起解码。
3.2 一份可直接抄的流式解析实现
下面这份代码是我在多个项目里打磨过的,处理了字节解码、事件切分、多行 data 拼接、异常捕获这些细节,可以直接拿去用:
async function streamChat(url, payload, { onMessage, onError, onDone, signal }) { const response = await fetch(url, { method: 'POST', headers: { 'Content-Type': 'application/json', 'Accept': 'text/event-stream', }, body: JSON.stringify(payload), signal, }); if (!response.ok) { throw new Error(`HTTP ${response.status}`); } const reader = response.body.getReader(); const decoder = new TextDecoder('utf-8'); let buffer = ''; try { while (true) { const { done, value } = await reader.read(); if (done) break; // stream: true 保证多字节字符不被截断 buffer += decoder.decode(value, { stream: true }); // 按空行切分事件 const parts = buffer.split('\n\n'); // 最后一段可能不完整,留在缓冲区 buffer = parts.pop() || ''; for (const part of parts) { if (!part.trim()) continue; const dataLines = []; for (const line of part.split('\n')) { if (line.startsWith('data:')) { dataLines.push(line.slice(5).trimStart()); } } if (dataLines.length === 0) continue; const data = dataLines.join('\n'); if (data === '[DONE]') { onDone && onDone(); return; } try { onMessage && onMessage(JSON.parse(data)); } catch { onMessage && onMessage(data); } } } onDone && onDone(); } catch (err) { if (err.name === 'AbortError') return; onError && onError(err); } finally { reader.releaseLock(); } }这段代码里有几个设计决策值得解释。buffer.split('\n\n')之后pop()出来的最后一段,是可能不完整的事件,必须留到下一轮,否则会丢数据。data:后面用slice(5)而不是slice(6),是因为规范里冒号后的空格是可选的,用trimStart()更稳妥。[DONE]是 OpenAI 风格的结束标记,很多服务端会发这个,收到就主动结束。
3.3 打字机效果的前端渲染策略
拿到流式消息之后,怎么渲染成打字机效果,这里也有讲究。最朴素的做法是每收到一个 token 就setState追加,但这样在高频推送下会触发大量重渲染,页面会卡。我一般用两种优化:
一是批量更新。用一个 ref 暂存待渲染的内容,用requestAnimationFrame或者一个 16ms 的定时器批量 flush 到 state,把渲染频率压到 60fps 以内。
二是光标动画。打字机效果的精髓其实在那个闪烁的光标,用一个 CSS 动画的::after伪元素就能实现,不需要额外的 DOM 节点:
.typing::after { content: '▋'; animation: blink 1s step-end infinite; } @keyframes blink { 50% { opacity: 0; } }流式结束时把typing类去掉,光标就消失了。这个细节虽小,但用户体验差别很大——有光标的时候用户知道"还在生成",没光标就以为"已经完了"。
常见问题:有时候流式输出到一半,前端突然收到一大坨内容,打字机效果变成"啪"一下全出来。这通常是中间有代理做了缓冲。排查顺序是:先看服务端有没有加
X-Accel-Buffering: no,再看 Nginx 的proxy_buffering,最后看 CDN 层有没有开压缩缓冲。逐层排除,基本都能定位到。
4. LangChain 结构化输出:让模型吐出能用的 JSON
4.1 为什么自然语言输出不够用
流式输出解决的是"体验"问题,但业务系统真正头疼的是"数据"问题。你让模型分析一段用户反馈,它给你回一段"这位用户对产品整体满意,但对物流速度有些不满……"——人看着挺好,但你要把它存进数据库、要做统计、要触发后续流程,这段自然语言就没法用了。你需要的是:
{ "sentiment": "positive", "issues": ["物流速度"], "score": 4 }这就是结构化输出要解决的问题。让模型直接返回符合特定 schema 的 JSON,程序拿到就能用,不用再做一轮解析。LangChain 在这方面提供了比较完整的工具链,核心思路是:定义好你想要的 schema,然后通过提示词约束或者模型原生能力,让输出严格符合这个 schema。
4.2 三种结构化输出的实现路径
LangChain 里实现结构化输出,大致有三条路,各有适用场景。
第一条路是提示词约束 + 输出解析器。这是最通用的方式,不依赖模型本身的能力。你在提示词里明确告诉模型"请以 JSON 格式返回,字段包括 xxx",然后用PydanticOutputParser或者JsonOutputParser把模型输出解析成对象。解析器会自动把 schema 的描述注入到提示词里,还会在解析失败时给出修复提示。
from langchain_core.output_parsers import PydanticOutputParser from langchain_core.prompts import ChatPromptTemplate from pydantic import BaseModel, Field class Feedback(BaseModel): sentiment: str = Field(description="情感倾向,positive/negative/neutral") issues: list[str] = Field(description="提到的问题列表") score: int = Field(description="评分,1-5") parser = PydanticOutputParser(pydantic_object=Feedback) prompt = ChatPromptTemplate.from_template( "分析以下反馈:{text}\n{format_instructions}" ).partial(format_instructions=parser.get_format_instructions()) chain = prompt | llm | parser result = chain.invoke({"text": "东西不错,就是发货太慢了"})第二条路是用模型原生的结构化输出能力。现在很多模型支持 function calling 或者 JSON mode,LangChain 封装成了with_structured_output方法,一行代码就能搞定:
structured_llm = llm.with_structured_output(Feedback) result = structured_llm.invoke("东西不错,就是发货太慢了")这种方式最省心,因为约束是在模型层面生效的,输出格式的稳定性比提示词约束高得多。但前提是你用的模型得支持这个能力。
第三条路是工具调用(Tool Calling)的变体。把 schema 定义成一个"工具",让模型"调用"这个工具来返回结构化数据。本质上和第二条路类似,只是实现机制不同。
4.3 三种路径的选型对比
| 路径 | 依赖模型能力 | 稳定性 | 灵活性 | 适用场景 |
|---|---|---|---|---|
| 提示词 + 解析器 | 否 | 中 | 高 | 模型不支持原生结构化输出,或需要复杂嵌套 |
| with_structured_output | 是 | 高 | 中 | 模型支持 function calling / JSON mode |
| 工具调用变体 | 是 | 高 | 中 | 需要和已有工具系统集成 |
我的经验是:能用原生能力就用原生能力,稳定性差距很明显。提示词约束在简单 schema 上还行,一旦字段多了、嵌套深了,模型就容易漏字段或者格式跑偏。但原生能力也不是万能的,有些模型对复杂嵌套 schema 的支持也不好,这时候还是得回到提示词约束,配合重试机制。
实操心得:不管用哪种方式,都要给解析加一层容错和重试。模型偶尔抽风返回个带 markdown 代码块包裹的 JSON(
json ...),或者多说了句"好的,这是结果:",解析器直接就崩了。我一般会写个清洗函数,先把代码块标记、前后缀文字剥掉,再尝试解析,失败就带着错误信息重试一次。
5. 流式与结构化:两个矛盾需求的共存方案
5.1 为什么这俩需求会打架
流式输出和结构化输出,本质上是有矛盾的。流式追求的是"尽快把内容推给用户",所以是逐 token 往外吐;结构化追求的是"输出必须符合 schema",而 schema 校验需要拿到完整的输出才能做。你不可能在只收到{"sentiment": "pos的时候就判断它合不合法。
这个矛盾在实际项目里经常表现为:产品经理说"我要打字机效果",后端说"我要结构化数据",两边一碰就发现没法同时满足。常见的妥协方案有几种,我逐个说说。
5.2 方案一:流式输出自然语言,结束后再结构化
这是最简单的方案。流式阶段正常推自然语言给用户看,等流结束后,把完整文本再走一遍结构化提取。好处是用户体验和数据结构都保住了,坏处是多了一次模型调用(或者一次本地解析),有额外延迟和成本。
如果结构化提取能用规则或者轻量模型做,这个方案其实很划算。比如从一段对话里提取订单号,用正则就够了,没必要再调一次大模型。
5.3 方案二:流式输出 JSON 片段,前端增量解析
这个方案更激进:让模型直接流式输出 JSON,前端边收边解析。但 JSON 是不完整的,{"a": 1, "b":这种状态没法直接JSON.parse。解决办法是用增量 JSON 解析器,比如partial-json这类库,它能解析不完整的 JSON,把已经完整的部分先返回。
import { parse } from 'partial-json'; let acc = ''; onMessage((chunk) => { acc += chunk; const partial = parse(acc, { allowPartial: true }); // partial 里已经能拿到部分字段 render(partial); });这个方案的好处是"边生成边结构化",用户能实时看到字段一个个填上,体验很酷。但坑也不少:一是模型得稳定输出 JSON,不能中途加解释文字;二是增量解析对嵌套结构的支持参差不齐;三是前端渲染逻辑复杂,字段顺序、缺失字段的处理都得考虑。
5.4 方案三:双通道,各走各的
还有一种思路是把两个需求彻底分开:一个通道流式推自然语言给用户看,另一个通道(或者同一响应的不同事件类型)推结构化数据给程序用。用 SSE 的event字段区分:
event: text data: {"content": "这位用户"} event: text data: {"content": "对物流"} event: structured data: {"sentiment": "negative", "issues": ["物流"]}前端监听text事件渲染打字机,业务逻辑监听structured事件处理数据。这个方案最灵活,但要求服务端能同时产出两种数据,实现复杂度最高。一般用在需要"实时展示 + 实时决策"的场景,比如客服系统里边聊边打标签。
5.5 选型建议
| 方案 | 实现复杂度 | 用户体验 | 数据实时性 | 推荐场景 |
|---|---|---|---|---|
| 流式 + 结束后结构化 | 低 | 好 | 低 | 大多数对话应用 |
| 流式 JSON 增量解析 | 高 | 极好 | 高 | 表单填充、实时抽取 |
| 双通道 | 最高 | 好 | 高 | 边聊边决策的复杂系统 |
我的建议是:先从方案一开始,它能覆盖 80% 的场景,实现简单、稳定。等确实有实时结构化的需求了,再考虑方案二或三。别一上来就追求最酷的方案,维护成本会让你怀疑人生。
6. 常见问题与排查技巧实录
6.1 流式连接相关的典型故障
流式接口的故障排查,最让人头疼的是"时好时坏"。本地测试一切正常,一上线就各种断连。我把踩过的坑整理成一张速查表:
| 现象 | 可能原因 | 排查方向 |
|---|---|---|
| 一次性返回全部内容 | 中间层缓冲 | 检查 Nginxproxy_buffering、X-Accel-Buffering |
| 连接几十秒后断开 | 读超时 | 调大proxy_read_timeout、网关超时 |
| 中文乱码 | 字节截断 | 确认TextDecoder用了stream: true |
| 前端收不到消息 | 事件格式错误 | 检查data:后是否有空行 |
| 偶发丢消息 | 缓冲区处理不当 | 检查split后是否保留了不完整片段 |
| 重连后重复内容 | 未用id字段 | 服务端发id,客户端带Last-Event-ID |
那个stream disconnected before completion: idle timeout waiting for sse的报错,几乎 90% 都是超时配置问题。模型思考时间长,中间又没有数据推送,网关就认为连接空闲把它掐了。解决办法有两个:一是调大超时,二是发心跳——服务端每隔 15 秒发一个注释行: keep-alive\n\n,注释行客户端会忽略,但能保持连接活跃。这个技巧非常实用,强烈建议加上。
6.2 结构化输出的稳定性问题
结构化输出最常见的失败是"格式不对"。模型返回了 JSON,但外面裹了层 markdown,或者字段名拼错了,或者该是数组的给了个字符串。我的处理策略是三层防御:
第一层是提示词里把 schema 描述清楚,字段含义、类型、取值范围都写明白,最好给个示例。
第二层是解析前做清洗,用正则把代码块标记剥掉,把前后多余的说明文字去掉。
第三层是解析失败后重试,把错误信息(比如"字段 score 应该是整数,你返回了字符串")作为反馈塞回提示词,让模型修正。LangChain 的OutputFixingParser就是干这个的,它会自动带着错误重试。
from langchain.output_parsers import OutputFixingParser fixing_parser = OutputFixingParser.from_llm(parser=parser, llm=llm)注意:重试不是万能的,如果模型本身能力不够,重试十次也修不好。这时候要么换模型,要么简化 schema。我见过有人非要用小模型做复杂嵌套抽取,折腾半天不如直接换个强点的模型,成本可能还更低。
6.3 性能与成本的平衡
流式 + 结构化这套组合,性能开销主要在模型调用上。几个优化点:
一是能本地解析就别调模型。简单的字段抽取用正则、用规则引擎,比调模型快几个数量级,还免费。
二是结构化提取用小模型。如果只是把自然语言转成 JSON,不需要太强的推理能力,用小模型完全够,成本能降一大截。
三是缓存。同样的输入如果会重复出现,把结构化结果缓存起来,别每次都重新算。
四是并行。流式输出和结构化提取如果互不依赖,可以并行发起,别串行等。
7. 一套可复用的工程化落地思路
把前面这些串起来,一个完整的流式 + 结构化系统大概长这样:前端用fetch+ReadableStream接收 SSE,自己维护缓冲区解析事件,用批量更新 + 光标动画实现打字机效果;服务端用 FastAPI 的StreamingResponse推流,加心跳保活,加X-Accel-Buffering防缓冲;结构化部分用 LangChain 的with_structured_output优先,不行就退回提示词 + 解析器 + 重试;两个需求用双通道或者"流式后结构化"的方式共存。
这套架构我在几个项目里跑下来,稳定性还不错。但要说最深的体会,其实是别过度设计。我一开始总想着把流式和结构化做到极致实时,结果代码复杂度飙升,bug 一堆。后来退回到"流式给用户看、结束后结构化给程序用"这个简单方案,反而跑得最稳。技术选型这事儿,够用就好,把省下来的精力放在业务逻辑上,收益更大。
另外分享一个小技巧:调试 SSE 的时候,别老盯着浏览器,用curl直接看原始响应最直观:
curl -N -H "Accept: text/event-stream" http://localhost:8000/stream-N参数关掉 curl 自己的缓冲,你能看到数据一条条实时打出来。如果 curl 这边是流式的,浏览器那边不是,那问题基本就在前端或者中间的代理层,排查范围一下就缩小了。这个命令我几乎每次调流式接口都会用,比开 DevTools 看 Network 面板高效得多。