Przenoszenie obciążeń wymagających silnego procesora w Node.js za pomocą worker_threads i Pools
Dowiedz się, w jaki sposób worker_threads w Node.js utrzymują pętlę zdarzeń w stanie reaktywnym: tworzenie pracowników, komunikacja typu żądanie-odpowiedź, buforzy przenoszone między procesami, zbiory zasobów oraz ich pułapki.
Node.js obsługuje tysiące jednoczesnych połączeń w ramach pojedynczego wątku pętli zdarzeń, ponieważ prawie cała praca polega na oczekiwaniu na odpowiedź z sieci lub dysku. Ten model przestaje być skuteczny, gdy żądanie wymaga rzeczywistych obliczeń: haszowanie dużych danych, zmiana rozdzielczości obrazu lub przetwarzanie zbiorów danych zajmuje dostępny wątek JavaScript, przez co wszystkie pozostałe żądania muszą czekać. Moduł worker_threads stanowi wbudowane rozwiązanie tego problemu. Po przeczytaniu tego przewodnika będziesz mógł przenosić kod wymagający intensywnego wykorzystania CPU do workerów, efektywnie wymieniać z nimi dane, uruchamiać je w grupach oraz rozpoznawać sytuacje, gdy worker nie jest odpowiednim narzędziem.
Dlaczego istnieją wątki robocze
worker_threads po raz pierwszy pojawił się w Node.js 10.5.0 pod flagą eksperymentalną i od wersji Node.js 12 jest uważany za stabilny. Umożliwia on pojedynczemu procesowi Node.js wykonywanie JavaScripta na kilku wątkach jednocześnie.
Różni się od child_process w istotny sposób. Proces potomny to zupełnie oddzielny proces Node.js z własną pamięcią oraz własną kopią środowiska wykonawczego. Pracownik funkcjonuje wewnątrz tego samego procesu, ale nie jest po prostu „inną nitką dzielącą się wszystkim”. Każdy pracownik ma swój własny izolat V8, własną stertę pamięci oraz własny pętlę zdarzeń. Nic nie jest dzielone w sposób domyślny; pracownicy komunikują się poprzez przekazywanie wiadomości, a pamięć dzielą tylko wtedy, gdy wyraźnie przekazujesz im SharedArrayBuffer. (Często powtarzanym twierdzeniem jest, że pracownicy dzielą instancję V8 oraz pamięć głównej nitki. To nie jest prawdą, a ta różnica tłumaczy większość aspektów projektu API.) W porównaniu z procesami, pracownicy są tańsi w uruchamianiu i komunikowaniu się z nimi, co czyni je naturalnym sposobem wykorzystania kilku rdzeni w ramach jednej instancji aplikacji.
Jakie koszty niesie z sobą zablokowana pętla zdarzeń
Głównym zadaniem workera jest utrzymanie pętli zdarzeń wolnej. Gdy na wątku głównym wykonywane jest obliczenie synchroniczne, serwer nie może przyjmować nowych żądań, nie może robić postępów w otwartych połączeniach, a nawet nie może odpowiedzieć na sprawdzenie stanu. Praktyczne konsekwencje wyglądają następująco:
- Z pogorszoną jakością doświadczenie użytkownika: odpowiedzi API przychodzą z opóźnieniem, a funkcje w czasie rzeczywistym, takie jak aktualizacje na żywo, działają nierównomiernie.
- Niższa przepustowość: jedno długotrwałe zadanie zajmuje jedyne wątki JavaScript, przez co liczba żądań na sekundę spada.
- Niestabilność: długotrwałe blokowanie może przekroczyć limity czasowe balansera obciążeń lub klienta, a te awarie mogą przenieść się na inne usługi.
Pеренiesienie intensywnych obliczeń do workera pozwala pętli głównej swobodnie obsługiwać sockety i timerzy, dzięki czemu aplikacja pozostaje responsywna nawet podczas wykonywania złożonych obliczeń.
Aby lepiej zrozumieć, jak współpracują pętla wydarzeń, libuv oraz ich zbiór wątków, sprawdź jak w rzeczywistości działa równoczesność w Node.js.
Wykorzystanie wszystkich rdzeni
Serwery zazwyczaj posiadają wiele rdzeni CPU. Bez workerów lub bez uruchamiania kilku procesów za pomocą child_process lub menedżera procesów takiego jak PM2, pojedynczy proces Node.js może wykorzystać mniej więcej jeden rdzeń do obsługi zadań JavaScript wymagających dużych obliczeń. Pozostałe rdzenie pozostają nieużywane. Workerzy umożliwiają jednemu aplikacji równoległe wykonywanie kilku złożonych zadań na dostępnych rdzeniach, co zwiększa ogólną wydajność obliczeniową.
Zadania, które stają się wykonalne
Obliczenia równoległe otwierają możliwości w dziedzinach, które tradycyjnie były przeznaczone dla języków z wątkami pierwszej klasy:
- Przetwarzanie danych: analiza w czasie rzeczywistym, modyfikacja obrazów i nagrań wideo oraz przekształcanie dużych zbiorów danych.
- Inferencja w uczeniu maszynowym: wykonywanie wyszkolonych modeli przy jednoczesnym zachowaniu interaktywności reszty aplikacji.
- Kryptografia: hashing, szyfrowanie i deszyfrowanie.
- Symulacje i obliczenia naukowe: symulacje równoległe oraz złożone modele matematyczne.
Kluczowym punktem jest to, że nic z tego nie kosztuje cię modelu I/O asynchronicznego. Procesory robocze dodają do niego możliwości obliczeń równoległych.
Bloki budulcowe modułu
worker_threads udostępnia niewielki zestaw prostych elementów:
Worker: klasa, której używa wątek główny do uruchomienia nowego procesora roboczego z skryptu.
isMainThread: wartość logiczna wskazująca, czy bieżący kod jest wykonywany na głównym wątku; przydatna, gdy jeden plik może pełnić obie role.parentPort: dostępna wewnątrz pracownika, stanowi kanał powrotny do wątku, który go stworzył. Na głównym wątku ma wartość null.workerData: dostępna wewnątrz pracownika, przechowuje klon wszystkich danych dostarczonych przez rodzica podczas tworzenia pracownika.MessagePort i MessageChannel: służą do tworzenia dodatkowych, niezależnych kanałów dwukierunkowych pomiędzy wątkami.SharedArrayBuffer i Atomics: prawdziwa pamięć współdzielona wraz z prostymi mechanizmami synchronizacji niezbędnymi do jej bezpiecznego używania. Są to najskuteczniejsze rozwiązanie, ale jednocześnie najbardziej złożone, ponieważ ponownie wprowadzają warunki wyścigu.Pierwszy pracownik: znajdowanie liczb pierwszych poza głównym wątkiem
Dobrym pierwszym ćwiczeniem jest celowo wolne zadanie: obliczanie wszystkich liczb pierwszych do dużego limitu. Przykład wykorzystuje dwa pliki znajdujące się w tej samej folderze.
Skrypt głównego wątku
Główny skrypt tworzy pracownika, przekazuje mu dane wejściowe i czeka na odpowiedź. Zanim je przeczytamy, zwróćmy uwagę na trzy rzeczy. Pracownik jest tworzony z oddzielnego ścieżki pliku. Dane wejściowe są przekazywane przez opcję workerData. A rodzicowy skrypt jest podpisany na trzy zdarzenia: message dla wyników, error dla wyjątków, których pracownik nie złapał, oraz exit w momencie zatrzymania wątku.
// 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);
}
Przejdźmy przez kluczowe linie:
- Sprawdzenie
isMainThreadchroni logikę tworzenia pracownika, dzięki czemu działa ona tylko na wątku głównym. Gałąźelseistnieje wyłącznie po to, by pokazać, co by się stało, gdyby ten sam plik został załadowany jako pracownik; w tym ustawieniu nigdy się nie wykonuje. largeNumberjest ustawiony na dość wysoką wartość (20 milionów), dzięki czemu obliczenia zajmują znaczną ilość czasu.
new Worker(path.join(__dirname, 'prime-worker.js'), { workerData: { limit: largeNumber } }) to kluczowa funkcja wywoławcza. Pierwszy argument wskazuje na skrypt, który ma być wykonywany przez nowy wątek – musi to być oddzielny plik. Drugi argument to obiekt opcji, którego pole workerData stanowi początkowe dane dostarczane pracownikowi.workerData jest kopiowany za pomocą algorytmu strukturalnego klonowania HTML. Pracownik otrzymuje własną kopię, a nie odwołanie do obiektu nadrzędnego.message otrzymuje wszystko, co przesyła pracownik – w tym przypadku tablicę liczb pierwszych. Słuchacz error aktywuje się w przypadku nieprzechwyconych wyjątków wewnątrz pracownika, natomiast słuchacz exit otrzymuje kod zakończenia, przy czym 0 oznacza normalne zakończenie.Skrypt pracownika
Pracownik zawiera jedynie kosztowne obliczenia oraz kod, który raportuje jego wynik. Funkcja findPrimes to prosty pętla dzielenia próbnego, celowo nieoptymalizowana, aby zużywać czas procesora.
// 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.");
}
Na co zwrócić uwagę:
- Pobiera
parentPortiworkerDataz modułu. - Warunek
if (parentPort)zapewnia, że logika zostanie wykonywana tylko wtedy, gdy plik jest ładowany jako pracownik; jeśli uruchomisz ją bezpośrednio za pomocąnode,parentPortbędzie miało wartośćnull, wtedy wyświetlana jest odpowiednia wiadomość. const { limit } = workerData;odczytuje dane przesłane przez rodzica.parentPort.postMessage(result)przekazuje liczby pierwsze do głównej wątku, gdzie obsługamessageje odbiera.
try...catch rejestruje każdą niepowodzenie podczas obliczeń i zgłasza je do procesu nadrzędnego, zamiast pozwolić wątkowi zostać zakończonemu.Jedna subtelność: w przypadku niepowodzenia wątek roboczy wysyła { error: ... } na ten sam kanał message, który jest używany do przekazywania wyników, ale główny skrypt traktuje każdą wiadomość jako tablicę liczb pierwszych. W rzeczywistym kodzie należy nadać wiadomościom wyraźną strukturę (na przykład pole status lub type), aby proces nadrzędny mógł odróżnić wynik od błędu.
Rozpoczęcie działania i odczyt wyników
Zapisz oba pliki obok siebie i uruchom program za pomocą node main.js. Zobaczysz, jak główny wątek się zgłasza, a wątek roboczy poinformuje o swoim uruchomieniu; po chwili pojawi się wynik wątka roboczego wraz z czasem trwania obliczeń.
Zwróć uwagę na kolejność wyświetlania linijki „Główny wątek kontynuuje wykonywanie innych zadań...”. W tym skrypcie jest ona wyświetlana wewnątrz obsługiwanego elementu message, więc pojawia się dopiero po zakończeniu pracy workera. Pokazuje to, że główny wątek był w stanie otrzymać i przetworzyć wynik, a nie że wykonywał w tym czasie inne zadania. Aby rzeczywiście zobaczyć, jak główny wątek pozostaje responsywny podczas obliczeń, uruchom coś na głównym wątku przed przybyciem wyniku, np. setInterval, który zapisuje informację o aktywności co kilkaset milisekund. Ten sygnał aktywności nadal jest wysyłany, podczas gdy worker wykonywa obliczenia, co dokładnie nie nastąpiłoby, gdyby findPrimes działało na głównym wątku.
Komunikacja dwukierunkowa przy użyciu protokołu żądanie-odpowiedź
Pierwszy przykład wysyła jeden wpis i otrzymuje jeden wynik. Większość rzeczywistych zastosowań wymaga długotrwałego pracownika, który obsługuje wiele żądań. Komunikacja jest w pełni dwukierunkowa: obiekt Worker na głównej nitce oraz parentPort wewnątrz pracownika dostarczają zarówno metodę .postMessage(...), jak i .on('message', ...).
Pracownik poniżej pozostaje aktywny i czeka na wiadomości. Gdy otrzymuje żądanie calculateFibonacci, oblicza wartość za pomocą prostej funkcji rekurencyjnej (celowo wolno) i odpowiada albo wiadomością fibonacciResult, albo fibonacciError, powtarzając identyfikator żądania, aby osoba wysyłająca mogła go dopasować do odpowiedzi.
// 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.');
}
W głównym wątku chodzi o przekształcenie przekazywania wiadomości w obietnice. Każde wezwanie generuje unikalny requestId, przechowuje funkcje resolve i reject obietnicy w strukturze Map pod tym identyfikatorem, a następnie wysyła żądanie. Gdy przychodzi odpowiedź, obsługa message wyszukuje ten identyfikator, usuwa odpowiedni wpis i realizuje odpowiadającą mu obietnicę. Przykład wykorzystuje pakiet uuid do generowania identyfikatorów; w aktualnych wersjach Node.js funkcja crypto.randomUUID() z wbudowanego modułu node:crypto pełni tę samą funkcję bez konieczności dodatkowych zależności.
// 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
}
})();
}
Podsumowanie wzorca:
- Główny wątek tworzy
requestIddla każdego żądania, dzięki czemu odpowiedzi mogą być powiązane z odpowiadającymi im obietnicami, nawet jeśli przychodzą w nieskoordynowanej kolejności.
pendingRequests przechowuje funkcje resolve i reject aż do momentu, gdy robot odpowie.calculateFibonacci i odpowiada wartością fibonacciResult lub fibonacciError, przenosząc oryginalny requestId.Uważnie obserwuj strukturę wiadomości, jeśli adaptujesz ten kod. Główny wątek umieszcza requestId na najwyższym poziomie wiadomości, obok payload, podczas gdy wątek roboczy dekonstruuje go z message.payload. W takiej formie wątek roboczy widzi requestId jako undefined, w związku z czym nie można dopasować odpowiedzi, a oczekiwana obietnica nigdy się nie realizuje. Zachowaj ID w jednym ustalonym miejscu po obu stronach (na przykład payload: { n, requestId }). Warto również dodać czas wygaśnięcia dla każdego oczekującego żądania, aby utracona odpowiedź zamieniła się w błąd zamiast powodować zatrzymanie procesu. Na koniec zwróć uwagę na blok finally: po zakończeniu pracy worker.terminate() zatrzymuje wątek, dzięki czemu proces może zostać zakończony.
Przenoszenie dużych danych bez ich kopiowania
Wszystko, co jest wysyłane za pomocą postMessage, domyślnie jest klonowane w formie strukturalnej, co oznacza, że odbiorca otrzymuje kopię. W przypadku małych wiadomości jest to w porządku. Jednak przy dużych danych binarnych sama kopia może stać się wąskim gardłem. Rozwiązaniem jest przeniesienie własności pamięci, zamiast jej kopiowania: obiekt staje się niewykorzystywalny dla nadawcy i natychmiast dostępny dla odbiorcy, bez żadnej duplikacji.
ArrayBuffer i MessagePort to najczęściej przekazywane obiekty. SharedArrayBuffer to odrębny przypadek: w ogóle nie jest przenoszony, ponieważ oba wątki mogą jednocześnie dostępować do tej samej pamięci.
Przykład poniżej zawiera dwa skrypty w jednej liście. Procesor (buffer-worker.js) otrzymuje bufor, podwaja każdy bajt bezpośrednio w tym buforze i wysyła go z powrotem jako obiekt przenoszalny. Główny skrypt (main-buffer.js) wypełnia bufor o pojemności 1 MB, przekazuje go do procesora i odczytuje przetworzone dane po jego powrocie.
// 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();
}
});
}
Kluczowe szczegóły:
- Drugi argument funkcji
postMessage, tutaj[myBuffer]z jednej strony i[buffer]z drugiej, to lista do przeniesienia. - Gdy wysyła się
myBuffer, jego własność przechodzi na procesora. W głównym wątku bufor staje się odłączony: jego wartośćbyteLengthspada do 0, a wszelkie widoki tablic typutyped-arrayodnoszące się do niego, takie jak wcześniejszyuint8, również mają długość 0.
Kompromis polega na tym, że po przekazaniu musisz przestać używać oryginalnego odniesienia. Jeśli oba wątki rzeczywiście potrzebują danych jednocześnie, potrzebujesz kopii lub SharedArrayBuffer.
Rozwiązywanie problemów i oczyszczanie
Pracownicy zawodzą tak jak każdy inny kod i podczas swojego działania zajmują pamięć oraz wątek. Dostępne narzędzia to:
worker.on('error', handler)w głównym wątku: odbiera wyjątki, które zostały rzucone w pracowniku i nigdy nie złapane.worker.on('exit', handler)w głównym wątku: uruchamia się za każdym razem, gdy pracownik przestaje działać, niezależnie od tego, czy zakończył pracę normalnie, czy uległ awarii. Kod wyjścia odróżnia te dwa przypadki, przy czym0oznacza sukces.process.on('uncaughtException', handler)wewnątrz pracownika: umożliwia zapisywanie lub raportowanie szczegółów bezpośrednio wewnątrz pracownika przed jego zatrzymaniem, dodatkowo do tego, co widzi proces nadrzędny.worker.terminate(): zatrzymuje pracownika tak szybko, jak to możliwe, i zwraca obietnicę, która zostanie rozwiązana z kodem wyjścia. Nie czeka na zakończenie bieżących zadań, więc jeśli potrzebujesz uporządkowanego wyłączenia, wysyłaj do pracownika wiadomość „stop” i pozwól mu samodzielnie się zakończyć. Zawsze upewnij się, że pracownicy, których już nie potrzebujesz, są zatrzymywani w jakiś sposób.
Wykonywanie wielu zadań za pomocą puli pracowników
Tworzenie nowego pracownika dla każdego zadań jest marnotrawstwem, a uruchamianie większej liczby pracowników wymagających dużo zasobów CPU, niż masz rdzeni, tylko zwiększa konflikty. Zbiór pracowników rozwiązuje oba te problemy: uruchamia określoną liczbę pracowników, umieszcza przychodzące zadania w kolejce i przekazuje każde zadanie następnemu wolnemu pracownikowi.
Społeczność uproszczona poniżej przechowuje tablicę rekordów pracowników, listę wolnych identyfikatorów pracowników oraz kolejkę z zadaniami w oczekiwaniu. Funkcja runTask zwraca obietnicę i umieszcza zadanie w kolejce; funkcja processQueue łączy najstarsze zadanie z wolnym pracownikiem; gdy pracownik zgłosi wynik, jego obietnica zostaje spełniona i wraca on na listę wolnych pracowników. Jeśli pracownik napotka błąd lub zakończy działanie w sposób nieprawidłowy, funkcja terminateWorker zastępuje go nowym, aby utrzymać pełną liczbę pracowników w zbiórce.
// 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;
Używany przez ten zbiór ogólny pracownik wykonywa funkcję obciążającą procesor dla każdej otrzymanej wiadomości i odpowiada statusem oraz payloadem. Pętla sumowania symuluje rzeczywiste obciążenie, takie jak przetwarzanie obrazów lub szyfrowanie.
// 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 });
}
});
}
Na koniec główny skrypt dostosowuje rozmiar zbiórku do liczby rdzeni procesora zgłaszanej przez os.cpus(), wysyła sześć zadań i czeka na ich ukończenie za pomocą Promise.allSettled. Obietnica każdego zadania jest mapowana na czytelny tekst o wyniku pomyślnym lub nieudanym.
// 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();
Traktuj ten zbiór wyłącznie jako przykład edukacyjny; ma kilka istotnych wad, które muszą zostać naprawione, zanim coś podobnego trafi do produkcji:
- Obsługa
messagew basenie rozwiązuje obietnicę zpayload, niezależnie odstatus, więc zadanie, które zawiodło wewnątrz pracownika, nadal jest uznawane za udane. Sprawdźstatusi odrzuć go przy wartości'error'. worker.terminate()powoduje wyjście pracownika z kodem innym niż zero. Ponieważ obsługaexitwywołujeterminateWorker, który tworzy zastępcę, wywołanieclose()uruchamia cykl nowych pracowników zamiast wyłączyć cały basen, a błąd w połączeniu z wydarzeniemexitmoże dwukrotnie zastąpić ten sam slot i dwa razy dodać jego ID do listy wolnych pozycji. Prawdziwy basen potrzebuje flagi „zamykania” i powinien zastępować tylko te pracowniki, które uległy nieoczekiwanemu awarii.
error; nienormalne zakończenie bez wydarzenia error pozostawia jego obietnicę w stanie oczekiwania.Number.MAX_SAFE_INTEGER, więc wyświetlane wyniki tracą precyzję. Użyj BigInt, jeśli ważne są dokładne wartości.Bardziej kompleksowe puli dodają również czas wyłączenia w stanie bezczynności, dynamiczne skalowanie oraz bardziej skrupulatną odbudowę po awariach. Dobrze utrzymywane biblioteki, takie jak Piscina, implementują te funkcje, a zazwyczaj stanowią lepszy wybór niż samodzielnie stworzona pula.
Kiedy wątki pracownicze są przydatne, a kiedy nie
Odpowiednie zastosowania
- Zadania obciążające CPU: wszystko, co zużywa znaczną ilość czasu procesora, takie jak skomplikowane obliczenia, kompresja, szyfrowanie lub przetwarzanie obrazów.
- Ochrona pętli zdarzeń: każda operacja, która w przeciwnym razie zablokowałaby główny wątek na dłużej, niż jest to do przyjęcia.
Nieodpowiednie zastosowania
- Zadania obciążające I/O: wywołania sieciowe, zapytania do bazy danych oraz operacje w systemie plików są już nieblokujące w Node.js. Umieszczanie ich w pracownikach dodaje nakład pracy bez żadnych korzyści.
- Mikroskali zadania: uruchamianie pracownika i przekazywanie mu wiadomości oba wymagają czasu. W przypadku szybkich obliczeń ten nakład może przewyższyć korzyści, co jest kolejnym powodem, by przechowywać pracowników o długim czasie życia w puli.
SharedArrayBuffer umożliwia to, ale skoordynowanie zapisów za pomocą Atomics jest niezwykle trudne do prawidłowego wykonania. Chyba że naprawdę tego potrzebujesz i rozumiesz konsekwencje, lepiej użyć przekazywania wiadomości.Praktyki zapewniające dobre funkcjonowanie workerów
- Zachowuj małe skrypty workerów. Ładuj tylko logikę i zależności niezbędne do wykonania zadania. Importowanie całego frameworka aplikacji do każdego workera wydłuża czas uruchamiania i zużywa więcej pamięci.
- Wybierz najtańszy sposób przekazywania danych. Strukturalne klonowanie za pomocą
postMessagenadaje się do małych wiadomości, dużeArrayBufferpowinny być przekazywane przez listę transferową, a głębokie lub złożone drzewa obiektów należy unikać, ponieważ ich klonowanie jest wolne.
error i exit na wątku głównym, a także otaczaj działania wewnątrz pracowników blokiem try...catch, aby błędy były przekazywane jako wyraźne komunikaty o błędzie.os.cpus().length to rozsądna wartość początkowa dla puli pracowników; w nowszych wersjach Node.js zalecany sposób na uzyskanie tej wartości to os.availableParallelism().message uruchamia długie, synchroniczne zadanie, kolejne wiadomości kumulują się za nim. Pracownik, który musi obsługiwać wiele żądań, powinien utrzymywać każdą jednostkę pracy w krótkim czasie lub, w rzadkich przypadkach, delegować zadania na swoich podpracowników.Błędy, które często sprawiają problemy
- Oczekiwanie na wspólne zmienne globalne. Dane przekazane do
new Worker()docierają do pracownika tylko poprzezworkerData(lub późniejsze wiadomości). Zmienne na poziomie modułu w głównej nitce nie są widoczne wewnątrz pracownika, ponieważ ten działa w swoim oddzielnym środowisku. Jedynym prawdziwym sposobem współdzielenia jestSharedArrayBuffer.
isMainThread. Określa ono kontekst wykonywania kodu, który jest obecnie uruchamiany. Jeden plik może być używany zarówno jako główny skrypt, jak i wątek roboczy, ale oddzielne pliki główne i robocze zazwyczaj są jaśniejsze.--inspect-brk, a w debuggerze konieczne jest przechodzenie pomiędzy kontekstami wątków. Edytory takie jak VS Code zajmują się tym w dużej mierze, jednak nadal jest to mniej bezpośrednie niż debugowanie głównego wątku.Główne wnioski
- Pracownicy uruchamiają JavaScript równolegle w oddzielnych izolatach w ramach jednego procesu; komunikują się poprzez przekazywanie wiadomości i dzielą pamięć wyłącznie za pomocą
SharedArrayBuffer. - Należy ich używać do zadań obciążających CPU, które w przeciwnym razie mogłyby zablokować pętlę zdarzeń, a nie do dostępu do sieci lub dysku, który Node.js już obsługuje asynchronicznie.
- Należy nadać wiadomościom jednoznaczny, spójny format, aby wyniki, błędy oraz identyfikatory zapytań były jednoznaczne po obu stronach.
- Należy przekazywać duże obiekty
ArrayBufferzamiast je klonować, pamiętając przy tym, że nadawca traci do nich dostęp po przekazaniu. - Należy ponownie wykorzystywać pracowników poprzez zbiór o wielkości odpowiadającej liczbie rdzeni procesora i upewnić się, że procedury wyłączania oraz przywracania po awarii są wyraźnie oddzielone, albo użyć sprawdzonej biblioteki do zarządzania takimi zbiorkami.
Literatura pokrewna
- Opanowanie konkurencji w Node.js: Unikanie zawieszeń API za pomocą p-map i Bottleneck — Dowiedz się, jak połączenie p-map i Bottleneck w Node.js zapobiega błędom związanym z ograniczeniami szybkości i przeciążeniem systemu poprzez kontrolę konkurencji oraz czasu realizacji zapytań.
- Wyjaśnienie konkurencji w Node.js: libuv, pętla zdarzeń i zbiór wątków — Dowiedz się, jak Node.js wykorzystuje prymitywy systemu operacyjnego dostarczane przez libuv oraz zbiór wątków do obsługi asynchronicznego I/O, a także o częstych problemach związanych z zbiorem wątków i wskazówkach dotyczących jego optymalizacji.