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

Node.js+ES6+TypeScript数据库并行化操作及替代方案咨询

基于ES6+TypeScript的Node.js数据库操作并行化方案

当然有,针对你的需求,以下是几种可行的替代/优化方案,均支持ES6模块与TypeScript,且能配合Prisma Client完成批量数据库处理:

方案1:修复worker_threads的ES模块支持问题

你之前遇到的worker_threads不支持ES导入的问题,其实在Node.js 12+版本已经得到解决,只需调整配置即可:

必要配置

  1. 在package.json中声明ES模块:
    {
      "type": "module"
    }
    
  2. 调整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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 13:18:15