前端开发··2 阅读·预计 16 分钟

GPT API 流式响应的工程化治理:背压控制、断线重连与消费端的 5 层防线

问题起点:一个普遍存在的「能跑就行」实现

// ❌ 最常见的写法——跑起来没问题,上线就炸
async function naiveStream(prompt: string) {
  const res = await fetch('/api/chat', {
    method: 'POST',
    body: JSON.stringify({ prompt }),
    headers: { 'Content-Type': 'application/json' }
  });

  const reader = res.body!.getReader();
  const decoder = new TextDecoder();
  let fullText = '';

  while (true) {
    const { done, value } = await reader.read();
    if (done) break;
    const chunk = decoder.decode(value, { stream: true });
    // 解析 SSE
    for (const line of chunk.split('\n')) {
      if (line.startsWith('data: ')) {
        const payload = JSON.parse(line.slice(6));
        const delta = payload.choices[0]?.delta?.content ?? '';
        fullText += delta;
        setMessages(prev => [...prev, { role: 'assistant', content: fullText }]);
      }
    }
  }
}

这段代码至少有 四个致命缺陷,我们逐层修补。


第一层:SSE 解析器——别用 split 处理流式协议

chunk.split('\n') 的最大问题:TCP 分包可能在任意字节位置截断。如果你收到的 chunk 是 "data: {"hello",下一帧才是 "}\n"JSON.parse 直接抛异常。

// ✅ 带缓冲区的增量 SSE 解析器
class SSEParser {
  private buffer = '';

  feed(chunk: string): SSEEvent[] {
    this.buffer += chunk;
    const events: SSEEvent[] = [];
    // 只在完整行边界切割
    const lines = this.buffer.split('\n');
    // 最后一行可能不完整,保留在 buffer
    this.buffer = lines.pop() ?? '';

    for (const line of lines) {
      if (line.startsWith('data: ')) {
        const raw = line.slice(6).trim();
        if (raw === '[DONE]') { events.push({ done: true }); continue; }
        try {
          events.push({ done: false, data: JSON.parse(raw) });
        } catch { /* 忽略畸形行,生产环境应上报 */ }
      }
    }
    return events;
  }
}

关键点:lines.pop() 把不完整的尾部还给缓冲区,等下一帧补齐后再处理。


第二层:背压控制——告诉下游「我吃不下了」

body.getReader().read() 是无背压的读取。如果网络极快而消费端渲染极慢(比如每次 setState 触发 React 全量 reconcile),未消费的 chunk 会堆积在内存中。

ReadableStream 内置的背压机制通过 reader.read() 的 Promise 未 resolve 来反向施压。但如果你在 read 之后立即进行耗时操作,背压就形同虚设。

// ❌ 背压失效:先读全了再慢慢处理,读取侧无任何等待
while (true) {
  const { done, value } = await reader.read(); // 瞬间读出
  if (done) break;
  hugeBuffer.push(value); // 堆在内存里
}
// 然后才慢慢渲染...

// ✅ 正确的背压消费:处理完一批才继续读下一批
async function* streamWithBackpressure(
  reader: ReadableStreamDefaultReader<Uint8Array>
): AsyncGenerator<string> {
  const decoder = new TextDecoder();

  while (true) {
    const { done, value } = await reader.read();
    if (done) break;
    // 在 yield 处暂缓读取:上游 TCP 窗口自动缩小
    yield decoder.decode(value, { stream: true });
  }
}

配合消费端的 批量 setState + requestAnimationFrame 节流

let pending = '';
let rafId = 0;

for await (const chunk of streamWithoutBackpressure(reader)) {
  pending += chunk;
  cancelAnimationFrame(rafId);
  rafId = requestAnimationFrame(() => {
    setMessages(prev => updateLastMessage(prev, pending));
    pending = '';
  });
}
// 兜底:最后一帧可能没有触发 rAF
if (pending) setMessages(prev => updateLastMessage(prev, pending));

rAF 将渲染频率自然限制在 60fps,大幅降低 React reconcile 次数。


第三层:AbortController + 竞态治理

用户在回复还没结束时快速切换话题,旧请求仍在跑。

// ✅ useRef 持有 AbortController,新请求到来时先 abort 旧请求
function useChatStream() {
  const abortRef = useRef<AbortController | null>(null);

  const send = useCallback(async (prompt: string) => {
    // 1. 杀死上一个还在跑的任务
    abortRef.current?.abort();

    const controller = new AbortController();
    abortRef.current = controller;

    try {
      const res = await fetch('/api/chat', {
        method: 'POST',
        body: JSON.stringify({ prompt }),
        signal: controller.signal, // 2. 接入 AbortSignal
        headers: { 'Content-Type': 'application/json' }
      });
      // ... 流式消费
    } catch (err) {
      if (err instanceof DOMException && err.name === 'AbortError') {
        return; // 3. 预期内的取消,不做任何处理
      }
      throw err;
    }
  }, []);

  // 组件卸载时清理
  useEffect(() => () => abortRef.current?.abort(), []);

  return { send };
}

