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

AWS Websocket调用Lambda时消息延迟发送问题求助

问题:Lambda通过Websocket发送消息时延迟一个版本,总是发送上一次的负载

问题现象

  • Lambda1以payload "test1"调用Lambda2 → Lambda2未向客户端发送任何内容
  • Lambda1以payload "test2"调用Lambda2 → Lambda2向客户端发送"test1"
  • Lambda1以payload "test3"调用Lambda2 → Lambda2向客户端发送"test2"

环境信息

  • Node.js 16
  • Websocket消息发送依赖AWS.ApiGatewayManagementApi.postToConnection
  • Websocket API未配置缓存,无报错日志,权限策略(Lambda调用、DynamoDB访问、Websocket消息发送)配置正确
  • Lambda2日志显示每次都能获取到当前调用的正确payload

相关代码

Lambda1 调用代码

const AWS = require("aws-sdk") ;
const lambda = new AWS.Lambda({ region: "eu-west-1" });

const invokeLambda = async(functionName, payload) => {
  try {
    await lambda
      .invoke({
        FunctionName: `${functionName}`,
        InvocationType: "Event",
        Payload: JSON.stringify(payload),
      })
      .promise();
  } catch (err) {
    console.error( err, `Error during the invokation of the lambda ${functionName}` );
  }
}

Lambda2 消息发送代码

import * as AWS from 'aws-sdk';
const apigwManagementApi = new AWS.ApiGatewayManagementApi({
  endpoint: process.env.WEBSOCKET_DOMAIN_NAME,
});
import { InvokeResponseDto } from "../dto";

export const handler = async (
  event: InvokeRequestDto
): Promise<APIGatewayProxyResultV2> => {
  const { userId, data } = event;

  // Retrieve all websocket connections
  const items = await dynamo.scan(....)

  // Send message on websocket to all connection
  items?.forEach(async (item) => {
    try {
      const params : InvokeResponseDto = {
        ConnectionId: item.connectionId,
        Data: JSON.stringify(data),
      }
      await apigwManagementApi.postToConnection(params).promise()
    } catch (error) {
      console.error(error, "ERROR at apigwManagementApi: ");
    }
  });
  console.info(`Sending the websocket message to ${items?.length} websockets`);

  return { statusCode: 200 };
};

问题根源

Lambda2的handler中使用了forEach遍历连接列表并执行异步的postToConnection操作,但forEach不会等待内部的异步函数执行完成。Lambda函数在handler返回{ statusCode: 200 }后,执行环境会被AWS Lambda冻结甚至复用,当后续调用触发Lambda时,之前未完成的异步发送任务才会继续执行,导致发送的是上一次调用的payload。

修复方案

需要确保所有Websocket消息发送的异步操作都在handler返回前完成,可通过以下两种方式修复:

方案1:使用for...of循环替代forEach

export const handler = async (
  event: InvokeRequestDto
): Promise<APIGatewayProxyResultV2> => {
  const { userId, data } = event;

  const items = await dynamo.scan(....)

  // 用for...of确保每个异步发送都完成后再继续
  if (items) {
    for (const item of items) {
      try {
        const params : InvokeResponseDto = {
          ConnectionId: item.connectionId,
          Data: JSON.stringify(data),
        }
        await apigwManagementApi.postToConnection(params).promise()
      } catch (error) {
        console.error(error, "ERROR at apigwManagementApi: ");
      }
    }
  }
  
  console.info(`Sending the websocket message to ${items?.length} websockets`);

  return { statusCode: 200 };
};

方案2:使用Promise.all并行发送(推荐,效率更高)

export const handler = async (
  event: InvokeRequestDto
): Promise<APIGatewayProxyResultV2> => {
  const { userId, data } = event;

  const items = await dynamo.scan(....)

  if (items) {
    // 生成所有发送任务的Promise,用Promise.all等待全部完成
    const sendPromises = items.map(async (item) => {
      try {
        const params : InvokeResponseDto = {
          ConnectionId: item.connectionId,
          Data: JSON.stringify(data),
        }
        await apigwManagementApi.postToConnection(params).promise()
      } catch (error) {
        console.error(error, "ERROR at apigwManagementApi: ");
      }
    });
    await Promise.all(sendPromises);
  }
  
  console.info(`Sending the websocket message to ${items?.length} websockets`);

  return { statusCode: 200 };
};

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 11:40:29