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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 07:36:42