第四层:断线重连——利用 lastCommitId 实现增量续传

流式传输中断(网络抖动、代理超时)后,重新请求整个 prompt 是巨大的 token 浪费。更合理的方案:把已收到的内容作为上下文回传,要求 API 从断点续写

// ✅ 带续传能力的流式请求
async function* resilientStream(prompt: string, previousText = '') {
  const body = previousText
    ? {
        prompt,
        // 关键:把已收到内容作为 assistant 消息追加到 messages 末尾
        continueFrom: previousText
      }
    : { prompt };

  const res = await fetch('/api/chat', {
    method: 'POST',
    body: JSON.stringify(body),
    headers: { 'Content-Type': 'application/json' }
  });
  // ... 解析并 yield
}

// 外层用指数退避重试
async function sendWithRetry(prompt: string) {
  let accumulated = '';

  for (let attempt = 0; attempt < 5; attempt++) {
    try {
      for await (const delta of resilientStream(prompt, accumulated)) {
        accumulated += delta;
        yield delta;
      }
      return; // 正常结束
    } catch (err) {
      if (attempt === 4) throw err;
      // 指数退避:1s, 2s, 4s, 8s
      await sleep(2 ** attempt * 1000);
    }
  }
}

注意:续传对服务端有要求——后端必须支持在 prompt 末尾拼接已收到的 assistant content,且续传请求返回的 delta 不包含已返回部分。如果你的后端不支持,至少要在 retry 时用 accumulated 保证 UI 不闪烁。


第五层:Unicode 边界——emoji 和中文不能切碎

TextDecoder{ stream: true } 选项已经处理了多字节 UTF-8 字符的跨 frame 问题。但还有一个隐蔽的坑:代理对(surrogate pairs)

某些 emoji 由两个 UTF-16 code unit 组成(如 👨‍👩‍👧 实际是 5 个 code unit)。如果恰好在代理对中间切割,decoder.decode(value) 会产出 \uFFFD(� 替换字符)。

// ✅ 安全的增量解码器
class SafeDecoder {
  private incomplete = new Uint8Array(0);
  private textDecoder = new TextDecoder();

  decode(chunk: Uint8Array): string {
    // 拼接上一轮未完成的字节
    const merged = new Uint8Array(this.incomplete.length + chunk.length);
    merged.set(this.incomplete);
    merged.set(chunk, this.incomplete.length);

    const text = this.textDecoder.decode(merged, { stream: true });

    // 末尾可能是一个不完整的代理对——保留最后 3 个字节做缓冲
    // 因为一个完整 Unicode 码点最多 4 字节,但代理对问题实际出在 2 字节的 UTF-16 code unit 上
    const keep = Math.min(3, chunk.length);
    this.incomplete = chunk.slice(-keep);

    return text.slice(0, -1 * keep || undefined); // 去掉尾部不完整部分
  }

  flush(): string {
    return this.textDecoder.decode(this.incomplete);
  }
}

实际上 { stream: true } 本身就能处理绝大部分场景。这个 SafeDecoder 的价值在于给你一个明确的「这部分还不完整,先别渲染」的信号,避免 UI 上出现闪烁的 �。


最终组装:五层防线协同工作

async function* productionStream(prompt: string): AsyncGenerator<string> {
  const controller = new AbortController();
  // 注册到全局以支持外部 abort

  for (let retry = 0; retry < 3; retry++) {
    try {
      const res = await fetch('/api/chat', {
        method: 'POST',
        body: JSON.stringify({ prompt, continueFrom: accumulated }),
        signal: controller.signal,
        headers: { 'Content-Type': 'application/json' }
      });

      const reader = res.body!.getReader();
      const parser = new SSEParser();
      const decoder = new TextDecoder();

      while (true) {
        const { done, value } = await reader.read();
        if (done) break;

        const chunk = decoder.decode(value, { stream: true });
        for (const event of parser.feed(chunk)) {
          if (event.done) return;
          const delta = event.data.choices[0]?.delta?.content ?? '';
          accumulated += delta;
          yield delta; // 背压点:消费者处理完才继续读
        }
      }
      return; // 流正常结束
    } catch (err) {
      if (retry === 2) throw err;
      await sleep(2 ** retry * 1000);
    }
  }
}

总结

防线解决的问题核心手段
SSE 解析器TCP 分包截断 JSON行缓冲区 + lines.pop()
背压控制内存溢出 / 频繁渲染async generator + rAF 节流
竞态治理快速切换导致状态错乱AbortController + useRef
断线重连网络抖动重复消费 token指数退避 + lastCommitId 续传
Unicode 安全emoji / 中文出现乱码TextDecoder stream 模式

这五层不是「最佳实践」的八股堆砌——每一条都是线上暴露过问题的真实防线。删掉任何一层,你最终都会在凌晨两点收到 PagerDuty 的报警。

0 评论

评论区

登录 后参与评论