流式传输断网重连与消息幂等去重
在大模型对话与实时数据看板落地过程中,Fetch EventSource 或 WebSocket 的长连接通信几乎是标配。很多刚接触流式传输的团队往往将注意力放在流式打字机的动效渲染上,直到弱网环境下用户频频抱怨“回答吞字”、“网络闪断后内容全丢”或者“重连后同一段话打印了两次”,才意识到流式管道底层的健壮性缺口。
流式响应本质上是一条连续的字节流序列。与传统单次 RPC 请求不同,流式传输在网络震荡或代理服务超时切断时,客户端拿到的只是一半的状态片段。如果简单粗暴地断线重连重新请求,服务端会从头生成,不仅浪费昂贵的 Token 开销,更让客户端视图陷入混乱。要构建一个工业级高可用的流式客户端,核心在于两件事:游标断点续传(Cursor Resumption)与消息粒度的幂等去重(Idempotent Deduplication)。
游标驱动的状态切片
标准的 SSE 协议原生支持Last-Event-ID标头,但在自定义的 POST 流式接口(例如携带复杂 prompt 或上下文结构的大模型会话)中,原生 EventSource 并不支持自定义 Header 和 Body,通常采用 Fetch API 配合ReadableStream手动消费。
这就要求前后端共同维护一套切片序列协议。服务端在吐出每个 Chunk 时,必须附带单调递增的chunkId或sequence,以及全局唯一的messageId。客户端维护一个本地滑动窗口,记录当前已成功渲染并落入持久化状态的最大sequence。
interface StreamChunkPayload { messageId: string; sequence: number; delta: string; status: 'streaming' | 'completed' | 'error'; } class ResilientStreamClient { private abortController: AbortController | null = null; private currentMessageId: string | null = null; private lastAcknowledgedSeq: number = -1; private retryAttempts: number = 0; private maxRetries: number = 5; private baseDelayMs: number = 1000; private receivedChunkIds: Set<string> = new Set(); constructor( private url: string, private payload: Record<string, any>, private onChunk: (delta: string, fullText: string) => void, private onStatusChange: (status: string) => void ) {} public async connect(): Promise<void> { this.abortController = new AbortController(); try { this.onStatusChange('connecting'); const response = await fetch(this.url, { method: 'POST', headers: { 'Content-Type': 'application/json', 'X-Last-Ack-Seq': String(this.lastAcknowledgedSeq), 'X-Message-Id': this.currentMessageId || '' }, body: JSON.stringify({ ...this.payload, resumeSeq: this.lastAcknowledgedSeq }), signal: this.abortController.signal }); if (!response.ok || !response.body) { throw new Error(`HTTP error! status: ${response.status}`); } this.retryAttempts = 0; this.onStatusChange('streaming'); await this.pumpStream(response.body.getReader()); } catch (err: any) { if (err.name === 'AbortError') { this.onStatusChange('aborted'); return; } this.handleReconnect(); } } private async pumpStream(reader: ReadableStreamDefaultReader<Uint8Array>): Promise<void> { const decoder = new TextDecoder('utf-8'); let buffer = ''; while (true) { const { done, value } = await reader.read(); if (done) break; buffer += decoder.decode(value, { stream: true }); const lines = buffer.split('\n\n'); buffer = lines.pop() || ''; for (const block of lines) { if (!block.trim()) continue; this.processEventBlock(block); } } } private processEventBlock(block: string): void { const lines = block.split('\n'); let dataStr = ''; for (const line of lines) { if (line.startsWith('data: ')) { dataStr = line.slice(6); } } if (!dataStr) return; try { const chunk: StreamChunkPayload = JSON.parse(dataStr); // 幂等去重防线:避免重连抖动下服务端重发已消费的消息段 const dedupeKey = `${chunk.messageId}:${chunk.sequence}`; if (this.receivedChunkIds.has(dedupeKey)) { return; } // 序列号对齐校验:若发现跳序则触发局部重同步 if (chunk.sequence <= this.lastAcknowledgedSeq) { return; } this.receivedChunkIds.add(dedupeKey); this.currentMessageId = chunk.messageId; this.lastAcknowledgedSeq = chunk.sequence; this.onChunk(chunk.delta, ''); if (chunk.status === 'completed') { this.onStatusChange('completed'); this.receivedChunkIds.clear(); } } catch (e) { console.error('Parse chunk error', e); } } private handleReconnect(): void { if (this.retryAttempts >= this.maxRetries) { this.onStatusChange('failed'); return; } this.onStatusChange('reconnecting'); this.retryAttempts++; // 指数退避与随机抖动,避免瞬时网络恢复时的雪崩请求 const jitter = Math.random() * 200; const delay = Math.min(this.baseDelayMs * Math.pow(2, this.retryAttempts) + jitter, 10000); setTimeout(() => { this.connect(); }, delay); } public abort(): void { if (this.abortController) { this.abortController.abort(); } } }重试风暴与退避策略
网络抖动往往不是单点现象,当网关层或者机房边缘节点发生秒级重启时,数以万计的在线客户端会同时触发重连逻辑。如果在断网捕获中直接调用重连函数,庞大的并发峰值会瞬间压垮服务端。
指数退避算法(Exponential Backoff)配合随机抖动(Jitter)是解决这一问题的经典工程实践。在代码实现中,退避时间随尝试次数呈 $2^n$ 增长,并附加一定的随机因子,使得不同终端的重发时间点均匀分布在时间轴上,为服务端的恢复留下缓冲空间。
前端状态一致性保护
断网重连不仅仅是网络层的重建,前端 UI 层的状态机同样需要严密防护。当流式中断时,UI 应该保持当前已有文本的展示,而不是闪烁清空;当重连握手成功后,新收到的首个 Chunk 必须紧接在已有字符的末尾追加,而非覆盖。
在生产环境中,还要防范由于网络延迟造成的“旧请求晚到”问题。使用 Fetch 的AbortController,在发起任何新的握手或放弃旧会话时立刻中断底层连接,确保同一个会话内只存在唯一活跃的 Reader 管道。通过游标序列号匹配、内存 Set 去重、退避重连与中止信号四大防线,流式交互才能在极其恶劣的弱网环境中保持如同行云流水般的稳定与从容。