1. 流式输出的本质:为什么我们需要 SSE
1.1 从“等一锅饭”到“边炒边上桌”
做过大模型应用的人都有一个共同体会:用户等一个完整回答的时间,往往比回答本身更让人焦虑。传统 HTTP 请求是“一锤子买卖”——客户端发请求,服务端算完所有 token,一次性返回。用户盯着转圈图标十几秒,体验极差。SSE(Server-Sent Events)解决的正是这个问题:服务端算出一个 token 就推一个,浏览器收到就渲染,这就是我们常说的“打字机效果”。
SSE 的本质其实非常朴素。它就是一个长连接,服务端持续往客户端写数据,格式遵循text/event-stream规范。每一帧数据以data:开头,以两个换行\n\n结束。浏览器端的EventSource会自动解析这些帧,触发onmessage回调。理解这一点很关键,因为后面所有的坑——断流、超时、粘包——都源于对这个格式的理解不够透彻。
我见过太多人把 SSE 和 WebSocket 混为一谈。简单说:SSE 是单向的(服务端到客户端),基于 HTTP,自动重连,实现简单;WebSocket 是双向的,需要协议升级,适合聊天室这类双向交互。大模型对话场景里,用户提问走普通 POST,回答走 SSE 推送,这种“一问一答”的模式用 SSE 完全够用,没必要上 WebSocket 增加复杂度。
1.2 一个最小可用的 SSE 服务端长什么样
先看一段 FastAPI 的最小实现,这是后面所有内容的地基:
from fastapi import FastAPI from fastapi.responses import StreamingResponse import asyncio app = FastAPI() async def event_generator(): for i in range(5): yield f"data: 第{i}条消息\n\n" await asyncio.sleep(0.5) @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 解析。X-Accel-Buffering: no是给 Nginx 看的,告诉它别缓冲,否则你推的数据会卡在反向代理层,用户看到的还是“一次性返回”。Cache-Control: no-cache防止中间层缓存。这三个 header 缺一个,打字机效果就可能变成“打字机卡顿”。
注意:
yield的字符串必须以\n\n结尾。少一个换行,浏览器就认为这一帧没结束,会一直等,表现就是“消息不显示”。
1.3 前端消费 SSE 的两种姿势
前端消费 SSE 有两条路。第一条是用原生EventSource:
const es = new EventSource('/stream'); es.onmessage = (e) => { console.log('收到:', e.data); }; es.onerror = (err) => { console.error('出错了', err); es.close(); };EventSource的优点是自动重连、代码极简。但它有个硬伤:只支持 GET 请求,不能自定义 header。这意味着你没法在请求头里带 Authorization token,也没法传复杂的 POST body。对于需要鉴权的大模型接口,这条路基本走不通。
第二条路是用fetch+ReadableStream手动解析,这也是我在生产环境里一直用的方案:
const response = await fetch('/stream', { method: 'POST', headers: { 'Content-Type': 'application/json', 'Authorization': `Bearer ${token}`, }, body: JSON.stringify({ prompt: '你好' }), }); 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 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; console.log('解析出:', JSON.parse(data)); } } }这段代码里最关键的是buffer的处理逻辑。网络传输是流式的,一次read()拿到的数据可能刚好把一帧切成两半。如果你直接对每次read()的结果做split,就会丢数据或者解析报错。正确做法是维护一个缓冲区,只处理完整的帧(以\n\n分隔),最后一段不完整的留到下一次拼接。这个细节我在项目里踩过坑,表现为“偶尔丢一条消息”,排查了半天才发现是分包问题。
2. LangChain 结构化输出:让模型吐出能用的 JSON
2.1 为什么“让模型返回 JSON”这么难
大模型天生是“话痨”。你让它返回 JSON,它可能给你返回:
好的,这是你要的 JSON: ```json {"name": "张三", "age": 25}希望对你有所帮助!
这段文本里混了自然语言和 markdown 代码块,直接 `JSON.parse` 必然报错。早期大家用正则去抠,写一堆 `replace` 和 `match`,脆弱得不行——模型换个措辞,正则就失效了。 LangChain 的结构化输出(Structured Output)就是来解决这个问题的。它的核心思路是:**不要靠事后解析,而是从生成阶段就约束模型**。具体来说,LangChain 提供了几种手段,从弱到强依次是:Prompt 约束、`with_structured_output`、以及底层的 Function Calling / JSON Mode。 ### 2.2 用 Pydantic 定义你的输出契约 LangChain 结构化输出的入口是 Pydantic 模型。你先定义好想要的数据结构: ```python from pydantic import BaseModel, Field from typing import List, Optional class Person(BaseModel): """人物信息""" name: str = Field(description="姓名") age: int = Field(description="年龄") skills: List[str] = Field(default_factory=list, description="技能列表") email: Optional[str] = Field(default=None, description="邮箱,可能没有")然后用with_structured_output把模型包一层:
from langchain_openai import ChatOpenAI llm = ChatOpenAI(model="gpt-4o-mini", temperature=0) structured_llm = llm.with_structured_output(Person) result = structured_llm.invoke("张三今年25岁,会Python和Go,邮箱是zhangsan@example.com") print(result) # Person(name='张三', age=25, skills=['Python', 'Go'], email='zhangsan@example.com')注意result直接就是Person对象,不是字符串,不需要你手动json.loads。这是with_structured_output最大的价值——它把“解析”这一步从你的代码里彻底拿掉了。
Field里的description不是摆设。模型在生成时会把 schema 和 description 一起看,description 写得越清楚,模型填错字段的概率越低。我一般会把每个字段的业务含义、格式要求、边界条件都写进去,比如description="年龄,必须是正整数,范围0-150"。
2.3 底层到底发生了什么:Function Calling 与 JSON Mode
with_structured_output不是魔法,它底层依赖两种机制,取决于你用的模型支持哪种。
第一种是Function Calling(也叫 Tool Calling)。LangChain 把你的 Pydantic schema 转成一个“工具定义”,发给模型。模型不直接生成文本,而是生成一个“调用这个工具”的指令,参数就是你要的 JSON。LangChain 拿到这个指令,解析参数,实例化成 Pydantic 对象。这种方式约束最强,因为模型是在“填参数”而不是“写作文”。
第二种是JSON Mode。模型被强制要求输出合法 JSON,但不保证符合你的 schema。LangChain 拿到 JSON 后再用 Pydantic 校验,不符合就报错或重试。这种方式约束弱一些,但兼容性好,很多模型都支持。
# 显式指定用哪种方式 structured_llm = llm.with_structured_output( Person, method="json_mode", # 或 "function_calling" )选哪种?我的经验是:能用 Function Calling 就用 Function Calling,准确率明显更高。只有当模型不支持 Function Calling 时,才退而求其次用 JSON Mode。判断方法很简单,看模型文档,或者直接试——如果with_structured_output报错说 method 不支持,就换。
注意:
temperature=0在结构化输出场景几乎是必须的。温度高了,模型容易“发挥创意”,字段名写错、类型写错、多塞字段,各种幺蛾子。
2.4 嵌套结构与复杂类型的处理
真实业务里,数据结构往往不是扁平的。比如一个订单,里面有商品列表,每个商品又有自己的属性:
class Product(BaseModel): name: str = Field(description="商品名") price: float = Field(description="单价") quantity: int = Field(description="数量") class Order(BaseModel): order_id: str = Field(description="订单号") customer: str = Field(description="客户名") products: List[Product] = Field(description="商品列表") total: float = Field(description="总金额") structured_llm = llm.with_structured_output(Order) result = structured_llm.invoke("订单A123,客户李四,买了2个苹果每个5块,3个香蕉每个2块")嵌套结构对模型来说难度更高,因为要同时保证外层和内层的字段都对。我的经验是:层级不要超过三层,字段不要超过十个。超过这个量级,模型出错的概率会显著上升。如果业务确实复杂,拆成多次调用,每次只抽一部分,比一次性抽一个大对象靠谱得多。
另外,Optional和默认值要慎用。模型看到Optional字段,可能会“偷懒”不填。如果某个字段业务上必须有,就别给默认值,让 Pydantic 校验时直接报错,逼模型填。
3. 打字机效果与结构化输出的冲突与调和
3.1 一个根本矛盾:流式是“半成品”,结构化是“成品”
这里有个很多人没意识到的矛盾。打字机效果要求边生成边推送,用户看到的是一个个 token 拼起来的半成品。而结构化输出要求最终结果是一个完整的、合法的 JSON。JSON 在没写完之前,是没法解析的——{"name": "张这种半截字符串,JSON.parse直接报错。
所以你不能简单地“把结构化输出的流直接推给前端”。前端拿到半截 JSON,解析不了,打字机效果就无从谈起。
那怎么办?业内有几种成熟的方案,我逐个说。
3.2 方案一:流式生成,结束后再结构化
最稳妥的方案是分两步:第一步,用普通流式输出让用户看到打字机效果;第二步,流结束后,把完整文本再喂给结构化输出链,得到 JSON。
async def stream_then_structure(prompt: str): # 第一步:流式输出给用户看 full_text = "" async for chunk in llm.astream(prompt): full_text += chunk.content yield f"data: {json.dumps({'type': 'text', 'content': chunk.content})}\n\n" # 第二步:结构化解析 structured_llm = llm.with_structured_output(Person) result = structured_llm.invoke(full_text) yield f"data: {json.dumps({'type': 'structured', 'content': result.model_dump()})}\n\n"这个方案的优点是简单、可靠,打字机效果和结构化输出都拿到了。缺点是多了一次模型调用,成本和延迟都翻倍。如果只是展示用,其实可以省掉第二步;如果后续逻辑需要结构化数据,这一步就省不了。
3.3 方案二:流式解析 JSON(增量解析)
如果你既想要打字机效果,又不想多调一次模型,那就得在流式过程中“增量解析” JSON。思路是:维护一个不断增长的 JSON 字符串,每收到一个 chunk 就尝试解析,能解析出多少算多少。
Python 里可以用ijson或者自己写一个简单的状态机。但说实话,自己写状态机很容易出 bug,尤其是遇到转义字符、嵌套对象的时候。我试过用partial-json-parser这类库,效果还行,但遇到复杂嵌套还是会翻车。
from partial_json_parser import loads as partial_loads buffer = "" async for chunk in structured_llm.astream(prompt): buffer += chunk try: partial = partial_loads(buffer) # partial 是当前能解析出的部分 yield f"data: {json.dumps(partial)}\n\n" except Exception: pass这个方案的优点是省一次调用,缺点是前端要处理“不完整对象”。比如{"name": "张三", "age":这种,解析出来可能是{"name": "张三"},age 字段还没出现。前端渲染时得考虑字段缺失的情况,逻辑会复杂不少。
我的建议是:如果前端只是展示,用方案一;如果前端需要实时根据结构化数据做交互(比如实时填表),才考虑方案二。方案二的复杂度,很多时候不值得。
3.4 方案三:双通道输出
还有一种折中方案:让模型同时输出自然语言和结构化数据,用特殊标记分隔。比如:
<text>张三今年25岁,会Python和Go。</text> <json>{"name": "张三", "age": 25, "skills": ["Python", "Go"]}</json>流式推送时,前端根据当前在哪个标记内,决定是渲染文本还是解析 JSON。这个方案的好处是一次调用搞定,坏处是模型不一定听话,标记可能写错、漏写、顺序颠倒。我在项目里用过,稳定性不如方案一,后来还是换回了两次调用。
4. 实战:封装一个可复用的 SSE 流式接口
4.1 整体架构设计
说了这么多原理,来点能直接抄的。我封装了一个通用的 SSE 流式接口,结构是这样的:
- 后端:FastAPI + LangChain,提供
/chat/stream接口 - 前端:Vue 3 + fetch,消费 SSE 并渲染打字机效果
- 协议:每一帧是一个 JSON,包含
type和content字段
协议设计很关键。我见过有人直接把模型输出的文本塞进data:,结果文本里有换行符,把 SSE 的帧格式搞乱了。永远用 JSON 包装你的数据,这样换行、特殊字符都被转义了,不会破坏帧结构。
# 每一帧的格式 {"type": "text", "content": "你"} {"type": "text", "content": "好"} {"type": "structured", "content": {"name": "张三"}} {"type": "done"} {"type": "error", "content": "模型调用失败"}4.2 后端完整实现
from fastapi import FastAPI, Request from fastapi.responses import StreamingResponse from langchain_openai import ChatOpenAI from pydantic import BaseModel, Field from typing import List, Optional import json import asyncio app = FastAPI() class Person(BaseModel): name: str = Field(description="姓名") age: int = Field(description="年龄") skills: List[str] = Field(default_factory=list, description="技能列表") class ChatRequest(BaseModel): prompt: str need_structured: bool = False llm = ChatOpenAI(model="gpt-4o-mini", temperature=0, streaming=True) async def sse_generator(prompt: str, need_structured: bool): full_text = "" try: async for chunk in llm.astream(prompt): if chunk.content: full_text += chunk.content yield f"data: {json.dumps({'type': 'text', 'content': chunk.content}, ensure_ascii=False)}\n\n" if need_structured: structured_llm = llm.with_structured_output(Person) result = await structured_llm.ainvoke(full_text) yield f"data: {json.dumps({'type': 'structured', 'content': result.model_dump()}, ensure_ascii=False)}\n\n" yield f"data: {json.dumps({'type': 'done'})}\n\n" except Exception as e: yield f"data: {json.dumps({'type': 'error', 'content': str(e)}, ensure_ascii=False)}\n\n" @app.post("/chat/stream") async def chat_stream(req: ChatRequest, request: Request): return StreamingResponse( sse_generator(req.prompt, req.need_structured), media_type="text/event-stream", headers={ "Cache-Control": "no-cache", "Connection": "keep-alive", "X-Accel-Buffering": "no", }, )几个关键点。ensure_ascii=False让中文正常显示,不然会变成\u4f60\u597d这种。ainvoke用异步版本,避免阻塞事件循环。异常要捕获并作为error帧推给前端,不然前端会一直等,直到超时。
4.3 前端 Vue 3 消费与渲染
// composables/useSSE.js import { ref } from 'vue'; export function useSSE() { const text = ref(''); const structured = ref(null); const loading = ref(false); const error = ref(''); async function start(prompt, needStructured = false) { text.value = ''; structured.value = null; error.value = ''; loading.value = true; try { const response = await fetch('/chat/stream', { method: 'POST', headers: { 'Content-Type': 'application/json' }, body: JSON.stringify({ prompt, need_structured: needStructured }), }); 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 frames = buffer.split('\n\n'); buffer = frames.pop(); for (const frame of frames) { if (!frame.startsWith('data: ')) continue; const payload = JSON.parse(frame.slice(6)); if (payload.type === 'text') { text.value += payload.content; } else if (payload.type === 'structured') { structured.value = payload.content; } else if (payload.type === 'error') { error.value = payload.content; } else if (payload.type === 'done') { loading.value = false; } } } } catch (e) { error.value = e.message; } finally { loading.value = false; } } return { text, structured, loading, error, start }; }这个 composable 可以直接在组件里用:
<template> <div> <button @click="start('介绍一下张三', true)">开始</button> <p>{{ text }}</p> <pre v-if="structured">{{ structured }}</pre> <p v-if="error" style="color: red">{{ error }}</p> </div> </template> <script setup> import { useSSE } from './composables/useSSE'; const { text, structured, loading, error, start } = useSSE(); </script>4.4 参数选择与性能调优
流式接口有几个参数直接影响体验,我列个表:
| 参数 | 推荐值 | 说明 |
|---|---|---|
temperature | 0 | 结构化输出场景必须为0,减少随机性 |
max_tokens | 按需 | 太小会截断,太大会浪费,一般2048够用 |
stream | True | 流式必须开 |
| 超时时间 | 60s | 前端 fetch 默认无超时,要手动加 AbortController |
| 心跳间隔 | 15s | 长时间无数据时发注释帧保活 |
超时这块特别说一下。浏览器 fetch 默认没有超时,如果服务端挂了,前端会一直等。我一般用AbortController加一个 60 秒的超时:
const controller = new AbortController(); const timeoutId = setTimeout(() => controller.abort(), 60000); fetch('/chat/stream', { signal: controller.signal, ... }) .finally(() => clearTimeout(timeoutId));心跳保活也很重要。有些反向代理会在连接空闲 30 秒后断开。解决办法是服务端定期发一个注释帧: heartbeat\n\n,这个帧不会被onmessage触发,但能保持连接活跃。
5. 常见问题与排查技巧实录
5.1 断流问题:idle timeout waiting for SSE
这是搜索热词里出现频率最高的问题。现象是:流式输出到一半,突然断了,控制台报stream disconnected before completion: idle timeout waiting for SSE。
原因通常有三个。第一,反向代理(Nginx)的proxy_read_timeout默认 60 秒,超过就断。解决办法是在 Nginx 配置里加:
location /chat/stream { proxy_pass http://backend; proxy_read_timeout 300s; proxy_buffering off; proxy_cache off; chunked_transfer_encoding on; }第二,模型本身生成太慢,两个 token 之间超过 60 秒。这种情况要么换更快的模型,要么加心跳帧。第三,客户端主动断开(用户关页面),这个属于正常,服务端捕获asyncio.CancelledError清理资源即可。
5.2 JSON 解析失败:模型返回了非法 JSON
即使有with_structured_output,偶尔还是会遇到解析失败。常见原因和排查方法:
| 现象 | 原因 | 解决 |
|---|---|---|
| 字段缺失 | 模型偷懒 | 字段设必填,不给默认值 |
| 类型错误 | 模型理解偏差 | description 写清楚类型 |
| 多出字段 | 模型自由发挥 | Pydantic 默认忽略多余字段,可设model_config = ConfigDict(extra='forbid')强制报错 |
| 中文乱码 | 编码问题 | 确保ensure_ascii=False |
| 截断 | max_tokens 太小 | 调大 max_tokens |
我一般会加一层重试:解析失败时,把错误信息拼回 prompt,让模型重新生成。LangChain 有with_retry可以配置:
structured_llm = llm.with_structured_output(Person).with_retry( stop_after_attempt=3, )5.3 前端渲染卡顿:频繁 setState 导致掉帧
打字机效果如果每个 token 都触发一次 Vue 的响应式更新,token 多了会卡。优化方法是批量更新:用一个缓冲区攒几个 token,再一次性更新。
let pending = ''; let rafId = null; function scheduleUpdate(content) { pending += content; if (rafId) return; rafId = requestAnimationFrame(() => { text.value += pending; pending = ''; rafId = null; }); }用requestAnimationFrame把更新对齐到浏览器刷新率,一帧最多更新一次,流畅度提升明显。这个技巧我在长文本场景下实测有效,从每秒卡顿几次变成丝滑。
5.4 踩坑清单:那些文档里不会写的事
最后分享几个我踩过的坑,都是文档里不会提但实际会遇到的。
坑一:StreamingResponse的 generator 里不能用同步阻塞调用。如果你在async def的 generator 里调了一个同步的llm.invoke(),整个事件循环会被阻塞,其他请求全部卡住。必须用ainvoke或astream。
坑二:Nginx 的proxy_buffering默认是 on。即使你服务端设了X-Accel-Buffering: no,有些 Nginx 版本还是不认,必须在 Nginx 配置里显式关掉。这个坑我排查了一下午,现象是“本地正常,线上不流式”。
坑三:EventSource不能跨域带 cookie。如果你的前端和后端不同域,EventSource默认不带 cookie,鉴权会失败。要么用fetch方案,要么配置withCredentials(但EventSource的withCredentials支持有限)。
坑四:结构化输出的 schema 别太复杂。我试过一个有 20 个字段、三层嵌套的 schema,模型准确率掉到 60% 以下。后来拆成三次调用,每次抽一部分,准确率回到 95% 以上。模型不是数据库,别指望它一次填完一个大表单。
坑五:流式输出的 token 不等于字符。一个中文字可能被切成多个 token,一个 token 也可能包含多个字符。前端做“逐字显示”时,如果按 token 渲染,中文会出现半个字的情况。解决办法是前端维护一个字符缓冲区,按字符粒度渲染,而不是按 token。
这些经验都是实际项目里一点点磨出来的,希望能帮你少走点弯路。流式输出和结构化输出这两个东西,单独看都不难,难的是把它们揉在一起还不打架。核心思路就一句话:该流式的地方流式,该结构化的地方结构化,别硬凑。