如何用fork()实现交错/拆分插入循环,让4个Worker各处理250条SQL插入?
用Node.js Cluster实现批量SQL插入的并行优化方案
当然可以!你的思路完全没问题——用4个Worker各处理250条数据的并行插入,绝对能大幅提升批量数据入库的效率。我之前做过类似的批量数据同步优化,正好给你拆解下Cluster和fork()的核心逻辑,再给你一套可落地的实现方案。
先搞懂Cluster和fork()的核心逻辑
Cluster是Node.js官方提供的多核CPU利用工具,本质上是通过fork()方法从主进程(Master)复制出多个独立的子进程(Worker)。每个Worker都是一个完整的Node.js实例,拥有自己的事件循环和内存空间,能并行执行任务,不会互相阻塞。
简单说:你可以把主进程当成「任务调度员」,负责拆分数据、分配任务;4个Worker就是「干活的工人」,各自拿一份数据独立执行插入操作,充分利用多核CPU的算力。
具体实现步骤
1. 主进程:拆分数据+调度Worker
主进程的核心工作是把1000条数据分成4份,然后fork出4个Worker,把对应的数据块发给每个Worker。
const cluster = require('cluster'); const numWorkers = 4; // 你指定的4个Worker const batchSize = 1000; // 每页数据量 // 模拟你的数据集(实际从页面获取) const mockData = Array.from({ length: batchSize }, (_, i) => ({ name: `User${i}`, height: 170 + Math.random() * 20, weight: 60 + Math.random() * 30 })); if (cluster.isPrimary) { console.log(`主进程 ${process.pid} 启动`); // 拆分数据为4个chunk,每个250条 const dataChunks = Array.from({ length: numWorkers }, (_, index) => { const start = index * (batchSize / numWorkers); const end = start + (batchSize / numWorkers); return mockData.slice(start, end); }); // 启动Worker并分配任务 dataChunks.forEach((chunk, index) => { const worker = cluster.fork(); worker.send({ taskId: index, data: chunk }); }); // 监听Worker完成/错误事件 cluster.on('exit', (worker, code, signal) => { console.log(`Worker ${worker.process.pid} 退出,代码:${code}`); }); cluster.on('message', (worker, message) => { if (message.type === 'success') { console.log(`Worker ${worker.process.pid} 完成 ${message.count} 条数据插入`); } else if (message.type === 'error') { console.error(`Worker ${worker.process.pid} 出错:${message.error}`); } }); } else { // Worker进程逻辑,下面单独写 }
2. Worker进程:处理插入任务
每个Worker收到主进程发来的数据后,要创建独立的SQL连接(重要:不能和主进程或其他Worker共享连接),然后执行批量插入——这里要注意,就算在Worker里,也不要逐行插,用SQL的批量插入语法效率更高。
// 接上面的else部分 const mysql = require('mysql2/promise'); // 用promise版本更方便异步处理 // Worker收到任务后执行 process.on('message', async (msg) => { const { taskId, data } = msg; let connection; try { // 创建独立的数据库连接 connection = await mysql.createConnection({ host: 'your-host', user: 'your-user', password: 'your-password', database: 'your-db' }); // 批量插入:把250条数据拼成一条INSERT语句 const placeholders = data.map(() => '(?, ?, ?)').join(','); const values = data.flatMap(item => [item.name, item.height, item.weight]); const sql = `INSERT INTO users (name, height, weight) VALUES ${placeholders}`; const [result] = await connection.execute(sql, values); // 通知主进程任务完成 process.send({ type: 'success', count: result.affectedRows }); } catch (err) { process.send({ type: 'error', error: err.message }); } finally { if (connection) await connection.end(); process.exit(0); // 任务完成后退出Worker } });
关键注意事项
- 独立数据库连接:每个Worker必须创建自己的连接,因为Node.js的TCP连接无法跨进程共享,共享连接会导致各种奇怪的错误。如果数据量更大,建议每个Worker用连接池代替单连接。
- 批量插入优化:就算用了Worker,也不要逐行执行
INSERT,用批量插入语法能把数据库IO次数从250次降到1次,效率提升非常明显。如果250条的SQL语句太长,可以再拆成更小的批次(比如每50条一批)。 - 灵活配置Worker数量:如果你的服务器CPU核数少于4,建议根据
os.cpus().length来设置Worker数量,避免过度抢占资源。 - 错误重试机制:可以给Worker加个简单的重试逻辑,如果插入失败,重试1-2次,提升任务稳定性。
- 进度跟踪:主进程可以维护一个计数器,收到所有Worker的完成通知后,标记整个批次任务完成,方便后续流程处理。
这样一套方案下来,相比原来的逐行插入,效率至少能提升3-4倍(具体取决于数据库性能和CPU核数),而且Cluster模块是Node.js官方维护的,稳定性有保障。
内容的提问来源于stack exchange,提问作者Kyle Dunne
相关产品推荐
相关产品推荐

