Node.js并行处理多JSON文件入库问题求助
解决Node.js多Worker处理JSON文件插入MongoDB的性能与连接问题
问题根源分析
- CPU使用率过高:按CPU核心数创建Worker会占满所有核心,加上每个Worker同步读取文件+同步数据映射,线程切换开销与IO阻塞导致CPU空转;若数据映射为CPU密集型,会直接拉满核心负载。
- MongoDB连接数过多:每个Worker单独创建MongoDB客户端,每个客户端默认维护100个连接的池,总连接数=Worker数×100,远超MongoDB默认最大连接数(1000),触发连接限制。
- 连接异常:连接数超限被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
相关产品推荐
相关产品推荐

