Node.jsda noldan thread pool yozamiz (worker threads)

Assalamu Alaykum bugun Nodejsda worker threadsdan foydalangan holda thread pool yaratamiz. Thread pool haqida ko'proq ma'lumot olmoqchi bo'lsangiz ushbu maqolani o'qishingiz mumkin.

Thread pool nima ?

Thread pool yaratishimiz uchun nimalar kerak ?

Worker worker.js faylida bu bizadagi vazifalarni bajaruvchi thread, pool manager bu thread pool logikasi implimentatsiyasi bo'ladi u tasklarni ishlayotgan va ishlamayotgan workerlar haqida ma'lumotni oladi va kelayotgan vazifalarni queueda saqlab qoladi va ularni ishga tushuradi. main bu fayl thread poolni ishga tushuradi va vazifalarni berib natijalarni olib sizga qaytaradi.

Worker.js fayl

Worker da biz fibonnaci raqamlarini hisoblaymiz chunki bu CPU intensive task bo'lib bera oladi va oxirida oddiy holatda va threadpooldagi yo'l bilan benchmark qilamiz.

const { parentPort } = require("worker_threads");
parentPort.on("message", ({ taskId, number }) => {
  const result = fib(number);
  parentPort.postMessage({ taskId, result });
});
function fib(n) {
  if (n <= 1) return n;
  let a = 0,
    b = 1;
  for (let i = 2; i <= n; i++) {
    [a, b] = [b, a + b];
  }
  return b;
}

Main.js fayl

Thread pool logikasini yozishdan oldin main.js ni yozib olamiz va pool o'zini qanday tutishini ko'rib olamiz .

const ThreadPool = require('./pool.js');
const os = require('os');
const poolSize = os.cpus().length;
const pool = new ThreadPool('./worker.js', poolSize);
async function main() {
  try {
    const inputs = [40, 41, 42, 43, 44];
    console.log('Fibonacci raqamlarni worker thread pooldan foydalangan holda hisoblash ...');
    const results = await Promise.all(
      inputs.map(n => pool.run(n))
    );
    inputs.forEach((n, index) => {
      console.log(`fib(${n}) = ${results[index]}`);
    });
  } catch (err) {
    console.error('Error:', err);
  } finally {
    if (pool.destroy) {
      await pool.destroy();
    }
  }
}
main();

Pool.js

Endi Threadpool logikasini yaratishni boshlaymiz . Hozir shunchakil class yaratamiz. Class 2 ta parametr qabul qiladi worker file joylashuvi hamda size nechta worker yaratishimiz kerakligi . Buni tepadagi main.jsdan ham ko'rib olishimiz mumkin.

const { Worker } = require('worker_threads');
const path = require('path');
class ThreadPool {
  constructor(workerFile, size) {
    this.workerFile = path.resolve(workerFile);
    this.size = size;
  }
}

Asosiy fieldlarni qo'shib chiqamiz yani worker threadlarni yaratamiz.

const { Worker } = require('worker_threads');
const path = require('path');
class ThreadPool {
  constructor(workerFile, poolSize = 4) {
    this.workerFile = workerFile;
    this.poolSize = poolSize;
    this.workers = [];
    this.idleWorkers = [];
    this._init();
  }
  _init() {
    for (let i = 0; i < this.poolSize; i++) {
      const worker = new Worker(this.workerFile);
      this.workers.push(worker);
      this.idleWorkers.push(worker);
    }
  }
}

Endi workerlardan kelgan ma'lumotlarni handle qilamiz.

_init() {
  for (let i = 0; i < this.poolSize; i++) {
    const worker = new Worker(this.workerFile);
    worker.on('message', ({ taskId, result }) => {
      const callback = this.callbacks.get(taskId);
      if (callback) {
        callback(result);
        this.callbacks.delete(taskId);
      }
      this.idleWorkers.push(worker);
      this._runNextTask();
    });
    worker.on('error', err => {
      console.error('Worker error:', err);
    });
    this.workers.push(worker);
    this.idleWorkers.push(worker);
  }
}

Endi tasklarni ishga tushuradigan va workerlarga beradigan metodni implimentatsiya qilamiz.

_runNextTask() {
  if (this.taskQueue.length === 0) return;
  if (this.idleWorkers.length === 0) return;
  const worker = this.idleWorkers.pop();
  const task = this.taskQueue.shift();
  worker.postMessage(task);
}

