SSE重连重复输出:游标续传与事件去重实战指南
发布时间:2026/10/4 2:06:26 锦皓数字建站

最近在做大模型对话产品的流式前端时被一个很不起眼的问题折腾了两天SSE 连接一断开重连之后聊天框像复读机一样开始重复吐字上一批已经渲染过的 token 又原封不动来了一遍。有时候只是重复几十个字有时候整段回答从头再来用户那边看到的就是满屏的复读。排查到最后发现根子不在大模型也不在单纯的网络抖动而在我们对 SSEServer-Sent Events重连语义的理解上——很多人包括我一开始都把重连当成了重新请求而不是断点续传。这篇文章就围绕这个场景把这套问题的机制、根因、解决方案和完整代码一步步讲清楚。适合正在接大模型流式输出的前端开发、负责流式接口封装的后端同学以及想要封装一套可靠 SSE 流式调用逻辑的团队参考。不管你是用 Vue、React 还是原生 JS也不管后端是 Python FastAPI 还是 Node.js思路都能直接套用。1. 先还原现场流式输出的复读机是怎么出现的1.1 最典型的复现路径我遇到的场景是这样的用户在手机端提问大模型已经输出了大概 800 个字此时用户从 WiFi 切换到移动网络TCP 连接瞬间断开。前端检测到断流后发起重连后端收到一个新的 HTTP 请求于是拿着原始 prompt 重新调用大模型从第一个 token 开始重新生成。前端这边的逻辑如果只是简单地把每次收到的内容追加到文本末尾那结果就是已经存在的 800 个字后面又追加了一遍从第 1 个字开始的新内容。另一个高频场景出现在长回答上。模型思考时间长或者两次 token 推送间隔超过代理的空闲超时时间连接被 Nginx 或云负载均衡静默掐断。前端等了一会儿没数据触发重连同样的问题又来一遍。你在 DevTools 里看到的网络状态往往是(failed) net::ERR_EMPTY_RESPONSE后端日志里则是stream disconnected before completion: idle timeout waiting for sse。1.2 为什么本地开发死活复现不出来这种问题最烦人的地方在于本地开发环境几乎不会触发。本地回环网络没有丢包没有代理层SSE 连接想断都难。我最初在本地用 Chrome 和 Postman 反复测试接口完全正常一次断流都没有。直到部署到测试环境开始有真机用户访问问题才集中爆发。所以如果你也在做流式输出务必在联调阶段就主动模拟断网、休眠恢复、代理超时别等到线上被用户截图挂出来。1.3 先分清两件容易混淆的事在动手排查前先厘清两个层面的重复。TCP 层的重传不背这个锅。TCP 协议本身会处理数据包丢失和重传应用层拿到的字节流是完整且不重复的这里不会出现同一个 token 两次到达的问题。真正的问题在应用层。要么是服务端把同一个请求重新执行了一遍要么是前端把已经渲染过的区间又追加了一次。另外如果你的前端请求库配置了自动重试HTTP 层面的重试也会造成同样的效果——每次重试都是一次全新的大模型调用。所以排查时第一步就是确认重复的内容到底是从服务端重新生成的还是同一个响应被重复渲染了。这两条路的修法完全不同。2. SSE 重连的机制缝隙问题为什么必然发生2.1 text/event-stream 的帧格式与续传锚点SSE 协议本身很简单服务端返回Content-Type: text/event-stream然后按固定格式往连接里写文本帧。一帧由若干字段和两个换行符组成id: 12 event: token data: {delta:你好}id字段是事件的全局序号event是事件类型data是实际内容空行表示这一帧结束。还有一个retry字段可以控制客户端重连的等待毫秒数形如retry: 3000。另外 SSE 允许发送注释行也就是以冒号开头的行: keep-alive这种注释行不会被解析成事件但它是实实在在的字节能起到保活的作用。后面讲心跳就会用到它。2.2 Last-Event-ID 是唯一可靠的续传凭据SSE 协议里有一个专门为重连设计的机制浏览器在重连时会自动携带一个Last-Event-ID请求头值就是客户端最后收到的那帧事件的id。服务端读取这个头就能知道客户端已经看到了第几条从而只回放后续内容。这是整个协议里唯一的续传锚点。可惜很多后端实现根本没读这个头。我的后端同事一开始就是直接忽略这个字段收到连接就重新跑 prompt效果等于把断线重连变成了重新生成。要解决重复吐字第一件必须做的事就是让服务端认识Last-Event-ID。2.3 原生 EventSource 与 fetch 流式方案的本质差别很多前端团队一开始会用浏览器原生的EventSource对象因为它自带自动重连而且重连时确实会自动带上Last-Event-ID。但用着用着就会发现一个硬伤EventSource不允许自定义请求头你没法往里塞Authorization: Bearer xxx也没法塞业务需要的 traceId。于是大家纷纷改用fetchReadableStream手动解析流式响应。这个切换本身就是重复吐字的温床。换成 fetch 后自动重连没了Last-Event-ID只能手动在下次请求的 header 里带上手动重连又引入了重连时机、重试次数、终止状态等问题。任何一个环节没做对都会导致重复。所以如果你想用 fetch 方案就得把协议原生提供的功能逐一重新实现一遍一个都不能漏。3. 根因剖析重复吐字的四个来源3.1 服务端把断线重连当成了新请求这是最普遍、最致命的一个来源。服务端无状态设计本身没问题问题在于你用什么标识来区分新请求和重连。如果服务端只认连接本身看到新连接就调用大模型那必然重新生成整个回答。这里有个容易被忽略的技术现实大模型推理接口本身不支持从第 N 个 token 继续生成。OpenAI 兼容的流式接口也好本地部署的模型也好一旦断掉你能做的就是重新传入完整上下文重新生成。所谓断点续传只能靠服务端自己在内存里把已经吐出来的 token 缓存住等客户端重连时再发出去。没有这层缓存续传就是空谈。3.2 前端缓冲区在断点处重复追加就算服务端做到了续传前端也可能制造重复。比如服务端在重连后把最后 20 个 token 重新发了一遍前端如果用的是content delta这种简单追加方式自然会把这 20 个 token 再渲染一遍。问题本质是前端不知道自己已经消费到哪个位置了。要解决前端必须有游标的概念。最简单的游标就是每帧事件的id。收到新帧时先判断id是否已经处理过处理过就跳过。这个判断必须放在渲染之前否则你更新了 UI 才发现是重复帧已经晚了。3.3 完成事件与断开信号的竞态还有一种隐蔽的重复大模型已经输出完了[DONE]帧刚发出来连接恰好在这时断开。前端的事件处理逻辑里错误处理优先于完成处理于是把这次断开当成流中断触发重连。服务端一看生成早结束了要么把整个缓冲从头再发一遍要么干脆重新生成。用户看到的现象就是回答已经完整显示闪了一下又开始从头吐字。解决办法是给前端连接加一个明确的状态机至少包含streaming、reconnecting、done、stopped四种状态。一旦进入done任何网络错误都不能再触发重连。这个状态切换的代码要放在 error handler 之前判断。3.4 代理空闲超时引发的假断开与假重复生产环境里SSE 连接要经过一层或多层代理。Nginx 默认的proxy_read_timeout是 60 秒云负载均衡通常也设置了空闲超时。所谓空闲不是指没有新 token而是指连接上没有任何字节流动。大模型偶尔思考时间超过 60 秒或者你的服务端在等待内部队列时一句话都不发代理就会判定连接空闲直接掐断。这会造成假的断开和假的重复。前端看到连接断了重连服务端重新生成——一切都是合理的但用户端出现的就是重复。更隐蔽的是如果 Nginx 开着proxy_buffering它会把 SSE 数据攒在缓冲区里直到超时才一次性吐给客户端不仅流式体验全毁还容易在重放时出现大段重复。后面配置章节我会专门讲怎么关掉这些坑。4. 解决思路游标续传、事件去重与幂等兜底4.1 方案一服务端事件缓冲 游标回放服务端为每次对话维护一个会话对象里面存两个东西已经生成的事件列表带递增的id以及一个是否已经完成的标记。新连接进来时读取Last-Event-ID如果这个会话已经存在就从last_event_id 1的位置开始把事件逐个回放回放结束后继续跟随实时生成的新事件。这里有一个设计要点事件列表里存的是已生成但待发送的 token不是原始 prompt。大模型每吐一个 token服务端就把它追加到列表里同时唤醒那个正在等待的 SSE 生成器。这样无论客户端什么时候重连都能从断点继续而不是从头再来。缓冲要设上限和过期时间。单次对话的 token 文本量并不大即使一万个 token 也就几十 KB正常业务撑得住但为了防止极端情况下内存膨胀事件数超过上限就丢弃最老的会话超过 10 分钟没人连接就整个清理。4.2 方案二前端按事件 ID 跳过已处理帧服务端做对游标回放后前端仍然要加一道防线按事件id去重。原因是服务端的缓冲可能因为重启、超时而被清空一旦清空回放逻辑失效客户端收到的还是从头开始的完整内容。这时候唯一能保护 UI 的就是前端自己记住我已经处理到哪个id了。具体实现是在前端解析 SSE 帧时先取id如果id小于等于本地已处理的processedId直接丢弃该帧否则更新processedId并渲染。这一步必须是渲染前的硬判断不能依赖字符串拼接后的检查。4.3 方案三字符级重叠消解兜底除了按帧去重还可以在渲染层做一个字符级兜底新内容追加前检查新内容的开头和现有文本的结尾是否有重叠有就去掉重叠部分再追加。function trimOverlap(existing, incoming) { const maxOverlap Math.min(existing.length, incoming.length, 200); for (let len maxOverlap; len 0; len--) { if (existing.endsWith(incoming.slice(0, len))) { return incoming.slice(len); } } return incoming; }但这个方案有两个明显的局限。第一它只能处理头部和尾部重叠的情况也就是重放恰好从你已渲染文本的尾部开始如果服务端从头完整重放重叠部分在文本中间endsWith匹配不上。第二JavaScript 字符串是 UTF-16 编码直接用slice可能会把一个 emoji 或生僻字从代理对中间切开造成乱码。真要按字符处理得先Array.from(str)转成码点数组再操作。所以这个方案只能作为辅助兜底不能作为主要手段。4.4 方案四请求幂等防止重复生成与重复计费大模型调用是要花钱的重复生成不仅体验差成本也直接翻倍。前端生成一个requestIdUUID每次重连都复用同一个 ID。服务端以requestId为键做幂等如果这个 ID 已经在生成中或已经完成直接返回已有的事件缓冲如果不存在才真正发起大模型调用。这有点像支付接口的幂等键设计——同一个业务请求无论重试多少次只被执行一次。加上这层之后即使前端因为 bug 重连十几次后端也只会产生一次模型调用钱不会白烧。4.5 取舍建议这四个方案不是四选一而是层层叠加上去的层级方案解决的问题代价服务端事件缓冲 游标回放重连后不重新生成内存占用、会话管理前端事件 ID 去重兜底服务端缓冲丢失维护 processedId渲染层字符重叠消解处理尾部重叠等边缘情况需注意字符边界请求层requestId 幂等防止重复调用与重复计费服务端需缓存状态我实际项目里就是全都要缺一层都有对应的翻车场景。5. 实操落地前后端完整改造5.1 后端实现FastAPI 事件缓冲下面是一个基于 FastAPI 的最小可用实现核心是SessionBuffer这个会话对象和event_source这个异步生成器。import asyncio import json import time from fastapi import FastAPI, Request from fastapi.responses import StreamingResponse app FastAPI() class SessionBuffer: def __init__(self): self.events [] # 事件列表每项 {id: int, data: dict} self.finished False self.notify asyncio.Event() self.updated_at time.time() self.max_size 5000 def append(self, data: dict, finished: bool False): self.events.append({id: len(self.events), data: data}) if len(self.events) self.max_size: self.events self.events[-self.max_size:] self.finished finished self.updated_at time.time() self.notify.set() sessions: dict[str, SessionBuffer] {} app.get(/api/chat/stream/{session_id}) async def chat_stream(session_id: str, request: Request): if session_id not in sessions: sessions[session_id] SessionBuffer() asyncio.create_task(run_generation(session_id)) buf sessions[session_id] try: last_event_id int(request.headers.get(last-event-id, -1)) except ValueError: last_event_id -1 async def event_source(): cursor last_event_id 1 while True: while cursor len(buf.events): ev buf.events[cursor] cursor 1 yield format_sse(ev[id], ev[data]) if buf.finished: yield event: done\ndata: [DONE]\n\n return try: await asyncio.wait_for(buf.notify.wait(), timeout15) buf.notify.clear() except asyncio.TimeoutError: yield : keep-alive\n\n return StreamingResponse( event_source(), media_typetext/event-stream, headers{ Cache-Control: no-cache, X-Accel-Buffering: no, }, ) def format_sse(event_id: int, data: dict) - str: return ( fid: {event_id}\n fevent: token\n fdata: {json.dumps(data, ensure_asciiFalse)}\n\n ) async def run_generation(session_id: str): buf sessions[session_id] # 实际项目这里换成大模型流式接口比如 openai.AsyncOpenAI 的 stream() for i, token in enumerate([你, 好, , 这, 是, 流, 式, 输, 出]): await asyncio.sleep(0.2) buf.append({delta: token, seq: i}) buf.append({delta: , seq: len(buf.events)}, finishedTrue)这里有四个要点。第一sessions字典是生成器的共享状态。第一次连接时创建会话并启动后台生成任务后面所有重连都复用同一个会话而不是重新跑生成。第二Last-Event-ID的读取用的是小写last-event-idFastAPI/Starlette 会把请求头统一成小写这一点容易写错。第三asyncio.wait_for(buf.notify.wait(), timeout15)是心跳的关键。15 秒内没有任何新 token就发一个注释行保活有新 token 就立即唤醒发送。这样代理层永远不会看到超过 20 秒的空闲。第四X-Accel-Buffering: no是告诉 Nginx 这个响应不要缓冲配合 Nginx 端的proxy_buffering off才能保证真正的逐 token 推送。需要额外说明的是实际接入大模型时建议不要直接把上游 API 的每个 chunk 原样转发。上游的 chunk 边界不一定和 token 边界对齐最好在run_generation里先累积、按完整事件组装好再buf.append保证每个 SSE 帧里的data都是完整可渲染的内容。另外记得加一个后台清理任务超过 10 分钟没有新事件的会话直接删掉避免内存泄漏。5.2 前端实现可复用的 SSE 流式客户端先说一下为什么不用原生 EventSource没法带Authorization头而且在很多场景下你也不想让浏览器自动重连而是希望自己控制重试策略。所以我封装了一个基于 fetch 的StreamClient把解析、游标、重连、状态机都收进去。class StreamClient { constructor({ buildUrl, onToken, onDone, onError }) { this.buildUrl buildUrl; this.onToken onToken; this.onDone onDone; this.onError onError; this.processedId null; // 已处理的最大事件 ID this.retryCount 0; this.maxRetry 8; this.controller null; this.finished false; } async start() { this.finished false; this.retryCount 0; this.processedId null; await this._open(); } async _open() { if (this.finished) return; const headers { Accept: text/event-stream }; if (this.processedId ! null) { headers[Last-Event-ID] String(this.processedId); } this.controller new AbortController(); try { const res await fetch(this.buildUrl(), { headers, signal: this.controller.signal, }); if (!res.ok || !res.body) { throw new Error(HTTP ${res.status}); } const reader res.body.getReader(); const decoder new TextDecoder(utf-8); let buffer ; while (true) { const { value, done } await reader.read(); if (done) break; buffer decoder.decode(value, { stream: true }); // SSE 事件以空行分隔切出完整帧留下可能残缺的尾部 const frames buffer.split(\n\n); buffer frames.pop(); for (const frame of frames) { this._handleFrame(frame); } } this._finish(); } catch (err) { if (err.name AbortError) return; this._retry(err); } } _handleFrame(frame) { let id null; let event message; let data ; for (const line of frame.split(\n)) { if (line.startsWith(id:)) id line.slice(3).trim(); else if (line.startsWith(event:)) event line.slice(6).trim(); else if (line.startsWith(data:)) data line.slice(5).trim() \n; } data data.trim(); // 游标去重必须放在渲染之前 if (id ! null) { const numericId Number(id); if (this.processedId ! null numericId this.processedId) { return; } this.processedId numericId; } if (event done || data [DONE]) { this._finish(); return; } if (event token data) { const payload JSON.parse(data); this.onToken(payload.delta || ); } } _finish() { if (this.finished) return; this.finished true; this.onDone this.onDone(); } _retry(err) { if (this.finished) return; // 完成后绝不重连 if (this.retryCount this.maxRetry) { this.onError this.onError(err); return; } const delay Math.min(1000 * 2 ** this.retryCount, 15000); this.retryCount 1; setTimeout(() this._open(), delay); } stop() { this.finished true; this.controller this.controller.abort(); } }调用侧很简单以对话页面为例const answerEl document.getElementById(answer); const streamClient new StreamClient({ buildUrl: () /api/chat/stream/${sessionId}, onToken: (delta) { // 这里的 delta 已经是服务端去重后的新内容 answerEl.textContent delta; }, onDone: () { // 收起 loading、停止录音、触发后续动作 }, onError: (err) { console.error(stream failed after retries, err); }, }); streamClient.start();这个类在 Vue 里可以放进onMounted在 React 里可以放进useEffect不需要依赖任何框架特性。关键是它把流式消息解析与封装这件事收敛到了一个类里业务层只关心onToken和onDone。这里有个细节要强调重连时发送的Last-Event-ID用的是processedId而不是最后收到的值的副本。这两个值在正常流里相等但如果你在_handleFrame里先渲染再更新游标会出现渲染了但游标没跟上的窗口重连后这帧会被服务端重发前端却又因为游标跳过了它反而丢内容。所以代码里我坚持先更新processedId再渲染。5.3 代理与心跳参数配置后端代码写得再好代理层配置不对也白搭。我的 Nginx 配置如下location /api/chat/stream/ { proxy_pass http://backend_upstream; proxy_http_version 1.1; proxy_set_header Connection ; proxy_set_header Host $host; proxy_buffering off; proxy_cache off; proxy_read_timeout 3600s; gzip off; }几个关键点proxy_buffering off是必须的。开着的话Nginx 会把上游的 SSE 数据攒着可能攒几秒甚至直到超时才一起吐给客户端流式意义直接消失。proxy_read_timeout设为 3600 秒只是兜底真正防空闲超时靠的是后端 15 秒一次的心跳注释行。Connection 配合proxy_http_version 1.1是为了让 Nginx 到上游的连接可以复用避免频繁握手。gzip off是因为压缩会引入缓冲影响逐 token 推送。如果你确实想压缩至少要把text/event-stream类型排除掉。前端这边我额外加了一个看门狗定时器如果超过 40 秒没有任何数据帧到达就主动abort()当前请求并触发重连。因为有些代理掐断连接时不会发 RSTreader.read()会一直 pending前端根本感知不到已经断线。看门狗让客户端拥有了假死检测能力。注意这个时间要比后端心跳间隔大一般取心跳间隔的 2~3 倍避免误杀。6. 踩坑记录与排查速查表6.1 我实际踩过的几个坑第一个坑是前面提到的 EventSource 换 fetch 之后忘了重连逻辑。当时只想着自定义 header换完发现断网后页面完全不恢复比重复吐字更尴尬。排查半天才反应过来原生 EventSource 自动做的重连和Last-Event-ID在 fetch 方案里全部要自己写。所以后来封装StreamClient时第一件事就是把这两个能力补回来。第二个坑是服务端回放和实时发送发生了竞争。最初实现里事件缓冲的追加和 SSE 生成器的发送是各自独立的重连时如果一边在回放旧事件、一边后台又在追加新事件游标可能被新事件干扰导致同一帧被发送两次。最后改成cursor只在生成器内部维护回放和追加串行化问题才消失。这个实际上就是并发控制的问题建议在服务端用单线程的事件循环配合条件变量不要用多线程共享可变列表。第三个坑是字符去重把 emoji 切碎了。当时前端兜底用了slice判断重叠用户输入里有个 去重后文本后半截直接出现乱码。后来才意识到 UTF-16 代理对的问题要么在去重逻辑里用Array.from要么干脆以事件 ID 去重为主、字符去重为辅。现在我的代码里字符去重只处理纯中文和英文场景遇到代理对就跳过。第四个坑是停止按钮引发的假重连。用户的停止操作会触发AbortController.abort()而 abort 会抛出AbortError。如果错误处理没区分AbortError和其他网络错误停止操作会被当成断线进而触发重连用户看到的是我明明点了停止它又开始从头回答。这一点在_open的 catch 里已经处理if (err.name AbortError) return;。6.2 排查速查表现象可能原因优先排查手段重连后从头完整重吐服务端忽略Last-Event-ID后端日志打印请求头确认 header 是否到达重连后尾部重复几十字服务端回放了已发送事件检查 cursor 是否从last_event_id 1开始连接 60 秒左右准时断开Nginx 或 LB 空闲超时确认是否有心跳注释行间隔是否小于超时断网后前端无任何反应fetch 方案没有重连逻辑检查StreamClient是否实现_retry点击停止后又开始输出abort 被当作网络错误区分AbortError进入finished状态回答结束后莫名又触发一次[DONE]和断线竞态状态机里done后禁止重连流式效果变攒包式Nginx 缓冲未关闭检查proxy_buffering off和X-Accel-Buffering重连后内容有乱码字符级去重切坏代理对改用事件 ID 去重或Array.from处理排查时还有一个很实用的习惯在前端把每个事件的id打个日志在后端把每次连接收到的Last-Event-ID打个日志两边对齐看是客户端没发还是服务端没用。这个对比能瞬间定位责任方比瞎猜快得多。最后再说一个我个人的体会。SSE 不是 WebSocket不要拿 WebSocket 那套连接即会话的思维去写它重连也不是重新请求而是带着游标的续传。把服务端缓冲、前端游标、请求幂等这三层都做扎实之后我再也没收到过重复吐字的反馈。如果你正在被这个 bug 折磨按这个顺序从上到下检查一遍基本都能解决。
锦
锦皓数字建站
深耕本土企业品牌数字化升级,专注原创端正雅致商务官网,从视觉设计到稳定运维全程保驾护航。