Node.js处理10GB CSV文件的代码内存泄漏排查求助
处理10GB CSV文件时Node.js流内存泄漏(移除transformer后恢复正常)
我编写了一个处理10GB大小CSV文件的Node.js脚本,核心逻辑为:
- 以流的形式打开文件
- 解析CSV对象
- 修改对象
- 将处理结果输出到新文件
但运行代码时出现内存泄漏问题,尝试多种调整均无效,移除管道中的transformer后泄漏现象消失。以下是我的代码:
'use strict'; import fs from 'node:fs'; import {parse, transform, stringify} from 'csv'; import lineByLine from 'n-readlines'; // big input file const inputFile = './input-data.csv'; // read headers first const linesReader = new lineByLine(inputFile); const firstLine = linesReader.next(); linesReader.close(); const headers = firstLine.toString() .split(',') .map(header => { return header .replace(/^"/, '') .replace(/"$/, '') .replace(/\s+/g, '_') .replace('(', '_') .replace(')', '_') .replace('.', '_') .replace(/_+$/, ''); }); // file stream const fileStream1 = fs.createReadStream(inputFile); // parser stream const parserStream1 = parse({delimiter: ',', cast: true, columns: headers, from_line: 1}); // transformer const transformer = transform(function(record) { return Object.assign({}, record, { SomeField: 'BlaBlaBla', }); }); // stringifier stream const stringifier = stringify({delimiter: ','}); console.log('Loading data...'); // chain of pipes fileStream1.on('error', err => { console.log(err); }) .pipe(parserStream1).on('error', err => {console.log(err); }) .pipe(transformer).on('error', err => { console.log(err); }) .pipe(stringifier).on('error', err => { console.log(err); }) .pipe(fs.createWriteStream('./_data/new-data.csv')).on('error', err => { console.log(err); }) .on('finish', () => { console.log('Loading data finished!'); });
问题原因
你使用的csv库中的transform方法可能存在背压处理缺陷,在处理大文件时,无法及时将转换后的记录传递给下游流,导致大量未处理数据堆积在内存中,引发泄漏。移除transformer后,数据直接从parser流向stringifier,没有中间环节的缓冲,因此内存恢复正常。
解决方案
方案1:改用through2实现转换流(推荐)
through2是Node.js生态中更可靠的流转换工具,对背压处理更严谨,能避免内存堆积问题:
- 安装依赖:
npm install through2
- 替换原transformer代码:
import through2 from 'through2'; // 替换原来的transformer const transformer = through2.obj(function(record, enc, callback) { const modifiedRecord = Object.assign({}, record, { SomeField: 'BlaBlaBla', }); this.push(modifiedRecord); callback(); });
方案2:升级csv库
可能是csv库旧版本的bug,升级到最新版可修复:
npm update csv
额外优化:统一用csv库读取表头
原代码用n-readlines单独读取表头,可改用csv库自身实现,避免重复打开文件:
'use strict'; import fs from 'node:fs'; import {parse, stringify} from 'csv'; import through2 from 'through2'; const inputFile = './input-data.csv'; let headers; // 先读取表头 const headerStream = fs.createReadStream(inputFile); const headerParser = parse({ delimiter: ',', from_line: 1, to_line: 1 }); headerStream .pipe(headerParser) .on('data', (row) => { headers = row.map(header => { return header .replace(/^"/, '') .replace(/"$/, '') .replace(/\s+/g, '_') .replace('(', '_') .replace(')', '_') .replace('.', '_') .replace(/_+$/, ''); }); }) .on('end', startProcessing) .on('error', err => console.log(err)); // 表头读取完成后启动主处理流程 function startProcessing() { const fileStream = fs.createReadStream(inputFile); const parserStream = parse({ delimiter: ',', cast: true, columns: headers, from_line: 2 }); const transformer = through2.obj(function(record, enc, callback) { this.push(Object.assign({}, record, { SomeField: 'BlaBlaBla' })); callback(); }); const stringifier = stringify({ delimiter: ',' }); console.log('Loading data...'); fileStream .on('error', err => console.log(err)) .pipe(parserStream).on('error', err => console.log(err)) .pipe(transformer).on('error', err => console.log(err)) .pipe(stringifier).on('error', err => console.log(err)) .pipe(fs.createWriteStream('./_data/new-data.csv')).on('error', err => console.log(err)) .on('finish', () => console.log('Loading data finished!')); }
内容的提问来源于stack exchange,提问作者Jakeroid
相关产品推荐
相关产品推荐

