前端SSE封装实战:从协议解析到断线重连,打造稳定的大模型流式输出
发布时间:2026/9/26 17:00:56来源:尧图网络
最近我在做大模型对话类功能服务端把生成结果用SSE一条条往浏览器推。刚开始图省事直接在组件里写了一坨连接逻辑结果断线重连、用户取消生成、页面切换时连接没释放……各种问题扎堆出现。后来我干脆封装了一整套SSE工具把连接管理、协议解析、断线重连、取消逻辑全部收拢到一层项目里所有依赖流式输出的功能都直接复用。这篇就把这套封装的完整思路和实战用法整理出来。SSE技术本身不难但生产环境要用得稳需要抠的细节远比文档里写的多。文章覆盖前端封装方案、协议解析、重连设计、大模型流式渲染以及一堆实测踩坑记录适合正在接SSE的前端或者想给AI接口做流式输出的同学参考。1. SSSE不是新东西先搞清报文和选型逻辑1.1 先看一组真实的SSE报文要理解SSE最直接的办法是看它到底在网络上传输什么。服务端如果开启了SSE返回的Content-Type通常是text/event-stream响应体不是一次性吐出的JSON而是一块一块往外推的文本。打开浏览器Network面板能看到它在很长一段时间里保持pending状态但同时不断有帧数据到达。注意这块调试方式特别直观一会儿第五部分专门讲。典型报文长这样data: 你好我是AI助手 data: 今天可以帮你完成两个任务 data: 第一是翻译第二是总结 event: done data: [DONE]字段行之间用两个换行符分隔。data: 是数据行event: 是自定义事件名。上面这组报文可以理解为先推送一条消息“你好我是AI助手”再推送一条跨两行的内容解析时要把两行data拼成一行最后推一条名为done的事件。这里最基础的规律是一条消息以空行结束消息内部按行解析字段多行data需要拼接。很多人一看这格式觉得简单但真到自己写解析器就会明白细节都在传输边界里。网络层不会保证一次read就能拿到完整消息经常这一块是半条、下一块是半条再加一条完整的。所以解析器必须维护buffer把没拼完的内容留到下一次继续处理。这个坑后面还会详细讲。1.2 和WebSocket放在一起看差别马上就清楚了WebSocket需要先做一次HTTP升级握手之后客户端和服务端在同一个TCP连接上双向收发消息。它是独立的协议在很多内网环境里需要单独开端口或者让网关放行而且断线之后要自己实现重连、心跳、消息补偿。SSE不一样它就是从普通HTTP请求演变来的服务端把响应头改一下、连接不关闭就能持续推数据。SSE是单向下行客户端没法用同一个连接主动给服务端发业务数据。这一点在选型时很重要。如果业务是聊天室、协同编辑这类需要双向频繁交互的场景WebSocket是正确选择。但如果只是“服务端有状态变化、需要主动推给客户端”比如消息通知、进度推送、AI流式输出SSE在成本、兼容性、可维护性上的优势非常明显。有一个容易忽略的点SSE基于HTTP公司内部的网关、负载均衡、防火墙对HTTP最熟悉接入时几乎没有额外成本。WebSocket则需要运维层面对超时、端口、协议升级都要额外关注。所以不是拿SSE硬扛所有场景而是大部分推送场景里SSE刚好够用又足够省事。1.3 为什么大模型接口清一色选SSE最近两年对接大模型接口会发现不同厂商在流式输出的协议选择上惊人一致几乎都基于SSE。原因其实非常简单大模型生成回答是逐token计算、边算边返回服务端要把每一次生成结果及时推给客户端。这类数据的流向是单向的客户端只需要接收不需要在同一连接上反向发数据SSE正好是为此设计的。更现实的因素是AI接口通常要经过网关、内容审核、负载均衡等多个中间层这些中间层对HTTP最友好。SSE直接建立在HTTP之上接入过程不引入新协议运维方便。反观WebSocket如果要用在AI场景里除了协议升级还得处理代理超时、跨端口通信、连接保活等一堆问题对大多数业务属于杀鸡用牛刀。所以只要不是双向实时交互我在前端项目里的默认选择都是SSE。2. 前端接入SSE的两条路线怎么选2.1 原生EventSource开箱即用但有明显边界浏览器自带EventSource用起来确实简单new EventSource(/api/stream)然后addEventListener监听message事件。很多教程把它当作SSE入门示例但真正拿到项目里用几个限制会很快让人头疼。第一EventSource只支持GET请求。现在大部分AI接口要求通过POST提交用户输入、系统参数、会话上下文这些内容用query参数塞都塞不进去长度和结构都受限。第二EventSource不能自定义Header。需要带Authorization或自定义鉴权字段时它直接没有这个能力。第三EventSource的断线重连是浏览器内置的断开后自动重连但策略不可控重试间隔、最大次数都没法按业务配置。第四用户手动取消生成时虽然可以调用close()但在复杂业务状态机里控制的颗粒度不够细。所以EventSource更适合快速验证、无需鉴权、单向通知的轻量场景。一旦项目上了鉴权、上了POST、上了复杂的错误处理或者你对接的是大模型接口基本上都得切换到第二种方案。2.2 fetch ReadableStream把控制权彻底拿回来fetch返回的response.body本身就是一个ReadableStream可以用getReader()拿到底层reader不断调用read()读取服务端推过来的数据块。SSE协议本质上只是HTTP响应体的编码方式解析规则前面说过就那几条自己解析并不难。换来的是对请求的完全控制任意HTTP方法、任意Header、带body、带跨域策略全是fetch的基础能力。这意味着项目里原有的请求封装的拦截器、鉴权逻辑、错误提示都能继续复用。我在公司项目里所有SSE请求都会先过统一鉴权层token过期统一跳登录这些用EventSource根本做不到。而且因为底层是标准fetch未来如果要把这套逻辑迁移到小程序、Node端、桌面端只要平台实现了fetch和ReadableStream核心代码基本不用改。两种方案怎么取舍我直接给个经验表维度原生EventSourcefetch ReadableStream请求方式仅GET任意HTTP方法自定义Header不支持支持断线重连内置且策略固定自行实现完全可控取消连接有close()粒度粗AbortController精准可控多事件分发支持event字段解析后自行分发鉴权兼容弱强适合场景快速原型、简单通知生产环境、AI流式交互对于大模型流式输出直接选第二条路线接下来要写的就是这套方案的完整封装。3. 一个能上生产的SSE封装怎么写3.1 封装前先列四张问题清单动手写代码之前我习惯先把要解决的问题列出来。SSE封装如果只做一个“能接收消息”的壳那跟不封装没有本质区别。真正值得封装的是四块连接生命周期、协议解析、重连策略、事件分发。连接生命周期最容易被忽略。第一版代码里很多人直接在组件里创建fetch和reader页面挂载就连接数据到了就渲染。问题在于页面切换、用户点击停止、重复点击发送时连接谁负责销毁销毁后是否重连正在进行的请求怎么取消。这些逻辑全散在业务层代码越写越乱。封装的核心价值是把生命周期统一收敛页面只需要知道“开始”、“取消”、“收到消息”三个动作底层怎么处理一律不关心。协议解析这块前面说过不只是几个if判断分包、多行data、事件名处理都要处理。重连策略是生产环境的核心网络抖动、网关超时、服务端重启任何一种情况都会中断连接没有重连的话用户只能刷新页面。事件分发则是把SSE协议里的event字段变成业务可见的订阅模型消息、完成、错误各走各的回调不让调用方接触协议细节。这段设计思路有点像面向对象里的“对接口编程而不是对实现编程”。封装里先定义好对外的事件接口内部是EventSource还是fetch实现都无所谓将来换WebSocket也一样能平替业务代码一行不用动。3.2 解析SSE报文必须注意的几个细节写解析器时我踩过不少坑先说几个关键点。第一一条SSE消息以空行作为结束标记。网络传输非常不可靠一次reader.read()拿到的字节可能是半条消息也可能是多条消息混在一起必须引入buffer把收到的字符串暂存起来找到分隔符就切出完整消息处理处理不了的就留在buffer里继续等。第二data行不一定只有一个。服务端发多段内容时会拆成多行data解析逻辑按换行拼接拼接顺序要保持原样。第三data: 后面可能有一个空格也可能没有解析时应该只去掉紧跟冒号的那个空格而不是直接trim整行。直接trim会把正常文本前后的空格误删这一点很容易被忽略。第四冒号开头的行是注释行客户端应当忽略。很多服务端会拿它做心跳每隔一段时间发一行“: ping”中间网络设备看到有数据传输就不会因为idle timeout把连接掐断。客户端解析时要跳过这些注释行不要当成业务数据往外推。3.3 封装实现完整代码与设计说明下面给出一份TypeScript版本的封装能力包括POST请求、自定义Header、协议解析、指数退避重连、主动取消、定时器清理。代码可以直接复制到项目里按需改一下配置项就行。type SSEOptions { url: string; method?: GET | POST; headers?: Recordstring, string; body?: string | FormData | null; maxRetries?: number; retryBaseInterval?: number; maxRetryInterval?: number; idleTimeout?: number; onOpen?: () void; onEvent?: (event: string, data: string) void; onError?: (error: Error) void; onDone?: () void; }; class SSEStream { private controller: AbortController; private retryCount 0; private closed false; private timers: number[] []; constructor(private options: SSEOptions) { this.controller new AbortController(); } start() { this.closed false; this.connect(); } private async connect() { if (this.closed) return; const controller new AbortController(); this.controller controller; try { const response await fetch(this.options.url, { method: this.options.method ?? GET, headers: this.options.headers, body: this.options.body, signal: controller.signal, }); if (!response.ok) { throw new Error(HTTP ${response.status} ${response.statusText}); } if (!response.body) { throw new Error(response.body is empty); } const reader response.body.getReader(); const decoder new TextDecoder(utf-8); let buffer ; let lastDataTime Date.now(); this.retryCount 0; this.options.onOpen?.(); const idleTimer setInterval(() { if (Date.now() - lastDataTime (this.options.idleTimeout ?? 35000)) { controller.abort(); clearInterval(idleTimer); } }, 5000); while (true) { const { done, value } await reader.read(); if (done) break; lastDataTime Date.now(); buffer decoder.decode(value, { stream: true }); buffer this.consumeBuffer(buffer, controller); } clearInterval(idleTimer); if (!this.closed) { this.options.onDone?.(); } } catch (error) { if (this.closed) return; this.options.onError?.(error as Error); this.scheduleReconnect(); } } private consumeBuffer(buffer: string, _controller: AbortController): string { let idx: number; while ((idx buffer.indexOf(\n\n)) ! -1) { const rawMessage buffer.slice(0, idx); buffer buffer.slice(idx 2); const parsed this.parseSseMessage(rawMessage); if (parsed) { this.options.onEvent?.(parsed.event, parsed.data); } } return buffer; } private parseSseMessage(raw: string) { const lines raw.split(\n); let event message; let data ; for (const line of lines) { if (line.startsWith(:)) continue; if (line.startsWith(event:)) { event line.slice(6).trim(); } else if (line.startsWith(data:)) { let value line.slice(5); if (value.startsWith( )) value value.slice(1); data data ? \n : ; data value; } } if (!data) return null; return { event, data }; } private scheduleReconnect() { if (this.closed) return; if (this.retryCount (this.options.maxRetries ?? 3)) return; const delay Math.min( (this.options.retryBaseInterval ?? 1000) * Math.pow(2, this.retryCount), this.options.maxRetryInterval ?? 15000 ); this.retryCount; const timer setTimeout(() { this.connect(); }, delay); this.timers.push(timer); } abort() { this.closed true; this.controller.abort(); this.timers.forEach((t) clearTimeout(t)); this.timers []; } } export function createSSE(options: SSEOptions) { const stream new SSEStream(options); stream.start(); return { abort: () stream.abort(), }; }代码里几个点值得单独说明。connect方法每次都会创建新的AbortController因为同一个controller不能跨请求复用。reader.read()循环里每收到一块数据就刷新lastDataTime然后交给consumeBuffer切分完整消息。idleTimeout的定时器用来对抗代理层空档超过约定时间没有数据就主动abort让catch走到重连逻辑。scheduleReconnect做了指数退避第一次等1秒第二次2秒第三次4秒封顶15秒最大重试次数默认3次。abort方法置closed为true同时清掉所有定时器这个布尔值相当关键它区分了“用户主动取消”和“异常断开”两种状态前者不允许重连后者必须重连。3.4 业务侧接入只写关心的事封装完之后业务侧的代码会变得非常干净const stream createSSE({ url: /api/chat, method: POST, headers: { Content-Type: application/json, Authorization: Bearer ${token}, }, body: JSON.stringify({ prompt: 帮我写一份周报, stream: true }), onEvent: (event, data) { if (event message) { appendAnswer(JSON.parse(data).delta); } else if (event done) { markAsFinished(); } }, onError: (err) { showToast(err.message); }, });业务方只需要关注事件名和data里的业务字段传输层的解析、重连、心跳、取消全部由封装接管。这里我建议把不同环境的API域名做成环境变量注入不要在代码里手写IP或域名避免测试环境、预发环境来回切的时候改错地方。以后如果要把大模型接入从SSE改成WebSocket只要保持onEvent、onError、onDone这套对外接口一致业务代码一行都不用动。4. 实战大模型流式回答实时渲染4.1 增量数据映射到页面状态对接大模型联调时各家返回的data结构不完全一样有的直接是字符串有的外包一层JSON。我们项目里统一约定为这种结构{id:chatcmpl_123,delta:你,created:1690000000}delta字段就是本次新增加的文本前端要做的就是把delta追加到上一次的结果后面。这里最忌讳的做法是每收到一条data就把整个结果字符串setState一次。服务端可能在几百毫秒内推几十条数据每次setState都触发React重新渲染页面会明显卡顿。我的做法是维护两个变量完整文本fullText和渲染文本renderText。onEvent里只更新fullText并打一个“待渲染”标记用requestAnimationFrame统一在下一帧把fullText赋给renderText触发一次setState。这样渲染频率和屏幕刷新率天然对齐看起来非常流畅调用方也只需要关心两三百毫秒一次的视觉更新。4.2 流式Markdown和代码高亮怎么处理大模型回答通常带Markdown格式。流式过程中直接解析Markdown有一个大问题半截语法会让页面疯狂闪烁。代码块刚渲染到一半被截断、表格格式错乱、列表符号乱飘整个界面不稳定。我踩过这个坑之后定了两条规则。第一条流式过程中只展示纯文本不做Markdown解析。视觉上虽然看不到标题和加粗但内容连贯稳定用户的注意力不会被闪烁打断。第二条流结束后再对完整文本做Markdown解析需要代码高亮的话再交给highlight.js之类的库处理。虽然结束时会有一次纯文本到富文本的切换但那是单次变化体验完全能接受。如果产品坚持要边输出边渲染Markdown可以做防抖停止收到数据200毫秒后再解析一次。即便如此代码块的闪烁也很难完全消除。我的建议是能不全量解析就不全量解析优先保证流式过程的视觉稳定性。4.3 多轮对话的竞态与并发控制多轮对话场景里有一个非常典型的问题用户快速发送第二条消息时第一条请求可能还没结束。如果两个SSE同时存在页面上到底该显示哪边的结果我之前遇到过线上事故两条连接同时回流数据回答内容互相覆盖、拼接错乱场面完全失控。解决方案不复杂前端约定同一时间只有一个活跃SSE连接新请求发起前先调用旧连接的abort方法再创建新连接。如果要支持多个会话并排展示就按会话维度独立维护连接但单个会话内部仍然保证单连接。这个约定后端同样要知道服务端收到新连接后要能识别这是同一会话主动断开旧连接不然旧连接还在消耗token。我在后端网关层加了一个会话版本号校验请求头带上会话ID和版本号旧版本号的连接会被强制终止。竞态问题还有一层组件卸载时后端可能还在推数据。封装里abort方法置了closed标记用户关闭页面或切换路由时主动调用就能保证回调不会再触发也避免了对已卸载组件执行setState导致的内存泄漏警告。5. 常见问题与排查实录5.1 连接被网关掐断idle timeout怎么处理这类报错在SSE生产环境极其常见报错文案里通常能看到idle timeout、connection closed、stream disconnected before completion之类的关键词。本质是中间设备比如Nginx、云负载均衡、网关在连接长时间没有数据传输时判定连接没人用主动断开。排查顺序我一般按四步走。第一步先确认客户端是不是长时间收不到数据用Network面板观察帧间隔如果一直有心跳帧在跳问题可能不在idle timeout。第二步确认服务端有没有发心跳没有的话让服务端每10到15秒发一行“: ping”注释行成本极低效果立竿见影。第三步检查代理层配置Nginx的proxy_read_timeout默认60秒如果SSE消息间隔超过这个值连接会被断掉要么调大超时时间要么让心跳间隔短于它。第四步客户端也要做兜底超过30多秒没收到任何数据就主动abort重连。这套策略我称为三层同时处理服务端保活、客户端兜底、代理放行。缺了任何一层都可能出现偶发断连。前端封装的idleTimeout参数就是为这个场景准备的。5.2 中文乱码TextDecoder的stream参数救场用ReadableStream接收中文字符时有一个细节特别容易被忽略多字节字符在传输中会被底层网络包切分。比如“你好”二字的字节序列可能被拆到相邻两个数据块里。如果每次decode都独立处理没有告诉解码器这是一段连续的流第一块数据就会因为字节不完整而出现乱码。英文和数字不受影响所以这个问题排查起来很隐蔽用户反馈往往只是“偶尔蹦出几个乱码”。正确做法是创建一次TextDecoder每次decode传入{stream: true}流结束时再调用一次decoder.decode()来冲刷剩余内容。我在第三节的封装里已经这样处理了。这个参数一定要写上不然后期上线就会收到一堆中文乱码类的反馈。5.3 长文本渲染的内存控制一次长对话下来完整回答的文本可能上万字符。前端如果把全部内容都放进DOM里持续渲染页面会越来越卡。我的建议是保留“原文全量”和“渲染截断”两个层次。原文可以交给后端存储或者放在不参与渲染的状态里给导出、复制、历史记录功能使用。渲染层只保留最近一段内容超过阈值就截断展示再提供一个“查看完整内容”的入口。如果产品上必须整段展示全文那至少要做虚拟滚动或者按段落懒加载不要让浏览器一次性排版上万行DOM节点。另外提醒一点代码高亮库对大段代码块的处理相当吃性能等流式结束后再高亮能省下主线程不少时间。5.4 生产环境SSE调试三板斧第一斧curl直连。用-N参数让curl不缓冲直接把SSE流实时打到终端确认接口返回的Content-Type和data内容是否正常curl -N http://localhost:8080/api/chat \ -H Content-Type: application/json \ -d {prompt:你好,stream:true}第二斧浏览器Network面板。SSE连接会一直保持pending状态点进去能看到Response里的帧记录一条条消息按时间排列帧间隔一目了然还能直接判断有没有心跳、消息是否被截断。第三斧服务端日志和客户端日志双向对齐。我在服务端每条SSE消息里带一个序号或时间戳客户端接收后把最后一条消息序号打印出来。两边一对立刻就能定位是服务端没发、还是客户端没收到。这个方法比两边各查各的日志高效得多建议直接安排在联调流程里。封装写好之后我给自己定了一个额外规矩凡是用到SSE的页面路由离开时必须主动abort并清理所有定时器。早期我犯过低级错误组件卸载没走完abort流程路由切走之后页面还挂着SSE连接后端Nginx连接数一路被拖满。从那以后每次发版前我都会先用Network面板扫一遍有没有遗留的pending连接。SSE本身不难难的是把生命周期和边界条件都处理干净。你如果正准备封装自己的SSE工具建议从第一天就把连接销毁当成一等公民设计进去后面能少踩很多坑。
网站建设高端定制企业官网