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

NodeJS中流导入大CSV至PostgreSQL(Drizzle ORM)遇背压内存溢出问题

解决大CSV流式插入PostgreSQL的背压与内存溢出问题

问题根源

你的代码中,data事件触发时直接调用异步的addDataToDB,但Node.js可读流不会等待异步操作完成——csv_parse会持续解析并推送数据,大量未完成的数据库插入Promise堆积在内存中,最终导致堆内存耗尽抛出FATAL ERROR: Reached heap limit Allocation failed。

解决方案

方法1:手动用pause/resume控制流

通过暂停/恢复流,确保每一行数据插入完成后再处理下一行:

const readStream = createReadStream("file_name.csv")
  .pipe(parse({ delimiter: ",", from_line: 2 }));

readStream.on("data", async (row) => {
  readStream.pause(); // 暂停流,避免继续推送数据
  try {
    await addDataToDB(row[2]);
  } catch (err) {
    console.error("单条插入失败:", err);
    // 可选:记录失败行到日志
  } finally {
    readStream.resume(); // 完成后恢复流
  }
});

readStream.on("end", () => console.log("所有数据处理完成"));
readStream.on("error", (err) => console.error("流处理错误:", err));

方法2:使用异步迭代器(更简洁)

Node.js可读流支持异步迭代,for await...of会自动处理背压,等待当前异步操作完成后再处理下一条数据:

async function processCSV() {
  const parser = createReadStream("file_name.csv")
    .pipe(parse({ delimiter: ",", from_line: 2 }));

  for await (const row of parser) {
    try {
      await addDataToDB(row[2]);
    } catch (err) {
      console.error("插入失败:", err);
    }
  }
  console.log("所有数据处理完成");
}

processCSV().catch(err => console.error("整体处理失败:", err));

方法3:批量插入(最优方案)

单条插入效率极低,攒一批数据再插入能大幅减少数据库请求次数,从根源降低内存压力:

async function processCSVInBatches(batchSize = 200) {
  const parser = createReadStream("file_name.csv")
    .pipe(parse({ delimiter: ",", from_line: 2 }));
  
  let batch = [];

  for await (const row of parser) {
    batch.push({ yourColumnName: row[2] }); // 按Drizzle表结构组装数据
    if (batch.length >= batchSize) {
      try {
        // Drizzle批量插入语法
        await db.insert(yourTable).values(batch);
        batch = []; // 清空批次
      } catch (err) {
        console.error("批量插入失败:", err);
        // 可选:将失败批次写入错误日志
      }
    }
  }

  // 处理剩余的不足一批的数据
  if (batch.length > 0) {
    try {
      await db.insert(yourTable).values(batch);
    } catch (err) {
      console.error("剩余数据插入失败:", err);
    }
  }

  console.log("所有数据处理完成");
}

processCSVInBatches(300).catch(err => console.error("处理失败:", err));

注意事项

  • 调整batchSize:根据数据库性能和本地内存情况测试最优值(通常100-500之间)
  • 错误处理:必须捕获插入异常,避免单个失败行中断整个流程
  • 临时堆内存调整:若需应急,可启动Node时加--max-old-space-size=4096(分配4GB堆内存),但这只是临时方案,核心解决还是控制流和批量插入

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 10:07:38