Node.js中ChatGPT流式API返回ReadableStream客户端接收异常排查
问题描述
我是Node.js新手,正在开发一个基于ChatGPT功能的程序,可根据主题生成笑话,使用的是https://api.openai.com/v1/chat/completions的流式版本。目前能看到服务器端返回的流包含多个数据块,但客户端无法正确接收:客户端的console.log({done, value});仅触发两次,调试发现服务器端的流有更多数据块,且客户端解码后的值为{}。请问在服务器端需要补充哪些配置才能正确实现流传输?
OpenAPI 工具函数
import { createParser, ParsedEvent, ReconnectInterval } from "eventsource-parser"; export const config = { runtime: "edge", }; export async function OpenAIStream(payload) { const encoder = new TextEncoder(); const decoder = new TextDecoder(); let counter = 0; const res = await fetch("https://api.openai.com/v1/chat/completions", { headers: { "Content-Type": "application/json", Authorization: `Bearer ${process.env.OPENAI_API_KEY}`, }, method: "POST", body: JSON.stringify(payload), }); const stream = new ReadableStream({ async start(controller) { function onParse(event: ParsedEvent | ReconnectInterval) { if (event.type === "event") { const data = event.data; if (data === "[DONE]") { controller.close(); return; } try { const json = JSON.parse(data); const text = json.choices[0].delta?.content || ""; if (counter < 2 && (text.match(/\n/) || []).length) { return; } console.log(text); const queue = encoder.encode(text); controller.enqueue(queue); counter++; } catch (e) { controller.error(e); } } } // stream response (SSE) from OpenAI may be fragmented into multiple chunks // this ensures we properly read chunks & invoke an event for each SSE event stream const parser = createParser(onParse); // https://web.dev/streams/#asynchronous-iteration for await (const chunk of res.body as any) { parser.feed(decoder.decode(chunk)); } }, }); return stream; }
Nest 控制器
import { Body, Controller, Post } from '@nestjs/common'; import { AppService } from './app.service'; import { OpenAIStream } from './helpers/openai'; import { ChatCompletionRequestMessage } from 'openai'; const MAX_RESPONSE_TOKENS = 200;//1024; @Controller() export class AppController { constructor(private readonly appService: AppService) { } @Post("joke") async generate(@Body() message: JokeTemplate) { let messages: Array<ChatCompletionRequestMessage> = [ { "role": "system", "content": "You are a joke engine." }, { "role": "user", "content": `Tell me a joke about ${message.subject}` }] const payload = { model: 'gpt-3.5-turbo', max_tokens: MAX_RESPONSE_TOKENS, temperature: 0, messages, stream: true }; const stream = await OpenAIStream(payload); return new Response(stream); } } interface JokeTemplate { subject: string; }
客户端请求触发代码
const triggerGPTRequest = async (e: any) => { setGptResponse(''); setLoading(true); const response = await fetch("/api/joke", { method: "POST", headers: { "Content-Type": "application/json", }, body: JSON.stringify({ 'subject': promptText }), }); if (!response.ok) { throw new Error(response.statusText); } const data = response.body; if (!data) { return; } const reader = data.getReader(); const decoder = new TextDecoder(); let done = false; while (!done) { const {value, done: doneReading} = await reader.read(); done = doneReading; const chunkValue = decoder.decode(value); console.log({done, value}); setGptResponse((prev) => prev + chunkValue); } setLoading(false); }
解决方案
1. 修改Nest控制器的响应处理逻辑
Nest默认响应机制会缓冲或合并流式数据,需要手动控制响应对象,配置必要头信息并将流管道到响应:
import { Body, Controller, Post, Res } from '@nestjs/common'; import { Response } from 'express'; // 保留其他原有导入 @Controller() export class AppController { constructor(private readonly appService: AppService) { } @Post("joke") async generate(@Body() message: JokeTemplate, @Res() res: Response) { let messages: Array<ChatCompletionRequestMessage> = [ { "role": "system", "content": "You are a joke engine." }, { "role": "user", "content": `Tell me a joke about ${message.subject}` }] const payload = { model: 'gpt-3.5-turbo', max_tokens: MAX_RESPONSE_TOKENS, temperature: 0, messages, stream: true }; const stream = await OpenAIStream(payload); // 配置流式响应头 res.setHeader('Content-Type', 'text/plain; charset=utf-8'); res.setHeader('Transfer-Encoding', 'chunked'); res.setHeader('Cache-Control', 'no-cache'); res.setHeader('Connection', 'keep-alive'); // 禁用压缩,避免流被合并 res.setHeader('Content-Encoding', 'identity'); // 将流管道到Express响应 (stream as any).pipe(res); } }
2. 移除OpenAIStream中的不必要过滤
当前代码中存在过滤逻辑,会跳过前两个含换行符的文本块,可能导致数据丢失,直接注释或删除:
// 注释或删除这段代码 // if (counter < 2 && (text.match(/\n/) || []).length) { // return; // }
3. 优化客户端解码方式
客户端解码时添加{ stream: true }参数,确保流式数据正确处理:
const chunkValue = decoder.decode(value, { stream: true });
内容的提问来源于stack exchange,提问作者4imble
相关产品推荐
相关产品推荐

