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

使用Node.js、fs和readline处理大CSV时内存泄漏求助

排查大型CSV导入SQL时的JavaScript内存溢出问题

问题描述

将大型CSV数据导入SQL时,执行到一半抛出JavaScript heap out of memory错误,日志显示最后一条记录停在PapaParse执行前。尝试在readline的line事件中使用await控制行处理顺序,但问题未解决。

原始代码

const readline = require('readline');
const fs = require('fs');

const postFromSarenetCSV = async (rq, res = response) => {
 try {
 const rl = readline.createInterface({
 input: fs.createReadStream('data/unparsedSarenet.csv',"UTF8"),
 crlfDelay: Infinity
 });
 await rl.on('line', async (line) => {
 //create object with papaparse library:
 let object = "column1,column2,column3\r\n"
 object = object.concat(line)
 let json = Papa.parse(object, {newline: "\r\n", header: true})
 //create sql request:
 var data = json.data[0]
 let {query, values} = createReplaceQuery('services', Object.entries(data))
 await dbConnection.query(query, values);
 });
 await events.once(rl, 'close');
 res.status(200).json({
 ok: true,
 msg: 'Database updated'
 })
 } catch (err) {
 console.error(err);
 }
} 

问题根源

  1. 错误的异步事件处理逻辑:rl.on('line')并非Promise,使用await无法等待所有行处理完成。每一行都会立即发起数据库请求,没有并发控制,当CSV行数较多时,成千上万个未完成的请求会堆积在事件循环中,快速耗尽内存。
  2. 低效的PapaParse使用方式:每次处理单行都拼接表头再调用PapaParse,重复创建解析实例,额外增加内存开销。
  3. 未处理的并发堆积:即使在line回调中用了await,readline仍会快速读取所有行并触发回调,导致所有异步请求几乎同时被初始化,内存占用瞬间飙升。

修复方案

方案1:手动控制并发数 + 简化行解析

const readline = require('readline');
const fs = require('fs');

// 限制同时执行的数据库请求数量,根据数据库性能调整
const CONCURRENCY_LIMIT = 5;
let activeRequests = 0;
let resolveImportComplete;
const importCompletePromise = new Promise(resolve => resolveImportComplete = resolve);

const postFromSarenetCSV = async (rq, res) => {
  try {
    const rl = readline.createInterface({
      input: fs.createReadStream('data/unparsedSarenet.csv', "UTF8"),
      crlfDelay: Infinity
    });

    rl.on('line', async (line) => {
      // 达到并发上限时等待
      while (activeRequests >= CONCURRENCY_LIMIT) {
        await new Promise(resolve => setTimeout(resolve, 100));
      }

      activeRequests++;
      try {
        // 直接分割行获取数据,替代PapaParse的冗余调用
        const [column1, column2, column3] = line.split(',');
        const data = { column1, column2, column3 };
        
        const { query, values } = createReplaceQuery('services', Object.entries(data));
        await dbConnection.query(query, values);
      } catch (rowError) {
        console.error(`处理行失败: ${line}`, rowError);
        // 可选:若需要终止整个导入流程,可关闭流并抛出错误
        // rl.close();
        // throw rowError;
      } finally {
        activeRequests--;
        // 流已关闭且无活跃请求时,标记导入完成
        if (activeRequests === 0 && rl.closed) {
          resolveImportComplete();
        }
      }
    });

    rl.on('close', () => {
      // 流关闭时仍有活跃请求,等待全部完成
      if (activeRequests === 0) {
        resolveImportComplete();
      }
    });

    await importCompletePromise;

    res.status(200).json({
      ok: true,
      msg: '数据库更新完成'
    });
  } catch (err) {
    console.error('导入失败', err);
    res.status(500).json({
      ok: false,
      msg: '导入失败'
    });
  }
};

方案2:使用PapaParse流式解析(更适合复杂CSV)

如果CSV存在含逗号的字段,手动分割会出错,推荐用PapaParse的流式解析:

const fs = require('fs');
const Papa = require('papaparse');

const CONCURRENCY_LIMIT = 5;
let activeRequests = 0;
let resolveImportComplete;
const importCompletePromise = new Promise(resolve => resolveImportComplete = resolve);

const postFromSarenetCSV = async (rq, res) => {
  try {
    const stream = fs.createReadStream('data/unparsedSarenet.csv', "UTF8");

    Papa.parse(stream, {
      header: true,
      step: async (result) => {
        while (activeRequests >= CONCURRENCY_LIMIT) {
          await new Promise(resolve => setTimeout(resolve, 100));
        }

        activeRequests++;
        try {
          const { query, values } = createReplaceQuery('services', Object.entries(result.data));
          await dbConnection.query(query, values);
        } catch (rowError) {
          console.error('处理行失败', rowError);
        } finally {
          activeRequests--;
        }
      },
      complete: () => {
        // 等待剩余活跃请求完成
        const waitForRequests = async () => {
          if (activeRequests > 0) {
            await new Promise(resolve => setTimeout(resolve, 100));
            await waitForRequests();
          }
          resolveImportComplete();
        };
        waitForRequests();
      },
      error: (parseError) => {
        console.error('CSV解析错误', parseError);
        throw parseError;
      }
    });

    await importCompletePromise;

    res.status(200).json({
      ok: true,
      msg: '数据库更新完成'
    });
  } catch (err) {
    console.error('导入失败', err);
    res.status(500).json({
      ok: false,
      msg: '导入失败'
    });
  }
};

额外优化建议

  • 批量插入:将多行数据打包成批量SQL请求(比如每100行执行一次),大幅减少数据库交互次数,降低内存占用。
  • 调整Node内存限制:临时方案可启动时增加内存,如node --max-old-space-size=4096 your-script.js,但核心还是优化并发逻辑。
  • 添加进度日志:记录已处理行数,方便监控导入进度和定位问题。

内容的提问来源于stack exchange,提问作者Alejandro Fabra Segarra

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 18:50:32