Startseite / Artikel / Auslagern von rechenintensiven Aufgaben in Node.js mit worker_threads und Pools

Auslagern von rechenintensiven Aufgaben in Node.js mit worker_threads und Pools

Erfahren Sie, wie Node.js worker_threads den Event-Loop reaktiv halten: Erstellen von Worker-Prozessen, Nachrichtenaustausch zwischen Anfragen und Antworten, übertragbare Puffer, Pools sowie ihre Fallstricke.

4814 Wörter

Node.js kann Tausende von gleichzeitigen Verbindungen in einem eindimensionalen Ereigniszyklus verarbeiten, weil fast die gesamte Arbeit auf das Netzwerk oder die Festplatte wartet. Das Modell bricht zusammen, sobald eine Anfrage echte Berechnungen erfordert: Das Hashen eines großen Datenträgers, das Anpassen der Größe eines Bildes oder die Verarbeitung einer Datensammlung halten den einzigen JavaScript-Thread beschäftigt, während alle anderen Anfragen warten müssen. Das worker_threads-Modul ist die eingebaute Lösung dafür. Nachdem Sie diesen Leitfaden durchgearbeitet haben, werden Sie in der Lage sein, CPU-intensiven Code in Worker-Prozesse zu verlagern, effizient mit ihnen Daten auszutauschen, sie als Pool zu betreiben und die Fälle zu erkennen, in denen Worker nicht das richtige Werkzeug sind.

Warum es Worker-Threads gibt

worker_threads tauchte erstmals in Node.js 10.5.0 hinter einem experimentellen Flag auf und gilt seit der Node.js 12-Version als stabil. Es ermöglicht es einem einzelnen Node.js-Prozess, JavaScript auf mehreren Threads gleichzeitig auszuführen.

Es unterscheidet sich in wichtiger Hinsicht von child_process. Ein Kindprozess ist ein völlig separater Node.js-Prozess mit eigenem Speicher und eigener Kopie der Laufzeitumgebung. Ein Worker existiert innerhalb desselben Prozesses, ist aber nicht einfach „ein weiterer Thread, der alles teilt“. Jeder Worker erhält sein eigenes V8-Isolat, seinen eigenen Heap sowie seine eigene Ereignisschleife. Nichts wird implizit geteilt; Worker kommunizieren durch das Übergeben von Nachrichten und teilen den Speicher nur dann, wenn man ihnen explizit einen SharedArrayBuffer übergibt. (Häufig wird behauptet, Worker teilen sich die V8-Instanz und den Speicher des Hauptthreads. Das ist nicht korrekt, und dieser Unterschied erklärt den größten Teil der folgenden API-Struktur.) Im Vergleich zu Prozessen sind Worker günstiger in der Startung und Kommunikation, was sie zur natürlichen Lösung macht, um mehrere Kerne einer Anwendungsinstanz zu nutzen.

Was Sie für eine blockierte Ereignisschleife bezahlen

Die Hauptaufgabe eines Workers besteht darin, den Event-Loop frei zu halten. Während eine synchrone Berechnung auf dem Hauptthread ausgeführt wird, kann der Server keine neuen Anfragen entgegennehmen, macht bei offenen Verbindungen keinen Fortschritt und kann nicht einmal eine Gesundheitsprüfung beantworten. Die praktischen Folgen sind wie folgt:

  • Verbesserte Benutzererfahrung: API-Antworten kommen verspätet an, und Echtzeitfunktionen wie Live-Updates funktionieren unregelmäßig.
  • Niedrigere Durchsatzrate: Eine lange Aufgabe belegt den einzigen JavaScript-Thread, wodurch die Anzahl der Anfragen pro Sekunde abnimmt.
  • Instabilität: Langes Blockieren kann die Timeout-Werte des Load Balancers oder der Clients überschreiten, und solche Fehler können sich auf andere Dienste auswirken.

Durch Verlegung der aufwendigen Berechnungen an einen Worker bleibt der Hauptloop frei, um Sockets und Timer weiter zu verwalten, sodass die Anwendung auch während intensiver Berechnungen reaktiv bleibt.

Für eine detailliertere Betrachtung davon, wie der Event-Loop, libuv und ihr Thread-Pool zusammenwirken, siehe wie die Konkurrenzfähigkeit von Node.js im Hintergrund funktioniert.

