如何在Node.js中去除超大型JSON文件中的重复数据?
在Node.js中处理超大型JSON文件去重(按empId)
核心结论:必须使用流处理
你的文件大小超过80GB、包含7亿条记录,直接将整个文件加载到内存会导致内存溢出(OOM),因此必须采用流式处理——边读取、边解析、边去重、边输出,全程不加载完整文件到内存。
实现方案
下面提供两种主流实现方式,分别适配不同内存情况:
方案1:内存足够时用Set去重(快速)
适用于唯一empId数量在内存可承受范围内(比如千万级),依赖JSONStream库实现流式JSON解析。
- 安装依赖
npm install jsonstream
- 代码实现
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,牺牲一点速度换取内存安全。
- 安装依赖
npm install jsonstream level
- 代码实现
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
相关产品推荐
相关产品推荐

