1. 从"能跑"到"敢上线":流式解析为什么必须工程化
流式解析这件事,第一次跑通的时候特别爽。后端一个接口推过来,前端EventSource一挂,字一个个往外蹦,感觉产品瞬间高级了。但真正把它放进生产环境,问题就来了:网络断了怎么办?用户切到后台再切回来,消息丢了一半怎么办?服务端推了个半截 JSON,前端JSON.parse直接抛异常,整个页面白屏怎么办?
这些问题的共同点是——它们都不是"流式"本身的问题,而是工程化缺失的问题。流式解析的工程化,核心就三件事:连接的生命周期管理、数据帧的边界处理、异常状态的可恢复性。把这三件事做扎实,流式功能才算从 demo 变成了可交付的能力。
这篇内容适合两类人看:一类是刚把 SSE 或 Web Streams 跑通、准备上线的前端/全栈同学;另一类是后端要封装流式接口、需要和前端约定协议边界的工程师。我会围绕SSE、Web Streams API、TransformStream这几个关键词,把流式解析从协议层到代码层拆开讲,重点放在那些文档里不写、但线上一定会遇到的坑。
先给一个整体判断:流式解析的难点从来不在"怎么读流",而在"怎么在流断掉、乱序、粘包、超时的情况下,让上层业务代码感觉不到这些破事"。工程化的目标就是造一层"防弹衣",把底层的不确定性挡在业务逻辑之外。
2. SSE 与 Web Streams:两套模型到底该怎么选
2.1 SSE 的本质是一个"长连接 + 文本协议"
很多人把 SSE 当成"服务端推送",这个理解不够准确。SSE 的全称是 Server-Sent Events,它建立在普通 HTTP 之上,本质是服务端保持一个 HTTP 响应不关闭,持续往响应体里写文本。浏览器端的EventSource只是帮你把这堆文本按协议解析成事件。
它的协议格式非常朴素,就是纯文本,用两个换行分隔一个事件块:
event: message data: {"id":1,"text":"hello"} event: message data: {"id":2,"text":"world"}注意几个细节:data:后面如果有多行,会被拼接成一行;event:字段决定事件类型;以:开头的行是注释,常被用作心跳保活。这些规则看起来简单,但前端手写解析时最容易在换行符上翻车——\n和\r\n混用、最后一个事件块没有结尾空行,都会导致解析错位。
SSE 最大的优势是浏览器原生支持、自动重连、走标准 HTTP 语义,对基础设施友好。劣势也很明显:只能服务端单向推、只支持文本、连接数受浏览器同域限制(HTTP/1.1 下同域大约 6 个)。所以它适合"通知类、进度类、AI 对话类"这种单向、低频、文本为主的场景。
2.2 Web Streams API 是更底层的"数据管道"
Web Streams API提供的是ReadableStream、WritableStream、TransformStream这套抽象。它不关心数据是 SSE、是 NDJSON 还是二进制,只关心"一块一块的数据怎么流动、怎么转换、怎么背压"。
fetch返回的response.body就是一个ReadableStream。你可以直接读它,也可以接一个TransformStream做中间处理。这就是为什么现在很多流式方案不再用EventSource,而是用fetch + ReadableStream——因为它给了你完全的控制权:可以带自定义 header、可以用 POST、可以中途 abort、可以插入任意转换逻辑。
TransformStream是这套模型里最值得说的东西。它由一对readable和writable组成,你往writable写,从readable读,中间经过你的transform函数。这天然就是一个"解析器"的位置:上游是原始字节流,下游是结构化事件流,中间用 TransformStream 做协议解析。
2.3 选型对照:别为了技术而技术
| 维度 | EventSource (SSE) | fetch + Web Streams |
|---|---|---|
| 请求方法 | 仅 GET | 任意方法,支持 POST |
| 自定义 Header | 不支持 | 完全支持 |
| 自动重连 | 内置 | 需自己实现 |
| 二进制支持 | 不支持 | 支持 |
| 解析控制粒度 | 黑盒,只能拿事件 | 完全可控 |
| 中断控制 | close() | AbortController |
| 适用场景 | 简单通知、进度 | AI 对话、复杂协议 |
我的经验是:如果只是"服务端推个通知",用 EventSource 最省事;只要涉及 POST 传参、鉴权 header、自定义协议、需要精细控制重连策略,一律上 fetch + Web Streams。现在主流的 AI 对话类产品,几乎都是后者,因为对话请求要带上下文、要 POST、要能随时中断生成。
3. 手写一个能扛住生产的流式解析器
3.1 为什么不能直接response.text()
新手最常见的写法是const text = await response.text(),然后等全部返回再解析。这在流式场景下等于把流式的意义完全抹掉了——用户要等所有内容生成完才看到第一个字。正确做法是拿到response.body这个ReadableStream,用getReader()逐块读取。
const response = await fetch('/api/chat', { method: 'POST', headers: { 'Content-Type': 'application/json' }, body: JSON.stringify({ prompt }), signal: controller.signal }); const reader = response.body.getReader(); const decoder = new TextDecoder('utf-8'); while (true) { const { done, value } = await reader.read(); if (done) break; const chunk = decoder.decode(value, { stream: true }); // 处理 chunk }这里有个极其关键的细节:decoder.decode(value, { stream: true })里的stream: true不能省。因为一个 UTF-8 中文字符占 3 个字节,网络分块完全可能把这三个字节切到两个 chunk 里。如果不加stream: true,每个 chunk 独立解码,遇到被切断的多字节字符就会解出乱码(经典的"锟斤拷"就是这么来的)。加上这个参数,TextDecoder会在内部缓存不完整的字节序列,等下一个 chunk 到了再拼起来解码。
3.2 用 TransformStream 把"字节流"变成"事件流"
直接在while循环里写解析逻辑,代码会越来越乱。更好的做法是把解析逻辑封装成一个TransformStream,让数据管道自己完成转换:
function createSSEParser() { let buffer = ''; return new TransformStream({ transform(chunk, controller) { buffer += chunk; const blocks = buffer.split('\n\n'); // 最后一段可能不完整,留在 buffer 里 buffer = blocks.pop(); for (const block of blocks) { const event = parseBlock(block); if (event) controller.enqueue(event); } }, flush(controller) { // 流结束时,处理残留的 buffer if (buffer.trim()) { const event = parseBlock(buffer); if (event) controller.enqueue(event); } } }); }这段代码里有三个工程化要点,值得单独拎出来说:
第一,buffer的存在是为了处理"粘包"和"半包"。网络传输不保证一次read()就返回一个完整的事件块。可能一次返回两个半事件,也可能一个事件被切成三次返回。所以必须维护一个缓冲区,只处理"完整的分隔符"之前的内容,剩下的留在 buffer 里等下一块。blocks.pop()就是干这个的——最后一段永远是不完整的,不能处理。
第二,flush不能忘。流结束的时候,buffer 里可能还残留最后一个事件块(因为服务端最后一个事件后面可能没有\n\n)。如果不在flush里处理,最后一个消息就会永远丢失。这个 bug 特别隐蔽,因为大部分时候服务端会补上结尾空行,只有特定情况下才暴露。
第三,parseBlock要能容错。一个事件块里可能有event:、data:、id:、retry:多个字段,data:还可能有多行。解析时要按行遍历,遇到不认识的行直接跳过,遇到data:就累加。永远不要假设服务端发来的格式百分百规范,防御性解析是流式解析的基本素养。
3.3 把管道串起来
有了 parser,整个链路就清爽了:
const response = await fetch('/api/chat', { /* ... */ }); const eventStream = response.body .pipeThrough(new TextDecoderStream()) .pipeThrough(createSSEParser()); const reader = eventStream.getReader(); while (true) { const { done, value } = await reader.read(); if (done) break; handleEvent(value); // value 已经是结构化对象 }注意这里用了TextDecoderStream,它是浏览器内置的、把TextDecoder包装成 TransformStream 的版本,效果和前面手写decoder.decode(value, {stream:true})一样,但更优雅。整条管道是:字节流 → 文本流 → 事件流,每一层职责单一,测试和替换都很方便。
4. 断线、超时、半包:线上真正会咬人的地方
4.1 "stream disconnected before completion" 到底在说什么
这个报错信息很多人见过,字面意思是"流在完成前断开了"。它可能来自几个完全不同的原因,排查时一定要区分:
- 服务端主动关闭:比如后端生成完了但没发结束标记,或者后端进程被重启。
- 中间层超时:反向代理、网关、负载均衡对空闲连接有超时限制,长时间没有数据流动就被掐断。
- 客户端网络抖动:移动端切网络、Wi-Fi 转 4G,TCP 连接直接断。
- 空闲超时(idle timeout):这是最常见的一种,连接建立后一段时间内没有任何数据往来,被判定为空闲连接强制关闭。
区分方法:看断开的时间点。如果是固定时长(比如 60 秒、120 秒)后断开,基本就是空闲超时;如果是随机时间断开,多半是网络或服务端问题。
4.2 心跳保活:让连接"看起来一直在忙"
对付空闲超时,最直接的办法是定期发送心跳。SSE 协议里,以:开头的行是注释,客户端会忽略,但足以让中间层认为连接是活跃的:
: heartbeat服务端每隔 15~30 秒发一次心跳,就能有效避开大多数空闲超时阈值。心跳间隔要小于中间层的最小超时时间,这个值需要和运维确认,不能拍脑袋。我一般会取超时阈值的 1/3 到 1/2,留足余量。
客户端这边也要配合:如果超过 N 个心跳周期没收到任何数据(包括心跳),就主动判定连接已死并重连。因为有些网络故障下,TCP 连接不会立刻报错,而是"假死"——你以为还连着,其实数据早就过不来了。这种"半开连接"是最坑的,只能靠应用层超时来兜底。
4.3 重连策略:指数退避 + 断点续传
重连不能无脑立刻重试,否则服务端一挂,所有客户端瞬间发起重连,直接把服务打垮(惊群效应)。标准做法是指数退避:
function getRetryDelay(attempt) { const base = 1000; const max = 30000; const delay = Math.min(base * Math.pow(2, attempt), max); // 加随机抖动,避免所有客户端同时重连 return delay + Math.random() * 1000; }Math.random()这个抖动很关键。如果所有客户端都按1s, 2s, 4s, 8s精确重连,它们会在同一时刻集体冲击服务端。加上随机抖动,重连请求就被打散了。
更进阶的是断点续传。SSE 协议支持id:字段和Last-Event-ID请求头:服务端给每个事件编号,客户端重连时带上最后收到的事件 ID,服务端从那个 ID 之后继续推。这样断线期间的消息不会丢。实现这个需要服务端维护一个短期的事件缓冲区,成本不低,但对"消息不能丢"的场景(比如订单状态、任务进度)是刚需。
4.4 半包与粘包的完整处理链路
前面提了 buffer 的思路,这里给一个更完整的排查链路。假设你发现前端偶尔解析出错,按这个顺序查:
- 先确认是不是编码问题:打印原始字节,看有没有被切断的多字节字符。如果是,检查
TextDecoder有没有加stream: true。 - 再确认分隔符:把原始文本打出来,看事件之间到底是
\n\n还是\r\n\r\n。有些服务端框架默认用\r\n,前端只按\n\n切就会切不开。 - 然后确认 buffer 逻辑:在
transform里打印每次的buffer和切出来的blocks,看有没有把不完整的块当成完整的处理了。 - 最后确认 flush:在流结束时打印残留 buffer,看最后一个事件有没有被丢掉。
这个顺序是从底层到上层,先排除编码,再排除协议,最后排除逻辑,能避免在错误的方向上浪费时间。
5. 跨语言协作:前后端在流式协议上的约定
5.1 后端封装流式接口的通用结构
不管后端是 Java、Python 还是 Node,封装流式接口的骨架都差不多:设置正确的响应头 → 拿到输出流 → 循环写数据 → 主动 flush → 处理客户端断开。
以 Java 为例(Spring 体系),核心是返回一个流式的响应体,并确保每次写入后立即 flush,否则数据会攒在缓冲区里,前端迟迟收不到:
@GetMapping(value = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public void stream(HttpServletResponse response) throws IOException { response.setContentType("text/event-stream"); response.setCharacterEncoding("UTF-8"); response.setHeader("Cache-Control", "no-cache"); response.setHeader("X-Accel-Buffering", "no"); // 关键:禁用代理缓冲 PrintWriter writer = response.getWriter(); for (String msg : generateMessages()) { writer.write("data: " + msg + "\n\n"); writer.flush(); // 关键:每次写完立即 flush } }这里有两个新手必踩的坑:
坑一:X-Accel-Buffering: no。如果前面有 Nginx 之类的反向代理,它默认会缓冲响应,导致你的流式数据被攒成一大块才发给客户端,流式效果完全消失。这个 header 就是告诉代理"别缓冲"。有些场景还需要在代理配置里显式关闭缓冲。
坑二:忘记 flush。很多输出流的默认缓冲区是 8KB,你不 flush,数据就一直躺在缓冲区里。表现就是"前端要等很久才一次性收到一大段",而不是逐字出现。
5.2 前后端必须对齐的协议细节
流式接口最容易出问题的地方,是前后端对协议的理解不一致。上线前一定要把下面这些点白纸黑字约定清楚:
| 约定项 | 建议值 | 说明 |
|---|---|---|
| 分隔符 | \n\n | 统一用\n,别混\r\n |
| 数据字段 | data: | 多行数据用多个data: |
| 结束标记 | event: done或[DONE] | 显式告诉前端流正常结束 |
| 错误传递 | event: error+ data | 业务错误也走流,别直接断连 |
| 心跳 | : ping | 注释行,客户端忽略 |
| 编码 | UTF-8 | 前后端统一 |
"结束标记"这一条特别重要。如果服务端生成完就直接关连接,前端无法区分"正常结束"和"异常断开"。加一个显式的结束事件,前端收到它就知道是正常完成,可以停止重连逻辑;没收到就断了,才触发重连。这个区分能避免"明明生成完了,前端还在傻傻重连"的尴尬。
5.3 错误也要走流
一个反直觉但很重要的设计:业务错误不要用 HTTP 状态码返回,而是用流内的事件返回。因为流一旦开始,HTTP 状态码早就发出去了(200),你没法再改。如果生成到一半出错,只能通过流内发一个event: error事件,把错误信息传给前端。
// 前端处理 if (event.type === 'error') { showError(event.data.message); // 注意:这里不要重连,因为是业务错误,重连也没用 return; }区分"可重试错误"(网络问题)和"不可重试错误"(业务错误、参数错误)非常关键。对不可重试错误做重连,只会浪费资源、放大问题。
6. 把流式解析封装成可复用的工程模块
6.1 抽象出一个 StreamClient
散落在各处的流式代码,维护起来是灾难。我的做法是封装一个StreamClient,把连接、解析、重连、中断全部收进去,业务层只关心"收到消息"和"收到错误"两个回调:
class StreamClient { constructor(url, options = {}) { this.url = url; this.options = options; this.controller = null; this.retryCount = 0; this.maxRetry = options.maxRetry ?? 5; } async start({ onMessage, onError, onDone }) { this.controller = new AbortController(); try { const response = await fetch(this.url, { ...this.options, signal: this.controller.signal }); if (!response.ok) throw new Error(`HTTP ${response.status}`); const stream = response.body .pipeThrough(new TextDecoderStream()) .pipeThrough(createSSEParser()); const reader = stream.getReader(); while (true) { const { done, value } = await reader.read(); if (done) break; if (value.type === 'done') { onDone?.(); return; } if (value.type === 'error') { onError?.(value.data); return; } onMessage?.(value.data); } onDone?.(); } catch (err) { if (err.name === 'AbortError') return; // 主动中断,不算错误 this.handleRetry({ onMessage, onError, onDone }); } } abort() { this.controller?.abort(); } }这个封装有几个设计取舍值得说:
AbortError要单独处理。用户主动点"停止生成"时,fetch会抛AbortError。这不是错误,不应该触发重连,也不应该弹错误提示。很多实现忘了这一条,导致用户一停止就弹个报错,体验很差。
重连要带上状态。如果支持断点续传,重连时要带上Last-Event-ID。这个 ID 应该在onMessage里记录,重连时作为 header 传回去。
重试次数要有上限。无限重连在服务端持续不可用时,会让客户端一直空转。设个上限(比如 5 次),超过就放弃并通知用户。
6.2 用状态机管理连接生命周期
流式连接的状态其实不少:idle、connecting、streaming、reconnecting、done、error。用状态机管理,能避免"在错误的状态做了错误的操作":
| 当前状态 | 触发事件 | 目标状态 | 动作 |
|---|---|---|---|
| idle | start | connecting | 发起请求 |
| connecting | 响应成功 | streaming | 开始读流 |
| streaming | 收到 done | done | 关闭连接 |
| streaming | 网络错误 | reconnecting | 退避后重连 |
| reconnecting | 重连成功 | streaming | 续传 |
| reconnecting | 超过上限 | error | 通知用户 |
| 任意 | abort | idle | 中断并清理 |
有了这张表,代码里的if/else就能收敛成清晰的状态转移,每个状态该做什么、不该做什么一目了然。这也是排查问题时的重要工具——出问题时先看当前状态,往往就能定位到问题。
6.3 内存与性能:别让流式拖垮页面
流式场景下有两个性能陷阱:
陷阱一:DOM 更新过于频繁。如果每收到一个字就更新一次 DOM,高频流式下页面会卡。正确做法是用requestAnimationFrame或定时器做批量更新,把短时间内的多个消息合并成一次渲染。
陷阱二:消息列表无限增长。长对话场景下,消息越堆越多,内存和渲染压力都会上来。需要做虚拟滚动或者定期裁剪历史消息。这个不是流式独有的问题,但流式场景下暴露得更快。
还有一个容易忽略的点:及时释放 reader 和 stream。流结束后,reader应该releaseLock(),避免资源泄漏。虽然现代浏览器 GC 会处理大部分情况,但显式释放是好习惯。
7. 上线前的自检清单与踩坑复盘
7.1 一份可以直接抄的自检清单
流式功能上线前,我会按这个清单过一遍:
- [ ]
TextDecoder是否加了stream: true(或用了TextDecoderStream) - [ ] 解析器是否有 buffer 处理半包
- [ ]
flush是否处理了残留数据 - [ ] 分隔符前后端是否统一(
\n\nvs\r\n\r\n) - [ ] 是否有显式的结束事件
- [ ] 服务端是否每次写入后 flush
- [ ] 代理层是否禁用了缓冲(
X-Accel-Buffering: no) - [ ] 是否有心跳保活
- [ ] 重连是否用了指数退避 + 抖动
- [ ]
AbortError是否被正确忽略 - [ ] 业务错误是否走流内事件而非 HTTP 状态码
- [ ] 是否有重试次数上限
- [ ] 高频更新是否做了批量渲染
这份清单里的每一条,背后都是一个真实踩过的坑。尤其是前三条,几乎每个手写流式解析的人都栽过。
7.2 几个印象深刻的坑
坑一:本地好好的,上线就断。本地直连后端,没有代理,流式一切正常。上线后前面挂了网关,60 秒空闲就断。原因是网关的空闲超时是 60 秒,而我们的心跳间隔设成了 90 秒。心跳间隔必须小于链路中所有中间层的最小超时值,这个值要挨个确认,不能想当然。
坑二:中文偶尔乱码。排查了很久,最后发现是某个 chunk 恰好把中文字符切开了,而解码时没加stream: true。这个 bug 复现概率低,但一旦出现就是乱码,非常影响观感。凡是处理 UTF-8 文本流,stream: true是标配。
坑三:最后一个消息丢失。服务端生成完最后一个消息后直接关连接,没有补结尾空行。前端解析器只处理\n\n之前的内容,最后一个消息就留在了 buffer 里,永远没被处理。加上flush逻辑后解决。这个坑的隐蔽性在于,它只在"最后一个消息"上出问题,前面的都正常,很容易被忽略。
坑四:用户点停止后弹错误。用户主动中断生成,结果弹了个"网络错误"。原因是没区分AbortError和真正的网络错误。主动中断是正常操作,不是错误,这个区分做不好,用户体验会大打折扣。
7.3 关于"工程化"的一点个人体会
流式解析这个领域,技术门槛其实不高,难的是把边界情况想全。我见过太多项目,happy path 跑得飞起,一遇到网络抖动、服务端重启、用户切后台就各种诡异问题。工程化的价值,恰恰体现在这些"不 happy"的路径上。
我的建议是:在写第一行流式代码之前,先把状态机和错误处理想清楚。哪些错误可重试、哪些不可重试、断线后消息能不能丢、用户中断怎么处理——这些问题想明白了,代码自然就稳了。反过来,如果一开始只想着"怎么把字蹦出来",后面补这些逻辑会非常痛苦,因为它们是穿插在整个流程里的,不是能"补"上去的。
另外,流式功能的测试不能只测正常流程。要专门构造断网、慢网、服务端中途关闭、发送畸形数据这些场景。Chrome DevTools 的网络限速和离线模式是很好的工具,能模拟出大部分异常情况。有条件的话,写一些针对解析器的单元测试,把各种半包、粘包的输入喂进去,验证输出是否正确。解析器是纯函数式的,非常适合单测,投入产出比很高。
最后说一句关于协议设计的:前后端的流式协议,越简单越好。不要试图在流里塞太复杂的结构,SSE 的文本协议本身就不适合承载复杂语义。把复杂逻辑放在应用层,让流只负责"搬运",这样出问题时排查范围小,替换实现也容易。我见过把整个业务状态机塞进流协议的,后期维护简直是噩梦。