如何在Node后端处理CustomGPT的流式SSE并转发给客户端
Node后端处理CustomGPT流式SSE并转发客户端方案
需求背景
开发Node后端服务,与第三方CustomGPT交互,接收其流式SSE响应并实时转发给客户端,以此隐藏CustomGPT的凭证及敏感配置。
CustomGPT API文档中关于stream参数的说明:
是否启用流式响应,若启用,响应会以纯数据的Server-Sent Events(SSE)形式实时发送,流以
status: "finish"消息终止。
官方提供的Node示例代码(仅做基础调用,未处理流式):
import sdk from '@api/customgpt'; sdk.auth('sdk-token-value'); sdk.postApiV1ProjectsProjectidConversationsSessionidMessages({ response_source: 'own_content', prompt: "Tell me how to consume your SSE's in Node" }, { stream: 'true', lang: 'en', projectId: '1234', sessionId: '1' }) .then(({ data }) => console.log(data)) .catch(err => console.error(err));
上述示例运行时会延迟许久后一次性返回完整响应,而非流式分段传输,返回格式符合MDN定义的SSE规范。
已尝试方案及问题
- 前端EventSource实现:仅适用于前端,且仅支持GET请求,无法携带认证信息,不适合后端中转场景
const evtSource = new EventSource("http://localhost:4000/event-source"); evtSource.onmessage = (event) => { if (event.data) { setData(JSON.parse(event.data)); } }; - 手动拆分SSE响应:尝试拆分
\n\n字符后用res.write发送,但所有消息会一次性发送,丢失流式实时性。
解决方案
核心思路是让后端直接获取CustomGPT的原始可读流,逐块接收后实时转发给客户端,以下是基于Express和Axios的实现:
1. 后端中转服务实现
import express from 'express'; import axios from 'axios'; const app = express(); app.use(express.json()); // 处理客户端请求,中转CustomGPT流式响应 app.post('/api/chat', async (req, res) => { const { prompt, projectId, sessionId } = req.body; const customGptToken = 'sdk-token-value'; // 隐藏在后端,不暴露给客户端 // 设置SSE响应头,告知客户端接收流式数据 res.writeHead(200, { 'Content-Type': 'text/event-stream', 'Cache-Control': 'no-cache', 'Connection': 'keep-alive', 'Access-Control-Allow-Origin': '*' // 根据实际需求配置跨域 }); try { // 向CustomGPT发起流式请求,获取可读流 const customGptResponse = await axios.post( `https://api.customgpt.ai/api/v1/projects/${projectId}/conversations/${sessionId}/messages`, { response_source: 'own_content', prompt }, { headers: { Authorization: `Bearer ${customGptToken}` }, params: { stream: 'true', lang: 'en' }, responseType: 'stream' // 关键配置:获取原始响应流 } ); // 监听流数据,实时转发给客户端 customGptResponse.data.on('data', (chunk) => { const chunkStr = chunk.toString('utf8'); // 拆分SSE事件并逐个转发 const events = chunkStr.split('\n\n').filter(event => event.trim()); events.forEach(event => res.write(`${event}\n\n`)); }); // 监听流结束,发送终止信号并关闭连接 customGptResponse.data.on('end', () => { res.write('data: {"status": "finish"}\n\n'); res.end(); }); // 处理流错误 customGptResponse.data.on('error', (err) => { console.error('CustomGPT stream error:', err); res.write(`data: {"error": "${err.message}"}\n\n`); res.end(); }); } catch (err) { console.error('Request to CustomGPT failed:', err); res.write(`data: {"error": "${err.message}"}\n\n`); res.end(); } }); app.listen(4000, () => console.log('中转服务运行在端口4000'));
2. 客户端接收流式响应
客户端可通过Fetch API处理POST流式请求(支持携带认证信息):
async function sendChat(prompt) { const response = await fetch('http://localhost:4000/api/chat', { method: 'POST', headers: { 'Content-Type': 'application/json' }, body: JSON.stringify({ prompt, projectId: '1234', sessionId: '1' }) }); const reader = response.body.getReader(); const decoder = new TextDecoder('utf8'); while (true) { const { done, value } = await reader.read(); if (done) break; const chunk = decoder.decode(value); const events = chunk.split('\n\n').filter(e => e.trim()); events.forEach(event => { if (event.startsWith('data:')) { const data = JSON.parse(event.slice(5).trim()); if (data.status === 'finish') { console.log('对话结束'); } else { console.log('收到内容:', data); // 这里更新UI展示实时内容 } } }); } }
关键注意事项
- 必须设置正确的SSE响应头,确保客户端能持续接收流式数据
- 启用
responseType: 'stream'是获取实时流的核心,避免SDK/HTTP库缓冲完整响应 - 不要对接收的Chunk做额外缓冲处理,实时转发才能保证流式效果
- 处理流的错误和结束事件,确保客户端能正确感知流的状态变化
内容的提问来源于stack exchange,提问作者Roland
相关产品推荐
相关产品推荐

