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

Node.js使用fast-csv读取CSV时,如何高效在.on("data")内执行数据库查询?

这个问题我太熟了——fast-csv的流式读取速度确实快,要是直接在data事件里怼数据库查询,分分钟把数据库连接池撑爆,要么请求排队超时,要么直接报错。下面给你两个最实用的高效解决方案,附带代码示例,你可以根据场景选:


方案一:控制并发数(通用场景首选)

这种方式限制同时执行的数据库操作数量,既不会让CSV读取阻塞,也不会给数据库造成太大压力。我一般用async库的queue来实现,它能帮你自动维护任务队列,控制并发数。

代码示例

const csv = require('fast-csv');
const async = require('async');
const mysql = require('mysql2/promise'); // 用Promise版本更方便异步处理

// 1. 创建数据库连接池(核心:复用连接,避免频繁创建销毁)
const pool = mysql.createPool({
  host: '你的数据库地址',
  user: '用户名',
  password: '密码',
  database: '目标库',
  waitForConnections: true,
  connectionLimit: 10, // 连接池最大连接数,根据数据库承载调整
  queueLimit: 0
});

// 2. 创建并发队列,限制同时执行5个DB任务(不要超过连接池的connectionLimit)
const dbTaskQueue = async.queue(async (task, callback) => {
  try {
    const { phone, full_name, first_name, last_name, date_time } = task;
    
    // 执行第一个插入:用户表
    const [userResult] = await pool.query(
      "INSERT INTO users SET phone = ?, name = ?, first_name = ?, last_name = ?",
      [phone, full_name, first_name, last_name]
    );
    
    // 用插入后的ID执行第二个插入:会话表
    await pool.query(
      "INSERT INTO conversation (user_one, user_two, ip, time) VALUES(?, ?, ?, ?)",
      ['1', userResult.insertId, "", date_time]
    );
    
    callback(null); // 标记任务成功完成
  } catch (err) {
    console.error(`处理任务失败: ${JSON.stringify(task)}`, err);
    callback(err); // 标记任务失败,可后续处理重试
  }
}, 5);

// 监听队列错误
dbTaskQueue.error((err, task) => {
  console.error(`任务 ${JSON.stringify(task)} 执行失败`, err);
});

// 3. 读取CSV并推送任务到队列
csv.fromPath(newpath)
  .on("data", (data) => {
    // 假设CSV列顺序是:phone, full_name, first_name, last_name
    const taskData = {
      phone: data[0],
      full_name: data[1],
      first_name: data[2],
      last_name: data[3],
      date_time: new Date().toISOString() // 或从CSV取对应时间列
    };
    
    // 把任务推到队列,队列会自动控制并发执行
    dbTaskQueue.push(taskData, (err) => {
      if (err) console.error('任务推送失败', err);
    });
  })
  .on("end", () => {
    console.log('CSV读取完成,等待所有DB任务结束...');
    // 所有任务执行完后关闭连接池
    dbTaskQueue.drain(() => {
      console.log('所有数据库操作完成');
      pool.end();
    });
  })
  .on("error", (err) => {
    console.error('CSV读取出错', err);
    pool.end();
  });

方案二:批量插入(纯数据导入场景效率最高)

如果只是单纯把CSV数据导入数据库,批量插入比单条插入效率高N倍——因为数据库处理批量请求的开销远小于多次单请求。这里一定要用事务,确保users和conversation的插入是原子性的,避免数据不一致。

代码示例

const csv = require('fast-csv');
const mysql = require('mysql2/promise');

const pool = mysql.createPool({ /* 连接池配置同上 */ });

const BATCH_SIZE = 100; // 每攒100条数据执行一次批量插入,可根据数据库调整
let userBatch = []; // 暂存用户数据

csv.fromPath(newpath)
  .on("data", (data) => {
    // 把CSV数据转成批量插入的格式
    userBatch.push([data[0], data[1], data[2], data[3]]);
    
    // 达到批量大小就执行插入
    if (userBatch.length >= BATCH_SIZE) {
      processBatch();
    }
  })
  .on("end", async () => {
    // 处理剩余的不足批量的数据
    if (userBatch.length > 0) await processBatch();
    console.log('所有数据导入完成');
    pool.end();
  })
  .on("error", (err) => {
    console.error('CSV读取出错', err);
    pool.end();
  });

// 批量处理核心函数
async function processBatch() {
  const connection = await pool.getConnection();
  try {
    await connection.beginTransaction(); // 开启事务,保证操作原子性
    
    // 1. 批量插入用户表
    const [userResult] = await connection.query(
      "INSERT INTO users (phone, name, first_name, last_name) VALUES ?",
      [userBatch]
    );
    
    // 2. 生成会话表的批量数据(用户ID是连续自增的,从insertId开始)
    const conversationBatch = [];
    const startUserId = userResult.insertId;
    for (let i = 0; i < userResult.affectedRows; i++) {
      conversationBatch.push(['1', startUserId + i, "", new Date().toISOString()]);
    }
    
    // 3. 批量插入会话表
    await connection.query(
      "INSERT INTO conversation (user_one, user_two, ip, time) VALUES ?",
      [conversationBatch]
    );
    
    await connection.commit(); // 提交事务
    console.log(`成功导入${userResult.affectedRows}条数据`);
    userBatch = []; // 清空批量数组
  } catch (err) {
    await connection.rollback(); // 出错回滚
    console.error('批量导入失败,已回滚', err);
    throw err;
  } finally {
    connection.release(); // 释放连接回池
  }
}

额外注意事项
  • 连接池是必须的:不要用单个数据库连接,连接池会帮你复用连接,避免频繁创建销毁的开销。
  • 调整参数适配你的数据库:比如连接池connectionLimit、并发数、批量大小,都要根据数据库的承载能力来试,别贪多。
  • 做好数据校验:CSV里可能有脏数据(空值、格式错误),建议在data事件里先做校验,避免一条坏数据卡断整个导入。

内容的提问来源于stack exchange,提问作者Sunny Kumar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 09:45:06