Jeder Kern nutzen

Server verfügen in der Regel über viele CPU-Kerne. Ohne Worker oder ohne Ausführung mehrerer Prozesse mithilfe von child_process oder einem Prozessmanager wie PM2 kann ein einzelner Node.js-Prozess für CPU-intensive JavaScript-Aufgaben in der Regel nur einen Kern nutzen. Die restlichen Kerne bleiben ungenutzt. Worker ermöglichen es einer Anwendung, mehrere aufwändige Aufgaben parallel über die verfügbaren Kerne auszuführen und somit die gesamte Rechenleistung zu erhöhen.

Arbeitslasten, die praktikabel werden

Parallelrechnung eröffnet Bereiche, die traditionell Sprachen mit erstklassigem Threading vorbehalten waren:

  • Datenverarbeitung: Echtzeit-Analytik, Manipulation von Bildern und Videos sowie Transformation großer Datensätze.
  • Maschinelles Lernen – Inferenz:Ausführung von trainierten Modellen, während der Rest der Anwendung weiterhin interaktiv bleibt.
  • Kryptographie: Hashing, Verschlüsselung und Entschlüsselung.
  • Simulation und wissenschaftliche Rechenverfahren: Parallelsimulationen sowie komplexe mathematische Modelle.

Der entscheidende Punkt ist, dass all das Ihnen das asynchrone I/O-Modell nicht kostet. Worker fügen dazu parallele Berechnungen hinzu.

Die Bausteine des Moduls

