如何用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
相关产品推荐
相关产品推荐

