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

Node.js 与 ChatGPT API 的流式响应实践:从阻塞轮询到 Server-Sent Events 的架构演进

一、两种接口范式的本质差异

ChatGPT API 提供了 stream: falsestream: true 两种模式。前者在一次 HTTP 往返中返回完整响应,后者通过 SSE(Server-Sent Events)逐 token 推送。这不仅仅是"快慢"的区别——它决定了整个服务端架构的内存模型和并发能力。

❌ 反例:阻塞式调用

// 收到用户消息 → 转发给 OpenAI → 等全量返回 → 再响应客户端
app.post('/api/chat', async (req, res) => {
  const completion = await openai.chat.completions.create({
    model: 'gpt-4o',
    messages: req.body.messages,
    stream: false,
  });

  res.json({ reply: completion.choices[0].message.content });
});

问题清单:

  • 用户盯着空白页面等 8-15 秒,体验极差
  • 每个请求占用一个 Node 事件循环槽位,1k 并发直接打满
  • 完整响应在内存中组装,长文本场景下 GC 压力陡增
  • 无法实现"打字机效果",产品竞争力掉一档

✅ 正例:全链路流式管道

app.post('/api/chat', async (req, res) => {
  res.setHeader('Content-Type', 'text/event-stream');
  res.setHeader('Cache-Control', 'no-cache');
  res.setHeader('Connection', 'keep-alive');
  res.flushHeaders(); // 立刻发送 HTTP 头,不要等 body

  const stream = await openai.chat.completions.create({
    model: 'gpt-4o',
    messages: req.body.messages,
    stream: true,
  });

  for await (const chunk of stream) {
    const delta = chunk.choices[0]?.delta?.content;
    if (delta) {
      res.write(`data: ${JSON.stringify({ content: delta })}\n\n`);
    }
  }

  res.write('data: [DONE]\n\n');
  res.end();
});

二、ReadableStream 不是 Response.body 的平替

很多开发者习惯用 fetchresponse.body 直接泵送。这在原型阶段可行,但生产环境中遇到背压(backpressure)就会丢数据。

❌ 直接用 response.body 泵送

const response = await fetch('https://api.openai.com/v1/chat/completions', {
  method: 'POST',
  headers: { 'Content-Type': 'application/json', Authorization: `Bearer ${key}` },
  body: JSON.stringify({ model: 'gpt-4o', messages, stream: true }),
});

const reader = response.body!.getReader();
const decoder = new TextDecoder();

while (true) {
  const { done, value } = await reader.read();
  if (done) break;
  res.write(decoder.decode(value, { stream: true }));
  // ⚠️ 没有检查 res.write() 的返回值
}

res.write() 返回 false 时表示内部缓冲区已满,下游消费速度跟不上。继续写入会导致内存无限膨胀,最终 OOM。

✅ Node.js 原生 Readable 配合背压控制

import { Readable } from 'node:stream';

function createOpenAIReadable(messages: Message[]): Readable {
  return new Readable({
    async read(this: Readable) {
      if (this._reading) return;
      this._reading = true;

      try {
        const stream = await openai.chat.completions.create({
          model: 'gpt-4o',
          messages,
          stream: true,
        });

        for await (const chunk of stream) {
          const delta = chunk.choices[0]?.delta?.content;
          if (!delta) continue;

          // push 返回 false → 缓冲区满,暂停上游消费
          if (!this.push(`data: ${JSON.stringify({ content: delta })}\n\n`)) {
            break;
          }
        }

        this.push('data: [DONE]\n\n');
      } finally {
        this._reading = false;
        this.push(null); // 结束流
      }
    },
  });
}

关键点:push() 返回 false 时,Node.js 流内部缓冲区达到 highWaterMark(默认 16KB),自动暂停,等下游排空后才会再次调用 _read()。这套机制是 libuv 层级的,不需要手动写 drain 事件。

三、生产级管道:OpenAI → Transform → SSE

真实场景中,你不会只做透传。中间需要一个 Transform 流做 token 计数、内容审核、自定义格式转换。

import { Transform } from 'node:stream';

class TokenCounterTransform extends Transform {
  private count = 0;
  private readonly maxTokens: number;

  constructor(options: { maxTokens: number }) {
    super({ objectMode: true });
    this.maxTokens = options.maxTokens;
  }

  _transform(
    chunk: { content: string },
    _encoding: string,
    callback: (err?: Error | null, data?: any) => void
  ) {
    this.count += chunk.content.length;

    if (this.count > this.maxTokens) {
      callback(new Error('Token overflow'));
      return;
    }

    callback(null, `data: ${JSON.stringify(chunk)}\n\n`);
  }
}

三流串联:

app.post('/api/chat', async (req, res) => {
  res.setHeader('Content-Type', 'text/event-stream');
  res.flushHeaders();

  const source = createOpenAIReadable(req.body.messages);
  const counter = new TokenCounterTransform({ maxTokens: 4096 });

  source.pipe(counter).pipe(res);

  // 连接断开时销毁上游,避免 token 浪费
  req.on('close', () => {
    source.destroy();
  });
});

四、一条经验法则

方案适用场景关键坑点
stream: false非实时批处理、短应答并发能力差,用户体验差
stream: true + 直接 res.write原型验证无背压控制,OOM 风险
stream: true + Readable + pipe生产环境需处理客户端断开、token 上限、异常恢复

五、客户端断开是你的付费陷阱

用户关掉浏览器标签页的那一刻,你的 Node 进程可能还在继续从 OpenAI 拉 token——每一 byte 都在计费。

req.on('close', () => {
  openaiStream.controller?.abort(); // 主动取消上游
  source.destroy();
});

req.on('close') 而不是 'end':客户端中断连接触发 close,正常结束触发 end,前者才是我们需要 abort 的信号。

六、总结

对接 ChatGPT API 的流式响应不是"加个 stream: true 就行"。完整的生产方案包含:

  1. 全链路流式传输,从 API 到客户端一个 pipe 走到底
  2. Native Readable 而非裸 fetch body,利用 Node.js 内置背压机制
  3. Transform 流做中间处理,保持管道模式而非回调地狱
  4. 客户端断开立即 abort,避免无效的 API 开销

流是 Node.js 的基因,不要用写阻塞代码的方式写 ChatGPT 代理。

0 评论

评论区

登录 后参与评论