Node.js+ES6+TypeScript数据库并行化操作及替代方案咨询
基于ES6+TypeScript的Node.js数据库操作并行化方案
当然有,针对你的需求,以下是几种可行的替代/优化方案,均支持ES6模块与TypeScript,且能配合Prisma Client完成批量数据库处理:
方案1:修复worker_threads的ES模块支持问题
你之前遇到的worker_threads不支持ES导入的问题,其实在Node.js 12+版本已经得到解决,只需调整配置即可:
必要配置
- 在
package.json中声明ES模块:{ "type": "module" } - 调整
tsconfig.json以适配ES模块编译:{ "compilerOptions": { "module": "ESNext", "target": "ES2020", "moduleResolution": "node", "outDir": "./dist", "esModuleInterop": true, "skipLibCheck": true } }
代码示例
主线程(CRON任务入口)
import { Worker } from 'worker_threads'; import cron from 'node-cron'; // CRON任务:每天凌晨1点执行 cron.schedule('0 1 * * *', async () => { const totalRecords = 10000; const threadCount = 4; const chunkSize = Math.ceil(totalRecords / threadCount); for (let i = 0; i < threadCount; i++) { const start = i * chunkSize; const end = Math.min(start + chunkSize - 1, totalRecords - 1); // 创建worker并传递索引范围 const worker = new Worker('./dist/worker.js', { workerData: { start, end } }); worker.on('message', (result) => { console.log(`Worker ${i}完成处理:${result.processed}条记录`); }); worker.on('error', (err) => { console.error(`Worker ${i}出错:`, err); }); } });
Worker线程(数据库处理逻辑)
import { parentPort, workerData } from 'worker_threads'; import { PrismaClient } from '@prisma/client'; const prisma = new PrismaClient(); async function processRecords() { const { start, end } = workerData; let processed = 0; // 处理指定索引范围的记录(示例:更新状态) const records = await prisma.someModel.findMany({ where: { id: { gte: start, lte: end } } }); for (const record of records) { await prisma.someModel.update({ where: { id: record.id }, data: { status: 'processed' } }); processed++; } await prisma.$disconnect(); parentPort?.postMessage({ processed }); } processRecords().catch(err => { console.error('处理失败:', err); parentPort?.postMessage({ error: err.message }); });
方案2:使用child_process.fork创建独立子进程
fork会生成独立的Node.js进程,天然支持ES模块,每个进程拥有独立的V8实例与Prisma连接池,适合IO密集型的数据库批量操作:
代码示例
主线程
import { fork } from 'child_process'; import cron from 'node-cron'; cron.schedule('0 1 * * *', () => { const totalRecords = 10000; const processCount = 4; const chunkSize = Math.ceil(totalRecords / processCount); for (let i = 0; i < processCount; i++) { const start = i * chunkSize; const end = Math.min(start + chunkSize - 1, totalRecords - 1); const child = fork('./dist/child-process.js'); child.send({ start, end }); child.on('message', (msg) => { console.log(`子进程${i}完成:处理${msg.count}条记录`); child.kill(); }); child.on('error', (err) => { console.error(`子进程${i}出错:`, err); }); } });
子进程
import { PrismaClient } from '@prisma/client'; import process from 'process'; const prisma = new PrismaClient(); process.on('message', async (data) => { const { start, end } = data as { start: number; end: number }; let count = 0; try { const records = await prisma.someModel.findMany({ where: { id: { gte: start, lte: end } } }); for (const record of records) { await prisma.someModel.update({ where: { id: record.id }, data: { status: 'processed' } }); count++; } process.send({ count }); } catch (err) { process.send({ error: (err as Error).message }); } finally { await prisma.$disconnect(); process.exit(0); } });
方案3:使用cluster模块实现多核并行
cluster模块专为利用多核CPU设计,主进程可将任务分配给多个工作进程,适合CPU密集型的数据库处理场景:
代码示例
import cluster from 'cluster'; import os from 'os'; import cron from 'node-cron'; import { PrismaClient } from '@prisma/client'; const numCPUs = os.cpus().length; const totalRecords = 10000; const chunkSize = Math.ceil(totalRecords / numCPUs); if (cluster.isPrimary) { // CRON任务触发时启动工作进程 cron.schedule('0 1 * * *', () => { console.log(`主进程 ${process.pid} 启动`); for (let i = 0; i < numCPUs; i++) { const worker = cluster.fork(); const start = i * chunkSize; const end = Math.min(start + chunkSize - 1, totalRecords - 1); worker.send({ start, end }); } cluster.on('exit', (worker, code, signal) => { console.log(`工作进程 ${worker.process.pid} 退出`); }); }); } else { const prisma = new PrismaClient(); process.on('message', async (data) => { const { start, end } = data as { start: number; end: number }; let count = 0; try { const records = await prisma.someModel.findMany({ where: { id: { gte: start, lte: end } } }); for (const record of records) { await prisma.someModel.update({ where: { id: record.id }, data: { status: 'processed' } }); count++; } console.log(`工作进程 ${process.pid} 完成:处理${count}条记录`); } catch (err) { console.error(`工作进程 ${process.pid} 出错:`, err); } finally { await prisma.$disconnect(); process.exit(0); } }); }
关键注意事项
- Prisma Client实例化:每个worker/子进程必须单独创建Prisma Client实例,禁止跨进程共享实例
- 数据库连接池:调整Prisma的
connection_limit配置(在prisma/schema.prisma的datasource块中),避免总连接数超过数据库上限 - 错误处理:务必在每个worker/子进程中捕获异常,避免单个进程崩溃导致整个CRON任务失败
内容的提问来源于stack exchange,提问作者Dudu Dudu
相关产品推荐
相关产品推荐

