Node.js中如何用流分片过滤S3上的JSON.gz对象数组
问题:流式处理S3上的大JSON数组时,如何按对象组分片而非字节分片?
背景与问题
S3存储着超大体积的JSON数组文件,格式如下:
[{"name":'A', "lastName":'A', age: 18}, {"name":'B', "lastName":'B', age: 20}, ...]
为了优化内存占用,希望通过Streams流式过滤数据,避免将整个文件加载到内存。但设置objectMode: true后,控制台依然输出字节字符串分片(比如"{name:'A', lastName:'"),而非预期的对象组(比如[{"name":'A', "lastName":'A'}, {"name":'B', "lastName":'B'}])。
原尝试代码:
const filteredData = []; const filterTransform = new Transform({ objectMode: true, transform(chunk, _, callback) { console.log("chunk : "+chunk); try { const filteredData = chunk.map((item: any) => ({ name: item.name, lastName: item.lastName, })); filteredData.push(JSON.stringify(filteredData)); } catch (err) { callback(err); } callback(); }, }); const client = getS3Client(); const command = new GetObjectCommand({ Bucket: bucket, Key: key, }); const data:GetObjectCommandOutput= await client.send(command); const readStream = (dataSingle.Body! as Readable) .pipe(zlib.createGunzip()) .pipe(filterTransform)
核心原因
你当前的流链中,zlib.createGunzip()输出的仍是原始字节/字符串流,并没有将JSON数据解析为JavaScript对象。objectMode: true仅告知Transform流可以处理对象,但上游未输出对象,因此你拿到的还是字节分片。
解决方案
步骤1:使用JSONStream库解析流式JSON
JSONStream是专门处理大JSON文件的流式解析库,能将JSON数组拆分为单个对象输出,让下游Transform流可以直接处理对象(此时objectMode才会生效)。
先安装依赖:
npm install jsonstream
步骤2:修改流链与处理逻辑
调整流顺序,在解压后先通过JSONStream.parse('*')解析JSON数组的每个元素,再传递给自定义Transform流。同时在Transform流中实现对象分组逻辑:
import JSONStream from 'JSONStream'; import { Transform } from 'stream'; import { S3Client, GetObjectCommand } from '@aws-sdk/client-s3'; import zlib from 'zlib'; // 存储待分组的过滤后对象 const pendingItems: Array<{name: string; lastName: string}> = []; const GROUP_SIZE = 2; // 每2个对象为一组 const filterTransform = new Transform({ objectMode: true, transform(item, _, callback) { try { // 过滤出需要的字段 const filteredItem = { name: item.name, lastName: item.lastName, }; pendingItems.push(filteredItem); // 达到分组数量时,推送分组JSON if (pendingItems.length === GROUP_SIZE) { this.push(JSON.stringify(pendingItems)); pendingItems.length = 0; // 清空待分组数组 } callback(); } catch (err) { callback(err as Error); } }, // 处理最后一批不足GROUP_SIZE的对象 flush(callback) { if (pendingItems.length > 0) { this.push(JSON.stringify(pendingItems)); } callback(); } }); async function processS3Json() { const client = new S3Client({ // 填写你的S3客户端配置,比如region等 }); const command = new GetObjectCommand({ Bucket: '你的存储桶名称', Key: 'JSON文件的Key(比如data.json.gz)', }); const data = await client.send(command); (data.Body as NodeJS.ReadableStream) .pipe(zlib.createGunzip()) .pipe(JSONStream.parse('*')) // 解析JSON数组的每个元素为单个对象 .pipe(filterTransform) .on('data', (chunk) => { // 这里拿到的就是你期望的对象组JSON字符串 console.log('chunk:', chunk); }) .on('end', () => { console.log('文件处理完成'); }) .on('error', (err) => { console.error('处理出错:', err); }); } processS3Json();
关键说明
JSONStream.parse('*'):该语法表示解析JSON数组的每一个元素,将数组拆分为单个对象逐个输出到下游流objectMode: true:此时上游输出的是JavaScript对象,Transform流的transform方法可以直接处理单个人员对象flush方法:流结束时触发,用于处理最后一批不足分组数量的对象,避免数据丢失
无第三方库的替代思路(不推荐)
如果不想依赖第三方库,需要手动拼接JSON字符串、识别数组元素的分隔符(,)、处理JSON的开头[和结尾],但这种方式容易出现解析错误(比如字符串中包含,的情况),维护成本高,因此更推荐使用成熟的JSONStream库。
内容的提问来源于stack exchange,提问作者Ananthalakshmi Sankar
相关产品推荐
相关产品推荐

