从 SSE 流式原理到 LangChain 结构化输出:打字机效果与 JSON 解析全方案实战
现在的 AI 应用,谁还没个打字机效果都不好意思上线。但说实话,我看过太多项目把“流式输出”做成了摆设——前端拿到一堆碎文本直接拼上去,后端一个yield扔出去就算完事,等到要让模型输出结构化 JSON 的时候,整个链路直接崩掉。今天这篇文章,咱们就把这条链路上的每一个环节都拆开揉碎:SSE 协议到底在传什么、EventSource 怎么接、LangChain 怎么配合流式输出、以及最关键的结构化 JSON 在流式场景下怎么解析才不会翻车。
适合谁看?如果你正在做 LLM 应用的落地开发,或者你只是听说过 SSE 和 LangChain 但一直没啃下源码,这篇文章能帮你省掉至少一周的试错时间。我会把前端、后端、协议层、解析层全部串起来讲,拿真实生产环境的方案说事。
1. 从轮询到推送:SSE 流式原理与选型逻辑
很多人一听“流式输出”就想到 WebSocket,这其实是个误区。Server-Sent Events(SSE)和 WebSocket 虽然都能做实时推送,但它们的定位完全不同。SSE 是单向的、基于 HTTP 的服务器推送协议,浏览器原生支持,一行EventSource就能接住;WebSocket 是双向全双工通信,适合聊天室、游戏这种需要频繁双向交互的场景。
选择 SSE 来做大模型流式输出的理由非常实际:大模型生成文本本来就是单向的——模型只管往客户端推数据,不需要客户端频繁往回发消息。你用一个 WebSocket 连接去做一件只需要单向通知的事,等于开着卡车去送快递,不是不行,但成本和复杂度完全不成比例。而且 SSE 走的是普通 HTTP,天然兼容各种网关、代理、负载均衡器,不存在 WebSocket 那种“连接要升级、防火墙要放行、长连接容易断”的麻烦。
1.1 SSE 的本质:就是一段不断输出的 HTTP 响应
SSE 这个协议本身没什么神秘感。它就是一个 HTTP 响应,只不过Content-Type设为text/event-stream,服务器端不关闭连接,持续往客户端写数据。每一帧数据遵循固定的格式规范:
data: 这是一行内容 data: 这是第二行内容事件帧之间用空行分隔,每条数据以data:前缀开头。如果一行放不下,可以分成多行data:,客户端会自动用换行符拼接起来。如果想要把某条数据单独拉到指定事件类型,还可以加event:字段,客户端监听对应事件名即可。
后端在 Python 里实现 SSE 端点其实非常简单——用 FastAPI 的StreamingResponse配合生成器,一行yield就是一次推送:
from fastapi import FastAPI from fastapi.responses import StreamingResponse app = FastAPI() def token_stream(): for token in ["你好", ",", "世界", "!"]: yield f"data: {token}\n\n" @app.get("/api/chat/stream") async def chat_stream(): return StreamingResponse( token_stream(), media_type="text/event-stream", headers={ "Cache-Control": "no-cache", "Connection": "keep-alive", "X-Accel-Buffering": "no" } )注意那个X-Accel-Buffering: no头,这个是我在生产环境踩过坑才加上的。Nginx 默认会缓冲后端响应,如果不关掉缓冲,前端收到的不是逐字打字机效果,而是等整个响应结束才一次性拿到全部数据——流式效果直接归零。遇到这种情况,第一反应不要怀疑代码,先查网关层是不是在“捣乱”。
1.2 为什么 EventSource 比 fetch 流式读取更省心
浏览器原生的 SSE 客户端是EventSource对象,用法极简:
const source = new EventSource('/api/chat/stream'); source.onmessage = (event) => { const token = event.data; // 把 token 追加到界面上 appendText(token); }; source.onerror = (err) => { // 处理断连逻辑 console.error('SSE 连接断开', err); };用 fetch 走流式读取也可以,但要注意一个关键差异:EventSource 自带自动重连机制,连接断开后浏览器会按retry:字段指定的间隔自动重新建立连接。而用 fetch 的ReadableStream做手动解析,一旦连接中断,你得自己写重连逻辑,还要处理各种边缘情况。
EventSource 还有一个限制:它只支持 GET 请求。如果你需要把用户的历史对话通过 POST 发给后端再开启流式响应,要么把参数拼在 URL 上,要么改用fetch加POST配合ReadableStream来手动解析 SSE 帧。后者更灵活,但需要自己实现帧解析逻辑,这个我们放到后面讲。
2. 前端打字机效果实现:不只是“把文字贴上去”
打字机效果的实现难度不在于“渲染文字”,而在于“渲染的节奏怎么和生成节奏对齐”。模型输出的 token 到达时间是随机的,最快的几个 token 可能几十毫秒就到了,最慢的可能要等好几秒,如果无脑把每个到达的 token 直接 append 到 DOM,页面会有一阵一阵的卡顿感——这其实是 render 频率超过了浏览器的刷新率,属于典型的性能问题。
2.1 节流渲染:让每帧最多只动一次 DOM
正确的做法是引入一个节流机制:把 SSE 收到的 token 先塞进一个缓冲区,然后通过requestAnimationFrame循环把缓冲内容定时写入 DOM。每帧最多渲染一次帧缓冲,浏览器会自然地在空闲时刷新,不会出现内容以肉眼可见的“一坨一坨”冲出来的情况。
export function createTypewriter(element) { let buffer = ''; let rendering = false; let canceled = false; function flush() { if (canceled) return; element.textContent += buffer; buffer = ''; rendering = false; } function scheduleFlush() { if (rendering) return; rendering = true; requestAnimationFrame(flush); } return { push(text) { buffer += text; scheduleFlush(); }, done() { canceled = true; flush(); } }; }这里最关键的是scheduleFlush里的防抖逻辑:同一帧里不管 push 进来多少 token,只会触发一次flush。这样即使后端一次塞给前端十几个 token,界面也是平滑地滚动更新,不会闪烁或跳动。
除了渲染节流,还需要处理 markdown 的渲染问题。如果直接用textContent把文本塞进去,接完流之后要做一次marked或markdown-it解析,把纯文本转成带样式的 HTML。这里有个细节要注意:流式过程中不要实时做 markdown 转 HTML,因为 markdown 语法往往是跨多个 token 的,比如代码块的三个反引号可能分几次传输,中途转换很容易解析出残次品。稳妥的做法是流式过程只展示纯文本,流结束后整体转一次 markdown。
2.2 断连重连与消息状态的精准控制
打字机的“结束态”比“进行态”更容易被忽视。前端跟用户沟通的方式不能只看文字有没有打完,还要知道后端是不是主动关闭了流。SSE 连接结束后浏览器会触发onerror,但这个onerror在正常关闭和异常断连时都会触发,区分两者需要靠后端发送一个特殊的结束标识。
惯例的做法是:后端在流式响应完所有 token 后,额外发送一个data: [DONE]帧,前端收到这个标志就正常关闭连接,不再触发重连逻辑。如果是中途超时、网关断连、服务器异常,前端收不到[DONE]就需要走自动重连,同时更新界面上的“重试”提示。
event: message data: 今天天气不错 data: [DONE]事件监听里判断event.data === '[DONE]'后调用source.close(),这是一个很容易被新手漏掉的操作——你不主动 close,EventSource 会一直按重试间隔尝试重新连接,产生一堆无意义的请求。
2.3 SSE 用 GET 也能带参数:EventSource 的兼容处理
刚才提到 EventSource 只支持 GET,很多人会被卡在这一步。如果业务场景要求用 POST 传参,最干净的方案是保留 EventSource 的自动重连能力,把复杂参数放在服务端会话里,前端只需要用 GET 请求一个临时会话标识:
// 先通过 POST 建立会话,拿到 session_id const { session_id } = await fetch('/api/chat/start', { method: 'POST', headers: { 'Content-Type': 'application/json' }, body: JSON.stringify({ messages: history }) }).then(res => res.json()); // 再开启 SSE 流式连接 const source = new EventSource(`/api/chat/stream?session_id=${session_id}`);后端拿到 session_id 后从内存或 Redis 中取对应用户的请求参数,再开始生成流。这样既绕过了 GET 的参数长度限制,又保留了对断线重连的友好性,算是我在多个项目里反复验证过的标准做法。
3. LangChain 流式输出实战:Stream 与回调的取舍
LangChain 做流式输出有两个层面的接口,一个是最外层的stream()方法,另一个是底层的callbacks机制。很多教程只讲stream(),但真正要把流式输出部署到生产环境的项目里,你大概率需要用回调回调来拿全链路的 token。
3.1 用.stream()拿到大模型 token 流
stream()的用法非常简单,直接迭代即可:
from langchain_openai import ChatOpenAI from langchain_core.messages import HumanMessage llm = ChatOpenAI( model="gpt-4o-mini", temperature=0, streaming=True ) for chunk in llm.stream([HumanMessage(content="讲个冷笑话")]): if hasattr(chunk, "content") and chunk.content: yield f"data: {chunk.content}\n\n"注意一点:llm.stream()在内部其实也是开启 streaming,然后逐块 yield AIMessageChunk。streaming=True这个参数不是必须的,但加上能确保底层调用过程使用流式接口,减少首 token 延迟。在 FastAPI 的生成器函数里逐块yield出去,前端打字机效果就有了。
3.2 用 callbacks 实现多路输出与日志采集
生产环境里普遍存在的需求是:同一个流式响应里,你既要给最终用户看完整回答,又要给运营看 token 用量、给调试者看中间推理过程。stream()只能让你拿到最终的 token 流,拿不到“模型在中间到底走了哪几步”。这种情况就要用 LangChain 的AsyncCallbackHandler:
from langchain_core.callbacks import AsyncCallbackHandler from langchain_core.agents import AgentFinish class TokenCollector(AsyncCallbackHandler): def __init__(self): self.token_buffer = [] self.agent_logs = [] async def on_llm_new_token(self, token: str, **kwargs): self.token_buffer.append(token) # 这里把 token 实时推送到 SSE async def on_agent_action(self, action, **kwargs): self.agent_logs.append(f"思考: {action.log}") async def on_agent_finish(self, finish: AgentFinish, **kwargs): self.agent_logs.append(f"完成: {finish.return_values}")回调机制特别适合 Agent 场景:用户在界面上看到的不只是最终回答,还能看到“Agent 正在搜索资料”“Agent 正在调用计算器”这样的中间状态,体验会非常像真的在看一个人思考过程在推进。
3.3 LangChain 版本差异的坑:接口更新比模型还快
LangChain 的接口变化非常频繁,尤其是 LangChain v0.1 到 v0.2 再到 v0.3 这个跨度里,很多老教程里的代码已经跑不起来了。比如langchain.chains.LLMChain在较新版本里被降级为 Legacy,官方更推荐直接用RunnableSequence或LangGraph。如果你照着网上的旧教程抄代码,很可能遇到DeprecationWarning甚至直接报错。
我的建议是:项目初始就锁定一个 LangChain 版本,不要用latest。比如langchain==0.2.x配langchain-openai的版本组合,功能相对稳定。升级某个核心包时,最好全量跑一遍链路的流式测试用例,因为回调签名、流式输出的 chunk 类型、事件触发顺序都可能被改动影响。
4. 结构化输出与 JSON 解析:流式场景下最大的坑
标题里说的“结构化输出”,指的是让 LLM 稳定地返回一段符合 JSON Schema 的数据,而不是一段格式可疑的文本。在这个基础上叠加流式传输,问题立刻复杂了一个量级:你收到的不是一个完整 JSON,而是被拆成几千个碎片的 JSON 字符串。如何从碎片流中稳定地解析出结构化数据,是这门手艺的真正核心。
4.1 LangChain 的with_structured_output:把格式控制交给模型
LangChain 提供了非常顺手的结构化输出能力。核心方法是with_structured_output(),配合 Pydantic 模型定义输出格式:
from pydantic import BaseModel, Field from langchain_openai import ChatOpenAI class WeatherReport(BaseModel): city: str = Field(description="城市名称") date: str = Field(description="日期,格式 YYYY-MM-DD") temperature: float = Field(description="气温,摄氏度") advice: str = Field(description="出行建议") llm = ChatOpenAI(model="gpt-4o", temperature=0) structured_llm = llm.with_structured_output(WeatherReport) result = structured_llm.invoke("明天北京天气怎么样?") print(result.city, result.temperature) # 返回的是 Pydantic 对象,不是原始字符串底层原理是:LangChain 会把 Pydantic 模型定义转换成 JSON Schema,通过 tool calling(函数调用)机制让模型选择一个合法的结构来响应。只要模型支持工具调用,这种做法的稳定性比“把人话转 JSON”高一个档次。result是WeatherReport实例,直接.city就能取到字段。
但这里有个容易忽略的问题:with_structured_output()默认不做流式输出,配合invoke()使用时会一次性返回完整结果。如果你想在流式过程中拿到结构化的片段,需要另外想办法。
4.2 流式场景下的 JSON 断帧问题:看似玄学,实则线性
当你开启 SSE 流式输出时,模型生成的 JSON 会被拆成无数个 token 一帧帧传过来。前端的event.data可能长这样:
{"ci ty": "北京 ", "tempe rature": 23 .5}这不是 JSON 坏了,这只是 JSON 被“截”了。你没法直接JSON.parse这种半截数据。粗暴的做法是攒着等全部数据到齐再解析——这样最稳,但同时失去了流式输出的意义;另一种不想放弃流式的做法是用“增量 JSON 解析器”。
增量 JSON 解析的核心思路是:维护一个不断累积的字符串缓冲区,每次收到新 token 就尝试解析一次缓冲区内容,解析成功就更新 UI,解析失败就继续等。由于 JSON 对象在结构上存在不完整状态,单纯靠JSON.parse的 try-catch 无法支持增量,所以需要判断“当前文本是否是一个合法 JSON 的前缀”。我常用的做法是加入一个“文本平衡校验”:
function hasBalancedBrackets(str) {