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

AWS Lambda处理DynamoDB流时创建表并绑定触发器失败问题

问题分析与修复:DynamoDB流触发Lambda中createEventSourceMapping未执行的原因

核心问题1:REMOVE分支的return语句直接终止了整个函数

你的代码中,当遇到REMOVE类型的记录时,直接执行return null,这会立即终止整个Lambda handler函数,导致后续所有记录(包括INSERT类型)都不会被处理。如果事件流中存在REMOVE记录且排在INSERT记录之前,createEventSourceMapping自然不会被调用。

修复:将return null改为continue,跳过当前REMOVE记录,继续处理后续记录:

if(record.eventName == 'REMOVE') {
    continue; // 跳过该记录,继续循环处理其他记录
}

核心问题2:混合使用async/await与.then()导致异步流程混乱

你在await dynamodb.createTable().promise()之后又链式调用.then()和.catch(),这种混合写法会让异步流程的执行顺序和错误处理变得不透明,内部的await lambda.createEventSourceMapping()可能因为未被正确追踪而出现执行异常,且错误信息可能被隐藏。

修复:统一使用async/await + try/catch的写法,让代码逻辑更清晰,便于调试:

exports.handler = async (event) => {
    for (const record of event.Records) {
        let tableStreamArn = '';

        if(record.eventName === 'REMOVE') {
            continue; // 跳过删除操作
        }

        if(record.eventName === 'INSERT') {
            console.log("Campaign Id here : ", record.dynamodb.NewImage.campaignId.N);
            const campaignId = record.dynamodb.NewImage.campaignId.N;

            const params = {
                AttributeDefinitions: [
                    { AttributeName: "userId", AttributeType: "S" },
                    { AttributeName: "externalDateTime", AttributeType: "S" }
                ],
                KeySchema: [
                    { AttributeName: "userId", KeyType: "HASH" },
                    { AttributeName: "externalDateTime", KeyType: "RANGE" }
                ],
                BillingMode: 'PAY_PER_REQUEST',
                TableName: `CampaignId-${campaignId}-RawVotes`,
                StreamSpecification: {
                    StreamEnabled: true,
                    StreamViewType: 'NEW_IMAGE'
                }
            };

            try {
                // 创建表
                const data = await dynamodb.createTable(params).promise();
                console.log("Successfully created Table", data);
                tableStreamArn = data.TableDescription.LatestStreamArn;
                console.log("Table stream ARN here : ", tableStreamArn);

                // 创建事件源映射
                const esmParams = {
                    FunctionName: process.env.AGGREGATION_FUNCTION_NAME,
                    Enabled: true,
                    EventSourceArn: tableStreamArn,
                    StartingPosition: 'LATEST',
                };
                const esmData = await lambda.createEventSourceMapping(esmParams).promise();
                console.log("Successfully attached event : ", esmData);
            } catch (err) {
                console.error("Error during table creation or event mapping:", err);
                // 可根据需求决定是否继续处理其他记录,或抛出错误终止
                // throw err;
            }
        }
    }
};

额外建议

  • 权限检查:确保当前Lambda角色拥有dynamodb:CreateTable和lambda:CreateEventSourceMapping权限,同时目标聚合Lambda函数的资源策略允许当前Lambda创建触发器。
  • 避免重复创建:可以在创建表前调用dynamodb.describeTable()检查表是否已存在,防止重复执行创建操作引发错误。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 14:15:41