如何使用NODEJS、AWS Lambda自动将S3存储桶JSON数据导入DynamoDB
实现方案:基于Node.js + AWS Lambda + DynamoDB 自动导入S3存储桶JSON数据
前置依赖准备
- 提前创建好用于存放JSON文件的S3存储桶
- 提前创建目标DynamoDB表,确认好表的分区键、排序键规则,和JSON数据的主键字段对应
- 配置Lambda执行角色,至少开通以下权限:S3存储桶的读取权限、DynamoDB表的写入权限、CloudWatch日志的写入权限(用于排查问题)
步骤1:配置Lambda触发器
进入Lambda控制台,给当前函数添加S3类型触发器:
- 触发事件选择
s3:ObjectCreated:*(所有文件新增事件触发) - 绑定目标S3存储桶
- 可按需添加后缀过滤,比如只匹配
.json后缀的文件,避免其他格式文件误触发
步骤2:编写Node.js版本Lambda代码
以下是可直接使用的代码示例,适配aws-sdk v3(Lambda Node.js 18.x及以上运行时默认内置,无需额外打包依赖):
import { S3Client, GetObjectCommand } from "@aws-sdk/client-s3"; import { DynamoDBClient, BatchWriteItemCommand } from "@aws-sdk/client-dynamodb"; import { marshall } from "@aws-sdk/util-dynamodb"; const s3Client = new S3Client({ region: process.env.AWS_REGION }); const dynamoClient = new DynamoDBClient({ region: process.env.AWS_REGION }); // 替换为你的DynamoDB表名 const DYNAMO_TABLE_NAME = "你的目标表名"; // BatchWriteItem单次最多写入25条,这里做分片阈值设置 const BATCH_WRITE_MAX_COUNT = 25; // 工具函数:将S3的Readable流转为字符串 const streamToString = (stream) => new Promise((resolve, reject) => { const chunks = []; stream.on("data", (chunk) => chunks.push(chunk)); stream.on("end", () => resolve(Buffer.concat(chunks).toString("utf8"))); stream.on("error", reject); }); // 工具函数:数组分片,适配批量写入限制 const sliceArray = (arr, size) => { const res = []; for (let i = 0; i < arr.length; i += size) { res.push(arr.slice(i, i + size)); } return res; }; export const handler = async (event) => { // 从触发事件中获取S3桶名和文件名 const s3Record = event.Records[0].s3; const bucketName = s3Record.bucket.name; const fileName = decodeURIComponent(s3Record.object.key.replace(/\+/g, " ")); try { // 1. 从S3读取JSON文件内容 const getObjectCmd = new GetObjectCommand({ Bucket: bucketName, Key: fileName, }); const s3Response = await s3Client.send(getObjectCmd); const fileContent = await streamToString(s3Response.Body); const jsonData = JSON.parse(fileContent); // 适配JSON内容为单条对象或者多条对象数组的情况 const dataList = Array.isArray(jsonData) ? jsonData : [jsonData]; // 2. 分片批量写入DynamoDB const batchGroups = sliceArray(dataList, BATCH_WRITE_MAX_COUNT); for (const group of batchGroups) { const writeRequests = group.map(item => ({ PutRequest: { Item: marshall(item) // 自动将JS对象转为DynamoDB支持的格式 } })); const batchWriteCmd = new BatchWriteItemCommand({ RequestItems: { [DYNAMO_TABLE_NAME]: writeRequests } }); await dynamoClient.send(batchWriteCmd); } console.log(`成功导入${dataList.length}条数据到DynamoDB`); return { statusCode: 200, body: "导入成功" }; } catch (err) { console.error("导入失败:", err); throw err; } };
注意事项
- 如果单个JSON文件数据量过大,可适当调高Lambda的内存配置和超时时间,避免执行中断
- JSON数据的每一条记录必须包含DynamoDB表定义的主键字段,否则对应记录会写入失败
- 如果需要导入历史S3文件,直接将对应S3事件参数传入Lambda手动触发即可
- 可根据业务需求添加写入失败重试、失败数据落盘等扩展逻辑
内容的提问来源于stack exchange,提问作者Naveenkumar K
相关产品推荐
相关产品推荐

