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

如何使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 14:54:03