Главная / Статьи / Перенос тяжелых задач, требующих мощности CPU, в Node.js с использованием worker_threads и Pools

Перенос тяжелых задач, требующих мощности CPU, в Node.js с использованием worker_threads и Pools

Узнайте, как worker_threads в Node.js обеспечивают отзывчивость цикла событий: создание рабочих процессов, обмен сообщениями между запросами и ответами, передаваемые буферы, пулы ресурсов и связанные с ними проблемы.

4814 слов

Node.js обрабатывает тысячи одновременных соединений в рамках однопоточного цикла событий, поскольку практически вся работа сводится к ожиданию ответов с сети или диска. Эта модель теряет эффективность, когда требуется реальная обработка данных: хэширование большого объема информации, изменение размера изображения или анализ набора данных занимают единственную JavaScript-поток, в результате чего все остальные запросы вынуждены ждать. Модуль worker_threads представляет собой встроенное решение этой проблемы. Прочитав этот гид, вы сможете переносить код, зависящий от производительности CPU, в рабочие потоки, эффективно обмениваться с ними данными, запускать их в виде пула и определять случаи, когда использование рабочих потоков нецелесообразно.

Почему существуют рабочие потоки

worker_threads впервые появился в Node.js 10.5.0 под экспериментальным флагом и с версии Node.js 12 считается стабильным. Он позволяет одному процессу Node.js выполнять JavaScript на нескольких потоках одновременно.

Он отличается от child_process по одному важному аспекту. Дочерний процесс — это полностью отдельный процесс Node.js с собственной памятью и собственной копией среды выполнения. Рабочий процесс находится внутри того же процесса, но это не просто «ещё одна нить, делящая всё». Каждый рабочий процесс имеет собственный изолированный интерпретатор V8, собственную кучу данных и собственный цикл событий. Ничего не делится автоматически; рабочие процессы обмениваются сообщениями, а память они делят только тогда, когда вам явно передаётся объект SharedArrayBuffer. (Часто утверждается, что рабочие процессы делят инстанцию V8 и память главной нити. Это неточно, и именно это различие объясняет большинство особенностей дизайна API ниже.) По сравнению с процессами, рабочие процессы легче запускать и с ними проще обмениваться данными, что делает их естественным способом использования нескольких ядер в рамках одной инстанции приложения.

Какова стоимость заблокированного цикла событий

Основная задача рабочего процесса — освобождать цикл событий. Пока на главном потоке выполняются синхронные вычисления, сервер не может принимать новые запросы, не может продвигаться в работе с открытыми соединениями и даже не может ответить на проверку работоспособности. Практические последствия выглядят следующим образом:

  • Ухудшение пользовательского опыта: ответы API приходят с задержкой, а такие функции в реальном времени, как обновления в прямом эфире, работают с перебоями.
  • Снижение пропускной способности: одна длительная задача занимает единственный JavaScript-поток, в результате чего количество запросов в секунду снижается.
  • Нестабильность: длительное блокирование может привести к нарушению тайм-аутов балансировщика нагрузки или клиента, и эти сбои могут распространиться на другие сервисы.

Перенос тяжелых вычислений в рабочий процесс позволяет главному циклу продолжать обслуживание сокетов и таймеров, благодаря чему приложение остается отзывчивым даже во время выполнения сложных вычислений.

Чтобы более подробно узнать о том, как взаимодействуют цикл событий, libuv и их пул потоков, ознакомьтесь с том, как на самом деле работает многозадачность в Node.js.

Использование всех ядер

У серверов обычно много ядер процессора. Без рабочих процессов или без запуска нескольких процессов с помощью child_process или менеджера процессов, такого как PM2, один процесс Node.js может использовать примерно одно ядро для выполнения JavaScript-кода, требующего высокой производительности CPU. Остальные ядра остаются неиспользуемыми. Рабочие процессы позволяют одному приложению параллельно выполнять несколько ресурсоемких задач на всех доступных ядрах, тем самым повышая общую производительность вычислений.

