1. 流式解析到底在解决什么问题
先把场景说清楚。你调用一个大模型接口,问它“帮我写一段快速排序”,如果走传统的请求-响应模式,客户端发一个 POST,服务端把整段回答在内存里拼完,再一次性返回。用户看到的就是:转圈、转圈、转圈,然后“啪”一下整段文字全出来。短回答还好,一旦回答上千字,等待时间可能十几秒,体验非常糟糕。
流式解析要解决的就是这个“等待焦虑”。它让服务端每生成一小段内容就立刻推给客户端,客户端边收边渲染,用户看到文字像打字机一样一个个蹦出来。这背后涉及三个层面的工程问题:传输协议怎么选、数据流怎么切分、前端怎么消费和渲染。这三个问题串起来,就是“流式解析工程化”这个标题真正要覆盖的范围。
热搜词里出现了 SSE、Web Streams API、TransformStream、OpenAI,这几个词基本勾勒出了当前主流方案的技术轮廓。SSE 负责传输,Web Streams API 负责在浏览器端处理数据流,TransformStream 负责在管道中做转换,OpenAI 的接口则是这套方案最典型的应用场景。我下面会把这几个环节拆开讲,每个环节都给出可复现的代码和踩坑记录。
这篇文章适合谁看?如果你正在做 AI 对话类产品的前端或全栈开发,或者你已经在用 SSE 但总觉得哪里不对劲——比如流断了不知道怎么恢复、中文乱码、多个流并发时状态混乱——那这篇内容应该能帮到你。我会从协议选型讲到代码落地,再讲到线上排查,尽量把每个决策背后的“为什么”说清楚。
2. 传输协议选型:为什么 SSE 成了默认答案
2.1 SSE、WebSocket、轮询三者的真实取舍
在流式场景里,可选的传输方式主要有三种:短轮询、WebSocket、SSE。很多人一上来就觉得 WebSocket 更“高级”,但实际上在 AI 对话这个场景里,SSE 才是更合适的选择。我把三者的关键差异列出来:
| 维度 | 短轮询 | WebSocket | SSE |
|---|---|---|---|
| 通信方向 | 客户端拉 | 全双工 | 服务端推 |
| 协议 | HTTP | 独立协议 | HTTP |
| 自动重连 | 需自己实现 | 需自己实现 | 浏览器内置 |
| 实现复杂度 | 低 | 高 | 低 |
| 代理兼容性 | 好 | 一般 | 好 |
| 适合场景 | 低频更新 | 双向实时 | 单向流式推送 |
AI 对话的本质是“客户端发一次请求,服务端持续推回答”,这是典型的单向推送。WebSocket 的全双工能力在这里是浪费的,反而带来了额外的连接管理成本。短轮询则会产生大量无效请求,延迟也不可控。SSE 基于 HTTP,天然穿透大多数代理和网关,浏览器还内置了重连机制,工程上最省心。
注意:SSE 是单向的,客户端不能通过同一个连接发消息。如果你需要中途打断生成,得用另一个 HTTP 请求去通知服务端,或者用 AbortController 直接断开连接。
2.2 SSE 协议格式的细节
SSE 的报文格式看起来简单,但有几个细节如果没注意,会导致解析失败。一个标准的 SSE 事件长这样:
data: {"choices":[{"delta":{"content":"你"}}]} data: {"choices":[{"delta":{"content":"好"}}]} data: [DONE]每条消息以data:开头,以两个换行符\n\n结束。注意这个双换行是必须的,它是事件的分隔符。如果服务端只发了一个换行,浏览器会认为事件还没结束,继续等待后续数据。
还有一个容易忽略的点:SSE 支持event:、id:、retry:等字段。id字段用于断线重连时告诉服务端从哪里继续,retry用于指定重连间隔。但在 AI 对话场景里,我们通常不需要这些,因为每次对话是独立的,断了就重新发起。
OpenAI 的流式接口返回的就是标准 SSE 格式,每个 chunk 是一个 JSON,最后以data: [DONE]结束。这个[DONE]不是 JSON,是一个特殊标记,解析的时候要单独处理,否则JSON.parse会直接抛异常。
2.3 为什么不用 fetch 的 responseType: 'stream'
有人会想,既然 SSE 就是 HTTP 流,那我直接用 fetch 拿 response.body 不就行了?确实可以,而且现代浏览器里 fetch 的 response.body 就是一个 ReadableStream。但这里有个关键区别:原生 EventSource 会自动帮你做事件切分、重连、状态管理,而 fetch 方案需要你自己处理这些。
那为什么很多项目还是选了 fetch 而不是 EventSource?因为 EventSource 有两个硬伤:第一,它只支持 GET 请求,没法带复杂的 POST body;第二,它不能自定义请求头,没法传 Authorization。而调用大模型接口通常需要 POST 加自定义 header,所以 fetch + ReadableStream 成了更实际的选择。
这就引出了下一个话题:拿到 ReadableStream 之后,怎么把它变成一个个可渲染的文本片段。
3. Web Streams API:把字节流变成可读文本
3.1 ReadableStream 的基本消费方式
fetch 返回的 response.body 是一个 ReadableStream,里面是 Uint8Array 类型的字节块。最原始的消费方式是用 getReader():
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'); while (true) { const { done, value } = await reader.read(); if (done) break; const text = decoder.decode(value, { stream: true }); console.log(text); }这里有个关键参数:decoder.decode(value, { stream: true })。如果不传stream: true,当一个多字节字符(比如中文)被切分到两个 chunk 里时,解码会出错,出现乱码。传了stream: true之后,TextDecoder 会缓存不完整的字节序列,等下一个 chunk 来了再一起解码。这个细节我在实际项目里踩过坑,中文回答偶尔出现“”就是因为这个。
3.2 用 TransformStream 做管道化处理
上面的写法能用,但把所有逻辑堆在一个 while 循环里,代码会越来越乱。更好的方式是用 TransformStream 把处理逻辑拆成独立的管道阶段。Web Streams API 提供了pipeThrough方法,可以把多个 TransformStream 串起来:
const decoder = new TextDecoder('utf-8'); const splitSSE = new TransformStream({ transform(chunk, controller) { const text = decoder.decode(chunk, { stream: true }); // 按双换行切分事件 const events = text.split('\n\n'); for (const event of events) { if (event.trim()) controller.enqueue(event); } } }); const parseJSON = new TransformStream({ transform(event, controller) { const line = event.trim(); if (line.startsWith('data: ')) { const data = line.slice(6); if (data === '[DONE]') { controller.enqueue({ done: true }); return; } try { const parsed = JSON.parse(data); const content = parsed.choices?.[0]?.delta?.content; if (content) controller.enqueue({ content }); } catch (e) { // 忽略解析失败的 chunk } } } }); const stream = response.body .pipeThrough(splitSSE) .pipeThrough(parseJSON); const reader = stream.getReader(); while (true) { const { done, value } = await reader.read(); if (done) break; if (value.content) { appendToDOM(value.content); } }这种管道化的写法有几个好处:每个 TransformStream 只负责一件事,方便单独测试;可以灵活增删中间环节,比如加一个统计 token 数的环节;代码可读性明显提升。
3.3 处理跨 chunk 边界的问题
上面的splitSSE有一个隐藏 bug:如果一个 SSE 事件被切分到两个 chunk 里,split('\n\n')会把不完整的事件也切出来。比如第一个 chunk 结尾是data: {"cho,第二个 chunk 开头是ices":...},直接切分会得到两个残缺的片段。
正确的做法是维护一个缓冲区,只处理完整的事件:
let buffer = ''; const splitSSE = new TransformStream({ transform(chunk, controller) { buffer += decoder.decode(chunk, { stream: true }); const parts = buffer.split('\n\n'); // 最后一段可能不完整,留在缓冲区 buffer = parts.pop(); for (const part of parts) { if (part.trim()) controller.enqueue(part); } }, flush(controller) { // 流结束时处理剩余数据 if (buffer.trim()) controller.enqueue(buffer); } });这个buffer的维护是流式解析里最容易出错的地方。我见过不少项目因为没处理跨 chunk 边界,导致长回答偶尔丢字或者 JSON 解析失败。判断标准很简单:如果服务端返回的 chunk 大小不固定,就必须做缓冲。
4. 前端渲染与状态管理
4.1 增量渲染的性能考量
拿到文本片段之后,最直接的做法是每来一个片段就innerHTML += content或者setState(prev => prev + content)。这在短回答上没问题,但回答长了之后,频繁的 DOM 操作和 React 重渲染会明显卡顿。
我的做法是用一个缓冲区加 requestAnimationFrame 做批量更新:
let pending = ''; let rafId = null; function appendToDOM(text) { pending += text; if (rafId) return; rafId = requestAnimationFrame(() => { document.getElementById('output').textContent += pending; pending = ''; rafId = null; }); }这样每帧最多更新一次 DOM,即使服务端每秒推几十个 chunk,渲染压力也可控。在 React 里可以用类似思路,把流式内容存在 ref 里,用 useSyncExternalStore 或者定时 flush 到 state。
4.2 中断生成与 AbortController
用户点了“停止生成”按钮,你得真的把请求断掉,不然服务端还在跑,白白消耗 token。AbortController 是标准做法:
const controller = new AbortController(); fetch('/api/chat', { signal: controller.signal, // ... }); // 用户点击停止 stopButton.onclick = () => controller.abort();abort 之后,reader.read() 会抛出一个 AbortError,需要在 catch 里单独处理,不要当成真正的错误上报。另外要注意,abort 只是断开了客户端连接,服务端是否停止生成取决于服务端的实现。如果服务端用的是 OpenAI 的流式接口,客户端断开后,服务端的写入会失败,通常也会跟着停止,但这不算强保证。
4.3 多轮对话的状态隔离
一个页面上可能有多个对话同时进行,或者用户快速切换对话。这时候要确保每个流的 reader 和 buffer 是独立的,不能共用全局变量。我习惯把每个流封装成一个类或者闭包:
function createStreamSession(onChunk, onDone, onError) { let buffer = ''; let controller = new AbortController(); async function start(payload) { const response = await fetch('/api/chat', { method: 'POST', signal: controller.signal, headers: { 'Content-Type': 'application/json' }, body: JSON.stringify(payload) }); // ... 消费流 } function stop() { controller.abort(); } return { start, stop }; }每个会话持有自己的 controller 和 buffer,切换对话时调用对应的 stop,状态就不会串。
5. 服务端实现要点
5.1 设置正确的响应头
服务端要返回 SSE,必须设置这几个响应头:
Content-Type: text/event-stream Cache-Control: no-cache Connection: keep-alive X-Accel-Buffering: noX-Accel-Buffering: no是给 Nginx 看的,告诉它不要缓冲这个响应。如果不设,Nginx 默认会缓冲,导致客户端收不到实时数据,等整个响应结束才一次性收到。这个坑我在线上环境遇到过,本地测试正常,一上 Nginx 就变成“假流式”。
5.2 转发上游流式响应
如果服务端是转发 OpenAI 的流,用 Node.js 的 fetch 拿到上游 response.body 之后,可以直接 pipe 给客户端:
const upstream = await fetch('https://api.openai.com/v1/chat/completions', { method: 'POST', headers: { 'Authorization': `Bearer ${process.env.OPENAI_API_KEY}`, 'Content-Type': 'application/json' }, body: JSON.stringify({ model: 'gpt-4', stream: true, messages }) }); res.writeHead(200, { 'Content-Type': 'text/event-stream', 'Cache-Control': 'no-cache', 'Connection': 'keep-alive', 'X-Accel-Buffering': 'no' }); const reader = upstream.body.getReader(); while (true) { const { done, value } = await reader.read(); if (done) break; res.write(value); } res.end();注意这里直接把上游的字节流转发给客户端,不做解析,让客户端去解析。这样服务端逻辑最简单,也最少出错。如果需要在服务端做内容过滤或者统计,再考虑加 TransformStream。
5.3 超时与心跳
SSE 连接如果长时间没有数据,中间的网络设备可能会主动断开。OpenAI 的接口在生成过程中会持续推数据,一般不会空闲太久。但如果你的服务端在两次 chunk 之间有较长的处理时间,建议加心跳:
const heartbeat = setInterval(() => { res.write(': heartbeat\n\n'); }, 15000); // 流结束时清理 res.on('close', () => clearInterval(heartbeat));以冒号开头的行是 SSE 的注释,客户端会忽略,但能保持连接活跃。热搜词里有一条 “stream disconnected before completion: idle timeout waiting for sse”,说的就是空闲超时导致流中断,加心跳是标准解法。
6. 常见问题与排查实录
6.1 流式解析问题速查表
| 现象 | 可能原因 | 排查方向 |
|---|---|---|
| 中文乱码 | TextDecoder 没传 stream: true | 检查 decode 参数 |
| 偶尔丢字 | 跨 chunk 边界没缓冲 | 检查 buffer 逻辑 |
| 假流式(一次性出) | Nginx 缓冲 | 加 X-Accel-Buffering: no |
| JSON 解析报错 | 把 [DONE] 当 JSON 解析 | 单独处理 [DONE] |
| 流中途断开 | 空闲超时 | 加心跳或调整超时 |
| abort 后报错 | 没区分 AbortError | catch 里单独判断 |
| 多对话串内容 | 全局 buffer 共用 | 每个会话独立封装 |
6.2 几个我踩过的坑
第一个坑是JSON.parse的容错。上游返回的 chunk 偶尔会有空行或者格式不标准的行,直接 parse 会抛异常。我的做法是 try-catch 包住,解析失败就跳过,不要让一个坏 chunk 打断整个流。
第二个坑是 React 的闭包问题。在 useEffect 里启动流,回调里更新 state,如果依赖数组没写对,会拿到旧的 state。我后来改成用 ref 存流式内容,或者用 useReducer 来管理,避免闭包陷阱。
第三个坑是服务端的 res.write 背压。如果客户端消费慢,res.write 返回 false,继续写会占用内存。生产环境要监听 drain 事件,或者用 pipeline 自动处理背压。
提示:调试 SSE 的时候,用 curl 加 -N 参数可以关闭缓冲,直接看到流式输出:
curl -N -X POST ...。这比在浏览器里看 Network 面板更直观。
6.3 关于 OpenAI 兼容接口的注意事项
现在很多模型服务都提供 OpenAI 兼容的接口,但流式返回的格式可能有细微差异。有的服务返回的 chunk 里delta.content是空字符串,有的会在最后一个 chunk 里带finish_reason。解析的时候要兼容这些情况,不要假设每个 chunk 都有 content。另外,API Key 一定要放在服务端,不要暴露在前端代码里,这是基本的安全底线。
7. 工程化封装的一点思路
把上面这些环节串起来,一个可复用的流式解析模块应该包含:传输层(fetch + AbortController)、解析层(TransformStream 管道)、渲染层(批量更新)、状态层(会话隔离)。我习惯把它封装成一个不依赖框架的核心类,然后在上层用 React/Vue 的适配器去对接。这样核心逻辑可以单独测试,换框架也不用重写。
测试的时候,用一个模拟的 ReadableStream 来构造各种边界情况:chunk 被切分、中文跨 chunk、[DONE] 标记、空 chunk、异常 chunk。把这些 case 都覆盖到,线上出问题的概率会小很多。
流式解析这件事,协议本身不复杂,难的是边界情况的处理。把 buffer 管理、编码处理、中断恢复这几个点做扎实,基本就能稳定运行了。