Décharger les tâches gourmandes en CPU dans Node.js à l’aide de worker_threads et de Pools
Découvrez comment les worker_threads de Node.js maintiennent le cycle d’événements réactif : création de travailleurs, communication demande-réponse, buffers transférables, pools, ainsi que leurs pièges.
Node.js gère des milliers de connexions simultanées grâce à un boucle d’événements à thread unique, car presque tout le travail consiste en des opérations d’attente réseau ou disque. Ce modèle cesse de fonctionner dès qu’une requête nécessite des calculs réels : le hachage d’un gros en-tête, la redimensionnement d’une image ou l’analyse d’un ensemble de données occupent le seul thread JavaScript, tandis que toutes les autres requêtes doivent attendre leur tour. Le module worker_threads constitue la solution intégrée à ce problème. Après avoir suivi ce guide, vous serez en mesure de déplacer le code gourmand en CPU dans des workers, d’échanger des données avec eux de manière efficace, de les exécuter en groupe, et de reconnaître les cas où un worker n’est pas l’outil adapté.
Pourquoi les threads workers existent
worker_threads est apparu pour la première fois dans Node.js 10.5.0 sous forme expérimentale, et est considéré comme stable depuis la version Node.js 12. Il permet à un seul processus Node.js d’exécuter du JavaScript sur plusieurs threads en même temps.
Celui-ci diffère de child_process d’une manière importante. Un processus enfant est un processus Node.js entièrement distinct, disposant de sa propre mémoire et de sa propre copie du moteur d’exécution. Un travailleur, quant à lui, vit à l’intérieur du même processus, mais il ne s’agit pas simplement d’« un autre thread qui partage tout ». Chaque travailleur dispose de son propre isolat V8, de sa propre pile mémoire et de son propre boucle d’événements. Rien n’est partagé implicitement ; les travailleurs communiquent en échangeant des messages, et ils ne partagent leur mémoire que lorsque vous leur transmettez explicitement un SharedArrayBuffer. (On affirme souvent que les travailleurs partagent l’instance V8 et la mémoire du thread principal. Cela n’est pas exact, et cette différence explique en grande partie la conception de l’API présentée ci-dessous.) Par rapport aux processus, les travailleurs sont moins coûteux à lancer et à faire communiquer, ce qui en fait le moyen naturel d’utiliser plusieurs cœurs au sein d’une même instance d’application.
Le coût pour vous d’une boucle d’événements bloquée
La tâche principale d’un worker est de maintenir le cycle d’événements libre. Tant qu’un calcul synchrone s’exécute sur le thread principal, le serveur ne peut pas accepter de nouvelles requêtes, ne peut pas faire avancer les connexions ouvertes, et ne peut même pas répondre à une vérification de santé. Les conséquences pratiques sont les suivantes :
- Expérience utilisateur dégradée : les réponses API arrivent en retard, et les fonctionnalités en temps réel telles que les mises à jour en direct sont perturbées.
- Rendement réduit : une tâche longue occupe le seul thread JavaScript, ce qui entraîne une diminution du nombre de requêtes par seconde.
- Instabilité : un blocage prolongé peut dépasser les temps d’attente du balanceur de charge ou du client, et ces pannes peuvent se propager à d’autres services.
En déplaçant les calculs lourds vers un worker, le cycle principal reste libre pour gérer les sockets et les temporiseurs, permettant ainsi à l’application de rester réactive même lorsqu’un calcul important est en cours.
Pour en savoir plus sur la manière dont le cycle d’événements, libuv et leur pool de threads s’interconnectent, consultez comment fonctionne la concurrence dans Node.js en profondeur.
Utilisation de chaque cœur
Les serveurs disposent généralement de nombreux cœurs CPU. Sans travailleurs, ou sans exécution de plusieurs processus via child_process ou un gestionnaire de processus comme PM2, un seul processus Node.js peut utiliser environ un cœur pour le JavaScript gourmand en ressources CPU. Les autres cœurs restent inutilisés. Les travailleurs permettent à une application d’exécuter plusieurs tâches lourdes en parallèle sur les cœurs disponibles, augmentant ainsi la capacité de calcul globale.
Charges de travail qui deviennent réalisables
Le calcul parallèle ouvre la voie à des domaines qui étaient traditionnellement réservés aux langages disposant de threads de première classe :
- Traitement des données : analyse en temps réel, manipulation d’images et de vidéos, ainsi que transformation de grands ensembles de données.
- Inference d’apprentissage automatique : exécution de modèles entraînés tout en maintenant l’interactivité du reste de l’application.
- Cryptographie : hachage, chiffrement et déchiffrement.
- Simulation et calcul scientifique : simulations parallèles et modèles mathématiques complexes.
Le point clé est que rien de tout cela ne vous fait perdre le modèle d’E/S asynchrone. Les workers ajoutent une computation parallèle à ce modèle.
Les éléments constitutifs du module
worker_threads expose un petit ensemble de primitives :
Worker: la classe que l’thread principal utilise pour démarrer un nouveau worker à partir d’un script.
isMainThread: un booléen qui indique si le code actuel s’exécute sur le thread principal, utile lorsque un même fichier peut jouer les deux rôles.parentPort: disponible à l’intérieur d’un worker, c’est le canal menant au thread qui l’a créé. Il vaut null sur le thread principal.workerData: disponible à l’intérieur d’un worker, il contient une copie des données fournies par le parent lors de la création du worker.MessagePort et MessageChannel: utilisés pour créer des canaux bidirectionnels indépendants supplémentaires entre les threads.SharedArrayBuffer et Atomics : une mémoire partagée véritable accompagnée des primitives de synchronisation nécessaires pour l’utiliser en toute sécurité. Ce sont l’option la plus efficace, mais aussi la plus complexe, car elles réintroduisent des conditions de course.Un premier travailleur : trouver des nombres premiers en dehors du thread principal
Un bon premier exercice consiste à effectuer une tâche délibérément lente : calculer tous les nombres premiers jusqu’à une grande limite. L’exemple utilise deux fichiers situés dans le même dossier.
Le script du thread principal
Le script principal crée le travailleur, lui fournit ses données d’entrée et attend sa réponse. Avant de la lire, notez trois points importants : le travailleur est créé à partir d’un chemin de fichier distinct ; les données d’entrée sont transmises via l’option workerData ; enfin, le script parent s’abonne à trois événements : message pour les résultats, error pour les exceptions non capturées par le travailleur, et exit lorsque la thread s’arrête.
// 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);
}
Décortiquons les lignes clés :
- La vérification
isMainThreadprotège la logique de création du travailleur en s’assurant qu’elle ne s’exécute que sur la thread principale. Le branchementelseexiste uniquement pour montrer ce qui se passerait si ce même fichier était chargé en tant que travailleur ; dans cette configuration, il n’est jamais exécuté. largeNumberest défini à une valeur suffisamment élevée (20 millions) afin que le calcul prenne un temps notable.
new Worker(path.join(__dirname, 'prime-worker.js'), { workerData: { limit: largeNumber } }) est l’appel principal. Le premier argument indique le script que la nouvelle thread exécutera, qui doit être un fichier distinct. Le deuxième est un objet d’options dont workerData constitue les données initiales du worker.workerData est copié à l’aide de l’algorithme de clonage structuré d’HTML. Le worker reçoit sa propre copie, et non une référence à l’objet parent.message reçoit tout ce que le worker envoie, ici l’array des nombres premiers. L’écouteur error est déclenché en cas d’exception non capturée à l’intérieur du worker, tandis que l’écouteur exit reçoit un code de sortie, où 0 signifie une fin normale.Le script du worker
Le worker ne contient que le calcul coûteux et le code qui rapporte son résultat. La fonction findPrimes est une boucle de division par essai simple, délibérément non optimisée afin d’exploiter le temps de traitement du CPU.
// 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.");
}
Points à noter :
- Elle récupère
parentPortetworkerDatadepuis le module. - La condition
if (parentPort)garantit que la logique ne s’exécute que lorsque le fichier est chargé en tant que worker ; si vous l’exécutez directement avecnode,parentPortvautnullet un message est affiché à la place. const { limit } = workerData;lit les données envoyées par le parent.parentPort.postMessage(result)transmet les nombres premiers au thread principal, où le gestionnaire demessageles récupère.
try...catch enregistre toute erreur survenue pendant le calcul et la rapporte au thread parent plutôt que de laisser le thread s’arrêter.Une subtilité : en cas d’échec, le thread travailleur publie { error: ... } sur le même canal message utilisé pour les résultats, mais le script principal traite chaque message comme un tableau de nombres premiers. Dans du code réel, donnez aux messages une structure explicite (par exemple un champ status ou type) afin que le thread parent puisse distinguer un résultat d’une erreur.
Lancer le programme et lire la sortie
Enregistrez les deux fichiers côte à côte et lancez le programme avec node main.js. Vous verrez le thread principal se déclarer, le thread travailleur indiquer qu’il a commencé à fonctionner, et quelque temps plus tard le résultat obtenu par ce dernier suivi du temps écoulé.
Faites attention à l’ordre de la ligne « Thread principal qui continue avec d’autres tâches... ». Dans ce script, elle est affichée à l’intérieur du gestionnaire message, ce qui signifie qu’elle n’apparaît que lorsque le thread d’exécution a terminé. Cela montre que le thread principal a pu recevoir et traiter le résultat, et non qu’il effectuait d’autres tâches en même temps. Pour voir réellement le thread principal rester réactif pendant le calcul, lancez une action sur ce dernier avant l’arrivée du résultat, par exemple un setInterval qui enregistre un signal de vie toutes les quelques centaines de millisecondes. Ce signal continue d’être émis tant que le thread d’exécution effectue les calculs, ce qui ne se produirait pas si findPrimes était exécuté sur le thread principal.
Messagerie bidirectionnelle avec un protocole demande-réponse
Le premier exemple envoie une entrée et reçoit un résultat. La plupart des utilisations réelles nécessitent un travailleur à long terme capable de gérer de nombreuses requêtes. La communication est entièrement bidirectionnelle : un objet Worker sur le thread principal ainsi que parentPort à l’intérieur du travailleur offrent tous deux les méthodes .postMessage(...) et .on('message', ...).
Le travailleur présent ci-dessous reste actif et attend des messages. Lorsqu’il reçoit une requête calculateFibonacci, il calcule la valeur à l’aide d’une fonction récursive naïve (intentionnellement lente) et répond soit par un message fibonacciResult, soit par un message fibonacciError, en incluant l’identifiant de la requête afin que l’appelant puisse correspondre la réponse à sa demande.
// 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.');
}
Dans la partie principale, l’astuce consiste à transformer le transfert de messages en promesses. Chaque appel génère un requestId unique, stocke les fonctions resolve et reject de la promesse dans un Map associé à cet ID, puis envoie la requête. Lorsqu’une réponse arrive, le gestionnaire de message recherche l’ID correspondant, supprime la entrée associée et résout la promesse correspondante. L’exemple utilise le package uuid pour générer les IDs ; dans les versions actuelles de Node.js, crypto.randomUUID() du module intégré node:crypto remplit la même fonction sans nécessiter de dépendance supplémentaire.
// 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
}
})();
}
En résumé, le schéma est le suivant :
- Le thread principal crée un
requestIdpour chaque requête afin que les réponses puissent être associées à leurs promesses, même si elles arrivent dans un ordre différent.
pendingRequests contient les fonctions resolve et reject en attendant que le travailleur réponde.fibonacciResult ou fibonacciError, en incluant l’requestId d’origine.Prêtez une attention particulière à la structure du message si vous adaptez ce code. Le thread principal place requestId au niveau le plus élevé du message, à côté de payload, tandis que le worker le récupère via message.payload. Tel quel, le worker considère donc requestId comme undefined, ce qui empêche de correspondre à la réponse et fait en sorte que la promesse attendue ne se résolve jamais. Conservez l’ID à un endroit convenu des deux côtés (par exemple payload: { n, requestId }). Il est également utile d’ajouter un délai d’expiration pour chaque demande en attente, afin qu’une réponse perdue se transforme en erreur plutôt que de provoquer une blocage. Enfin, notez le bloc finally : une fois le travail terminé, worker.terminate() arrête le thread permettant ainsi au processus de se terminer.
Déplacer de grandes quantités de données sans les copier
Tout ce qui est envoyé via postMessage est par défaut cloné de manière structurée, ce qui signifie que le destinataire reçoit une copie. Cela convient pour de petits messages. Mais pour de gros ensembles de données binaires, la copie elle-même peut devenir un goulot d’étranglement. La solution consiste à transférer la propriété de la mémoire sous-jacente au lieu de la copier : l’objet devient inutilisable pour l’expéditeur et immédiatement utilisable pour le destinataire, sans duplication.
ArrayBuffer et MessagePort sont les objets le plus souvent transférés. SharedArrayBuffer constitue un cas particulier : il n’est pas du tout transféré, car les deux threads peuvent accéder à la même mémoire simultanément.
L’exemple ci-dessous contient deux scripts dans une même liste. Le travailleur (buffer-worker.js) reçoit un buffer, double chaque octet sur place et renvoie le buffer sous forme de transférable. Le script principal (main-buffer.js) remplit un buffer de 1 MB, le transfère au travailleur et lit les données traitées lorsqu’il est renvoyé.
// 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();
}
});
}
Les détails clés :
- Le deuxième argument de
postMessage, ici[myBuffer]d’un côté et[buffer]de l’autre, correspond à la liste de transfert. - Lorsque
myBufferest envoyé, sa propriété passe au travailleur. Dans le thread principal, le buffer devient détaché : sa valeurbyteLengthpasse à 0 et toute vue de type tableau binaire sur lui, comme leuint8précédemment mentionné, a également une longueur de 0.
L’inconvénient est qu’il faut cesser d’utiliser la référence originale après le transfert. Si les deux threads ont réellement besoin des données en même temps, il faut une copie ou un SharedArrayBuffer.
Gestion des erreurs et nettoyage
Les travailleurs échouent comme n’importe quel autre code, et ils conservent de la mémoire ainsi qu’un thread tant qu’ils sont actifs. Les outils disponibles :
worker.on('error', handler)sur le thread principal : reçoit les exceptions qui ont été levées dans le travailleur et n’ont jamais été capturées.worker.on('exit', handler)sur le thread principal : est déclenché chaque fois que le worker s’arrête, qu’il ait terminé normalement ou qu’il soit tombé en panne. Le code de sortie permet de distinguer les deux cas,0indiquant une réussite.process.on('uncaughtException', handler)à l’intérieur du worker : vous offre un endroit pour enregistrer ou signaler des détails depuis l’intérieur du worker avant qu’il ne cesse de fonctionner, en plus de ce que voit le processus parent.worker.terminate(): arrête le worker dès que possible et renvoie une promesse qui se résout avec le code de sortie. Il ne patiente pas la fin des tâches en cours, donc si vous avez besoin d’une fermeture propre, envoyez au worker un message « stop » et laissez-le quitter de lui-même. Assurez-vous toujours que les workers qui ne sont plus nécessaires soient arrêtés d’une manière ou d’une autre.
Exécution de nombreuses tâches via un pool de workers
Créer un nouveau travailleur pour chaque tâche est gaspilleur, et exécuter plus de travailleurs gourmands en CPU que le nombre de cœurs disponibles ne fait qu’augmenter les conflits. Un pool de travailleurs résout ces deux problèmes : il lance un nombre fixe de travailleurs, place les tâches arrivantes dans une file d’attente et attribue chaque tâche au prochain travailleur disponible.
Le pool simplifié présenté ci-dessous contient un tableau de fiches de travailleurs, une liste des IDs de travailleurs libres et une file d’attente des tâches en attente. runTask renvoie une promesse et insère la tâche dans la file ; processQueue associe la tâche la plus ancienne à un travailleur libre ; lorsque le travailleur rapporte son état, sa promesse est résolue et il rejoint à nouveau la liste des travailleurs libres. Si un travailleur rencontre une erreur ou s’arrête de manière anormale, terminateWorker le remplace par un nouveau travailleur afin de maintenir la taille maximale du pool.
// 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;
Le travailleur générique utilisé par le pool exécute une fonction gourmande en CPU pour chaque message qu’il reçoit et répond avec un status ainsi qu’un payload. La boucle de sommation représente n’importe quel travail réel, tel que le traitement d’images ou le chiffrement.
// 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 });
}
});
}
Finalement, le script principal définit la taille du pool en fonction du nombre de cœurs CPU indiqué par os.cpus(), soumet six tâches et attend leur exécution complète à l’aide de Promise.allSettled. La promesse associée à chaque tâche est convertie en une chaîne lisible indiquant un succès ou un échec.
// 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();
Considérez ce pool strictement comme un exemple pédagogique ; il présente plusieurs défauts importants avant qu’un système similaire ne puisse être mis en production :
- Le gestionnaire de
messagedu pool résout la promesse avec lepayloadquel que soit lestatus; ainsi, une tâche qui a échoué à l’intérieur du worker est toujours considérée comme réussie. Vérifiez lestatuset rejetez la promesse en cas de valeur'error'. worker.terminate()fait en sorte qu’un worker s’arrête avec un code différent de zéro. Comme le gestionnaire d’arrêt appelleterminateWorker, ce qui crée un worker de remplacement, l’appel àclose()déclenche un cycle de nouveaux workers au lieu d’arrêter le pool. De plus, une erreur suivie de son événementexitpeut occuper deux fois la même place et inscrire son ID sur la liste des slots libres plus d’une fois. Un vrai pool a besoin d’un indicateur de « fermeture » et ne devrait remplacer que les workers qui ont planté de manière inattendue.
error ; une sortie anormale sans événement error laisse sa promesse en attente.Number.MAX_SAFE_INTEGER, ce qui fait perdre en précision les résultats affichés. Utilisez BigInt si des valeurs exactes sont importantes.Les pools plus complets ajoutent également des délais de timeout pour les threads inactifs, une mise à l’échelle dynamique et des mécanismes de récupération plus rigoureux. Des bibliothèques bien entretenues comme Piscina implémentent ces fonctionnalités, et elles constituent généralement un meilleur choix qu’un pool conçu manuellement.
Lorsque les threads d’exécution sont utiles et lorsqu’ils ne le sont pas
Cas d’usage appropriés
- Tâches gourmandes en CPU : tout ce qui consomme beaucoup de temps processeur, comme des calculs intensifs, la compression, le chiffrement ou le traitement d’images.
- Protection du cycle d’événements : toute opération qui bloquerait sinon le thread principal plus longtemps que ce que l’on peut tolérer.
Adaptations insuffisantes
- Travaux gourmands en E/S : les appels réseau, les requêtes de base de données et les opérations sur le système de fichiers sont déjà non bloquants dans Node.js. Les envelopper dans un worker ajoute de la charge inutile sans aucun avantage.
- Tâches mineures : démarrer un worker et envoyer des messages coûtent du temps. Pour des calculs rapides, cette charge supplémentaire peut annuler les gains, ce qui est une autre raison de conserver des workers à longue durée de vie dans un pool.
SharedArrayBuffer le rend possible, mais coordonner les écritures avec Atomics est extrêmement difficile à gérer correctement. À moins que vous en ayez vraiment besoin et de comprendre les conséquences, préférez le passage de messages.Pratiques pour maintenir les workers en bon état
- Rendez les scripts des workers petits. Chargez uniquement la logique et les dépendances nécessaires à la tâche. Importer l’ensemble du framework de votre application dans chaque worker allonge le temps de démarrage et consomme plus de mémoire.
- Choisissez le chemin de transmission des données le moins coûteux. La clonage structuré via
postMessageconvient bien pour de petits messages ; les grandsArrayBufferdoivent être transmis via la liste de transfert, et il vaut mieux éviter les graphes d’objets profonds ou complexes car leur clonage est lent.
error et exit sur le thread principal, et enveloppez les tâches exécutées par les travailleurs dans des blocs try...catch afin que les échecs soient rapportés sous forme de messages d’erreur explicites.os.cpus().length constitue une valeur de départ raisonnable pour déterminer la taille d’un pool ; sur les versions plus récentes de Node.js, os.availableParallelism() est la méthode recommandée pour obtenir ce chiffre.message exécute une tâche synchrone longue, les messages suivants s’accumulent derrière elle. Un worker qui doit gérer de nombreuses requêtes devrait garder chaque unité de travail courte ou, dans des cas rares, déléguer à ses propres sous-workers.Erreurs fréquentes
- Attendre des variables globales partagées. Les données fournies à
new Worker()ne parviennent au worker que viaworkerData(ou des messages ultérieurs). Les variables au niveau de module dans le thread principal ne sont pas visibles à l’intérieur du worker, car celui-ci s’exécute dans son propre environnement isolé. Le seul véritable partage est possible grâce àSharedArrayBuffer.
isMainThread. Ce paramètre décrit le contexte d’exécution du code qui est actuellement en cours d’exécution. Un seul fichier peut être utilisé pour fonctionner soit comme script principal, soit comme thread d’arbeitant, mais il est généralement plus clair d’utiliser des fichiers distincts pour le script principal et les threads d’arbeitant.--inspect-brk et passer entre les contextes de thread dans le débogueur. Des éditeurs comme VS Code gèrent en grande partie cette fonctionnalité, mais cela reste moins direct que le débogage du thread principal.Points clés
- Les workers exécutent du JavaScript en parallèle dans des isolats distincts au sein d’un même processus ; ils communiquent par l’envoi de messages et ne partagent la mémoire qu’à travers
SharedArrayBuffer. - Utilisez-les pour des tâches gourmandes en CPU qui, sinon, bloqueraient la boucle d’événements, et non pour des accès réseau ou disque que Node.js gère déjà de manière asynchrone.
- Donnez aux messages une structure explicite et cohérente afin que les résultats, les erreurs et les identifiants de requête soient univoques des deux côtés.
- Transférez de grands
ArrayBufferplutôt que de les cloner, et n’oubliez pas que l’expéditeur perd tout accès à eux par la suite. - Réutilisez les workers grâce à un pool dont la taille correspond au nombre de cœurs du système, et assurez-vous que les opérations de fermeture et de récupération en cas de panne soient clairement séparées, ou utilisez une bibliothèque de pool éprouvée.
Lectures complémentaires
- Maitriser la concurrentise de Node.js : Éviter les crashes d’API avec p-map et Bottleneck — Découvrez comment combiner p-map et Bottleneck dans Node.js pour prévenir les erreurs de limitation de débit et la surcharge du système en contrôlant la concurrentise et le timing des requêtes.
- Explication de la concurrentise de Node.js : libuv, la boucle d’événements et le pool de threads — Apprenez comment Node.js utilise les primitives du système d’exploitation de libuv ainsi que son pool de threads travailleurs pour gérer les I/O asynchrone, ainsi que les pièges courants liés au pool de threads et des conseils d’optimisation.