Accueil / Articles / Décharger les tâches gourmandes en CPU dans Node.js à l’aide de worker_threads et de Pools

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.

4814 mots

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 isMainThread protège la logique de création du travailleur en s’assurant qu’elle ne s’exécute que sur la thread principale. Le branchement else existe 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é.
    • largeNumber est 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.
  • L’écouteur 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 parentPort et workerData depuis 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 avec node, parentPort vaut null et 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 de message les récupère.
  • Le bloc 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 requestId pour chaque requête afin que les réponses puissent être associées à leurs promesses, même si elles arrivent dans un ordre différent.
  • Le tableau pendingRequests contient les fonctions resolve et reject en attendant que le travailleur réponde.
  • calculateFibonacci et répond par 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 myBuffer est envoyé, sa propriété passe au travailleur. Dans le thread principal, le buffer devient détaché : sa valeur byteLength passe à 0 et toute vue de type tableau binaire sur lui, comme le uint8 précédemment mentionné, a également une longueur de 0.
  • Le travailleur modifie les données et les renvoie, de sorte que le thread principal reçoit à nouveau un buffer.
  • Les 1 MB de données ne sont jamais copiés dans un sens ou dans l’autre.
  • 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, 0 indiquant 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 message du pool résout la promesse avec le payload quel que soit le status ; ainsi, une tâche qui a échoué à l’intérieur du worker est toujours considérée comme réussie. Vérifiez le status et 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 appelle terminateWorker, 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énement exit peut 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.
  • Une tâche en cours lorsque son exécuteur plante est rejetée uniquement via le chemin error ; une sortie anormale sans événement error laisse sa promesse en attente.
  • Dans la démonstration, les sommes atteignant deux milliards dépassent Number.MAX_SAFE_INTEGER, ce qui fait perdre en précision les résultats affichés. Utilisez BigInt si des valeurs exactes sont importantes.
  • Le journal « Doing other stuff » s’exécute après que toutes les tâches ont été attendues, de sorte qu’il ne montre pas la concurrence avec ces tâches, comme dans le premier exemple.
  • 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.
  • État mutable partagé : 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

    1. 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.
    2. Choisissez le chemin de transmission des données le moins coûteux. La clonage structuré via postMessage convient bien pour de petits messages ; les grands ArrayBuffer doivent être transmis via la liste de transfert, et il vaut mieux éviter les graphes d’objets profonds ou complexes car leur clonage est lent.
  • Mettez les travailleurs en pause. Terminez leur exécution ou permettez-leur de quitter lorsqu’ils ont terminé leur travail. Les travailleurs inactifs consomment encore de la mémoire.
  • Gérez les erreurs des deux côtés. Écoutez les événements 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.
  • Surveillez l’utilisation des ressources. Trop de travailleurs entraînent un changement fréquent de contexte et une consommation accrue de mémoire, ce qui peut ralentir les opérations au lieu de les accélérer. 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.
  • Ne bloquez pas la boucle propre du worker. Un worker dispose également d’une boucle d’événements. Si son gestionnaire de 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 via workerData (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.
  • Mauvaise interprétation de 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.
  • Présupposition que le débogage fonctionne comme d’habitude. L’inspecteur de Node.js peut être associé aux threads d’arbeitant, mais il faut lancer le processus avec des paramètres d’inspecteur tels que --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.
  • Laisser les threads d’arbeitant en cours d’exécution. Les threads d’ Arbeitant oubliés accumulent de la mémoire et des ressources au fil du temps, ce qui peut épuiser les ressources dans les services à long terme.
  • 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 ArrayBuffer plutô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

  • Terminer les chutes de demandes dans SvelteKit avec des fonctions de chargement parallèles — Découvrez pourquoi les appels séquentiels à awaits ralentissent les pages, et comment SvelteKit, les fonctions de chargement, Promise.all, des appels prudents à parent() ainsi que les promesses en flux éliminent ces chutes de demandes.