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
相关产品推荐
相关产品推荐

