You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.17 12:37:03