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

Node.js并行处理多JSON文件入库问题求助

解决Node.js多Worker处理JSON文件插入MongoDB的性能与连接问题

问题根源分析

  1. CPU使用率过高:按CPU核心数创建Worker会占满所有核心,加上每个Worker同步读取文件+同步数据映射,线程切换开销与IO阻塞导致CPU空转;若数据映射为CPU密集型,会直接拉满核心负载。
  2. MongoDB连接数过多:每个Worker单独创建MongoDB客户端,每个客户端默认维护100个连接的池,总连接数=Worker数×100,远超MongoDB默认最大连接数(1000),触发连接限制。
  3. 连接异常:连接数超限被MongoDB拒绝,Worker退出时未关闭连接导致泄漏,频繁创建销毁连接引发不稳定。

具体解决方案

1. 优化Worker调度与CPU负载

  • 限制Worker数量:用os.cpus().length - 1(留1个核心给主线程处理消息和DB操作),避免系统过载。
  • 异步IO处理:Worker内部用fs.promises.readFile异步读取文件,避免线程阻塞导致CPU空转。
  • 批量任务分发:给每个Worker分配一批文件(而非单个),减少线程切换开销;用任务队列实现限流,Worker完成一批任务后再领取下一批。

2. 统一MongoDB连接池(核心优化)

不要让Worker直接操作MongoDB,改为主线程维护单个连接池,Worker仅负责文件读取与数据映射:

  • 主线程初始化MongoDB客户端,设置合理的连接池参数(maxPoolSize建议设为50-100,根据Worker数量调整)。
  • Worker读取并处理完JSON数据后,将结果发送给主线程,由主线程统一执行批量插入。
  • 若必须让Worker操作DB,每个Worker复用单个MongoDB客户端,且将maxPoolSize设为总连接数上限 / Worker数,避免连接池溢出。

3. 解决连接异常

  • 正确关闭连接:Worker退出前调用client.close(),主线程在所有任务完成后关闭客户端。
  • 配置连接参数:设置connectTimeoutMS=30000、socketTimeoutMS=30000、minPoolSize=5,提升连接稳定性。
  • 添加重试机制:对DB插入操作包裹重试逻辑,处理临时连接异常(自定义重试或用p-retry库)。

代码调整示例

主线程(readFiles.js)

const { Worker } = require('worker_threads');
const fs = require('fs/promises');
const path = require('path');
const { MongoClient } = require('mongodb');
const os = require('os');

// 配置
const WORKER_COUNT = os.cpus().length - 1;
const MONGO_URI = 'mongodb://localhost:27017';
const DB_NAME = 'test';
const COLLECTION_NAME = 'data';
const JSON_DIR = './json-files';

// 主线程维护MongoDB连接
let mongoClient;
let collection;

// 初始化DB连接
async function initDB() {
  mongoClient = new MongoClient(MONGO_URI, {
    maxPoolSize: 80, // 合理设置连接池大小
    connectTimeoutMS: 30000,
    socketTimeoutMS: 30000
  });
  await mongoClient.connect();
  collection = mongoClient.db(DB_NAME).collection(COLLECTION_NAME);
}

// 分发任务给Worker
async function dispatchTasks(fileList) {
  const chunkSize = Math.ceil(fileList.length / WORKER_COUNT);
  const chunks = [];
  for (let i = 0; i < fileList.length; i += chunkSize) {
    chunks.push(fileList.slice(i, i + chunkSize));
  }

  const workerPromises = chunks.map((chunk, index) => {
    return new Promise((resolve) => {
      const worker = new Worker('./Worker.js', { workerData: { fileChunk: chunk, jsonDir: JSON_DIR } });
      worker.on('message', async (processedData) => {
        // 主线程批量插入DB
        if (processedData.length > 0) {
          await collection.insertMany(processedData, { ordered: false });
        }
      });
      worker.on('error', (err) => console.error(`Worker ${index} error:`, err));
      worker.on('exit', (code) => {
        console.log(`Worker ${index} exited with code ${code}`);
        resolve();
      });
    });
  });

  await Promise.all(workerPromises);
}

// 主流程
async function main() {
  try {
    await initDB();
    // 获取所有JSON文件名
    const files = await fs.readdir(JSON_DIR);
    const jsonFiles = files.filter(file => path.extname(file) === '.json');
    await dispatchTasks(jsonFiles);
    // 返回合并后的结果(可从DB查询或收集Worker返回的数据)
    const allData = await collection.find({}).toArray();
    console.log('All data inserted successfully:', allData.length);
  } catch (err) {
    console.error('Main process error:', err);
  } finally {
    if (mongoClient) await mongoClient.close();
  }
}

main();

Worker线程(Worker.js)

const { parentPort, workerData } = require('worker_threads');
const fs = require('fs/promises');
const path = require('path');

// 仅负责读取文件和数据映射
async function processFiles(fileChunk, jsonDir) {
  const processedData = [];
  for (const file of fileChunk) {
    const filePath = path.join(jsonDir, file);
    try {
      const rawData = await fs.readFile(filePath, 'utf8');
      const jsonData = JSON.parse(rawData);
      // 数据映射逻辑(示例)
      const mappedData = {
        ...jsonData,
        processedAt: new Date(),
        sourceFile: file
      };
      processedData.push(mappedData);
    } catch (err) {
      console.error(`Error processing file ${file}:`, err);
    }
  }
  // 发送处理后的数据给主线程
  parentPort.postMessage(processedData);
}

processFiles(workerData.fileChunk, workerData.jsonDir);

额外优化建议

  • 批量插入优化:用bulkWrite代替insertMany,支持更灵活的插入/更新逻辑,减少DB请求次数。
  • 任务队列限流:若文件数量极大,用p-queue等库实现任务队列,控制并发处理的文件数量,避免瞬间IO过载。
  • CPU密集型映射优化:若数据映射逻辑复杂,可改用cluster模块(多进程),避免Worker线程抢占主线程的V8资源。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 12:15:43