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
相关产品推荐
相关产品推荐

