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

如何使用@aws-sdk/client-s3流式读取S3中的大JSON Lines文件

使用@aws-sdk/client-s3流式处理S3中的JSON Lines文件

实现方案

直接基于S3响应的流式能力,结合Node.js内置的readline模块逐行解析JSON Lines内容,全程无需加载整个文件到内存。

完整代码示例

const { S3Client, GetObjectCommand } = require("@aws-sdk/client-s3");
const readline = require("readline");

async function processJsonLinesFromS3(region, bucketParams) {
  // 初始化S3客户端
  const s3Client = new S3Client({ region });
  // 获取S3对象响应
  const s3Response = await s3Client.send(new GetObjectCommand(bucketParams));

  // S3响应体本身就是可读流(Node.js环境下)
  const s3Stream = s3Response.Body;

  // 创建逐行解析接口
  const lineReader = readline.createInterface({
    input: s3Stream,
    crlfDelay: Infinity // 兼容所有换行符格式
  });

  // 逐行处理数据
  for await (const line of lineReader) {
    // 跳过空行
    if (!line.trim()) continue;

    try {
      // 解析单行JSON
      const record = JSON.parse(line);
      // 替换为你的DynamoDB写入逻辑
      // 示例:单条写入
      // await dynamoClient.send(new PutItemCommand({
      //   TableName: "your-target-table",
      //   Item: record
      // }));
      // 若数据量大,推荐用BatchWriteItem批量提交
      console.log("已处理记录:", record);
    } catch (error) {
      console.error("行解析失败,内容:", line, "错误:", error);
    }
  }

  console.log("所有JSON Lines处理完成");
}

// 调用示例
processJsonLinesFromS3("us-east-1", {
  Bucket: "your-bucket",
  Key: "data/your-file.jsonl"
}).catch(err => console.error("流程出错:", err));

关键细节

  • 流式特性:s3Response.Body在Node.js环境下是标准的Readable流,无需额外转换即可直接用于流式处理。
  • 内存优化:逐行读取解析,内存占用始终保持在较低水平,适合处理超大文件。
  • 容错处理:跳过空行、捕获解析错误,避免单个无效行导致整个流程终止。
  • 性能优化:大量数据导入时,建议将记录分批,用BatchWriteItemCommand批量写入DynamoDB,减少API调用次数。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 14:35:22