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

如何用NestJS+Fastify实现类OpenAI的Post流式SSE API?

问题描述

我尝试使用NestJS和Fastify实现类似OpenAI的Create chat completion API的Server-Sent Events(SSE)流式API,希望客户端能像官方API一样通过设置responseType: "stream"的axios接收流式数据。以下是我尝试的简化代码:

controller.ts

@Post("/chatgpt/chat")
@UseGuards(AuthGuard)
@HttpCode(200)
@Header("Content-type", "text/event-stream")
async chatWithPostStream(@Body() body: ChatGPTChatDto, @Req() fastifyRequest: FastifyRequest, @Res() fastifyResponse: FastifyReply) {
  return await this.chatGPTService.chatPostStream(body);
}

service.ts

async chatWithPostStream(body: ChatGPTChatDto) {
  const openaiKey = "sk-xxxxxxxxxxxxxxxxxxxxxx";
  let chatGPTResponse;
  
  try {
    chatGPTResponse = await axios.post("https://api.openai.com/v1/chat/completions", {
      model: "gpt-3.5-turbo",
      messages: [{"role": "user", "content": "Hello!"}],//仅测试用,原应为`body.content`
      stream: true
    }, {
      headers: {
        "Content-Type": "application/json",
        "Authorization": `Bearer ${openaiKey}`
      },
      responseType: "stream"
    });
  } catch (e) {
    console.error("ChatGPT API Error:", e);
  }

  let outputText = "";

  chatGPTResponse.data.on("data", (chunk) => {
    //TODO: 将相同的chunk发送给客户端

    if (!chunk.toString("utf-8").includes("data: [DONE]")) {
      try {
        console.log("chunk text: ", chunk.toString("utf-8"));
        let choice = JSON.parse(chunk.toString("utf-8").substring(5)).choices[0];
        outputText += choice.delta.content ? choice.delta.content : "";
      } catch (e) {
        console.log("json parse error: ", e.message, "Original JSON:", chunk.toString("utf-8").substring(5));
      }
    } else {
      console.log("Final output text", outputText);
    }
  });

  chatGPTResponse.data.on("error", (e) => {
    console.error("ChatGPT Error", e);
  });

  chatGPTResponse.data.on("end", () => {
    console.log("ChatGPT End");
  });
}

过去几天我尝试了多种方式实现该功能,但均未成功。请问如何创建类似官方API的流式API,让Web客户端能以相同方式通过axios监听流式数据?


解决方案

1. 调整Controller逻辑

需要放弃NestJS默认的返回机制,直接操控Fastify的响应流,同时补充SSE必需的响应头,避免浏览器缓存或中断连接。

修改后的controller.ts:

@Post("/chatgpt/chat")
@UseGuards(AuthGuard)
@HttpCode(200)
async chatWithPostStream(
  @Body() body: ChatGPTChatDto,
  @Res({ passthrough: true }) fastifyResponse: FastifyReply
) {
  // 设置SSE核心响应头
  fastifyResponse.header('Content-Type', 'text/event-stream');
  fastifyResponse.header('Cache-Control', 'no-cache');
  fastifyResponse.header('Connection', 'keep-alive');
  fastifyResponse.header('Transfer-Encoding', 'chunked');

  // 调用Service处理流转发
  await this.chatGPTService.chatPostStream(body, fastifyResponse);
  
  // 无需返回值,直接通过响应流发送数据
}

2. 修改Service实现流转发

将Fastify响应对象传入Service,收到OpenAI的流式数据后,直接原样转发给客户端,同时处理错误和流结束事件。

修改后的service.ts:

async chatPostStream(body: ChatGPTChatDto, fastifyResponse: FastifyReply) {
  const openaiKey = "sk-xxxxxxxxxxxxxxxxxxxxxx";
  
  try {
    const chatGPTResponse = await axios.post(
      "https://api.openai.com/v1/chat/completions",
      {
        model: "gpt-3.5-turbo",
        messages: body.messages, // 使用客户端传入的实际对话数据
        stream: true
      },
      {
        headers: {
          "Content-Type": "application/json",
          "Authorization": `Bearer ${openaiKey}`
        },
        responseType: "stream"
      }
    );

    // 监听OpenAI数据流,直接转发给客户端
    chatGPTResponse.data.on("data", (chunk) => {
      const chunkStr = chunk.toString("utf-8");
      // 保持OpenAI原始SSE格式转发
      fastifyResponse.write(chunkStr);
      
      if (chunkStr.includes("data: [DONE]")) {
        console.log("流式响应结束");
      }
    });

    // 处理OpenAI API错误
    chatGPTResponse.data.on("error", (err) => {
      console.error("OpenAI API错误:", err);
      fastifyResponse.write(`data: ${JSON.stringify({ error: err.message })}\n\n`);
      fastifyResponse.end();
    });

    // 数据流结束时关闭客户端响应
    chatGPTResponse.data.on("end", () => {
      console.log("OpenAI响应流结束");
      fastifyResponse.end();
    });

  } catch (e) {
    console.error("请求OpenAI失败:", e);
    fastifyResponse.status(500).send(`data: ${JSON.stringify({ error: e.message })}\n\n`);
  }
}

3. 客户端axios调用示例

客户端代码与调用OpenAI官方API逻辑完全兼容,只需设置responseType: "stream"并监听数据流:

import axios from 'axios';

async function callStreamAPI() {
  const response = await axios.post('/chatgpt/chat', {
    messages: [{ role: 'user', content: 'Hello!' }]
  }, {
    responseType: 'stream',
    headers: {
      'Authorization': 'Bearer YOUR_AUTH_TOKEN' // 替换为实际认证令牌
    }
  });

  response.data.on('data', (chunk) => {
    const chunkStr = chunk.toString('utf-8');
    // 解析SSE格式数据
    const lines = chunkStr.split('\n').filter(line => line.trim() !== '');
    for (const line of lines) {
      if (line.startsWith('data: ')) {
        const data = line.substring(6);
        if (data === '[DONE]') {
          console.log('流式响应结束');
          return;
        }
        try {
          const result = JSON.parse(data);
          console.log('收到内容:', result.choices[0].delta.content);
        } catch (e) {
          console.log('解析数据失败:', e);
        }
      }
    }
  });

  response.data.on('error', (err) => {
    console.error('请求错误:', err);
  });

  response.data.on('end', () => {
    console.log('响应结束');
  });
}

callStreamAPI();

关键注意事项

  • 必须设置Transfer-Encoding: chunked和Cache-Control: no-cache,确保流式数据实时传输且不被缓存。
  • 不要在Controller中返回任何值,所有数据通过fastifyResponse.write()直接发送。
  • 严格保持OpenAI的原始SSE格式转发,不要修改数据结构,保证客户端兼容性。
  • 错误场景下要发送符合SSE格式的错误信息,并调用end()关闭响应流。

内容的提问来源于stack exchange,提问作者hash070

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 19:09:59