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

如何将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)方式接收数据,比如浏览器端用EventSource API。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 17:05:31