Node.js中如何实现CSV流单行数据的条件编辑修改
Node.js 大体积CSV低内存逐行处理方案
核心思路是基于Node.js Stream的背压机制,搭建「读流-逐行解析-单行处理-格式转换-写流」的全链路管道,内存中始终仅保留当前处理的单条数据,内存占用稳定在百KB级,可支持GB级CSV文件处理,不会出现OOM问题。
核心实现逻辑
全程不缓存全量CSV数据,每解析出一行数据就立刻完成判定、业务执行、字段修改、写入操作,处理完的行数据直接释放内存。
可直接复用的实现代码
const fs = require('fs'); const { Transform } = require('stream'); const csvtojson = require('csvtojson'); const { Parser } = require('json2csv'); // 配置参数,根据实际业务修改 const TARGET_CSV_PATH = './data.csv'; const TEMP_FILE_PATH = './data_temp.csv'; const BOOLEAN_MARK_FIELD = 'processed'; // 替换为实际的布尔标记位字段名 const CHUNK_SIZE = 64 * 1024; // 每次读写块大小设为64KB,平衡性能和内存占用 // 初始化CSV序列化工具,仅第一行写入表头 const csvSerializer = new Parser({ header: true }); let hasWriteHeader = false; // 自定义逐行处理转换流 const lineProcessTransform = new Transform({ writableObjectMode: true, readableObjectMode: false, highWaterMark: CHUNK_SIZE, async transform(singleRow, _, callback) { try { // 兼容CSV解析出的字符串类型布尔值 const flagValue = singleRow[BOOLEAN_MARK_FIELD]; if (flagValue === false || flagValue === 'false') { // 执行自定义业务逻辑,支持异步操作 await customBusinessOperation(singleRow); // 修改标记位为true singleRow[BOOLEAN_MARK_FIELD] = true; } // 序列化当前行为CSV格式 let lineContent; if (hasWriteHeader) { lineContent = '\n' + csvSerializer.parse(singleRow, { header: false }); } else { lineContent = csvSerializer.parse(singleRow); hasWriteHeader = true; } // 推送到下游写入流 callback(null, lineContent); } catch (err) { callback(err); } } }); // 组装流管道 const readStream = fs.createReadStream(TARGET_CSV_PATH, { highWaterMark: CHUNK_SIZE }); const writeStream = fs.createWriteStream(TEMP_FILE_PATH, { highWaterMark: CHUNK_SIZE }); readStream // 配置csvtojson为逐行输出模式,不缓存全量数据 .pipe(csvtojson({ downstreamFormat: 'line' })) .pipe(lineProcessTransform) .pipe(writeStream); // 处理完成后原子替换原文件 writeStream.on('finish', async () => { await fs.promises.rename(TEMP_FILE_PATH, TARGET_CSV_PATH); console.log('处理完成'); }); // 统一错误处理,出错时清理临时文件 [readStream, lineProcessTransform, writeStream].forEach(stream => { stream.on('error', (err) => { fs.unlink(TEMP_FILE_PATH, () => {}); console.error('处理失败:', err); process.exit(1); }); }); /** * 自定义业务操作,替换为实际逻辑 * @param {object} rowData 当前行数据 */ async function customBusinessOperation(rowData) { // 示例:存库、调用接口、数据计算等操作 // 异步操作会阻塞当前行处理流程,不会出现并发乱序、流背压问题 }
避坑要点
- 禁止使用会拉取全量数据的API:不要调用csvtojson的
then()方法获取全量数组,不要额外定义数组缓存行数据,所有操作必须在单行进入转换流时即时处理 - 异步业务必须适配异步流:如果自定义业务包含IO、接口调用等异步逻辑,必须等异步操作执行完成后再调用
callback推送数据,否则会出现内存泄漏、数据顺序错乱 - 必须用临时文件中转:禁止直接读写同一个文件,否则读流未完成时原文件会被写流覆盖导致数据丢失,临时文件写入完成后通过
fs.rename做原子替换是最安全的方案 - 注意布尔值类型兼容:CSV解析出的字段默认是字符串类型,判定标记位时要同时匹配字符串
'false'和布尔值false,避免漏处理 - 不要随意调大高水位线:64KB~128KB是兼顾性能和内存的合理区间,过大的块大小不会明显提升处理速度,反而会抬升内存占用
内容的提问来源于stack exchange,提问作者abhi
相关产品推荐
相关产品推荐