Endi tasklarni main threaddan qabul qilib oladigan run metodini yozamiz

run(number) {
  return new Promise(resolve => {
    const taskId = ++this.taskId;
    this.callbacks.set(taskId, resolve);
    this.taskQueue.push({ taskId, number });
    this._runNextTask();
  });
}

Oxirida ammallar bajarilgandan so'ng worker threadlarni o'chirib yuboramiz.

async destroy() {
  await Promise.all(this.workers.map(w => w.terminate()));
}

Bizda quyidagi to'liq kode kelib chiqadi :

const { Worker } = require('worker_threads');
const path = require('path');
class ThreadPool {
  constructor(workerFile, poolSize = 4) {
    this.workerFile = workerFile;
    this.poolSize = poolSize;
    this.workers = [];
    this.idleWorkers = [];
    this.taskQueue = [];
    this.taskId = 0;
    this.callbacks = new Map();
    this._init();
  }
  _init() {
    for (let i = 0; i < this.poolSize; i++) {
      const worker = new Worker(this.workerFile);
      worker.on('message', ({ taskId, result }) => {
        const callback = this.callbacks.get(taskId);
        if (callback) {
          callback(result);
          this.callbacks.delete(taskId);
        }
        this.idleWorkers.push(worker);
        this._runNextTask();
      });
      worker.on('error', err => {
        console.error('Worker error:', err);
      });
      this.workers.push(worker);
      this.idleWorkers.push(worker);
    }
  }
  _runNextTask() {
    if (this.taskQueue.length === 0) return;
    if (this.idleWorkers.length === 0) return;
    const worker = this.idleWorkers.pop();
    const task = this.taskQueue.shift();
    worker.postMessage(task);
  }
  run(number) {
    return new Promise(resolve => {
      const taskId = ++this.taskId;
      this.callbacks.set(taskId, resolve);
      this.taskQueue.push({ taskId, number });
      this._runNextTask();
    });
  }
  async destroy() {
    await Promise.all(this.workers.map(w => w.terminate()));
  }
}
module.exports = ThreadPool;

Natija :

Natija
Natija

Benchmark

Endi oddiy holatda va har bir task uchun alohida worker yaratgan holda , threadpooldan foydalangan holda fibonnachi raqamlarini hisoblaymiz va qanaqa natija bo'lishini ko'ramiz.

single thread oddiy uslub:

function fib(n) {
  const isBigInt = typeof n === "bigint";
  if (!isBigInt && !Number.isInteger(n)) {
    throw new TypeError("n must be an integer");
  }
  if (n <= 1) return n;
  let a = isBigInt ? 0n : 0;
  let b = isBigInt ? 1n : 1;
  for (let i = isBigInt ? 2n : 2; i <= n; i++) {
    [a, b] = [b, a + b];
  }
  return b;
}
const numbers = [
  4448n,
  6440n,
];
console.time("single-thread");
const results = numbers.map((n) => fib(n));
console.timeEnd("single-thread");

har bir task uchun worker yaratib :

const { Worker } = require("worker_threads");
const path = require("path");
function runFib(number) {
  return new Promise((resolve, reject) => {
    const worker = new Worker(path.resolve(__dirname, "worker.js"));
    worker.once("message", (result) => {
      resolve(result);
      worker.terminate();
    });
    worker.once("error", reject);
    worker.postMessage({ number });
  });
}
(async () => {
  const numbers = [
    23340n,
    4441n,
    44442n,
    4443n
  ];
  console.time("one-off-workers");
  const results = await Promise.all(numbers.map((n) => runFib(n)));
  console.timeEnd("one-off-workers");
  console.log(results);
})();
module.exports = runFib;

Note: Benchmark inputni repostoridan topishingiz mumkin.

benchmark natijalari
benchmark natijalari

Xulosa

Ko'rib turganingizdek threadpool hammasidan tezroq chunki u CPUdan to'liq foydalangan holda ishlamoqda, har bir task uchun yangi thread ochgan uslub ham tezroq ammo threadpooldan sekin ,chunki u har bir thread almashishi va yaratib o'chirilishi uchun qo'shimcha vaqt oladi va input hajmi oshgan sari oradagi farq sekinlashadi. Oddiy usul esa ko'rib turganingizdek 4 barobar ko'proq vaqt yo'qotdi chunki u faqatgina bitta threaddan foydalanmoqda.