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

Node.js使用pg库批量插入CSV到PostgreSQL:如何正确返回Promise

解决方法

你的问题核心在于Promise的作用域不对,pool.connect回调里返回的Promise无法被外部获取,且手动维护计数器跟踪异步操作的方式不够可靠。下面是修正后的代码,同时优化了连接管理和异步流程:

const parse = require('csv-parse'); // 确保已导入csv-parse库

let parsedata = [];
let csvStream = parse()
  .on("data", (data) => {
    parsedata.push(data);
  })
  .on("end", async () => { // 将回调改为async函数,方便用await处理异步逻辑
    parsedata.shift(); // 移除CSV表头行
    if (parsedata.length === 0) return;

    const query = "INSERT INTO Test_Stage (EmployeeID, UID) VALUES ($1, $2)";
    
    // 使用pool.connect的Promise版本替代回调形式
    const client = await pool.connect();
    try {
      // 将所有插入操作转为Promise,用Promise.all等待全部完成
      const insertPromises = parsedata.map(row => 
        client.query(query, row)
      );
      await Promise.all(insertPromises);
      console.log('所有数据插入完成');
      // 此处可直接执行后续函数
      // yourFollowUpFunction();
    } catch (err) {
      console.error('插入失败:', err.stack);
      throw err; // 按需处理错误,比如终止流程或记录日志
    } finally {
      client.release(); // 必须释放客户端,避免连接池耗尽
    }
  });

stream.pipe(csvStream);

关键修正说明:

  1. Promise作用域修正:把end事件回调改为async函数,内部用await处理所有异步操作,确保异步流程可被外部追踪。
  2. 替换计数器为Promise.all:用parsedata.map将每个client.query转换为Promise,Promise.all会等待所有插入操作完成,比手动维护计数器更可靠,还能避免异步操作顺序问题。
  3. 连接池正确管理:用await pool.connect()获取客户端,通过try/finally确保客户端被释放,防止连接泄漏导致的资源耗尽。
  4. 统一错误处理:用try/catch捕获插入过程中的所有错误,避免单个插入失败导致整个流程崩溃(若需忽略单个错误,可在map中单独处理每个Promise的异常)。

如果需要将整个CSV解析+插入流程包装成可复用的异步函数,方便后续链式调用,可参考以下代码:

async function importCsvToDb(stream) {
  return new Promise((resolve, reject) => {
    let parsedata = [];
    const csvStream = parse()
      .on("data", (data) => parsedata.push(data))
      .on("error", reject) // 捕获CSV解析阶段的错误
      .on("end", async () => {
        parsedata.shift();
        if (parsedata.length === 0) return resolve();

        const query = "INSERT INTO Test_Stage (EmployeeID, UID) VALUES ($1, $2)";
        const client = await pool.connect();
        try {
          await Promise.all(parsedata.map(row => client.query(query, row)));
          resolve('所有数据插入完成');
        } catch (err) {
          reject(err);
        } finally {
          client.release();
        }
      });
    stream.pipe(csvStream);
  });
}

// 使用示例
importCsvToDb(stream)
  .then(() => {
    console.log('导入完成,执行后续逻辑');
    // yourFollowUpFunction();
  })
  .catch(err => console.error('导入失败:', err));

内容的提问来源于stack exchange,提问作者Roger Holland

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 17:56:20