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

Node.js 长时流式响应的工程化治理:从超时防线到背压控制的 4 层守卫

Node.js 长时流式响应的工程化治理:从超时防线到背压控制的 4 层守卫

ChatGPT 的流式响应在本地跑 Demo 只需 30 行代码。但部署到生产环境——经过反向代理、负载均衡、移动弱网——超时断开、SSE 卡死、内存暴涨三个问题轮番轰炸。本文梳理 4 层工程化防线,并给出可落地的代码方案。

第 1 层:分层超时 —— 读、写、空闲三把锁

最常见的翻车是设置了单一超时。但实际上一个流式连接有三种截然不同的超时语义:

// ❌ 坏:单一超时糊弄所有人
const response = await fetch('https://api.openai.com/v1/chat/completions', {
  signal: AbortSignal.timeout(30000) // 流式输出 29 秒后 → 直接 500
});

// ✅ 好:HTTP 读超时 + 流读取空闲超时 + 总时长上限 三把锁
async function createStreamWithTimeout(body, options = {}) {
  const {
    connectTimeout = 10_000,   // 握手超时
    readTimeout = 5_000,       // chunk 间最大静默
    maxDuration = 120_000,     // 全流程硬上限
  } = options;

  const controller = new AbortController();
  const connectTimer = setTimeout(() => controller.abort(), connectTimeout);

  const response = await fetch('https://api.openai.com/v1/chat/completions', {
    method: 'POST',
    headers: { 'Content-Type': 'application/json', Authorization: `Bearer ${apiKey}` },
    body: JSON.stringify(body),
    signal: controller.signal,
  });
  clearTimeout(connectTimer);

  const maxTimer = setTimeout(() => controller.abort(), maxDuration);
  const reader = response.body.getReader();

  return (async function* readChunks() {
    let idleTimer = null;
    const resetIdle = () => {
      if (idleTimer) clearTimeout(idleTimer);
      idleTimer = setTimeout(() => controller.abort(), readTimeout);
    };
    resetIdle();

    try {
      while (true) {
        const { done, value } = await reader.read();
        if (done) break;
        resetIdle();
        yield value;
      }
    } finally {
      clearTimeout(maxTimer);
      if (idleTimer) clearTimeout(idleTimer);
      reader.releaseLock();
    }
  })();
}

AbortSignal.timeout(totalMs) 无法区分"握手慢"还是"流暂停"。经验数据:OpenAI 在生成长文本时 chunk 间隔通常 < 3s。超过 5s 无数据可视为上游异常。

第 2 层:连接断线检测 —— req.on('close') 而非猜

流式中最隐蔽的 bug:客户端断开了,Node.js 这边的写入仍在继续,直到后续 write()ERR_STREAM_WRITE_AFTER_END

// ❌ 缺检测:客户端断开后继续喂给 OpenAI,浪费 token
app.post('/chat', async (req, res) => {
  const stream = await openai.chat.completions.create({ stream: true, ... });
  for await (const chunk of stream) {
    res.write(`data: ${JSON.stringify(chunk)}\n\n`);
  }
  res.end();
});

// ✅ 客户端断开立即取消上游
app.post('/chat', async (req, res) => {
  const abortController = new AbortController();

  req.on('close', () => {
    abortController.abort(); // 通知 OpenAI 中止
  });

  const stream = await openai.chat.completions.create({
    stream: true,
    ...body,
  }, { signal: abortController.signal });

  for await (const chunk of stream) {
    if (abortController.signal.aborted) break;
    const ok = res.write(`data: ${JSON.stringify(chunk)}\n\n`);
    if (!ok) {
      // res 写缓冲区满 → 等待 drain
      await new Promise(resolve => res.once('drain', resolve));
    }
  }
  res.end();
});

关键点:不仅要监听 close 取消 OpenAI 请求,还要检查 res.write() 的返回值。返回值 false 意味着内部缓冲区已满——硬写会导致内存暴涨。

第 3 层:背压控制 —— 不是消费快,是控制生产慢

res.write() 返回 false 且 Node.js 内部缓冲区累积时,for await 循环仍从 OpenAI 高速读取数据填入内存。

// ❌ 忽略背压 → 1 分钟流式堆内存 2GB+
for await (const chunk of stream) {
  res.write(`data: ${JSON.stringify(chunk)}\n\n`); // 无视 Boolean 返回值
}

// ✅ 感知背压,暂停上游读取
async function* drainableStream(readable, writable) {
  for await (const chunk of readable) {
    if (!writable.write(chunk)) {
      // 内部 buffer 满,等待下游排空
      await new Promise(resolve => writable.once('drain', resolve));
    }
  }
}

// 用法
app.post('/chat', async (req, res) => {
  res.writeHead(200, { 'Content-Type': 'text/event-stream' });
  const stream = await openai.chat.completions.create({ stream: true, ... });
  yield* drainableStream(stream, res);
  res.end();
});

第 4 层:Reader 生命周期 —— releaseLockcancel

response.body.getReader() 获取的 ReadableStreamDefaultReader 被锁定后,其他消费者无法访问。中断时忘记 releaseLock() 会导致潜在的内存泄漏。

// ✅ 无论成功/失败/中断,始终释放锁
async function readSSEStream(response, onChunk, signal) {
  const reader = response.body.getReader();
  const decoder = new TextDecoder();
  let buffer = '';

  signal.addEventListener('abort', () => {
    reader.cancel().catch(() => {}); // 取消底层流
  }, { once: true });

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

      buffer += decoder.decode(value, { stream: true });
      const lines = buffer.split('\n');
      buffer = lines.pop() || ''; // 保留未完成行

      for (const line of lines) {
        if (line.startsWith('data: ')) {
          const data = line.slice(6);
          if (data === '[DONE]') return;
          onChunk(JSON.parse(data));
        }
      }
    }
  } finally {
    reader.releaseLock();
  }
}

reader.cancel() 的三层作用:

  • 告诉上游"我不读了,TCP 可以 RST"
  • 阻止后续 read() 返回新数据(抛 AbortError
  • 允许 GC 回收底层缓冲区

总结

防线手段不加的后果
分层超时三把锁分别限制握手/静默/总时长弱网抖动误杀正常连接
连接断线req.on('close') + AbortController客户端走了,OpenAI 还在烧 token
背压控制检查 write() 返回值 + drain 等待1 分钟跑出 OOM
Reader 释放finally { releaseLock() } + cancel()内存泄漏 + 连接句柄泄漏

四层防线写完你会发现:最大的工程难点不是调用 API,而是你永远预测不到用户的网络状况。

0 评论

评论区

登录 后参与评论