如何将OpenAI流式响应通过Lambda与API Gateway传输至客户端?
实现Lambda流式转发OpenAI响应到客户端的方案
核心前提
AWS API Gateway需支持流式响应,可选择以下任一方式:
- 使用HTTP API(原生支持流式,推荐)
- 为现有REST API开启集成响应的流式模式
- 使用Lambda函数URL(配置更简单,可绕过API Gateway直接暴露接口)
同时Lambda需使用Node.js 14+版本,支持流式响应API。
方案1:直接管道转发OpenAI流(最优,低开销)
无需手动解析每段数据,直接将OpenAI返回的流通过Lambda流式响应传递给客户端,减少中间处理开销。
修改后的Lambda代码:
import { pipeline } from 'stream/promises'; export const handler = awslambda.streamifyResponse(async (event, responseStream, context) => { try { const res = await openai.createCompletion({ ...params, stream: true, }, { responseType: 'stream' }); // 设置SSE格式响应头,匹配OpenAI流式输出规范 responseStream.setContentType('text/event-stream'); responseStream.setHeader('Cache-Control', 'no-cache'); responseStream.setHeader('Connection', 'keep-alive'); // 直接管道转发OpenAI的流到客户端响应流 await pipeline(res.data, responseStream); // 同步收集完整内容用于存储到DynamoDB let content = ''; res.data.on('data', (chunk) => { const lines = chunk.toString().split('\n').filter(line => line.trim() !== ''); for (const line of lines) { const message = line.replace(/^data: /, ''); if (message !== '[DONE]') { try { const parsed = JSON.parse(message); content += parsed.choices[0].text; } catch (e) {} } } }); res.data.on('end', () => { storeRecord(content); }); } catch (error) { console.error('请求错误:', error); responseStream.setStatusCode(500); responseStream.write('data: {"error": "请求失败"}\n\n'); responseStream.end(); } });
方案2:手动逐段处理后发送(适合需加工内容的场景)
如果需要对OpenAI返回的内容修改后再发送,可沿用你原有的解析逻辑,逐段处理后推送给客户端。
修改后的Lambda代码:
export const handler = awslambda.streamifyResponse(async (event, responseStream, context) => { try { // 设置SSE响应头 responseStream.setContentType('text/event-stream'); responseStream.setHeader('Cache-Control', 'no-cache'); responseStream.setHeader('Connection', 'keep-alive'); let content = ''; const res = await openai.createCompletion({ ...params, stream: true, }, { responseType: 'stream' }); res.data.on('data', (data) => { const lines = data.toString().split('\n').filter(line => line.trim() !== ''); for (const line of lines) { const message = line.replace(/^data: /, ''); if (message === '[DONE]') { // 发送结束标记 responseStream.write('data: [DONE]\n\n'); // 存储完整内容到DynamoDB storeRecord(content); responseStream.end(); return; } try { const parsed = JSON.parse(message); const text = parsed.choices[0].text; content += text; // 以SSE格式推送解析后的文本到客户端 responseStream.write(`data: ${JSON.stringify({ text })}\n\n`); } catch (error) { console.error('解析流数据失败:', message, error); } } }); res.data.on('error', (err) => { console.error('流传输错误:', err); responseStream.setStatusCode(500); responseStream.write('data: {"error": "流传输失败"}\n\n'); responseStream.end(); }); } catch (error) { console.error('请求初始化错误:', error); responseStream.setStatusCode(500); responseStream.write('data: {"error": "请求初始化失败"}\n\n'); responseStream.end(); } });
关键配置要点
- HTTP API集成:创建HTTP API并集成Lambda时,在集成设置中开启“流式响应”。
- REST API集成:在集成请求中设置“Content Handling”为“Passthrough”,阶段设置中开启“Enable Streaming”,并确保响应头为
text/event-stream。 - Lambda函数URL:直接为Lambda配置函数URL,在设置中开启“流式响应”,无需额外API Gateway配置。
注意事项
- 确保Lambda执行角色拥有调用OpenAI API、写入DynamoDB的权限。
- Lambda超时时间需覆盖OpenAI生成响应的最长耗时,避免连接提前中断。
- 客户端需使用Server-Sent Events(SSE)方式接收数据,比如浏览器端用
EventSourceAPI。
内容的提问来源于stack exchange,提问作者good1492
相关产品推荐
相关产品推荐

