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

如何在Node.js中去除超大型JSON文件中的重复数据?

在Node.js中处理超大型JSON文件去重(按empId)

核心结论:必须使用流处理

你的文件大小超过80GB、包含7亿条记录,直接将整个文件加载到内存会导致内存溢出(OOM),因此必须采用流式处理——边读取、边解析、边去重、边输出,全程不加载完整文件到内存。

实现方案

下面提供两种主流实现方式,分别适配不同内存情况:


方案1:内存足够时用Set去重(快速)

适用于唯一empId数量在内存可承受范围内(比如千万级),依赖JSONStream库实现流式JSON解析。

  1. 安装依赖
npm install jsonstream
  1. 代码实现
const fs = require('fs');
const JSONStream = require('JSONStream');
const { Transform } = require('stream');

// 存储已处理的empId,避免重复
const seenEmpIds = new Set();

// 创建转换流:过滤重复记录
const deduplicateStream = new Transform({
  objectMode: true,
  transform(chunk, _, callback) {
    const empId = chunk.empId;
    if (!seenEmpIds.has(empId)) {
      seenEmpIds.add(empId);
      this.push(chunk);
    }
    callback();
  }
});

// 输入/输出流初始化
const inputStream = fs.createReadStream('你的超大文件.json');
const outputStream = fs.createWriteStream('去重后的文件.json');
const jsonParser = JSONStream.parse('rows.*'); // 定位到rows数组的每个元素

// 手动拼接输出的JSON结构(避免流式输出格式错误)
outputStream.write('{"rows": [');
let isFirstRecord = true;

// 串联流处理逻辑
jsonParser.pipe(deduplicateStream)
  .on('data', (record) => {
    const recordStr = JSON.stringify(record);
    outputStream.write(isFirstRecord ? recordStr : `,${recordStr}`);
    isFirstRecord = false;
  })
  .on('end', () => {
    outputStream.write(']}');
    outputStream.end();
    console.log('去重完成');
  })
  .on('error', (err) => {
    console.error('处理出错:', err);
    outputStream.end();
  });

inputStream.pipe(jsonParser);

方案2:内存不足时用磁盘存储去重(低内存占用)

如果唯一empId数量接近7亿,内存无法承载Set,改用磁盘键值库level存储已处理的empId,牺牲一点速度换取内存安全。

  1. 安装依赖
npm install jsonstream level
  1. 代码实现
const fs = require('fs');
const JSONStream = require('JSONStream');
const { Transform } = require('stream');
const level = require('level');

// 创建磁盘数据库,存储已处理的empId
const empIdDB = level('./empid-store', { valueEncoding: 'utf8' });

// 异步转换流:查询磁盘数据库判断是否重复
const deduplicateStream = new Transform({
  objectMode: true,
  async transform(chunk, _, callback) {
    try {
      const empId = chunk.empId;
      // 检查empId是否已存在,不存在则写入数据库并保留记录
      await empIdDB.get(empId).catch(() => {
        empIdDB.put(empId, '1');
        this.push(chunk);
      });
      callback();
    } catch (err) {
      callback(err);
    }
  }
});

// 输入/输出流初始化
const inputStream = fs.createReadStream('你的超大文件.json');
const outputStream = fs.createWriteStream('去重后的文件.json');
const jsonParser = JSONStream.parse('rows.*');

// 拼接JSON结构
outputStream.write('{"rows": [');
let isFirstRecord = true;

// 串联流并处理收尾
jsonParser.pipe(deduplicateStream)
  .on('data', (record) => {
    const recordStr = JSON.stringify(record);
    outputStream.write(isFirstRecord ? recordStr : `,${recordStr}`);
    isFirstRecord = false;
  })
  .on('end', async () => {
    outputStream.write(']}');
    outputStream.end();
    await empIdDB.close();
    console.log('去重完成');
  })
  .on('error', async (err) => {
    console.error('处理出错:', err);
    await empIdDB.close();
    outputStream.end();
  });

inputStream.pipe(jsonParser);

关键注意事项

  • 流式JSON解析:必须用JSONStream或stream-json这类库,原生JSON.parse会加载整个文件到内存,直接导致崩溃。
  • 格式拼接:流式输出需要手动拼接外层的{"rows": [ ... ]}结构,避免出现语法错误。
  • 错误处理:必须监听流的error事件,防止程序意外崩溃,同时及时关闭数据库/文件流。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 15:57:19