Transferir tareas que consumen mucho la CPU en Node.js con worker_threads y Pools
Aprenda cómo los worker_threads de Node.js mantienen el bucle de eventos receptivo: creación de trabajadores, comunicación solicitud-respuesta, buffers transferibles, pools y sus problemas.
Node.js gestiona miles de conexiones concurrentes en un bucle de eventos de un solo hilo porque casi todo ese trabajo consiste en esperar a que la red o el disco respondan. Este modelo deja de funcionar en cuanto una solicitud requiere cálculos reales: el hasheo de un gran volumen de datos, el redimensionado de una imagen o el procesamiento intensivo de un conjunto de datos mantienen ocupado al único hilo de JavaScript, mientras que todas las demás solicitudes esperan su turno. El módulo worker_threads es la solución integrada para este problema. Después de estudiar esta guía, podrá trasladar código que depende intensamente del procesador a los workers, intercambiar datos con ellos de manera eficiente, ejecutarlos en grupos y reconocer los casos en los que un worker no es la herramienta adecuada.
Por qué existen los threads de trabajo
worker_threads apareció por primera vez en Node.js 10.5.0 bajo una bandera experimental y se ha considerado estable desde la línea Node.js 12. Permite que un único proceso de Node.js ejecute JavaScript en varios hilos al mismo tiempo.
Se diferencia de child_process de una manera importante. Un proceso hijo es un proceso completamente separado de Node.js con su propia memoria y su propia copia del entorno de ejecución. Un worker vive dentro del mismo proceso, pero no es simplemente “otro hilo que comparte todo”. Cada worker cuenta con su propio aislamiento V8, su propio montón de memoria y su propio bucle de eventos. Nada se comparte de forma implícita; los workers se comunican mediante el envío de mensajes y solo comparten memoria cuando se les proporciona explícitamente un SharedArrayBuffer. (Una afirmación frecuente es que los workers comparten la instancia V8 y la memoria del hilo principal. Eso no es cierto, y esa diferencia explica la mayor parte del diseño de la API que se describe a continuación.) En comparación con los procesos, los workers son más económicos de iniciar y comunicar, lo que los convierte en la forma natural de utilizar varios núcleos desde una única instancia de aplicación.
El costo que te impone un bucle de eventos bloqueado
La tarea principal de un worker es mantener libre el bucle de eventos. Mientras se ejecuta un cálculo síncrono en el hilo principal, el servidor no puede aceptar nuevas solicitudes, no puede avanzar en las conexiones abiertas e incluso no puede responder a una verificación de estado. Las consecuencias prácticas son las siguientes:
- Experiencia de usuario deteriorada: las respuestas de la API llegan tarde y las funciones en tiempo real, como las actualizaciones en vivo, se interrumpen.
- Menor rendimiento: una tarea larga ocupa el único hilo de JavaScript, por lo que disminuyen las solicitudes por segundo.
- Inestabilidad: el bloqueo prolongado puede superar los tiempos de espera del balanceador de carga o del cliente, y esas fallas pueden propagarse a otros servicios.
Al trasladar los cálculos intensivos a un worker, el bucle principal queda libre para seguir gestionando sockets y temporizadores, de modo que la aplicación sigue siendo receptiva incluso mientras se realizan cálculos complejos.
Para conocer en mayor profundidad cómo se integran el bucle de eventos, libuv y su piscina de hilos, consulte cómo funciona la concurrencia en Node.js en el fondo.
Uso de todos los núcleos
Los servidores suelen tener muchos núcleos de CPU. Sin workers, o sin ejecutar varios procesos mediante child_process o un gestor de procesos como PM2, un único proceso de Node.js puede utilizar aproximadamente un núcleo para el JavaScript con carga intensiva en la CPU. El resto permanece inactivo. Los workers permiten que una aplicación ejecute varias tareas pesadas en paralelo a través de los núcleos disponibles, aumentando así el rendimiento computacional total.
Cargas de trabajo que se vuelven prácticas
El cómputo paralelo abre campos que tradicionalmente quedaban reservados para lenguajes con hilos de primera clase:
- Procesamiento de datos: análisis en tiempo real, manipulación de imágenes y videos, y transformaciones de grandes conjuntos de datos.
- Inferencia de aprendizaje automático: ejecución de modelos entrenados mientras el resto de la aplicación permanece interactiva.
- Criptografía: hashing, cifrado y descifrado.
- Simulación e informática científica: simulaciones paralelas y modelos matemáticos complejos.
Lo importante es que nada de esto le hace perder el modelo de E/S asíncrona. Los workers añaden cálculos paralelos sobre él.
Los componentes básicos del módulo
worker_threads expone un pequeño conjunto de primitivas:
Worker: la clase que el hilo principal utiliza para iniciar un nuevo worker desde un script.
isMainThread: un booleano que indica si el código actual se está ejecutando en el hilo principal, útil cuando un mismo archivo puede desempeñar ambos roles.parentPort: disponible dentro de un worker, es el canal que conecta con el hilo que lo creó. Tiene el valor null en el hilo principal.workerData: disponible dentro de un worker, alberga una copia de los datos que el padre proporcionó al crearlo.MessagePort y MessageChannel: se utilizan para crear canales bidireccionales e independientes adicionales entre hilos.SharedArrayBuffer y Atomics: memoria compartida real junto con los primitivos de sincronización necesarios para utilizarla de forma segura. Son la opción más eficiente, pero también la más compleja, ya que reintroducen condiciones de carrera.Un primer trabajador: encontrar números primos fuera del hilo principal
Un buen ejercicio inicial es una tarea deliberadamente lenta: calcular todos los números primos hasta un límite grande. El ejemplo utiliza dos archivos que se encuentran en la misma carpeta.
El script del hilo principal
El script principal crea al trabajador, le proporciona los datos de entrada y espera la respuesta. Antes de leerla, observe tres cosas: el trabajador se crea a partir de una ruta de archivo separada; los datos de entrada se envían mediante la opción workerData; y el padre se suscribe a tres eventos: message para los resultados, error para las excepciones que el trabajador no capturó, y exit cuando el hilo se detiene.
// 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);
}
Analizando las líneas importantes:
- La comprobación
isMainThreadprotege la lógica de creación del trabajador para que solo se ejecute en el hilo principal. La ramaelseexiste únicamente para ilustrar qué ocurriría si este mismo archivo se cargara como trabajador; en esta configuración nunca se ejecuta. largeNumberse establece en un valor lo suficientemente alto (20 millones) como para que el cálculo tarde tiempo considerable en completarse.
new Worker(path.join(__dirname, 'prime-worker.js'), { workerData: { limit: largeNumber } }) es la llamada principal. El primer argumento indica el script que ejecutará el nuevo hilo, el cual debe ser un archivo independiente. El segundo es un objeto de opciones cuyo workerData se convierte en la entrada inicial del worker.workerData se copia utilizando el algoritmo de clonación estructurada de HTML. El worker recibe su propia copia, y no una referencia al objeto del padre.message recibe todo lo que envíe el worker, en este caso el array de números primos. El listener error se activa ante excepciones no capturadas dentro del worker, y el listener exit recibe un código de salida, donde 0 indica una finalización normal.El script del worker
El worker solo contiene el cálculo costoso y el código que informa sobre su resultado. La función findPrimes es un bucle simple de división por prueba, intencionalmente sin optimizar para que consuma tiempo de 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.");
}
Aspectos a tener en cuenta:
- Obtiene
parentPortyworkerDatadel módulo. - La condición
if (parentPort)asegura que la lógica solo se ejecute cuando el archivo se carga como worker; si se ejecuta directamente connode,parentPortesnully en su lugar se imprime un mensaje. const { limit } = workerData;lee la entrada enviada por el padre.parentPort.postMessage(result)envía los números primos al hilo principal, donde el manejador demessagelos recibe.
try...catch registra cualquier fallo durante el cálculo y lo informa al proceso padre en lugar de permitir que el hilo se detenga.Una sutileza: en caso de fallo, el proceso trabajador publica { error: ... } en el mismo canal message utilizado para los resultados, pero el script principal trata cada mensaje como un array de números primos. En código real, hay que dar a los mensajes una estructura explícita (por ejemplo, un campo status o type) para que el proceso padre pueda distinguir entre un resultado y un error.
Ejecutarlo y leer la salida
Guarde ambos archivos uno al lado del otro e inicie el programa con node main.js. Verá cómo el hilo principal se anuncia, cómo el proceso trabajador informa que ha comenzado y, después de algún tiempo, el resultado del proceso trabajador seguido del tiempo transcurrido.
Preste atención al orden de la línea “Hilo principal continuando con otras tareas...”. En este script se imprime dentro del manejador message, por lo que solo aparece después de que el trabajador finaliza. Esto demuestra que el hilo principal pudo recibir y procesar el resultado, no que estuviera realizando otras tareas mientras tanto. Para ver realmente cómo el hilo principal permanece receptivo durante el cálculo, inicie alguna acción en él antes de que llegue el resultado, como un setInterval que registre una señal de vida cada pocos cientos de milisegundos. La señal de vida seguirá activa mientras el trabajador realiza los cálculos, lo cual es exactamente lo que no ocurriría si findPrimes se ejecutara en el hilo principal.
Comunicación bidireccional con un protocolo de solicitud-respuesta
El primer ejemplo envía una entrada y recibe un resultado. La mayoría de los usos reales requieren un proceso que permanezca activo durante mucho tiempo y pueda manejar muchas solicitudes. La comunicación es completamente bidireccional: un objeto Worker en el hilo principal y parentPort dentro del proceso trabajador ofrecen tanto .postMessage(...) como .on('message', ...).
El proceso trabajador mostrado a continuación permanece activo esperando mensajes. Cuando recibe una solicitud calculateFibonacci, calcula el valor con una función recursiva simple (intencionalmente lenta) y responde con un mensaje fibonacciResult o fibonacciError, incluyendo el identificador de la solicitud para que quien la envió pueda hacerla coincidir con la respuesta.
// 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.');
}
En el lado principal, la estrategia consiste en convertir el intercambio de mensajes en promesas. Cada llamada genera un requestId único, almacena las funciones resolve y reject de la promesa en un Map bajo ese ID y envía la solicitud. Cuando llega una respuesta, el manejador de message busca el ID, elimina la entrada correspondiente y resuelve la promesa asociada. El ejemplo utiliza el paquete uuid para los IDs; en las versiones actuales de Node.js, crypto.randomUUID() del módulo integrado node:crypto realiza la misma función sin necesidad de dependencias.
// 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 resumen, el patrón es el siguiente:
- El hilo principal crea un
requestIdpara cada solicitud, de modo que las respuestas puedan asociarse a sus promesas incluso si llegan fuera de orden.
pendingRequests almacena las funciones resolve y reject hasta que el worker responda.calculateFibonacci y responde con fibonacciResult o fibonacciError, incluyendo el requestId original.Si adapta este código, preste atención cuidadosamente a la estructura del mensaje. El hilo principal coloca requestId en el nivel más alto del mensaje, junto a payload, mientras que el worker lo obtiene desde message.payload. Tal como está escrito, el worker considera requestId como undefined, por lo que no se puede asociar la respuesta y la promesa esperada nunca se resuelve. Mantenga el ID en un lugar acordado por ambas partes (por ejemplo, payload: { n, requestId }). También es recomendable agregar un tiempo de espera para cada solicitud pendiente, de modo que una respuesta perdida se convierta en un error en lugar de causar un bloqueo. Finalmente, observe el bloque finally: una vez completada la tarea, worker.terminate() detiene el hilo para que el proceso pueda finalizar.
Mover grandes cantidades de datos sin copiarlos
Todo lo que se envía a través de postMessage se clona de forma estructurada por defecto, lo que significa que el destinatario recibe una copia. Para mensajes pequeños eso está bien. Pero con grandes volúmenes de datos binarios, la propia copia puede convertirse en un cuello de botella. La solución es transferir la propiedad de la memoria subyacente en lugar de copiarla: el objeto deja de ser utilizable para el remitente y se vuelve inmediatamente utilizable para el destinatario, sin duplicación alguna.
ArrayBuffer y MessagePort son los objetos que se transfieren con mayor frecuencia. SharedArrayBuffer es un caso distinto: no se transfiere en absoluto, ya que ambos hilos pueden acceder a la misma memoria simultáneamente.
El ejemplo a continuación contiene dos scripts en una sola lista. El worker (buffer-worker.js) recibe un buffer, duplica cada byte en el mismo lugar y envía el buffer de vuelta como transferible. El script principal (main-buffer.js) llena un buffer de 1 MB, lo transfiere al worker y lee los datos procesados cuando este lo devuelve.
// 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();
}
});
}
Los detalles clave:
- El segundo argumento de
postMessage, aquí[myBuffer]por un lado y[buffer]por el otro, es la lista de transferencia. - Cuando se envía
myBuffer, su propiedad pasa al worker. En el hilo principal, el buffer queda desacoplado: subyteLengthdisminuye a 0 y cualquier vista de tipo array sobre él, como eluint8mencionado anteriormente, también tiene longitud 0.
El compromiso es que debes dejar de usar la referencia original después del transferencia. Si ambos hilos realmente necesitan los datos al mismo tiempo, necesitas una copia o un SharedArrayBuffer.
Manejo de errores y limpieza
Los trabajadores fallan como cualquier otro código, y mantienen memoria y un hilo mientras están activos. Las herramientas disponibles son:
worker.on('error', handler)en el hilo principal: recibe las excepciones que se lanzaron en el trabajador y que nunca fueron capturadas.worker.on('exit', handler)en el hilo principal: se dispara cada vez que el worker deja de funcionar, ya sea porque finalizó normalmente o se cayó. El código de salida distingue entre ambos casos, siendo0el valor para un éxito.process.on('uncaughtException', handler)dentro del worker: ofrece un lugar para registrar o informar detalles desde el interior del worker antes de que deje de funcionar, además de lo que ve el proceso padre.worker.terminate(): detiene el worker lo antes posible y devuelve una promesa que se resuelve con el código de salida. No espera a que termine el trabajo en curso, por lo que si necesita un cierre ordenado, envíe al worker un mensaje de “stop” y déjelo salir por sí mismo. Siempre asegúrese de detener los workers que ya no se necesitan de alguna manera.
Ejecutar muchas tareas a través de un grupo de workers
Crear un trabajador nuevo para cada tarea es un desperdicio, y ejecutar más trabajadores que los núcleos disponibles solo genera competencia por los recursos. Un grupo de trabajadores resuelve ambos problemas: inicia un número fijo de trabajadores, coloca las tareas recibidas en una cola y asigna cada tarea al siguiente trabajador disponible.
El grupo simplificado que se muestra a continuación mantiene un array de registros de trabajadores, una lista de IDs de trabajadores libres y una cola de tareas pendientes. runTask devuelve una promesa y agrega la tarea a la cola; processQueue empareja la tarea más antigua con un trabajador libre; cuando un trabajador envía una respuesta, su promesa se resuelve y vuelve a formar parte de la lista de trabajadores libres. Si un trabajador presenta errores o se cierra de forma anormal, terminateWorker lo reemplaza por uno nuevo para mantener el tamaño total del grupo.
// 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;
El trabajador genérico utilizado por el pool ejecuta una función que depende de la CPU para cada mensaje que recibe y responde con un status y un payload. El bucle de sumación simula cualquier carga de trabajo real, como el procesamiento de imágenes o la encriptación.
// 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 });
}
});
}
Finalmente, el script principal ajusta el tamaño del pool al número de núcleos de CPU indicados por os.cpus(), envía seis tareas y espera a que todas se completen con Promise.allSettled. La promesa de cada tarea se traduce en una cadena legible que indica si tuvo éxito o falló.
// 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();
Trate este pool estrictamente como un ejemplo didáctico; presenta varias deficiencias importantes antes de que algo similar pueda utilizarse en producción:
- El manejador de
messagede la piscina resuelve la promesa con elpayloadindependientemente delstatus; por lo tanto, una tarea que falló dentro del worker sigue contando como un éxito. Verifique elstatusy rechace la operación si es'error'. worker.terminate()hace que un worker salga con un código distinto de cero. Dado que el manejador deexitllama aterminateWorker, lo cual genera uno nuevo en su lugar, llamar aclose()desencadena un ciclo de workers nuevos en lugar de cerrar la piscina. Además, unerrorseguido de su eventoexitpuede reemplazar el mismo slot dos veces e insertar su ID en la lista de disponibles más de una vez. Una piscina real necesita una bandera de “cierre” y solo debería reemplazar a los workers que fallen inesperadamente.
error; una salida anormal sin un evento error deja su promesa en estado pendiente.Number.MAX_SAFE_INTEGER, por lo que los resultados impresos pierden precisión. Utilice BigInt si son importantes los valores exactos.Los pools más completos también incluyen tiempos de espera para procesos inactivos, escalado dinámico y una recuperación más cuidadosa. Las bibliotecas bien mantenidas como Piscina implementan estos detalles, y suelen ser una mejor opción que un pool creado manualmente.
Cuándo ayudan los hilos de trabajo y cuándo no
Casos adecuados
- Tareas dependientes de la CPU: todo aquello que consume una cantidad significativa de tiempo de procesamiento, como cálculos intensivos, compresión, encriptación o procesamiento de imágenes.
- Proteger el bucle de eventos: cualquier operación que, de lo contrario, bloquearía el hilo principal por más tiempo del que se puede tolerar.
Casos inadecuados
- Trabajo dependiente de E/S: las llamadas a red, las consultas a bases de datos y las operaciones del sistema de archivos ya son no bloqueantes en Node.js. Envolverlas en un worker añade sobrecarga sin ningún beneficio.
- Tareas pequeñas: iniciar un worker y pasar mensajes ambos requieren tiempo. En cálculos rápidos, la sobrecarga puede superar al beneficio, lo cual es otra razón para mantener los workers de larga duración en un pool.
SharedArrayBuffer lo hace posible, pero coordinar las escrituras con Atomics es extremadamente difícil de hacer bien. A menos que realmente lo necesites y entiendas las consecuencias, prefiere el intercambio de mensajes.Prácticas que mantienen sanos a los workers
- Mantén los scripts de los workers pequeños. Carga solo la lógica y las dependencias que necesite la tarea. Importar todo el framework de tu aplicación en cada worker aumenta el tiempo de inicio y el consumo de memoria.
- Elegir la vía de datos más económica. La clonación estructurada a través de
postMessagees adecuada para mensajes pequeños; los grandesArrayBufferdeben pasar por la lista de transferencias, y los gráficos de objetos profundos o complejos deben evitarse ya que su clonación es lenta.
error y exit en el hilo principal, y envuelve la tarea dentro del trabajador con try...catch para que los fallos se manifiesten como mensajes de error explícitos.os.cpus().length es un valor inicial razonable para un grupo de trabajadores; en versiones más recientes de Node.js, os.availableParallelism() es la forma recomendada para obtener ese número.message ejecuta una tarea síncrona prolongada, los mensajes siguientes se acumulan detrás de ella. Un worker que debe manejar muchas solicitudes debería mantener cada unidad de trabajo breve o, en casos excepcionales, delegarla a sus propios subworkers.Errores que suelen causar problemas
- Expectar variables globales compartidas. Los datos proporcionados a
new Worker()llegan al worker únicamente a través deworkerData(o mensajes posteriores). Las variables a nivel de módulo en el hilo principal no son visibles dentro del worker, ya que este se ejecuta de forma aislada. El único verdadero mecanismo de compartición esSharedArrayBuffer.
isMainThread. Este parámetro describe el contexto de ejecución del código que se está ejecutando en ese momento. Un mismo archivo puede utilizarse tanto como script principal como como hilo de trabajo, pero suele ser más claro tener archivos separados para cada uno.--inspect-brk y cambiar entre los contextos de hilo en el depurador. Editores como VS Code gestionan gran parte de esto, pero sigue siendo menos directo que depurar el hilo principal.Puntos clave
- Los workers ejecutan JavaScript en paralelo en entornos aislados dentro de un mismo proceso; se comunican mediante el intercambio de mensajes y comparten memoria únicamente a través de
SharedArrayBuffer. - Úsalos para tareas que consumen mucho recursos de la CPU y que de lo contrario bloquearían el bucle de eventos, no para accesos a red o disco que Node.js ya maneja de forma asíncrona.
- Dale a los mensajes una estructura explícita y consistente para que los resultados, errores e IDs de solicitud sean inequívocos en ambos lados.
- Transfiere grandes
ArrayBufferen lugar de clonarlos, y recuerda que el remitente pierde acceso a ellos posteriormente. - Reutiliza los workers a través de un pool del tamaño adecuado según los núcleos de la máquina, y asegúrate de que el cierre del sistema y la recuperación en caso de fallo estén claramente separados, o utiliza una biblioteca de pool probada.
Lecturas relacionadas
- Domar la concurrencia de Node.js: Evitando colapsos de API con p-map y Bottleneck — Aprenda cómo combinar p-map y Bottleneck en Node.js para prevenir errores de límite de velocidad y sobrecarga del sistema al controlar la concurrencia y el tiempo de las solicitudes.
- Explicación de la concurrencia en Node.js: libuv, el bucle de eventos y el pool de hilos — Aprenda cómo Node.js utiliza las primitivas del sistema operativo de libuv y su pool de hilos de trabajo para manejar E/S asíncrona, además de los problemas comunes en el pool de hilos y consejos para su ajuste.