Задачи, которые становятся выполнимыми

Параллельные вычисления открывают возможности в тех областях, которые ранее оставались за языками с первоклассной поддержкой потоков:

  • Обработка данных: анализ в реальном времени, манипуляция изображениями и видео, а также преобразование крупных наборов данных.
  • Инференс в машинном обучении: выполнение обученных моделей при сохранении интерактивности остальной части приложения.
  • Криптография: хеширование, шифрование и дешифрование.
  • Симуляции и научные вычисления: параллельные симуляции и сложные математические модели.

Ключевой момент заключается в том, что ничто из этого не лишает вас модели асинхронного ввода-вывода. Рабочие потоки добавляют к ней возможности параллельных вычислений.

Компоненты модуля

worker_threads предоставляет небольшой набор примитивов:

  • Worker: класс, который используется основным потоком для запуска нового рабочего потока из скрипта.
  • isMainThread: булево значение, указывающее на то, выполняется ли текущий код на главном потоке; полезно тогда, когда один файл может выполнять обе роли.
  • parentPort: доступен внутри рабочего потока и представляет собой канал обратной связи с потоком, создавшим его; на главном потоке он равен null.
  • workerData: доступен внутри рабочего потока и содержит копию всех данных, предоставленных родительским потоком при его создании.
  • MessagePort и MessageChannel: используются для создания дополнительных независимых двунаправленных каналов между потоками.
  • SharedArrayBuffer и Atomics: настоящая общедоступная память вместе с примитивами синхронизации, необходимыми для её безопасного использования. Это самый эффективный вариант, но также и самый сложный, поскольку он вновь создаёт ситуации конкуренции.
  • Первый рабочий процесс: поиск простых чисел вне основного потока

    Хорошей первой задачей является намеренно медленная операция: вычисление всех простых чисел до большого предела. В примере используются два файла, находящихся в одной папке.

    Скрипт основного потока

    Основной скрипт создаёт рабочий процесс, передаёт ему входные данные и ждёт ответа. Прежде чем его прочитать, обратите внимание на три момента: рабочий процесс создаётся из отдельного пути к файлу, входные данные передаются через параметр workerData, а родительский скрипт подписывается на три события — message для получения результатов, error в случае исключений, которые рабочий процесс не смог поймать, и exit при остановке потока.

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

    Рассмотрим ключевые строки кода:

    • Проверка isMainThread обеспечивает защиту логики запуска, гарантируя, что она выполняется только на основном потоке. Ветка else существует исключительно для демонстрации того, что произойдёт, если этот же файл будет загружен в качестве рабочего процесса; в данной конфигурации она никогда не выполняется.
    • Значение largeNumber установлено достаточно высоким (20 миллионов), чтобы вычисления занимали значительное время.
  • new Worker(path.join(__dirname, 'prime-worker.js'), { workerData: { limit: largeNumber } }) — это основной вызов. Первый аргумент указывает на скрипт, который будет выполняться новым потоком; это должен быть отдельный файл. Второй аргумент — объект с настройками, где поле workerData служит начальными данными для потока.
  • workerData копируется с использованием алгоритма структурированной копирования в HTML. Поток получает собственную копию, а не ссылку на объект родителя.
  • Слушатель события message принимает любые данные, отправленные потоком; в данном случае это массив простых чисел. Слушатель события error активируется при возникновении необработанных исключений внутри потока, а слушатель события exit получает код завершения, причем значение 0 означает нормальное завершение.
  • Скрипт потока

    Рабочий процесс содержит лишь дорогостоящие вычисления и код, отвечающий за вывод результата. Функция findPrimes представляет собой простой цикл проверки на делители; она намеренно не оптимизирована, чтобы потреблять максимальное количество ресурсов процессора.

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

    Что стоит обратить внимание:

    • В ней запрашиваются переменные parentPort и workerData из модуля.
    • Условие if (parentPort) гарантирует, что логика будет выполняться только тогда, когда файл загружается в качестве рабочего процесса; если его запустить напрямую с помощью node, значение parentPort будет равно null, и в этом случае будет выведено соответствующее сообщение.
    • const { limit } = workerData; считывает данные, отправленные родительским процессом.
    • parentPort.postMessage(result) передает найденные простые числа в главный поток, где обработчик message их принимает.
  • Блок try...catch фиксирует любые сбои во время вычислений и сообщает об этом родительскому процессу, вместо того чтобы позволить потоку завершиться.
  • Есть одна нюанс: при сбое рабочий поток отправляет { error: ... } по тому же каналу message, что и для результатов, но основной скрипт рассматривает каждое сообщение как массив простых чисел. В реальном коде необходимо придать сообщениям четкую структуру (например, поля status или type), чтобы родительский процесс мог отличить результат от ошибки.

    Запуск и просмотр вывода

    Сохраните оба файла рядом друг с другом и запустите программу командой node main.js. Вы увидите, как основной поток сообщит о себе, рабочий поток — о начале работы, а спустя некоторое время появятся результаты рабочего потока вместе с указанием затраченного времени.

    Обратите внимание на порядок вывода строки «Основной поток продолжает выполнять другие задачи...». В этом скрипте она выводится внутри обработчика message, поэтому появляется только после завершения работы рабочего потока. Это демонстрирует, что основной поток смог получить и обработать результат, а не то, что он выполнял другую работу в тот же момент. Чтобы действительно увидеть, что основной поток остается отзывчивым во время вычислений, запустите что-то в основном потоке до поступления результата, например setInterval, который записывает информацию о состоянии каждые несколько сотен миллисекунд. Эта запись будет происходить продолжительное время, пока рабочий поток выполняет вычисления, чего бы не произошло, если бы findPrimes выполнялся в основном потоке.

    Двусторонняя переписка с использованием протокола запрос-ответ

    В первом примере отправляется один входной запрос и получается один результат. В большинстве реальных случаев требуется работник с длительным сроком жизни, способный обрабатывать множество запросов. Общение полностью двунаправленное: объект Worker на главном потоке и parentPort внутри работника оба предоставляют методы .postMessage(...) и .on('message', ...).

    Работник, описанный ниже, остается активным и ждет сообщений. При получении запроса calculateFibonacci он вычисляет значение с помощью примитивной рекурсивной функции (намеренно медленно) и отвечает либо сообщением fibonacciResult, либо fibonacciError, включая идентификатор запроса, чтобы вызывающий код мог сопоставить ответ.

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

    На основном потоке секрет заключается в преобразовании передачи сообщений в обещания. Каждый вызов генерирует уникальный requestId, сохраняет функции resolve и reject обещания в структуре Map под этим идентификатором, а затем отправляет запрос. Когда поступает ответ, обработчик message находит соответствующий ID, удаляет запись и реализует действие, связанное с этим обещанием. В примере для генерации ID используется пакет uuid; в текущих версиях Node.js та же функция выполняется с помощью crypto.randomUUID() из встроенного модуля node:crypto без необходимости в дополнительных зависимостях.

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

    Краткое описание паттерна:

    • Основной поток генерирует requestId для каждого запроса, что позволяет сопоставлять ответы с соответствующими обещаниями, даже если они поступают в неправильном порядке.
  • Карта pendingRequests хранит функции resolve и reject до тех пор, пока рабочий процесс не даст ответ.
  • Рабочий процесс обрабатывает сообщения calculateFibonacci и отвечает с данными fibonacciResult или fibonacciError, включающими исходный requestId.
  • Внимательно следите за структурой сообщения при адаптации этого кода. Основной поток размещает requestId на верхнем уровне сообщения рядом с payload, в то время как рабочий поток извлекает его из message.payload. В текущей реализации рабочий поток воспринимает requestId как undefined, поэтому ответ невозможно сопоставить, и ожидаемое обещание так и не будет выполнено. Храните ID в одном и том же месте у обеих сторон (например, payload: { n, requestId }). Также стоит добавить таймаут для каждого задержанного запроса, чтобы потерянный ответ приводил к ошибке, а не к застою. Наконец, обратите внимание на блок finally: как только работа завершена, worker.terminate() останавливает поток, позволяя процессу завершиться.

    Передача больших данных без их копирования

    По умолчанию всё, что отправляется с помощью postMessage, копируется в структурированном виде, что означает, что получатель получает копию. Для небольших сообщений это не проблема. Однако для крупных бинарных данных сама копия может стать препятствием. Решение заключается в передаче прав собственности на соответствующую память вместо её копирования: объект становится непригодным для отправителя и сразу же доступен для получателя без дублирования.

    ArrayBuffer и MessagePort — это наиболее часто передаваемые объекты. SharedArrayBuffer представляет собой особый случай: его вообще не передают, поскольку оба потока могут одновременно доступаться к одной и той же памяти.

    В приведённом ниже примере в одном списке представлены два скрипта. Рабочий процесс (buffer-worker.js) получает буфер, удваивает каждый байт непосредственно в нём и возвращает этот буфер в виде объекта, предназначенного для передачи. Основной скрипт (main-buffer.js) заполняет буфер объёмом 1 МБ, передаёт его рабочему процессу и считывает обработанные данные по возвращении.

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

    Ключевые детали:

    • Второй аргумент функции postMessage — здесь [myBuffer] с одной стороны и [buffer] с другой — представляет собой список для передачи.
    • При отправке myBuffer его собственность переходит к рабочему процессу. В основном потоке буфер становится отсоединённым: его значение byteLength снижается до 0, а любые представления типизированных массивов над ним, такие как ранее использовавшийся uint8, также имеют длину 0.
  • Рабочий процесс изменяет данные и возвращает их обратно, так что основной поток снова получает буфер.
  • 1 МБ данных никогда не копируется в обоих направлениях.
  • Подобный подход имеет свою цену: после передачи данных необходимо прекратить использование исходного объекта. Если оба потока действительно нуждаются в данных одновременно, требуется копия или объект SharedArrayBuffer.

    Обработка ошибок и очистка

    Рабочие процессы могут выдавать ошибки, как и любой другой код, и пока они существуют, они занимают память и поток. Доступные инструменты:

    • worker.on('error', handler) в основном потоке: принимает исключения, возникшие в рабочем процессе и не пойманные.
    • worker.on('exit', handler) на главном потоке: срабатывает каждый раз, когда рабочий процесс останавливается, независимо от того, завершился ли он нормально или вышел из строя. Код завершения позволяет различить эти случаи, причем 0 означает успех.
    • process.on('uncaughtException', handler) внутри рабочего процесса: предоставляет возможность записывать или сообщать детали изнутри рабочего процесса перед его остановкой, помимо того, что видит родительский процесс.
    • worker.terminate(): останавливает рабочий процесс как можно скорее и возвращает обещание, которое решается с кодом завершения. Он не ждет завершения текущих операций, поэтому если вам нужно плавное завершение, отправьте рабочему процессу сообщение «stop» и позвольте ему самостоятельно выйти. Всегда убедитесь, что рабочие процессы, которые больше не нужны, останавливаются каким-либо образом.

    Выполнение множества задач с помощью пула рабочих процессов

    Создание нового рабочего процесса для каждой задачи является расточительством, а запуск большего количества процессов, требующих мощности CPU, чем имеется ядер, лишь усиливает конкуренцию. Пул рабочих процессов решает оба этих проблемы: он запускает фиксированное количество процессов, помещает поступающие задачи в очередь и передаёт каждую задачу следующему свободному процессу.

    Упрощённый пул, описанный ниже, содержит массив записей о процессах, список идентификаторов свободных процессов и очередь задач, ожидающих обработки. Функция runTask возвращает обещание и добавляет задачу в очередь; функция processQueue соединяет самую старую задачу с свободным процессом; когда процесс сообщает о завершении работы, его обещание реализуется, и он снова попадает в список свободных процессов. Если процесс выдаёт ошибку или завершает работу ненормально, функция terminateWorker заменяет его на новый, чтобы поддерживать полный размер пула.

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

    Генерический рабочий процесс, используемый пулом, выполняет функцию, зависящую от мощности CPU, для каждого полученного сообщения и возвращает status вместе с payload. Цикл суммирования имитирует любую реальную нагрузку, такую как обработка изображений или шифрование.

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

    В конце основной скрипт настраивает размер пула в соответствии с количеством ядер CPU, указанных функцией os.cpus(), отправляет шесть задач и ждёт их завершения с помощью Promise.allSettled. Обещание каждой задачи преобразуется в читаемую строку с информацией о успехе или неудаче.

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

    Рассматривайте этот пул исключительно как учебный пример; в нём имеется несколько серьезных недостатков, которые должны быть устранены перед тем, как подобные решения будут использоваться в производстве:

    • Обработчик message пула реализует выполнение обещания с использованием payload, независимо от значения status; поэтому задача, которая потерпела неудачу внутри рабочего процесса, всё равно считается успешной. Необходимо проверять значение status и отклонять запрос при значении 'error'.
    • Вызов worker.terminate() приводит к выходу рабочего процесса с кодом, отличным от нуля. Поскольку обработчик exit вызывает функцию terminateWorker, которая создаёт заменяющий рабочий процесс, вызов close() приводит к циклу создания новых рабочих процессов вместо закрытия пула. Кроме того, событие error, следующее за событием exit, может дважды занять ту же ячейку и более одного раза добавить её ID в список свободных ячеек. Настоящему пулу требуется флаг «закрытия», и он должен заменять рабочие процессы только в том случае, если они вышли из строя неожиданно.
  • Задача, которая находилась в процессе выполнения при сбое рабочей нити, отклоняется только через путь error; ненормальный выход без события error оставляет её обещание в состоянии ожидания.
  • В демонстрации суммы, превышающие два миллиарда, выходят за пределы Number.MAX_SAFE_INTEGER, поэтому выводимые результаты теряют точность. Используйте BigInt, если важны точные значения.
  • Лог «Выполнение других операций» записывается после того, как все задачи будут завершены, поэтому, как и в первом примере, он не отражает одновременное выполнение задач.
  • Более полноценные пулы также включают таймауты бездействия, динамическое масштабирование и более тщательную рекуперацию. Хорошо поддерживаемые библиотеки, такие как Piscina, реализуют эти функции, и они обычно являются лучшим выбором, чем самостоятельно созданный пул.

    Когда рабочие нити полезны, а когда — нет

    Подходящие случаи применения

    • Задачи, зависящие от производительности CPU: всё, что потребляет значительное количество времени процессора, такое как сложные вычисления, сжатие, шифрование или обработка изображений.
    • Защита цикла событий: любые операции, которые иначе могли бы заблокировать основной поток на длительное время, превышающее допустимые пределы.

    Несоответствующие случаи применения

    • Работа, зависящая от ввода-вывода: сетевые запросы, операции над базами данных и действия в файловой системе уже являются неблокирующими в Node.js. Использование рабочих процессов для их выполнения добавляет нагрузку без какой-либо пользы.
    • Короткие задачи: запуск рабочего процесса и передача сообщений оба занимают время. Для быстрых вычислений эта дополнительная нагрузка может превзойти получаемую выгоду, что ещё одна причина хранить долгоживущих рабочих процессов в пуле.
  • Общее изменяемое состояние: SharedArrayBuffer позволяет это сделать, но синхронизация записей с помощью Atomics чрезвычайно сложна для правильной реализации. Если у вас нет настоящей необходимости в этом и вы не понимаете последствий, лучше использовать передачу сообщений.
  • Практики, сохраняющие работоспособность воркеров

    1. Сохраняйте скрипты воркеров небольшими. Загружайте только ту логику и зависимости, которые требуются задаче. Импортирование всей фреймворковой структуры приложения в каждый воркер увеличивает время запуска и потребление памяти.
    2. Выбирайте самый экономичный способ передачи данных. Структурированное копирование с помощью postMessage подходит для небольших сообщений; большие объекты типа ArrayBuffer следует передавать через список передачи, а глубокие или сложные объектные структуры лучше избегать, так как их копирование медленно.
  • Останавливайте рабочие процессы. Завершайте их или позволяйте им выходить, когда работа завершена. Бездействующие процессы продолжают занимать память.
  • Обрабатывайте ошибки с обеих сторон. Наблюдайте за сигналами error и exit на главном потоке, а также оборачивайте выполнение кода внутри рабочих процессов в блоки try...catch, чтобы сбои возвращались в виде четких сообщений об ошибках.
  • Контролируйте использование ресурсов. Чрезмерное количество рабочих процессов приводит к частой смене контекста и увеличению потребления памяти, что может замедлить работу вместо того, чтобы ускорить ее. os.cpus().length — разумное начальное значение для количества процессов в пуле; в более новых версиях Node.js рекомендуется использовать os.availableParallelism() для получения этого числа.
  • Не блокируйте собственный цикл обработки задач рабочего процесса. У каждого рабочего процесса тоже есть цикл событий. Если обработчик message выполняет длительную синхронную задачу, последующие сообщения накапливаются в очереди. Рабочему процессу, который должен обрабатывать множество запросов, следует делать каждую единицу работы короткой или, в редких случаях, делегировать её своим подрабочим процессам.
  • Частые ошибки, которые мешают разработчикам

    • Ожидание совместного использования глобальных переменных. Данные, передаваемые в new Worker(), доходят до рабочего процесса только через workerData (или последующие сообщения). Переменные уровня модуля на основном потоке недоступны внутри рабочего процесса, поскольку тот работает в своей изолированной среде. Единственный способ истинного обмена данными — через SharedArrayBuffer.
  • Неправильное понимание isMainThread. Этот параметр описывает контекст выполнения кода, который в данный момент выполняется. Один файл может использоваться как основной скрипт, так и в качестве рабочей нити, однако отдельные файлы для основного и рабочего кода обычно делают структуру более понятной.
  • Предположение, что отладка работает как обычно. Инспектор Node.js может подключаться к рабочим нитям, но для этого процесс запускается с флагами инспектора, такими как --inspect-brk, и требуется переключаться между контекстами нитей в отладчике. Редакторы вроде VS Code автоматизируют большую часть этого процесса, однако это всё равно менее удобно, чем отладка основной нити.
  • Оставление рабочих нитей в активном состоянии. Забытые рабочие нити со временем накапливают память и ресурсы, что может привести к исчерпанию ресурсов в долгосрочно работающих сервисах.
  • Основные выводы

    • Рабочие процессы выполняют JavaScript параллельно в отдельных изоляциях внутри одного процесса; они обмениваются данными с помощью передачи сообщений и делят память только через SharedArrayBuffer.
    • Используйте их для задач, требующих больших вычислительных ресурсов CPU, которые иначе заблокируют цикл событий, а не для работы с сетью или диском, которую Node.js уже обрабатывает асинхронно.
    • Дайте сообщениям четкую, единообразную структуру, чтобы результаты, ошибки и идентификаторы запросов были однозначными с обеих сторон.
    • Передавайте большие объекты ArrayBuffer вместо их копирования, помня при этом, что отправитель теряет к ним доступ после этого.
    • Повторно используйте рабочие процессы через пул, размер которого соответствует количеству ядер машины, и убедитесь, что процессы выключения и восстановления после сбоев четко разделены, либо используйте проверенную библиотеку для управления пулом.

    Связанная литература

  • Прекращение последовательности запросов в SvelteKit с использованием функций параллельной загрузки — Узнайте, почему последовательные вызовы awaits замедляют загрузку страниц, и как функции загрузки SvelteKit, Promise.all, тщательные вызовы parent() и потоковые обещания устраняют эту проблему.