如何使用@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
相关产品推荐
相关产品推荐

