ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

SSE流式解析工程化实战:从协议选型到中断恢复的完整落地指南

SSE流式解析工程化实战:从协议选型到中断恢复的完整落地指南 流式解析这件事表面上看只是把接口返回的数据一段段吐到页面上但真到工程里落地问题会一个接一个冒出来SSE 连接莫名其妙断了、中文被切成半个字、用户点了停止按钮后台还在烧 token、页面切走了流还在跑。我在几个 AI 交互类项目里反复踩过这些坑之后慢慢总结出一套相对稳的工程化做法。这篇就围绕流式解析的工程化落地展开从协议选择、数据流管道搭建、中断控制、异常恢复到性能优化把每个环节的为什么和怎么做都讲透。不管你是刚接触 SSE 的前端还是想理清整条链路的全栈看完应该都能直接抄作业。1. 为什么流式解析值得单独做工程化1.1 从能跑到跑得稳之间的鸿沟很多人第一次接大模型接口写法都很朴素fetch拿到response.bodygetReader()读TextDecoder解码然后按\n\n切分取data:后面的内容JSON.parse一下把delta.content拼到页面上。这套代码在本地跑 Demo 完全没问题网速好的时候体验也顺滑。但只要上线问题就来了。我印象最深的一次是某个下午用户反馈回答卡住不动了。排查半天发现是网络抖动导致 SSE 连接被中间层掐断但前端没有任何重连或收尾逻辑界面就永远停在半句话上。还有一次是用户输入了一段带 emoji 的 prompt结果流式渲染时 emoji 变成了乱码方块——原因是TextDecoder每次decode一个 chunk而 UTF-8 的多字节字符恰好被切在了两个 chunk 中间。这些问题的共同点是它们都不在主流程里而在边界条件里。Demo 只跑主流程工程要处理所有边界。这就是为什么流式解析值得单独拎出来做工程化——它不是调个接口那么简单而是一整条数据管道的可靠性设计。1.2 流式解析到底在解决什么问题先明确一个前提大模型的输出是逐 token 生成的如果等全部生成完再一次性返回用户要盯着空白屏幕等好几秒甚至十几秒。流式解析的核心价值就是把这十几秒的等待摊开——模型生成一个字前端就渲染一个字用户感知到的首字延迟从十几秒降到几百毫秒。这个体验差异是巨大的。心理学上有个说法用户对等待的容忍度取决于是否有反馈。一个转圈的 loading 和一个正在逐字蹦出的回答感受完全不同。所以流式解析不是锦上添花而是 AI 交互产品的体验底线。但逐字渲染这个目标落到工程上要解决一串问题数据怎么传协议、怎么切分帧、怎么解编码、怎么渲染性能、怎么停中断、断了怎么办恢复。每一个都是独立的工程点串起来才是一条完整的管道。1.3 一条完整的流式管道长什么样我习惯把整条链路拆成五层来看这样排查问题时能快速定位是哪一层出了岔子层级职责常见技术典型故障传输层建立长连接、传输字节流SSE / WebSocket / fetch stream连接被掐断、idle timeout分帧层把字节流切成一条条消息按\n\n切分、事件边界识别半包、粘包解码层字节转字符串、JSON 解析TextDecoder、JSON.parse中文乱码、解析报错状态层维护消息列表、增量拼接状态管理、增量更新重复渲染、状态错乱渲染层把增量内容画到界面虚拟 DOM、节流批量更新卡顿、掉帧这五层里传输层和分帧层是最容易出问题的因为它们直接面对不可靠的网络。解码层次之主要坑在多字节字符。状态层和渲染层更多是性能问题。后面我会逐层展开。2. SSE 与 Web Streams API 的选型逻辑2.1 为什么大模型流式输出普遍选 SSE先说结论绝大多数大模型接口的流式输出用的是 SSEServer-Sent Events而不是 WebSocket。这不是随便选的背后有很实在的理由。WebSocket 是全双工协议客户端和服务端可以随时互发消息适合聊天室、协同编辑这种双向高频通信场景。但大模型的流式输出本质上是一问一答客户端发一次请求服务端持续推一段时间的响应推完就结束。这是典型的单向流用 WebSocket 属于杀鸡用牛刀。SSE 的优势在于它基于普通 HTTP服务端返回Content-Type: text/event-stream然后持续往响应体里写数据。对客户端来说就是一个永远读不完的 HTTP 响应。它天然支持断线重连浏览器原生EventSource有重连机制实现简单穿透代理和网关也比 WebSocket 友好——很多企业网关对 WebSocket 的升级握手有额外限制但对普通 HTTP 长响应基本放行。注意虽然浏览器原生有EventSource但实际项目里我几乎不用它。原因后面 2.3 会讲。2.2 SSE 的数据格式到底长什么样SSE 的协议格式其实很简单一条消息由若干字段行组成字段行之间用换行分隔消息之间用空行\n\n分隔。常见字段有data:消息内容可以有多行多行会被拼接event:事件类型客户端可以据此区分不同消息id:消息 ID用于断线重连时定位retry:重连等待时间毫秒大模型接口返回的典型长这样data: {choices:[{delta:{content:你}}]} data: {choices:[{delta:{content:好}}]} data: [DONE]注意最后那个data: [DONE]这是 OpenAI 风格的结束标记不是标准 SSE 的一部分而是接口约定。解析时看到它就该收尾了。这里有个容易忽略的点data:后面通常有一个空格但规范里这个空格是可选的。有些实现会写成data:{...}如果你的解析代码硬编码了data:带空格遇到不带空格的就会解析失败。稳妥的做法是切掉前缀后trim()一下。2.3 为什么我放弃 EventSource 改用 fetch ReadableStreamEventSource看起来很美但它有几个硬伤导致在真实项目里基本不可用第一它只支持 GET 请求没法带请求体。而大模型接口通常需要 POST 一个包含 messages、model、temperature 等参数的 JSON body。这一条就直接判了死刑。第二它没法自定义请求头。很多接口需要Authorization头传 API KeyEventSource加不了。第三它的重连是自动的、不可控的。断线后它会自己重连但重连时不会带上你原来的请求体等于重连了个寂寞还可能造成重复请求。所以实际项目里主流做法是用fetch发 POST 请求然后手动读取response.body这个ReadableStream。这样请求方法、请求头、请求体全都可控中断也能通过AbortController精确控制。代价是要自己处理分帧和重连但这部分逻辑封装一次就能复用。const controller new AbortController(); const response await fetch(/api/chat, { method: POST, headers: { Content-Type: application/json, Accept: text/event-stream }, body: JSON.stringify({ messages, model: gpt-4 }), signal: controller.signal }); const reader response.body.getReader();这段代码是整个流式解析的起点。controller.signal挂上去之后任何时候调controller.abort()都能立刻中断请求这是后面中断控制的基础。2.4 Web Streams API 提供的管道能力response.body是一个ReadableStream这是 Web Streams API 的核心对象。它提供了getReader()方法让你能一段段地读数据。但 Web Streams API 的能力不止于此它还有TransformStream和WritableStream三者可以串成一条管道ReadableStream → TransformStream → TransformStream → WritableStreamTransformStream是这条管道的灵魂。它接收上游的 chunk处理后传给下游。你可以把字节转字符串做成一个 Transform字符串切分成消息做成另一个 TransformJSON 解析再做一层。每一层职责单一组合起来就是一条清晰的数据处理流水线。这种管道式设计的最大好处是可测试、可复用。每个 Transform 都是纯函数式的转换单独写单测很容易。而且pipeThrough的写法比手写while(true) { await reader.read() }循环要优雅得多也更符合流式处理的思维模型。3. 手写一条可靠的流式解析管道3.1 用 TransformStream 拆解处理阶段我一般把整条管道拆成三个 Transform对应前面说的分帧、解码、解析三个阶段。先看分帧这一层它的任务是把字节流按 SSE 的消息边界切开function createSSEFramer() { let buffer ; const decoder new TextDecoder(utf-8); return new TransformStream({ transform(chunk, controller) { // stream: true 是关键见 3.2 buffer decoder.decode(chunk, { stream: true }); let boundary; while ((boundary buffer.indexOf(\n\n)) ! -1) { const rawEvent buffer.slice(0, boundary); buffer buffer.slice(boundary 2); controller.enqueue(rawEvent); } }, flush(controller) { // 流结束时把残留数据吐出去 buffer decoder.decode(); if (buffer.trim()) { controller.enqueue(buffer); } } }); }这段代码有两个关键点。第一是buffer的存在网络传过来的 chunk 不保证按消息边界对齐可能一条消息被切成两半也可能两条消息挤在一个 chunk 里。所以必须用一个缓冲区累积只有遇到完整的\n\n才切出一条消息。第二是flush方法流正常结束时缓冲区里可能还有没遇到\n\n的残留数据得在flush里补一刀否则最后一条消息会丢。3.2 TextDecoder 的 stream 参数中文乱码的根源decoder.decode(chunk, { stream: true })里的stream: true是很多人会漏掉的关键参数。它的作用是告诉解码器这个 chunk 可能不是完整的如果结尾有多字节字符被切断了先别急着报错把不完整的字节缓存起来等下一个 chunk 来了再拼。如果不加这个参数遇到 UTF-8 多字节字符被 chunk 边界切开的情况解码器会直接输出替换字符那个菱形问号中文就变乱码了。一个中文字符在 UTF-8 里占 3 个字节一个 emoji 可能占 4 个字节网络传输时被切开是家常便饭。我踩过一次这个坑当时测试用的是纯英文 prompt一切正常。上线后用户输入中文偶尔出现乱码还很难复现。后来才定位到是TextDecoder没开stream模式。这个坑的隐蔽性在于它只在特定网络分片情况下触发本地测试很难稳定复现。提示flush阶段调用decoder.decode()不带参数是为了把解码器内部缓存的最后几个字节吐出来。这一步也不能省。3.3 消息解析从 data 字段到业务对象分帧之后每条消息还是原始字符串需要进一步解析。这一层要做三件事识别字段、提取 data、JSON 解析。function createMessageParser() { return new TransformStream({ transform(rawEvent, controller) { const lines rawEvent.split(\n); let dataPayload ; for (const line of lines) { if (line.startsWith(data:)) { // 去掉前缀后 trim兼容带空格和不带空格两种写法 dataPayload line.slice(5).trim(); } } if (!dataPayload) return; // 结束标记 if (dataPayload [DONE]) { controller.enqueue({ type: done }); return; } try { const parsed JSON.parse(dataPayload); const delta parsed.choices?.[0]?.delta?.content; if (delta) { controller.enqueue({ type: delta, content: delta }); } } catch (err) { // 解析失败不要中断整条流记录后跳过 console.warn(SSE 消息解析失败:, dataPayload, err); } } }); }这里有个工程上的取舍JSON 解析失败时我选择记录日志后跳过而不是抛错中断整条流。原因是流式场景下一条消息解析失败不应该影响后续消息的接收。如果直接抛错用户会看到回答戛然而止体验更差。当然如果失败率异常高说明接口格式变了那要另做告警。3.4 把管道串起来完整的数据流三个 Transform 准备好之后用pipeThrough串起来最后接一个WritableStream或者手动读const response await fetch(/api/chat, { /* ... */ }); const stream response.body .pipeThrough(createSSEFramer()) .pipeThrough(createMessageParser()); const reader stream.getReader(); while (true) { const { done, value } await reader.read(); if (done) break; if (value.type delta) { appendToUI(value.content); } else if (value.type done) { finalizeUI(); break; } }这条管道的好处是每一层都可以单独测试。比如我想验证分帧逻辑就构造一个把消息切碎的字节流喂给createSSEFramer看它能不能正确还原。想验证解析逻辑就直接喂字符串。这种可测试性在排查线上问题时价值巨大。4. 中断、超时与异常恢复的实战处理4.1 AbortController让停止按钮真正停下来用户点停止生成按钮时如果只是前端停止渲染后台的请求还在跑token 还在烧这是很浪费的。正确做法是用AbortController把请求真正掐断。const controller new AbortController(); // 发起请求时挂上 signal fetch(/api/chat, { signal: controller.signal, /* ... */ }); // 用户点停止 function handleStop() { controller.abort(); finalizeUI(); // 把已收到的内容定稿 }abort()调用后fetch的 promise 会以AbortError拒绝reader.read()也会抛出异常。所以读取循环要包一层 try-catch把AbortError单独处理——它不是错误而是用户主动取消的正常流程。try { while (true) { const { done, value } await reader.read(); // ... } } catch (err) { if (err.name AbortError) { // 用户主动停止正常收尾 finalizeUI(); } else { // 真正的错误 handleError(err); } }这里有个细节abort()之后已经收到的内容要保留不能清空。用户点了停止是想就到这里而不是全部不要。所以finalizeUI要把当前已渲染的内容定稿并可能补一个已停止的标记。4.2 idle timeout那个让人头疼的连接断开线上最常见的一类报错是stream disconnected before completion: idle timeout waiting for sse。字面意思是SSE 等待超时流在完成前断开了。这个问题的根源通常不在客户端而在中间的代理层或网关。很多反向代理如 Nginx有默认的读超时比如 60 秒。如果服务端在这段时间内没有往连接里写任何数据代理就会认为连接空闲主动掐断。大模型在生成很长的回答时如果中间有较长的思考停顿比如推理模型就可能触发这个超时。解决思路有几个方向服务端定期发送心跳注释行: ping\n\n保持连接活跃调整代理的读超时时间客户端检测到异常断开后用已收到的上下文发起续写请求心跳是最常用的手段。SSE 规范里以冒号开头的行是注释客户端会忽略但能起到占位作用让代理认为连接是活跃的: keep-alive data: {choices:[{delta:{content:...}}]}服务端每隔 15 到 30 秒发一次心跳就能有效避免 idle timeout。4.3 断线续传用已收到的内容做恢复如果连接真的断了最优雅的处理是续传——把已经收到的内容作为上下文让模型接着往下写。这需要服务端支持客户端把已生成的部分回传服务端在 prompt 里说明请接着以下内容继续。不过续传有个坑模型可能会重复已经说过的内容。所以续传时要明确告诉模型不要重复直接从断点继续。实测下来续传的成功率不是 100%有时候模型会重新组织语言导致前后文风格不一致。所以我的做法是短回答直接重试长回答才用续传并且给用户一个继续生成的按钮让用户自己决定。4.4 错误分类与用户提示流式场景下的错误五花八门但给用户的提示不能都是出错了。我一般把错误分成几类分别给不同的提示错误类型判断依据用户提示处理方式用户取消AbortError无提示保留已生成内容网络断开TypeError / 网络错误网络不稳定请重试提供重试按钮服务端错误HTTP 4xx/5xx服务暂时不可用展示错误码解析失败JSON 解析异常无提示静默跳过记录日志超时idle timeout生成超时可继续提供续写按钮这种分类处理能让用户知道发生了什么而不是面对一个笼统的错误。尤其是用户取消这一类绝对不能弹错误提示否则用户会以为是自己操作导致的故障。5. 渲染性能与状态管理的细节5.1 高频增量更新为什么会卡流式渲染时模型可能每秒吐出几十个 token如果每个 token 都触发一次 React 的setState就会造成高频重渲染。在回答很长的时候页面会明显卡顿甚至掉帧。问题的本质是渲染的频率超过了屏幕刷新的频率。屏幕一般 60Hz也就是每 16.7ms 刷新一次而 token 到达的频率可能远高于此。多出来的渲染都是浪费。解决办法是批量更新——把短时间内的多个增量合并成一次渲染。常见做法是用一个缓冲区累积增量然后用requestAnimationFrame或者定时器批量 flushlet pending ; let rafId null; function appendToUI(text) { pending text; if (rafId) return; rafId requestAnimationFrame(() { setContent(prev prev pending); pending ; rafId null; }); }这样无论 token 来得多快每帧最多只渲染一次。实测下来这个改动能让长回答的渲染帧率稳定在 60fps。5.2 增量拼接 vs 全量替换另一个性能点是状态更新的方式。有两种做法第一种是增量拼接每次把新内容 append 到已有内容后面。第二种是全量替换每次用完整的新内容替换旧内容。增量拼接的性能更好因为它只处理新增的部分。但它的前提是状态里存的是累积后的完整字符串每次 append 都要做一次字符串拼接。在 JavaScript 里字符串是不可变的每次拼接都会创建新字符串。如果回答有几万字频繁拼接会有内存压力。全量替换则相反如果服务端每次都返回完整内容有些接口是这样设计的那客户端直接替换即可逻辑简单但传输量大。我一般用增量拼接因为大模型接口返回的就是 delta。但要注意拼接时不要用在超长字符串上反复操作可以考虑用一个数组存片段渲染时join。不过实测下来几万字的字符串拼接在现代浏览器里性能完全够用没必要过度优化。5.3 滚动跟随与用户打断的平衡流式渲染时内容不断增长页面需要自动滚动到底部让用户看到最新内容。但如果用户手动往上滚去看历史这时候自动滚动就会把用户拽回底部体验很糟。正确的做法是检测用户是否在底部附近。如果在底部就自动跟随如果用户往上滚了就停止自动滚动并显示一个回到底部的按钮。function isNearBottom(el, threshold 50) { return el.scrollHeight - el.scrollTop - el.clientHeight threshold; }这个threshold是关键不能设成 0。因为滚动位置是浮点数用户滚到几乎底部时scrollHeight - scrollTop - clientHeight可能是个很小的正数。设个 50px 的容差体验会自然很多。5.4 状态管理把流式状态和消息状态分开在 React 里我习惯把流式进行中的状态和消息列表的状态分开管理。消息列表是持久化的业务数据流式状态是临时的 UI 状态。const [messages, setMessages] useState([]); // 消息列表 const [streaming, setStreaming] useState(false); // 是否正在流式 const [streamingText, setStreamingText] useState(); // 当前流式内容流式进行时增量内容先写到streamingText渲染时把它作为一个临时消息显示在列表末尾。流式结束后再把streamingText作为一个正式消息 push 进messages清空streamingText。这样做的好处是消息列表的更新频率很低只在流结束时更新一次而高频的增量更新只影响streamingText这一个状态。React 的 diff 范围小性能自然好。而且这种分离让停止生成的逻辑很清晰——停止时把streamingText定稿进messages即可。6. 几个真实踩坑与排查链路6.1 中文乱码从现象到根因的完整排查现象用户反馈偶尔出现乱码本地无法复现。第一步我先怀疑是接口返回的编码问题。抓包看响应头Content-Type: text/event-stream; charsetutf-8编码声明没问题。第二步怀疑是TextDecoder的用法。检查代码发现写的是decoder.decode(chunk)没有stream: true。这就是根因——多字节字符被 chunk 边界切开时解码器无法正确处理。第三步验证。我写了个测试手动构造一个把中文字符从中间切开的字节流喂给不带stream的解码器果然输出乱码。加上stream: true后正常。这个坑的教训是涉及字节流解码的地方一定要考虑数据可能被任意切分这个前提。TextDecoder的stream参数就是为这个场景设计的别漏。6.2 最后一条消息丢失flush 的重要性现象流式回答的最后几个字偶尔丢失。排查时我盯着分帧逻辑看发现buffer里累积的数据只有在遇到\n\n时才被切出去。如果流结束时最后一条消息后面没有\n\n有些实现最后一条不带空行那这条消息就永远留在 buffer 里随着流结束被丢弃。修复方法就是在TransformStream的flush里把 buffer 里剩余的内容也 enqueue 出去。这个坑很隐蔽因为大部分时候最后一条消息后面是有\n\n的只有特定实现才会漏。6.3 重复请求自动重连的副作用现象偶尔出现同一条消息被回答两次。排查发现是用了EventSource的自动重连。当连接因为某种原因断开时EventSource会自动重连但重连时它不知道上次请求的 body因为EventSource只支持 GET于是服务端把它当成一个新请求重新生成了一遍回答。这个坑的根因是自动重连和请求体的矛盾。解决办法就是前面说的放弃EventSource用fetch手动控制重连重连时带上完整的请求上下文并且用消息 ID 去重。6.4 内存泄漏忘记取消订阅现象页面切换后控制台还在打印流式日志。排查发现是组件卸载时没有中断流。React 组件卸载后fetch的读取循环还在跑回调还在执行造成内存泄漏。修复方法是在useEffect的清理函数里调用controller.abort()useEffect(() { const controller new AbortController(); startStream(controller.signal); return () { controller.abort(); // 组件卸载时中断 }; }, []);这个坑在单页应用里很常见尤其是用户快速切换页面时。养成发起请求就想着怎么取消的习惯能避免很多这类问题。7. 工程化落地的几点个人体会流式解析的工程化说到底是在处理不确定性。网络是不确定的数据分片是不确定的用户操作是不确定的。工程化的价值就是把这些不确定性都收敛到可控的范围内。我个人的几条经验第一永远假设数据会被任意切分无论是字节层面还是消息层面都要有缓冲和边界处理。第二永远给用户一个停止的出口并且这个停止要真正掐断后台请求而不是只停前端渲染。第三错误要分类用户取消不是错误网络抖动要能重试服务端故障要能提示。第四性能优化要抓主要矛盾高频渲染用批量更新长列表用虚拟滚动别过早优化。还有一点是关于测试的。流式逻辑很难用常规的单元测试覆盖因为它的行为依赖时序。我的做法是把每个 Transform 单独测构造各种刁钻的输入——被切碎的消息、不带空行的结尾、超长消息、非法 JSON——看管道能不能正确处理。这些测试用例往往就是线上问题的预演。最后说个容易被忽略的点日志。流式场景出问题时光看前端报错很难定位。我一般会在关键节点打日志——连接建立、每条消息解析、流结束、异常抛出——并且带上时间戳和消息序号。这样出问题时能快速判断是没收到数据还是收到了但解析失败排查效率能提升一大截。
返回列表