本地正常的Server-Sent Event处理器部署到Cloud Run后响应极慢
流式获取OpenAI Chat Completion API数据的Cloud Functions Gen2部署问题
我用以下代码通过Server-Sent Events从OpenAI的Chat Completion API流式获取数据,本地在Firebase Functions模拟器运行时完全正常,但部署到基于Cloud Run的Cloud Functions Gen2后调用,会出现10-15秒延迟,随后所有文本一次性返回。
原代码:
export async function promptStream( apiKey: string, userId: string, prompt: Prompt, handleNewChunk: (chunk: string) => Promise<void> ) { const req = https.request( { hostname: "api.openai.com", port: 443, path: "/v1/chat/completions", method: "POST", headers: { "Content-Type": "application/json", Authorization: "Bearer " + apiKey, }, }, function (res) { res.on("data", async (data) => { handleNewChunk(data) } }); res.on("end", () => { console.log("No more data in response."); }); } ); const body = JSON.stringify({ ...prompt, stream: true, }); req.write(body); req.end(); }
问题原因
Cloud Run(以及基于它的Gen2函数)默认会缓冲响应内容,直到函数处理完成才一次性发送给客户端,直接破坏了流式传输的实时性。另外原代码没有正确处理OpenAI返回的SSE格式数据,也未配置必要的响应头禁用缓冲。
解决方法
调整代码以支持流式响应,正确处理OpenAI的SSE数据,同时配置响应头禁用缓冲:
修改后的代码示例:
import * as https from 'https'; import type { Request, Response } from 'express'; // 假设Prompt类型定义 type Prompt = { model: string; messages: Array<{ role: string; content: string }>; // 其他自定义字段 }; export async function promptStream( apiKey: string, userId: string, prompt: Prompt, res: Response ) { // 配置SSE响应头,强制禁用缓冲 res.setHeader('Content-Type', 'text/event-stream'); res.setHeader('Cache-Control', 'no-cache'); res.setHeader('Connection', 'keep-alive'); res.flushHeaders(); // 立即发送响应头,避免客户端等待 const req = https.request( { hostname: "api.openai.com", port: 443, path: "/v1/chat/completions", method: "POST", headers: { "Content-Type": "application/json", Authorization: "Bearer " + apiKey, }, }, function (openaiRes) { let buffer = ''; openaiRes.on("data", (chunk) => { buffer += chunk.toString(); // 按SSE的分隔符分割数据 const lines = buffer.split('\n\n'); buffer = lines.pop() || ''; // 保留未完成的最后一行 for (const line of lines) { if (!line.startsWith('data: ')) continue; const dataStr = line.slice(6).trim(); // 处理流结束标记 if (dataStr === '[DONE]') { res.write('data: [DONE]\n\n'); res.end(); return; } try { const data = JSON.parse(dataStr); const content = data.choices[0]?.delta?.content; if (content) { // 以SSE格式发送内容片段 res.write(`data: ${JSON.stringify({ content })}\n\n`); res.flush(); // 立即推送数据到客户端 } } catch (err) { console.error('解析OpenAI响应失败:', err); } } }); openaiRes.on("end", () => { console.log("流传输结束"); res.end(); }); openaiRes.on("error", (err) => { console.error('OpenAI请求出错:', err); res.status(500).end('请求OpenAI服务失败'); }); } ); req.on('error', (err) => { console.error('创建请求出错:', err); res.status(500).end('创建请求失败'); }); const body = JSON.stringify({ ...prompt, stream: true, }); req.write(body); req.end(); } // HTTP触发的Cloud Functions Gen2入口函数 export const streamChat = async (req: Request, res: Response) => { const apiKey = process.env.OPENAI_API_KEY!; const userId = req.body.userId; const prompt: Prompt = req.body.prompt; await promptStream(apiKey, userId, prompt, res); };
关键调整点
- 配置SSE响应头:必须指定
text/event-stream类型,同时设置no-cache和keep-alive,并调用flushHeaders()立即发送头信息,防止Cloud Run缓冲响应。 - 解析OpenAI的SSE格式:OpenAI返回的流式数据以
\n\n分隔,需提取每个data:字段的内容,过滤结束标记[DONE],只处理包含内容片段的响应。 - 实时推送数据:每次发送内容片段后调用
res.flush(),确保数据立即推送到客户端,避免被缓冲。 - 部署配置:部署时设置足够的超时时间(例如
gcloud functions deploy streamChat --gen2 --runtime nodejs20 --timeout 360s --trigger-http),防止长连接提前断开。
内容的提问来源于stack exchange,提问作者skillshot
相关产品推荐
相关产品推荐

