Node.js 与 ChatGPT API 的流式响应实践:从阻塞轮询到 Server-Sent Events 的架构演进
一、两种接口范式的本质差异
ChatGPT API 提供了 stream: false 和 stream: 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 的平替
很多开发者习惯用 fetch 的 response.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 就行"。完整的生产方案包含:
- 全链路流式传输,从 API 到客户端一个 pipe 走到底
- Native Readable 而非裸
fetchbody,利用 Node.js 内置背压机制 - Transform 流做中间处理,保持管道模式而非回调地狱
- 客户端断开立即 abort,避免无效的 API 开销
流是 Node.js 的基因,不要用写阻塞代码的方式写 ChatGPT 代理。
0 评论
评论区
登录 后参与评论