worker_threads stellt eine kleine Anzahl von Primitiven zur Verfügung:

  • Worker: Die Klasse, die der Hauptthread verwendet, um einen neuen Worker aus einem Skript zu starten.
  • isMainThread: ein Boolean, der anzeigt, ob der aktuelle Code auf dem Hauptthread ausgeführt wird – nützlich, wenn eine Datei beide Rollen übernehmen kann.
  • parentPort: innerhalb eines Workers verfügbar, es handelt sich dabei um den Kanal zurück zum Thread, der ihn erstellt hat. Auf dem Hauptthread ist dieser Wert null.
  • workerData: innerhalb eines Workers verfügbar, er enthält eine Kopie der Daten, die der Elternthread bei der Erstellung des Workers bereitgestellt hat.
  • MessagePort und MessageChannel: dienen dazu, zusätzliche, unabhängige bidirektionale Kanäle zwischen Threads zu erstellen.
  • SharedArrayBuffer und Atomics: eches gemeinsames Speichermedium zusammen mit den Synchronisierungsprimitiven, die zur sicheren Nutzung erforderlich sind. Sie stellen die effizienteste Option dar, sind aber auch am komplexesten, da sie erneut Rennbedingungen hervorrufen.
  • Ein erster Worker: Primzahlen außerhalb des Hauptthreads finden

    Eine gute erste Übung ist eine absichtlich langsame Aufgabe – nämlich die Berechnung aller Primzahlen bis zu einer hohen Grenze. Das Beispiel verwendet zwei Dateien, die sich im selben Verzeichnis befinden.

    Das Skript für den Hauptthread

    Die Hauptskript erstellt den Worker, gibt ihm die Eingabe und wartet auf die Antwort. Bevor man sie liest, sollten drei Dinge beachtet werden: Der Worker wird aus einem separaten Dateipfad erstellt. Die Eingabe wird über die Option workerData übermittelt. Zudem abonniert der Elternteil drei Ereignisse: message für die Ergebnisse, error für Ausnahmen, die der Worker nicht gefangen hat, sowie exit für den Fall, dass der Thread gestoppt wird.

    // 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);
    }
    

    Überblick über die wichtigen Zeilen:

    • Die Überprüfung mit isMainThread schützt die Spawning-Logik, sodass sie nur auf dem Hauptthread ausgeführt wird. Der else-Zweig dient ausschließlich dazu, zu veranschaulichen, was passieren würde, wenn dieselbe Datei als Worker geladen würde; in dieser Konfiguration wird er niemals ausgeführt.
    • largeNumber wird so hoch gewählt (20 Millionen), dass die Berechnung erhebliche Zeit in Anspruch nimmt.
  • new Worker(path.join(__dirname, 'prime-worker.js'), { workerData: { limit: largeNumber } }) ist der zentrale Aufruf. Das erste Argument weist auf das Skript hin, das der neue Thread ausführen wird – es muss eine eigene Datei sein. Das zweite Argument ist ein Options-Objekt, dessen workerData als anfängliche Eingabe für den Worker dient.
  • workerData wird mithilfe des HTML-structured-clone-Algorithmus kopiert. Der Worker erhält somit seine eigene Kopie und nicht einen Verweis auf das Objekt des Elterntasks.
  • Der message-Listener erhält alles, was der Worker sendet – in diesem Fall den Array mit Primzahlen. Der error-Listener triggert bei nicht gefangenen Ausnahmen innerhalb des Workers, und der exit-Listener erhält einen Abbruchcode, wobei 0 ein normales Beenden bedeutet.
  • Das Worker-Skript

    Der Worker enthält lediglich die aufwändige Berechnung sowie den Code, der das Ergebnis ausgibt. Die findPrimes-Funktion ist ein einfacher Versuchsteilungszyklus, der absichtlich unoptimiert ist, damit er CPU-Zeit verbraucht.

    // 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.");
    }
    

    Was man beachten sollte:

    • Sie holt parentPort und workerData aus dem Modul ab.
    • Die Bedingung if (parentPort) stellt sicher, dass die Logik nur dann ausgeführt wird, wenn die Datei als Worker geladen wird; wenn man sie direkt mit node ausführt, ist parentPort null und es wird stattdessen eine Nachricht ausgegeben.
    • const { limit } = workerData; liest die vom Elternteil gesendeten Eingaben ein.
    • parentPort.postMessage(result) übermittelt die Primzahlen an den Hauptthread, wo der message-Handler sie entgegennimmt.
  • Der try...catch-Block protokolliert jeden Fehler während der Berechnung und meldet ihn an den Elternprozess, anstatt den Thread absterben zu lassen.
  • Eine Nuance: Im Falle eines Fehlers sendet der Worker { error: ... } über denselben message-Channel, der für Ergebnisse verwendet wird, doch die Hauptskript behandelt jede Nachricht als Array von Primzahlen. In echtem Code sollte den Nachrichten eine explizite Struktur gegeben werden (zum Beispiel ein status- oder type-Feld), damit der Elternprozess zwischen Ergebnissen und Fehlern unterscheiden kann.

    Ausführen und Ausgabe lesen

    Speichern Sie beide Dateien nebeneinander und starten Sie das Programm mit node main.js. Sie werden sehen, wie der Hauptthread sich selbst ankündigt, der Worker meldet, dass er gestartet ist, und nach einer Weile folgt das Ergebnis des Workers zusammen mit der verstrichenen Zeit.

    Achten Sie auf die Reihenfolge der Zeile „Hauptthread setzt mit anderen Aufgaben fort...“. In diesem Skript wird sie innerhalb des message-Handlers ausgegeben, sodass sie erst nach Abschluss der Arbeit des Workers erscheint. Damit wird gezeigt, dass der Hauptthread das Ergebnis empfangen und verarbeiten konnte – nicht jedoch, dass er in der Zwischenzeit andere Aufgaben ausgeführt hat. Um tatsächlich zu sehen, dass der Hauptthread während der Berechnung weiterhin reaktiv bleibt, starten Sie etwas auf dem Hauptthread, bevor das Ergebnis eintrifft – beispielsweise einen setInterval, der alle paar Hundert Millisekunden ein „Heartbeat“-Signal protokolliert. Dieses Signal wird weiterhin gesendet, solange der Worker berechnet, was genau nicht der Fall wäre, wenn findPrimes auf dem Hauptthread ausgeführt würde.

    Zweiwege-Kommunikation mit einem Request-Response-Protokoll

    Das erste Beispiel sendet eine Eingabe und erhält ein Ergebnis. Die meisten tatsächlichen Anwendungen benötigen einen lang lebenden Worker, der viele Anfragen verarbeitet. Die Kommunikation ist vollständig zweireihig: Ein Worker-Objekt auf dem Hauptthread sowie parentPort innerhalb des Workers bieten sowohl .postMessage(...) als auch .on('message', ...) bereit.

    Der untenstehende Worker bleibt aktiv und wartet auf Nachrichten. Wenn er eine calculateFibonacci-Anfrage erhält, berechnet er den Wert mithilfe einer naiven rekursiven Funktion (absichtlich langsam) und antwortet mit einer fibonacciResult- oder fibonacciError-Nachricht, wobei er den Identifikator der Anfrage mitteilt, damit der Aufrufer die Antwort zuordnen kann.

    // 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.');
    }
    

    Auf der Hauptseite besteht die Lösung darin, den Nachrichtenaustausch in Promises umzuwandeln. Jeder Aufruf erzeugt eine eindeutige requestId, speichert die Funktionen resolve und reject des Promises in einem Map unter dieser ID und sendet anschließend die Anfrage. Wenn eine Antwort eintrifft, sucht der message-Handler nach der ID, löscht den Eintrag und erfüllt das entsprechende Promise. Im Beispiel wird das uuid-Paket für die IDs verwendet; in aktuellen Node.js-Versionen erledigt crypto.randomUUID() aus dem eingebauten node:crypto-Modul dieselbe Aufgabe ohne zusätzliche Abhängigkeiten.

    // 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
            }
        })();
    }
    

    Zusammenfassung des Musters:

    • Der Hauptthread erstellt für jede Anfrage eine requestId, damit Antworten auch dann ihren Promises zugeordnet werden können, wenn sie aus unterschiedlicher Reihenfolge eintreffen.
  • Der pendingRequests-Map enthält die resolve- und reject-Funktionen, bis der Worker antwortet.
  • Der Worker verarbeitet calculateFibonacci-Nachrichten und antwortet mit fibonacciResult oder fibonacciError, wobei die ursprüngliche requestId übergeben wird.
  • Achten Sie sorgfältig darauf, wie die Nachricht strukturiert ist, wenn Sie diesen Code anpassen. Der Hauptthread platziert requestId auf der obersten Ebene der Nachricht neben payload, während der Worker es über message.payload abruft. In dieser Form sieht der Worker requestId als undefined, wodurch keine Antwort abgeglichen werden kann und die erwartete Promise niemals abgeschlossen wird. Bewahren Sie die ID an einem einheitlich festgelegten Ort auf beiden Seiten auf (zum Beispiel payload: { n, requestId }). Es lohnt sich außerdem, einen Timeout für jede ausstehende Anfrage hinzuzufügen, damit eine verlorene Antwort zu einem Fehler wird statt zum Absturz des Prozesses. Beachten Sie schließlich den finally-Block: Sobald die Arbeit abgeschlossen ist, stoppt worker.terminate() den Thread, damit der Prozess beendet werden kann.

    Große Daten verschieben, ohne sie zu kopieren

    Alles, was über postMessage gesendet wird, wird standardmäßig strukturiert geklont, was bedeutet, dass der Empfänger eine Kopie erhält. Für kleine Nachrichten ist das in Ordnung. Bei großen binären Daten kann die Kopie selbst jedoch zum Engpass werden. Die Lösung besteht darin, statt zu kopieren das Eigentum am zugrunde liegenden Speicher zu übertragen: Das Objekt wird für den Sender unbrauchbar und sofort für den Empfänger nutzbar – ohne Duplikation.

    ArrayBuffer und MessagePort sind die am häufigsten übertragenen Objekte. SharedArrayBuffer ist ein Sonderfall: Es wird überhaupt nicht übertragen, da beide Threads gleichzeitig auf denselben Speicher zugreifen können.

    Das untenstehende Beispiel enthält zwei Skripte in einer Auflistung. Der Worker (buffer-worker.js) erhält einen Puffer, verdoppelt jeden Byte vor Ort und sendet den Puffer als übertragbares Objekt zurück. Das Hauptskript (main-buffer.js) füllt einen 1 MB großen Puffer, überträgt ihn an den Worker und liest die verarbeiteten Daten ab, sobald dieser zurückkehrt.

    // 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();
            }
        });
    }
    

    Die wichtigsten Details:

    • Der zweite Argument des postMessage-Funkts, hier [myBuffer] auf der einen Seite und [buffer] auf der anderen, ist die Übertragungsliste.
    • Wenn myBuffer gesendet wird, wechselt sein Besitzrecht an den Worker. Im Hauptthread wird der Puffer getrennt: seine byteLength-Wertigkeit sinkt auf 0 und jede Darstellung als typisierter Array, wie zum Beispiel der zuvor verwendete uint8, hat ebenfalls eine Länge von 0.
  • Der Worker ändert die Daten und überträgt sie zurück, sodass der Hauptthread wieder einen Puffer erhält.
  • Die 1 MB an Daten werden in keiner Richtung kopiert.
  • Der Kompromiss besteht darin, dass Sie nach der Übertragung die ursprüngliche Referenz nicht mehr verwenden dürfen. Wenn beide Threads die Daten tatsächlich gleichzeitig benötigen, brauchen Sie eine Kopie oder ein SharedArrayBuffer.

    Fehlerbehandlung und Aufräumen

    Worker versagen genauso wie jeder andere Code, und solange sie aktiv sind, belegen sie Speicher sowie einen Thread. Zu den verfügbaren Werkzeugen gehören:

    • worker.on('error', handler) im Hauptthread: empfängt Ausnahmen, die im Worker ausgelöst wurden und nicht gefangen wurden.
    • worker.on('exit', handler) im Hauptthread: wird ausgelöst, sobald der Worker stoppt – egal ob er normal abgeschlossen wurde oder abstürzte. Der Abbruchcode unterscheidet zwischen diesen Fällen, wobei 0 für Erfolg steht.
    • process.on('uncaughtException', handler) innerhalb des Workers: bietet die Möglichkeit, Details aus dem Worker selbst vor seinem Absturz zu protokollieren oder zu melden, zusätzlich zu dem, was der Elternteil sieht.
    • worker.terminate(): stoppt den Worker so schnell wie möglich und gibt einen Promise zurück, der mit dem Abbruchcode gelöst wird. Dabei wird nicht auf das Abschließen laufender Aufgaben gewartet; falls ein sanfter Abstieg erforderlich ist, senden Sie dem Worker eine „stop“-Nachricht und lassen Sie ihn selbst beenden. Stellen Sie stets sicher, dass Worker, die nicht mehr benötigt werden, auf irgendeine Weise gestoppt werden.

    Ausführung vieler Aufgaben über einen Worker-Pool

    Das Erstellen eines neuen Arbeiters für jede Aufgabe ist verschwenderisch, und das Laufen mehrerer CPU-intensiver Arbeitnehmer, als Kerne vorhanden sind, führt nur zu zusätzlicher Konkurrenz. Ein Worker-Pool löst beide Probleme: Er startet eine feste Anzahl von Arbeitnehmern, legt die eingehenden Aufgaben in eine Warteschlange und weist jede Aufgabe dem nächsten freien Arbeiter zu.

    Der unten dargestellte vereinfachte Pool enthält ein Array mit Aufzeichnungen der Arbeitnehmer, eine Liste der freien Worker-IDs sowie eine Warteschlange mit ausstehenden Aufgaben. runTask gibt eine Promise zurück und legt die Aufgabe in die Warteschlange; processQueue verbindet die älteste Aufgabe mit einem freien Arbeiter; wenn ein Arbeiter zurückmeldet, wird seine Promise erfüllt und er tritt wieder der Liste der freien Arbeitnehmer bei. Falls ein Arbeiter Fehler aufweist oder abnorm beendet wird, ersetzt terminateWorker ihn durch einen neuen, um die Größe des Pools beizubehalten.

    // 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;
    

    Der generische Worker, der vom Pool verwendet wird, führt für jede empfangene Nachricht eine CPU-intensive Funktion aus und antwortet mit einem status sowie einem payload. Die Summierungsloipe stellt eine beliebige reale Arbeitslast wie z. B. Bildverarbeitung oder Verschlüsselung dar.

    // 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 });
            }
        });
    }
    

    Schließlich passt das Hauptskript die Größe des Pools an die Anzahl der CPU-Kerne an, die von os.cpus() gemeldet werden, übergibt sechs Aufgaben und wartet mit Promise.allSettled auf deren Abschluss. Jede Promise einer Aufgabe wird mit einem verständlichen Erfolgs- oder Fehlertext versehen.

    // 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();
    

    Betrachten Sie diesen Pool streng als Lehrbeispiel; er weist mehrere Mängel auf, die vor dem Einsatz in der Produktion behoben werden müssen:

    • Der message-Handler des Pools erfüllt die Promise mit dem payload, unabhängig vom status; daher zählt eine im Worker fehlgeschlagene Aufgabe dennoch als Erfolg. Überprüfen Sie den status und lehnen Sie die Promise ab, wenn der Wert ‚error‘ ist.
    • worker.terminate() veranlasst einen Worker, mit einem nicht-nullen Code abzuschließen. Da der exit-Handler die Funktion terminateWorker aufruft, die einen Ersatz-Worker erstellt, führt ein Aufruf von close() zu einem Zyklus neuer Worker anstelle des Herunterfahrens des Pools. Zudem kann ein error zusammen mit dem zugehörigen exit-Event denselben Slot zweimal ersetzen und seine ID mehrfach auf die freie Liste setzen. Ein echter Pool benötigt ein „Schließungs“-Flag und sollte nur solche Worker ersetzen, die unerwartet abgestürzt sind.
  • Eine Aufgabe, die noch in Bearbeitung war, als ihr Worker abstürzte, wird nur über den error-Pfad abgelehnt; ein ungewöhnlicher Abbruch ohne error-Ereignis lässt deren Promise ausstehen.
  • In der Demo überschreiten die Summen von bis zu zwei Milliarden den Wert von Number.MAX_SAFE_INTEGER, wodurch die ausgegebenen Ergebnisse an Präzision verlieren. Verwenden Sie BigInt, wenn genaue Werte wichtig sind.
  • Das Protokoll „Doing other stuff“ wird erst nachdem alle Aufgaben abgewartet wurden ausgeführt, sodass es, wie im ersten Beispiel, keine Darstellung der Konkurrenz mit den Aufgaben zeigt.
  • Komplettere Pools fügen außerdem Inaktivitätszeitlimits, dynamisches Skalieren sowie eine sorgfältigere Wiederherstellung hinzu. Gut gepflegte Bibliotheken wie Piscina implementieren diese Funktionen, und sie sind in der Regel eine bessere Wahl als selbst erstellte Pools.

    Wann Worker-Threads hilfreich sind und wann nicht

    Gute Anwendungsfälle

    • CPU-intensive Aufgaben: alles, was erhebliche CPU-Zeit verbraucht, wie zum Beispiel aufwändige Berechnungen, Komprimierung, Verschlüsselung oder Bildverarbeitung.
    • Schutz des Event Loops: jede Operation, die sonst den Hauptthread länger blockieren würde, als man ertragen kann.

    Ungeeignete Anwendungen

    • I/O-intensive Arbeiten: Netzwerkaufrufe, Datenbankabfragen und Dateisystemoperationen sind in Node.js bereits nicht blockierend. Ihr Einbetten in einen Worker führt nur zu zusätzlichem Overhead ohne Vorteil.
    • Kleine Aufgaben: Sowohl das Starten eines Workers als auch das Übermitteln von Nachrichten kosten Zeit. Bei schnellen Berechnungen kann der Overhead die Vorteile überwiegen – das ist ein weiterer Grund, lang lebende Worker in einem Pool zu halten.
  • Gemeinsamer, veränderlicher Zustand: SharedArrayBuffer macht dies möglich, doch die Koordination von Schreibvorgängen mit Atomics ist äußerst schwierig richtig umzusetzen. Es sei denn, Sie benötigen es wirklich und verstehen die Folgen, sollten Sie lieber Nachrichtenaustausch verwenden.
  • Praktiken, die Worker gesund halten

    1. Halten Sie die Worker-Skripte klein. Laden Sie nur die Logik sowie Abhängigkeiten, die die Aufgabe benötigt. Das Importieren des gesamten Anwendungsframeworks in jeden Worker verlängert die Startzeit und erhöht den Speicherverbrauch.
    2. Wählen Sie den kostengünstigsten Datentransferweg. Strukturiertes Klonen über postMessage eignet sich für kleine Nachrichten; große ArrayBuffer-Objekte sollten über die Übertragungsliste gesendet werden. Tiefgehende oder komplexe Objektgraphen sollten vermieden werden, da ihr Klonen zeitaufwendig ist.
  • Stellen Sie die Worker still. Beenden Sie sie oder lassen Sie sie gehen, sobald ihre Arbeit abgeschlossen ist. Inaktive Worker belegen weiterhin Speicher.
  • Behandeln Sie Fehler auf beiden Seiten. Achten Sie im Hauptthread auf error und exit, und umschließen Sie die Arbeit innerhalb der Worker mit try...catch, damit Fehler als explizite Fehlermeldungen zurückgegeben werden.
  • Überwachen Sie die Ressourcennutzung. Zu viele Worker verursachen Kontextwechsel und einen höheren Speicherverbrauch, was die Leistung eher verlangsamen kann. os.cpus().length ist eine angemessene Ausgangswert für einen Pool; in neueren Node.js-Versionen ist os.availableParallelism() die empfohlene Methode, um diese Zahl zu ermitteln.
  • Blockieren Sie nicht die eigene Schleife des Workers. Ein Worker verfügt ebenfalls über eine Event-Schleife. Wenn sein message-Handler eine lange, synchrone Aufgabe ausführt, werden weitere Nachrichten hinter ihr in der Warteschlange angehäuft. Ein Worker, der viele Anfragen bearbeiten muss, sollte jede Arbeitseinheit kurz halten oder in seltenen Fällen Aufgaben an eigene Unterworker delegieren.
  • Häufige Fehlerquellen

    • Erwartung gemeinsamer Globale. Daten, die an new Worker() übergeben werden, erreichen den Worker nur über workerData (oder spätere Nachrichten). Variablen auf Modulebene im Hauptthread sind innerhalb des Workers nicht sichtbar, da dieser in seiner eigenen Isolation läuft. Die einzige echte Datenaustauschmöglichkeit besteht über SharedArrayBuffer.
  • Falsche Interpretation von isMainThread. Es beschreibt den Ausführungskontext des derzeit laufenden Codes. Eine einzige Datei kann dazu verwendet werden, entweder als Hauptskript oder als Worker zu fungieren, doch getrennte Haupt- und Worker-Dateien sind in der Regel klarer.
  • Ausgangsannahme, dass das Debuggen wie gewohnt funktioniert. Der Node.js-Inspektor kann an Worker-Threads angehängt werden, doch man startet den Prozess mit Inspektionsflaggen wie --inspect-brk und muss im Debugger zwischen den Threadkontexten wechseln. Editor wie VS Code übernehmen einen Großteil dieser Aufgabe, bleiben aber dennoch weniger direkt als das Debuggen des Hauptthreads.
  • Lässt Worker weiterlaufen. Vergessene Worker sammeln im Laufe der Zeit Speicher und Threads an und können bei langlaufenden Diensten Ressourcen erschöpfen.
  • Kernpunkte

    • Worker führen JavaScript parallel in separaten Isolaten innerhalb eines Prozesses aus; sie kommunizieren durch Nachrichtenaustausch und teilen den Speicher nur über SharedArrayBuffer.
    • Verwenden Sie sie für rechenintensive Aufgaben, die sonst den Event-Loop blockieren würden, nicht für Netzwerk- oder Festplattenzugriffe, die von Node.js bereits asynchron abgewickelt werden.
    • Geben Sie den Nachrichten eine explizite, konsistente Struktur, damit Ergebnisse, Fehler und Anfragenummern auf beiden Seiten eindeutig sind.
    • Übertragen Sie große ArrayBuffer-Objekte anstelle dessen, sie zu klonen, und bedenken Sie, dass der Sender danach keinen Zugriff mehr hat.
    • Wiederverwenden Sie Worker über einen Pool, dessen Größe den Kernzahlen des Rechners entspricht, und stellen Sie sicher, dass Herunterfahren und Fehlerbehebung klar voneinander getrennt sind – oder verwenden Sie eine bewährte Pool-Bibliothek.

    Weitere Literatur

  • Das Beenden von Anfrage-Wasserfällen in SvelteKit mit parallelen Ladefunktionen — Erfahren Sie, warum sequentielle awaits Seiten verlangsamen, und wie SvelteKit-Ladefunktionen, Promise.all, sorgfältige Aufrufe von parent() sowie gestreamte Promises diese Wasserfälle beseitigen.