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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 03:07:13