Перенесення важких обчислень на CPU у Node.js за допомогою worker_threads та Pools
Дізнайтеся, як worker_threads у Node.js підтримують реактивність циклу подій: створення працівників, обмін повідомленнями між запитами та відповідями, передаванні буферів, використання пулів та проблеми, пов’язані з ними.
Node.js обробляє тисячі одночасних з’єднань у монопотоковому циклі подій, оскільки майже вся робота полягає у очікуванні на мережі чи диску. Ця модель припиняє працювати, як тільки запит вимагає реальних обчислень: хешування великого об’єму даних, зміна розміру зображення чи обробка набору даних займають єдиний потік JavaScript, тож усі інші запити змушені чекати. Модуль worker_threads є вбудованим рішенням. Пройшовши цей посібник, ви зможете переміщувати код, що обмежений потужністю CPU, у робочі потоки, ефективно обмінюватися з ними даними, запускати їх у пулі та розпізнавати випадки, коли робочий потік є неправильним інструментом.
Чому існують робочі потоки
worker_threads вперше з’явився у Node.js 10.5.0 під експериментальним флагом та вважається стабільним з версії Node.js 12. Він дозволяє одному процесу Node.js виконувати JavaScript на кількох потоках одночасно.
Він відрізняється від child_process у важливому аспекті. Дочірній процес — це повністю окремий процес Node.js із власною пам’яттю та власною копією середовища виконання. Робітник існує всередині того самого процесу, але це не просто „ще одна нитка, яка ділить усе“. Кожен робітник має власний ізольований модуль V8, власну купу та власний цикл подій. Нічого не ділиться імпліцитно; робітники спілкуються шляхом передачі повідомлень, а пам’ять вони ділять лише тоді, коли ви явно передаєте їм SharedArrayBuffer. (Часто повторюється твердження, що робітники ділять інстанцію V8 та пам’ять основної нитки. Це не є правдивим, і саме ця різниця пояснює більшу частину дизайну API нижче.) Порівняно з процесами, робітників легше запускати та з ними спілкуватися, що робить їх природним способом використання кількох ядер у одній інстанції програми.
Яку ціну вам коштує заблокований цикл подій
Основне завдання працівника — зберігати цикл подій вільним. Поки на основному потоці виконується синхронний обчислення, сервер не може приймати нові запити, не може просуватися у виконанні відкритих з’єднань та навіть не може відповісти на перевірку стану. Практичні наслідки виглядають так:
- Погіршення користувацького досвіду: відповіді API надходять з запізненням, а функції в реальному часі, такі як живі оновлення, працюють з перервами.
- Знижена продуктивність: одне тривале завдання займає єдиний JavaScript-потік, тому кількість запитів за секунду зменшується.
- Нестабільність: тривале блокування може перевищити таймаути балансувальника навантаження чи клієнта, і ці збої можуть поширитися на інші сервіси.
Перенесення важких обчислень на працівника звільняє основний цикл для обслуговування сокетів та таймерів, тож додаток залишається реактивним навіть під час виконання складних обчислень.
Щоб детальніше дізнатися про те, як взаємодіють цикл подій, libuv та їхній пул потоків, перегляньте як працює конкурентність у Node.js зсередини.
Використання кожного ядра
Сервери зазвичай мають багато ядер CPU. Без працівників або без запуску кількох процесів через child_process чи менеджер процесів, такий як PM2, один процес Node.js може використовувати приблизно одне ядро для виконання JavaScript-коду, що обмежений ресурсами CPU. Решта ядер залишаються без дії. Працівники дозволяють одному застосунку виконувати кілька важких завдань паралельно на всіх доступних ядрах, тим самим підвищуючи загальну продуктивність обробки даних.
Завдання, які стають практичними
Паралельні обчислення відкривають можливості у сферах, які традиційно залишалися для мов із першокласним підтримкою потоків:
- Обробка даних: аналітика в реальному часі, маніпуляція зображеннями та відео, а також перетворення великих наборів даних.
- Виконання моделей машинного навчання: запуск навчених моделей під час того, як решта додатку залишається інтерактивною.
- Криптографія: хеширування, шифрування та декодування.
- Симуляція та наукові обчислення: паралельні симуляції та складні математичні моделі.
Ключовий момент полягає у тому, що жоден з цих елементів не позбавляє вас асинхронної моделі введення-виведення. Робочі потоки додають до неї можливості паралельних обчислень.
Компоненти модуля
worker_threads надає невеликий набір примітивів:
Worker: клас, який використовується основним потоком для запуску нового робочого потоку зі скрипту.
isMainThread: булеве значення, яке вказує, чи виконується поточний код на головній смузі, що корисно тоді, коли один файл може виконувати обидві ролі.parentPort: доступний усередині працівника, це канал зворотного зв’язку до смуги, яка його створила. На головній смузі він дорівнює null.workerData: доступний усередині працівника, він містить копію будь-яких даних, які надала батьківська смуга під час створення працівника.MessagePort та MessageChannel: використовуються для створення додаткових незалежних двосторонніх каналів між смугами.SharedArrayBuffer та Atomics: справжня спільна пам’ять разом із примітивами синхронізації, необхідними для її безпечного використання. Це найефективніший варіант, але водночас і найскладніший, оскільки він знову створює умови конкуренції.Перший працівник: пошук простих чисел поза основною ниткою
Гарним першим завданням є навмисно повільна робота – обчислення всіх простих чисел до великого ліміту. У прикладі використовуються два файли, які знаходяться в одній папці.
Скрипт основної нитки
Основний скрипт створює працівника, надає йому вхідні дані та чекає на відповідь. Перш ніж її прочитати, зверніть увагу на три моменти. Працівник створюється з окремого шляху до файлу. Вхідні дані передаються через опцію workerData. А батьківський скрипт підписується на три події: message для отримання результатів, error у разі виникнення винятків, яких працівник не спіймав, та exit коли поток закінчує свою роботу.
// main.js
const { Worker, isMainThread, workerData } = require('worker_threads');
const path = require('path');
if (isMainThread) {
console.log(`Main thread started. PID: ${process.pid}`);
const largeNumber = 20_000_000; // A large number for prime calculation
console.log(`Starting CPU-intensive task (finding primes up to ${largeNumber})...`);
const startTime = Date.now();
const worker = new Worker(path.join(__dirname, 'prime-worker.js'), {
workerData: { limit: largeNumber }
});
worker.on('message', (result) => {
const endTime = Date.now();
console.log(`Worker finished. Found ${result.length} primes.`);
console.log(`Time taken by worker: ${endTime - startTime}ms`);
// console.log('Primes found:', result.slice(0, 10), '...'); // Log first 10 for brevity
// Demonstrate a simple non-blocking task on the main thread
console.log('Main thread continuing with other tasks...');
setTimeout(() => {
console.log('Main thread completed a separate non-blocking task.');
}, 100);
});
worker.on('error', (err) => {
console.error('Worker error:', err);
});
worker.on('exit', (code) => {
if (code !== 0)
console.error(`Worker stopped with exit code ${code}`);
else
console.log('Worker thread exited normally.');
});
} else {
// This part will not be executed in our scenario as worker.js runs separately.
// It's here for illustrative purposes if main.js itself was to be imported as a worker.
console.log('This code runs inside a worker thread if main.js was spawned as one.');
console.log('WorkerData:', workerData);
}
Розглянемо ключові рядки:
- Перевірка
isMainThreadзахищає логіку створення працівника, щоб вона виконувалася лише на основному потоці. Веткаelseіснує виключно для ілюстрації того, що станеться, якщо цей самий файл буде завантажений як працівник; у цій конфігурації вона ніколи не виконується. largeNumberвстановлюється на достатньо високе значення (20 мільйонів), щоб обчислення займало помітний час.
new Worker(path.join(__dirname, 'prime-worker.js'), { workerData: { limit: largeNumber } }) — це основний виклик. Перший аргумент вказує на скрипт, який буде виконуватися новою потоковою одиницею, і це має бути окремий файл. Другий аргумент — це об’єкт параметрів, у якому поле workerData стає початковим даним для робочої одиниці.workerData копіюється за допомогою алгоритму структурованої клонування HTML. Робоча одиниця отримує власну копію, а не посилання на об’єкт батька.message отримує будь-які дані, які надсилає робоча одиниця, у цьому випадку — масив простих чисел. Слухач error активується у разі виникнення неконтрольованих винятків усередині робочої одиниці, а слухач exit отримує код завершення, де 0 означає нормальне завершення.Скрипт робочої одиниці
Робітник містить лише дорогі обчислення та код, який повідомляє про їхній результат. Функція findPrimes — це простий цикл перебору дільників, навмисно неоптимізований, щоб він споживав час процесора.
// prime-worker.js
const { parentPort, workerData } = require('worker_threads');
// A simple (not highly optimized) function to find prime numbers
function findPrimes(limit) {
const primes = [];
for (let i = 2; i <= limit; i++) {
let isPrime = true;
for (let j = 2; j <= Math.sqrt(i); j++) {
if (i % j === 0) {
isPrime = false;
break;
}
}
if (isPrime) {
primes.push(i);
}
}
return primes;
}
// Ensure this code only runs if it's indeed a worker thread being executed
if (parentPort) {
const { limit } = workerData;
console.log(`Worker thread started to find primes up to ${limit}.`);
// Simulate an error for demonstration purposes sometimes
// if (Math.random() < 0.2) {
// throw new Error('Simulated worker error!');
// }
try {
const result = findPrimes(limit);
parentPort.postMessage(result); // Send the result back to the main thread
} catch (error) {
console.error('Error during prime calculation in worker:', error);
parentPort.postMessage({ error: error.message }); // Send error back
}
} else {
console.log("This script is intended to be run as a worker thread.");
}
Що варто зазначити:
- Вона отримує
parentPortтаworkerDataз модуля. - Умова
if (parentPort)гарантує, що логіка виконуватиметься лише тоді, коли файл завантажується як робітник; якщо виконати її безпосередньо за допомогоюnode,parentPortбуде дорівнюватиnull, і тоді буде виведено повідомлення. const { limit } = workerData;зчитує дані, надіслані батьківським процесом.parentPort.postMessage(result)передає прості числа на головну схему, де обробникmessageїх отримує.
try...catch фіксує будь-яку невдачу під час обчислень та повідомляє про неї батьківському процесу, замість того щоб дозволити потоку завершитися.Одна нюанс: у разі невдачі працівник публікує { error: ... } у тому самому каналі message, що використовується для результатів, але основний скрипт розглядає кожне повідомлення як масив простих чисел. У реальному коді необхідно надати повідомленням чітку структуру (наприклад, поля status або type), щоб батьківський процес міг відрізнити результат від помилки.
Запуск та читання результату
Збережіть обидві файли поруч один з одним та запустіть програму за допомогою node main.js. Ви побачите, як основний потік повідомляє про себе, як працівник інформує про початок роботи, а через деякий час — результат працівника разом із кількістю спливлого часу.
Зверніть увагу на порядок розташування рядка „Основний потік продовжує виконання інших завдань...“. У цьому скрипті він виводиться всередині обробника message, тому з’являється лише після завершення роботи працівника. Це демонструє, що основний потік зміг отримати та обробити результат, а не те, що він виконував іншу роботу протягом цього часу. Щоб справді побачити, як основний потік залишається реактивним під час обчислень, запустіть щось у основному потоці до отримання результату, наприклад setInterval, який записує інформацію про стан кожні кілька сотень мілісекунд. Цей сигнал продовжуватиме надходити, поки працівник виконує обчислення, що саме не сталося б, якби findPrimes виконувався у основному потоці.
Двостороння комунікація за допомогою протоколу запит-відповідь
У першому прикладі надсилається один вхідний даний та отримується один результат. У більшості реальних випадків потрібен працівник, який довго залишається активним та обробляє багато запитів. Комунікація є повністю двосторонньою: об’єкт Worker у основній нитці та parentPort всередині працівника обидва надають методи .postMessage(...) та .on('message', ...).
Процес, описаний нижче, залишається активним та чекає на повідомлення. Коли він отримує запит calculateFibonacci, він обчислює значення за допомогою наївної рекурсивної функції (навмисно повільно) та відповідає або повідомленням fibonacciResult, або fibonacciError, вказуючи ідентифікатор запиту, щоб викликаючий міг його ідентифікувати.
// request-response-worker.js
const { parentPort } = require('worker_threads');
function fibonacci(n) {
if (n <= 1) return n;
return fibonacci(n - 1) + fibonacci(n - 2);
}
if (parentPort) {
parentPort.on('message', (message) => {
if (message.type === 'calculateFibonacci') {
const { n, requestId } = message.payload;
console.log(`Worker: Calculating Fibonacci(${n}) for request ID ${requestId}`);
try {
const result = fibonacci(n);
parentPort.postMessage({ type: 'fibonacciResult', requestId, payload: result });
} catch (error) {
console.error(`Worker: Error calculating Fibonacci(${n}):`, error);
parentPort.postMessage({ type: 'fibonacciError', requestId, payload: error.message });
}
}
});
console.log('Worker ready to receive Fibonacci requests.');
}
У основному потоці секрет полягає у перетворенні передачі повідомлень на обіцянки. Кожен виклик генерує унікальний requestId, зберігає функції resolve та reject обіцянки у Map під цим ідентифікатором та надсилає запит. Коли надходить відповідь, обробник message знаходить ідентифікатор, видаляє запис та виконує відповідну обіцянку. У прикладі для генерації ідентифікаторів використовується пакет uuid; у поточних версіях Node.js функція crypto.randomUUID() з вбудованого модуля node:crypto виконує ту саму роботу без додаткових залежностей.
// main-request-response.js
const { Worker, isMainThread } = require('worker_threads');
const path = require('path');
const { v4: uuidv4 } = require('uuid'); // npm install uuid
if (isMainThread) {
const worker = new Worker(path.join(__dirname, 'request-response-worker.js'));
const pendingRequests = new Map();
worker.on('message', (message) => {
const { type, requestId, payload } = message;
if (pendingRequests.has(requestId)) {
const { resolve, reject } = pendingRequests.get(requestId);
pendingRequests.delete(requestId);
if (type === 'fibonacciResult') {
resolve(payload);
} else if (type === 'fibonacciError') {
reject(new Error(payload));
}
}
});
worker.on('error', (err) => {
console.error('Worker error:', err);
});
worker.on('exit', (code) => {
if (code !== 0) console.error(`Worker stopped with exit code ${code}`);
});
async function calculateFibonacciInWorker(n) {
const requestId = uuidv4();
return new Promise((resolve, reject) => {
pendingRequests.set(requestId, { resolve, reject });
worker.postMessage({ type: 'calculateFibonacci', payload: { n }, requestId });
});
}
(async () => {
console.log('Main: Sending Fibonacci requests to worker...');
try {
const fib40 = await calculateFibonacciInWorker(40);
console.log(`Main: Fibonacci(40) = ${fib40}`);
const fib35 = await calculateFibonacciInWorker(35);
console.log(`Main: Fibonacci(35) = ${fib35}`);
// This will block the main thread if done directly, but not here
// console.log(`Main: Local Fibonacci(40) = ${fibonacci(40)}`);
} catch (error) {
console.error('Main: Failed to get Fibonacci result from worker:', error);
} finally {
worker.terminate(); // Terminate worker when done
}
})();
}
Огляд паттерну:
- Основний потік створює
requestIdдля кожного запиту, щоб відповіді можна було прив’язати до їхніх обіцянок, навіть якщо вони надходять у неправильному порядку.
pendingRequests зберігає функції resolve та reject доки робітник не відповість.calculateFibonacci та відповідає з fibonacciResult або fibonacciError, передаючи початковий requestId.Якщо ви адаптуєте цей код, уважно стежте за структурою повідомлення. Основний потік розміщує requestId на найвищому рівні повідомлення поруч із payload, тоді як працівник отримує його з message.payload. У такому вигляді працівник бачить requestId як undefined, тому відповідь не може бути знайдена, а очікувана обіцянка ніколи не буде виконана. Треба зберігати ID у одному узгодженому місці з обох сторін (наприклад, payload: { n, requestId }). Також варто додати таймаут для кожного поточного запиту, щоб втрачена відповідь перетворилася на помилку, а не призвела до зациклення. Нарешті, зверніть увагу на блок finally: як тільки робота завершена, worker.terminate() зупиняє потік, щоб процес міг завершитися.
Переміщення великих даних без їх копіювання
Усе, що надсилається через postMessage, за замовчуванням копіюється у структурованому вигляді, що означає, що отримувач отримує копію. Для невеликих повідомлень це нормально. Однак для великих бінарних даних сама копія може стати причиною уповільнення. Рішенням є передача прав власності на внутрішню пам’ять замість її копіювання: об’єкт стає недоступним для відправника та одразу доступним для отримувача без дублювання.
ArrayBuffer та MessagePort — це найпоширеніші об’єкти, які передаються. SharedArrayBuffer — це окремий випадок: його зовсім не передають, оскільки обидві потоки можуть одночасно доступатися до тієї самої пам’яті.
У наведеному нижче прикладі містяться два скрипти в одному списку. Робочий процес (buffer-worker.js) отримує буфер, подвоює кожен байт без змін та надсилає цей буфер назад у вигляді об’єкта, придатного до передачі. Основний скрипт (main-buffer.js) заповнює буфер об’ємом 1 МБ, передає його робочому процесу та читає оброблені дані після його повернення.
// buffer-worker.js
const { parentPort } = require('worker_threads');
if (parentPort) {
parentPort.on('message', (message) => {
if (message.type === 'processBuffer') {
const { buffer } = message.payload; // This is now the ArrayBuffer
const uint8 = new Uint8Array(buffer);
// Modify the buffer in the worker
for (let i = 0; i < uint8.length; i++) {
uint8[i] = uint8[i] * 2;
}
console.log('Worker: Buffer processed. First 5 elements:', uint8.slice(0, 5));
// Send it back as a transferable
parentPort.postMessage({ type: 'bufferProcessed', payload: buffer }, [buffer]);
}
});
}
// main-buffer.js
const { Worker, isMainThread } = require('worker_threads');
const path = require('path');
if (isMainThread) {
const worker = new Worker(path.join(__dirname, 'buffer-worker.js'));
const bufferSize = 1024 * 1024; // 1MB
let myBuffer = new ArrayBuffer(bufferSize);
let uint8 = new Uint8Array(myBuffer);
// Initialize buffer
for (let i = 0; i < uint8.length; i++) {
uint8[i] = i % 256;
}
console.log('Main: Original buffer (first 5 elements):', uint8.slice(0, 5));
worker.postMessage({ type: 'processBuffer', payload: { buffer: myBuffer } }, [myBuffer]);
// After postMessage with transfer, myBuffer becomes detached/empty in main thread
// Attempting to access it will result in an error or empty view.
// console.log('Main: Buffer after transfer (should be detached):', uint8.slice(0, 5)); // This would likely show zeros or error.
worker.on('message', (message) => {
if (message.type === 'bufferProcessed') {
const receivedBuffer = message.payload;
const receivedUint8 = new Uint8Array(receivedBuffer);
console.log('Main: Received processed buffer (first 5 elements):', receivedUint8.slice(0, 5));
worker.terminate();
}
});
}
Ключові деталі:
- Другий аргумент функції
postMessage– тут[myBuffer]з одного боку та[buffer]з іншого – є списком для передачі. - Коли надсилається
myBuffer, його власність переходить до робочого процесу. У основному потоці буфер втрачає зв’язок: його значенняbyteLengthстає 0, а будь-який погляд на нього у вигляді типового масиву, наприкладuint8, також має довжину 0.
Компроміс полягає у тому, що після передачі необхідно припинити використання початкового посилання. Якщо обидві потокові лінії справді потребують даних одночасно, потрібна копія або SharedArrayBuffer.
Обробка помилок та очищення
Робітники можуть зламатися, як і будь-який інший код, і поки вони існують, вони займають пам’ять та потокову лінію. Доступні інструменти:
worker.on('error', handler)у основній потоковій лінії: отримує винятки, які були кинуті у робітнику та так і не спіймані.worker.on('exit', handler)у основному потоці: викликається щоразу, коли працівник зупиняється, незалежно від того, завершив він роботу нормально чи зламався. Код вихіду дозволяє розрізнити ці два випадки, причому0означає успіх.process.on('uncaughtException', handler)всередині працівника: надає можливість записувати журнал або повідомляти деталі безпосередньо з працівника перед його зупинкою, окрім того, що бачить батьківський процес.worker.terminate(): зупиняє працівника якомога швидше та повертає обіцянку, яка реалізується з кодом вихіду. Він не чекає на завершення поточних операцій, тому якщо вам потрібне грайсфулне завершення, надішліть працівнику повідомлення „stop“ та дозвольте йому вийти самостійно. Завжди переконуйтесь, що працівники, які більше не потрібні, зупиняються будь-яким способом.
Виконання багатьох завдань через пул працівників
Створення нового працівника для кожного завдання є марнотратством, а запуск більшої кількості працівників, які сильно навантажують процесор, ніж у вас є ядер, лише посилює конкуренцію. Пул працівників вирішує обидві проблеми: він запускає фіксовану кількість працівників, поміщає надходження завдань у чергу та передає кожне завдання наступному вільному працівнику.
У спрощеному пулі нижче зберігається масив записів працівників, список вільних ідентифікаторів працівників та черга завдань, що очікують обробки. Функція runTask повертає обіцянку та додає завдання до черги; функція processQueue поєднує найстаріше завдання з вільним працівником; коли працівник надсилає повідомлення про завершення роботи, його обіцянка виконується та він знову потрапляє до списку вільних працівників. Якщо працівник стикається з помилкою або завершує роботу ненормально, функція terminateWorker замінює його на нового, щоб пул залишався повнорозмірним.
// workerPool.js
const { Worker } = require('worker_threads');
const path = require('path');
class WorkerPool {
constructor(workerPath, numWorkers) {
this.workerPath = workerPath;
this.numWorkers = numWorkers;
this.workers = [];
this.freeWorkers = [];
this.queue = [];
this.initWorkers();
}
initWorkers() {
for (let i = 0; i < this.numWorkers; i++) {
const worker = new Worker(this.workerPath);
worker.id = i; // Assign an ID for easier debugging
worker.on('message', (message) => {
// Resolve the promise associated with this worker's task
this.workers[worker.id].resolve(message.payload);
this.returnWorkerToPool(worker.id);
this.processQueue();
});
worker.on('error', (err) => {
this.workers[worker.id].reject(err);
console.error(`Worker ${worker.id} error:`, err);
this.terminateWorker(worker.id); // Re-initialize or handle as appropriate
});
worker.on('exit', (code) => {
if (code !== 0) console.error(`Worker ${worker.id} exited with code ${code}`);
this.terminateWorker(worker.id); // Handle crashed worker
});
this.freeWorkers.push(worker.id);
this.workers.push({ instance: worker, busy: false, resolve: null, reject: null });
}
console.log(`Worker Pool initialized with ${this.numWorkers} workers.`);
}
runTask(taskData) {
return new Promise((resolve, reject) => {
this.queue.push({ taskData, resolve, reject });
this.processQueue();
});
}
processQueue() {
if (this.queue.length === 0 || this.freeWorkers.length === 0) {
return;
}
const workerId = this.freeWorkers.shift();
const workerInfo = this.workers[workerId];
workerInfo.busy = true;
const { taskData, resolve, reject } = this.queue.shift();
workerInfo.resolve = resolve;
workerInfo.reject = reject;
workerInfo.instance.postMessage(taskData);
}
returnWorkerToPool(workerId) {
const workerInfo = this.workers[workerId];
workerInfo.busy = false;
workerInfo.resolve = null;
workerInfo.reject = null;
this.freeWorkers.push(workerId);
}
terminateWorker(workerId) {
const workerInfo = this.workers[workerId];
if (workerInfo.instance) {
workerInfo.instance.terminate();
}
// Remove from current workers list and potentially replace
this.workers[workerId] = { instance: null, busy: false, resolve: null, reject: null };
// Optionally, re-create the worker to maintain pool size
console.log(`Worker ${workerId} terminated. Re-initializing...`);
const newWorker = new Worker(this.workerPath);
newWorker.id = workerId;
newWorker.on('message', (message) => {
this.workers[newWorker.id].resolve(message.payload);
this.returnWorkerToPool(newWorker.id);
this.processQueue();
});
newWorker.on('error', (err) => {
this.workers[newWorker.id].reject(err);
console.error(`Worker ${newWorker.id} error:`, err);
this.terminateWorker(newWorker.id);
});
newWorker.on('exit', (code) => {
if (code !== 0) console.error(`Worker ${newWorker.id} exited with code ${code}`);
this.terminateWorker(newWorker.id);
});
this.workers[workerId] = { instance: newWorker, busy: false, resolve: null, reject: null };
this.freeWorkers.push(workerId); // Add the new worker to the pool
}
close() {
for (const workerInfo of this.workers) {
if (workerInfo.instance) {
workerInfo.instance.terminate();
}
}
console.log('Worker Pool closed.');
}
}
module.exports = WorkerPool;
Генеричний працівник, який використовується пулом, виконує функцію, залежну від потужності CPU, для кожного отриманого повідомлення та відповідає з статусом та даними навантаження. Цикл підсумування імітує будь-яку реальну роботу, таку як обробка зображень чи шифрування.
// pool-worker.js
const { parentPort } = require('worker_threads');
function performHeavyCalculation(data) {
// Example heavy calculation: sum of numbers up to 'limit'
// This could be anything CPU-bound: image processing, encryption, etc.
let sum = 0;
for (let i = 0; i <= data.limit; i++) {
sum += i;
}
return sum;
}
if (parentPort) {
parentPort.on('message', (taskData) => {
try {
const result = performHeavyCalculation(taskData);
parentPort.postMessage({ status: 'success', payload: result });
} catch (error) {
parentPort.postMessage({ status: 'error', payload: error.message });
}
});
}
Нарешті, основний скрипт налаштовує розмір пулу відповідно до кількості ядер CPU, повідомленої функцією os.cpus(), надсилає шість завдань та чекає на їх виконання за допомогою Promise.allSettled. Кожна обіцянка завдання перетворюється на читабельний рядок про успіх або невдачу.
// main-pool.js
const WorkerPool = require('./workerPool');
const path = require('path');
const numCores = require('os').cpus().length;
const pool = new WorkerPool(path.join(__dirname, 'pool-worker.js'), numCores);
async function runExample() {
const tasks = [
{ limit: 1_000_000_000 },
{ limit: 500_000_000 },
{ limit: 1_500_000_000 },
{ limit: 750_000_000 },
{ limit: 2_000_000_000 },
{ limit: 250_000_000 }
];
console.log('Main: Submitting tasks to the worker pool...');
const results = await Promise.allSettled(tasks.map((task, index) =>
pool.runTask(task)
.then(res => `Task ${index} completed with result: ${res}`)
.catch(err => `Task ${index} failed: ${err.message}`)
));
results.forEach(res => console.log(res.value));
// Demonstrate main thread responsiveness
console.log('Main: Tasks submitted. Doing other stuff...');
await new Promise(resolve => setTimeout(resolve, 100)); // Simulate async work
console.log('Main: Other stuff done.');
pool.close();
console.log('Main: Worker pool closed.');
}
runExample();
Розглядайте цей пул суто як навчальний приклад; у ньому є кілька серйозних недоліків, які мають бути виправлені перед тим, як щось подібне потрапить у продакшн.
- Обробник
messageбасейну виконує обіцянку зpayload, незалежно відstatus, тому завдання, яке зазнало невдачі всередині працівника, все одно вважається успішним. Перевіртеstatusта відхиліть його, якщо значення дорівнює'error'. worker.terminate()змушує працівника завершити роботу з кодом, відмінним від нуля. Оскільки обробникexitвикликаєterminateWorker, який створює нового працівника, викликclose()спричиняє цикл створення нових працівників замість закриття басейну, а помилка разом із подієюexitможе двічі заповнити ту саму позицію та додати її ID до списку вільних місць більше одного разу. Справжній басейн потребує флага „закриття“ та має замінювати лише працівників, які несподівано зламалися.
error; ненормальний вихід без події error залишає його обіцянку у стані очікування.Number.MAX_SAFE_INTEGER, тому надруковані результати втрачають точність. Використовуйте BigInt, якщо важливі точні значення.Більш повні пули також додають таймаути для неактивних елементів, динамічне масштабування та більш ретельну відновлювальну функцію. Добре підтримувані бібліотеки, такі як Piscina, реалізують ці деталі, і вони зазвичай є кращим вибором, ніж самостійно створений пул.
Коли робочі потоки допомагають, а коли — ні
Підходящі випадки
- Завдання, що залежать від CPU: все, що споживає значну кількість часу процесора, таке як складні обчислення, стиснення, шифрування чи обробка зображень.
- Захист циклу подій: будь-яка операція, яка інакше могла б заблокувати основний потік на довший час, ніж це допустимо.
Невідповідні варіанти
- Робота, що залежить від вводу-виводу: мережеві запити, запити до баз даних та операції файлової системи в Node.js вже є неблокуючими. Обгортання їх у робітника створює додаткове навантаження без жодної користі.
- Дрібні завдання: запуск робітника та передача повідомлень — обидва процеси займають час. Для швидких обчислень це додаткове навантаження може перевищити отриману вигоду, що є ще однією причиною тримати довгоживучих робітників у пулі.
SharedArrayBuffer дозволяє це зробити, але координація записів за допомогою Atomics є надзвичайно складною для правильного виконання. Якщо ви справді цього не потребуєте та не розумієте наслідків, краще використовувати передачу повідомлень.Практики, які підтримують ефективність працівників
- Зробіть скрипти працівників мінімальними. Завантажуйте лише ту логіку та залежності, які потрібні завданню. Імпортування всього фреймворку вашого додатку до кожного працівника подовжує час запуску та споживає більше пам’яті.
- Оберіть найдешевший шлях передачі даних. Структуроване клонування через
postMessageпідходить для невеликих повідомлень; великіArrayBufferслід передавати через список передач, а глибокі чи складні дерева об’єктів варто уникати, оскільки їх клонування є повільним.
error та exit у головному потоці, а також обгортайте роботу всередині працівника за допомогою try...catch, щоб помилки поверталися у вигляді чітких повідомлень про помилку.os.cpus().length — це розумний початковий показник для пулу працівників; у новіших версіях Node.js рекомендується використовувати os.availableParallelism() для отримання цього значення.message виконує тривалу синхронну роботу, наступні повідомлення накопичуються позаду неї. Працівник, який мусить обробляти багато запитів, повинен робити кожну одиницю роботи короткою або, у рідкісних випадках, делегувати її своїм підпрацівникам.Поширені помилки, які заважають
- Очікування спільних глобальних змінних. Дані, передані через
new Worker(), потрапляють до працівника лише черезworkerData(або подальші повідомлення). Змінні рівня модуля у основному потоці не є видимими всередині працівника, оскільки той працює у своєму власному ізольованому середовищі. Єдиний справжній спосіб обміну даними — черезSharedArrayBuffer.
isMainThread. Цей параметр описує контекст виконання коду, який зараз працює. Один файл може використовуватися як основний скрипт, так і як робочий процес, але окремі файли для основного та робочого коду зазвичай є зрозумілішими.--inspect-brk, і потрібно перемикатися між контекстами потоків у налагоджувачі. Редактори на кшталт VS Code багато в чому автоматизують цей процес, проте це все одно менш зручно, ніж налагоджування основного потоку.Основні висновки
- Робітники виконують JavaScript паралельно у окремих ізольованих середовищах всередині одного процесу; вони спілкуються шляхом передачі повідомлень та ділять пам’ять лише через
SharedArrayBuffer. - Використовуйте їх для завдань, які інтенсивно використовують CPU та можуть заблокувати цикл подій, а не для мережевого чи дискового доступу, які Node.js вже обробляє асинхронно.
- Надавайте повідомленням чітку, послідовну структуру, щоб результати, помилки та ідентифікатори запитів були однозначними з обох сторін.
- Передавайте великі об’єкти
ArrayBufferзамість їх клонування та пам’ятайте, що відправник втрачає до них доступ після цього. - Повторно використовуйте робітників через пул, розмір якого відповідає кількості ядер машини, та переконайтеся, що процеси вимкнення та відновлення після збою чітко розділені, або використовуйте перевірену бібліотеку пулів.
Пов’язана література
- Контроль конкуруентності в Node.js: як уникнути збоїв API за допомогою p-map та Bottleneck — Дізнайтеся, як поєднання p-map та Bottleneck у Node.js допомагає запобігати помилкам обмеження частоти запитів та перевантаженню системи шляхом керування конкуруентністю та часом обробки запитів.
- Пояснення конкуруентності в Node.js: libuv, цикл подій та пул потоків — Дізнайтеся, як Node.js використовує примітиви операційної системи libuv та пул робочих потоків для обробки асинхронних операцій вводу-виводу, а також про поширені проблеми пулу потоків та поради щодо його налаштування.