Перейти к содержимому

Многопоточность: cluster, worker_threads, child_process

Node.js исполняет твой JS в одном потоке. Это фича, а не баг: Event Loop отлично обслуживает тысячи I/O-соединений, пока колбеки короткие. Но упираешься в CPU — и «неблокируемость» заканчивается. На 8-ядерной машине один процесс Node использует одно ядро, и остальные семь простаивают. Плюс любой тяжёлый расчёт в колбеке замораживает цикл для всех клиентов одновременно.

Краткая версия перечислила три инструмента. Здесь — когда и какой брать, как передавать данные между потоками без оверхеда, как организовать пул воркеров и почему exec с пользовательским вводом — это CVE в твоём коде.

Инструмент Единица Изоляция Для чего
cluster Процесс Полная (память, GC) Занять все ядра HTTP-сервером
worker_threads Поток (свой Event Loop) Собственный heap, общая память доступна CPU-нагрузка: хеши, сжатие, парсинг
child_process Процесс ОС Максимальная Внешние программы: ffmpeg, python, git

Правило выбора: HTTP-сервис → реплики (cluster/PM2/Kubernetes); тяжёлый расчёт внутри своего кода → worker_threads; нужна чужая программа → child_process.

Модуль cluster поднимает N процессов-воркеров, которые слушают один порт: мастер принимает соединение и отдаёт его воркеру (на Linux по умолчанию round-robin). Память у каждого воркера своя — утечка в одном не убивает остальные; воркер упал — мастер поднимает новый.

import cluster from 'node:cluster';
import { availableParallelism } from 'node:os';
import http from 'node:http';
if (cluster.isPrimary) {
const n = availableParallelism();
console.log(`мастер ${process.pid}: поднимаю ${n} воркеров`);
for (let i = 0; i < n; i++) cluster.fork();
cluster.on('exit', (worker, code, signal) => {
console.log(`воркер ${worker.process.pid} умер (${signal || code}), поднимаю новый`);
cluster.fork();
});
} else {
// каждый воркер — полноценный сервер на общем порту
http.createServer((req, res) => {
res.end(`ответ от воркера ${process.pid}\n`);
}).listen(3000);
}

Важные нюансы: воркеры не делят состояние — сессии в памяти работать не будут (нужен Redis); долгоживущие соединения (WebSocket) привязываются к конкретному воркеру; graceful shutdown делается в каждом воркере отдельно (сигналы сначала мастеру — он должен ретранслировать).

Окно терминала
# PM2 — то же самое без кода в проекте
pm2 start server.js -i max # по числу ядер
pm2 start server.js -i 4 # фиксированно 4 реплики
pm2 reload server.js # zero-downtime перезапуск

Воркер — это отдельный поток ОС со своим V8-изолятом, своим Event Loop и своей кучей, живущий внутри одного процесса (обзор API — в документации worker_threads). Общего состояния «по умолчанию» нет — общение идёт через сообщения.

// hasher.js — файл воркера
import { parentPort, workerData } from 'node:worker_threads';
import { createHash } from 'node:crypto';
import { readFileSync } from 'node:fs';
const file = readFileSync(workerData.file);
const hash = createHash('sha256').update(file).digest('hex');
parentPort.postMessage(hash); // отсылаем результат в мейн
main.js
import { Worker } from 'node:worker_threads';
const worker = new Worker('./hasher.js', {
workerData: { file: 'big.iso' }, // данные на старте — клонируются
});
worker.on('message', (hash) => console.log('sha256:', hash));
worker.on('error', (err) => console.error('воркер упал:', err));
worker.on('exit', (code) => {
if (code !== 0) console.error(`воркер завершился с кодом ${code}`);
});

Ключевой момент — как передаются данные:

  1. workerData и postMessage — объекты сериализуются структурно клонированием (structured clone). Подходит для JSON-подобных структур. Большой Buffer будет скопирован — 100 МБ через postMessage = 100 МБ копирования.
  2. transferList — передача владения без копирования: ArrayBuffer «переезжает» в воркер, в мейне становится невалидным (byteLength → 0). Единственный способ передать большой буфер быстро.
// передаём гигабайт без копирования
const ab = new ArrayBuffer(1024 ** 3);
const worker = new Worker('./processor.js');
worker.postMessage({ buffer: ab }, [ab]); // transfer list
console.log(ab.byteLength); // 0 — память больше не наша
  1. SharedArrayBuffer — общая память, видимая обоим потокам одновременно, без передачи вообще. Ноль копирований, но появляется всё, что есть в многопоточном программировании: гонки, видимость, атомарность. Координация — через Atomics (Atomics.wait, Atomics.notify, Atomics.add).
main.js
const sab = new SharedArrayBuffer(4);
const worker = new Worker('./counter.js', { workerData: sab });
new Int32Array(sab)[0] = 42;
worker.postMessage('go'); // воркер прочитает 42 из общей памяти

Создавать Worker на каждый запрос дорого (порождение изолята — десятки миллисекунд). Правильный паттерн — пул фиксированного размера, задания раздаются по очереди. Не пиши свой пул вручную (очередь, освобождение, ретраи — там подводных камней хватает), возьми Piscina:

import Piscina from 'piscina';
const pool = new Piscina({
filename: new URL('./image-resize.js', import.meta.url).pathname,
maxThreads: 4, // обычно = число ядер минус один под Event Loop
});
// API как у обычного вызова функции
const thumbnails = await pool.run({ input: 'photo.jpg', width: 200 });

Piscina следит за очередью, перезапускает упавших воркеров, поддерживает transferList и Named Tasks. Впрочем, для обучения стоит один раз написать пул из 30 строк на Worker + массиве, чтобы понять, что Piscina делает за тебя.

Три основных способа, разница — в том, как заворачивается команда и как возвращается результат (детали и опции — в документации child_process):

Функция Шелл Вывод Когда
execFile Нет (массив аргументов) Буфером целиком Короткие команды, файлы, ffmpeg
spawn Нет (массив аргументов) Стримами Долгие процессы, большой вывод
exec Да (/bin/sh -c) Буфером целиком Только доверенные команды без ввода
fork Нет IPC-канал (postMessage) Запуск другого Node-скрипта
import { execFile } from 'node:child_process';
import { promisify } from 'node:util';
const execFileAsync = promisify(execFile);
// безопасно: argv-массив, шелл не участвует
const { stdout } = await execFileAsync('ffmpeg', [
'-i', 'input.mp4', '-vn', '-acodec', 'libmp3lame', 'output.mp3',
]);

spawn — для процессов, живущих долго или пишущих много (сборка, рендер, стриминг): вывод приходит чанками через события stdout/stderr, а не копится в буфере. exec и execFile копят весь вывод в память — дефолтный лимит 1 МБ, превышение убивает процесс с ERR_CHILD_PROCESS_STDIO_MAXBUFFER.

fork — частный случай для Node: запускает другой JS-файл в новом процессе и даёт канал IPC (process.send / process.on('message')) — вместо worker_threads, когда нужна настоящая изоляция (падающий форк не роняет мейн) или другая версия рантайма.

  1. CPU-работа в HTTP-хендлере. Тяжёлый цикл/хеширование/парсинг большого JSON в колбеке — цикл забит, все соединения ждут. Класс ошибки «работало на ноутбуке». Лечение: worker_threads или вынос в очередь (см. System Design).
  2. Воркер на запрос. new Worker() внутри HTTP-хендлера под нагрузкой: каждый запрос порождает изолят — CPU уходит на создание/уничтожение потоков, а не на работу. Лечение: пул (Piscina) с размером ~числу ядер.
  3. exec с шаблонными строками. Даже «безобидный» exec(\git log ${branch}`)— инъекция, еслиbranchприходит извне. Лечение:execFile(‘git’, [‘log’, branch])`.
  4. Игнор error и exit у воркеров. Упавший воркер без обработчика — необработанное исключение в мейн-процессе или молча завершившийся поток. Всегда: worker.on('error'), worker.on('exit') с проверкой кода.
  5. Большие данные через postMessage. Передача 100 МБ Buffer через postMessage копирует 100 МБ — под нагрузкой это заметнее самой обработки. Лечение: transferList для разовых передач, SharedArrayBuffer для интенсивного обмена.
  6. Считать availableParallelism свободными ресурсами. На shared-хостинге и в контейнерах с лимитами CPU это число может быть равно числу ядер хоста, а не твоей квоте. В Kubernetes лучше управлять репликами через replicas, а не кластером внутри пода.
  7. Хранение состояния между воркерами cluster. Счётчики, кэши, rate-limit-окна в памяти работают «на каждый воркер» — лимит 100 rps превращается в 400 на 4 воркерах. Лечение: внешнее состояние (Redis).
  1. Чем worker_threads отличаются от процессов cluster? Воркеры — потоки в одном процессе: свой Event Loop и heap, общение через сообщения, дешевле создание, изоляция слабее (упавший воркер роняет процесс при необработанном исключении). Cluster — отдельные процессы с общим listen-портом, полная изоляция памяти.
  2. Как передать большой буфер воркеру без копирования? Через transferList в postMessage — владение ArrayBuffer передаётся, в отправителе буфер обнуляется. Для постоянного обмена — SharedArrayBuffer + Atomics.
  3. Когда spawn, когда execFile? execFile — короткие команды, результат буфером целиком, удобен в связке с promisify. spawn — долгие процессы и большой вывод: стримы событий, не копит в памяти, умеет stdio пайпы.
  4. Почему exec опасен? Прогоняет строку через /bin/sh -c: интерполяция пользовательского ввода даёт произвольное выполнение команд. execFile/spawn с argv-массивом обходят шелл — ввод остаётся аргументом, а не кодом.
  5. Что такое Piscina и зачем она нужна? Пул worker_threads: фиксированное число потоков, очередь задач, перезапуск упавших, поддержка transferList. Нужна, чтобы не писать свой пул и не порождать воркер на каждый запрос.
  6. Как Node использует все ядра CPU? Один процесс Node = одно ядро для JS. Для полного использования: несколько процессов — cluster, PM2 или реплики за балансировщиком/Kubernetes. worker_threads добавляют потоки для CPU-задач, но Event Loop мейн-потока всё равно один.
  7. Что будет, если воркер worker_threads завершится с ненулевым кодом? Мейн получит событие exit с кодом; сам процесс не умрёт. Но необработанное исключение в воркере — событие error в мейне; без обработчика оно уронит весь процесс.
  1. PM2-кластер с убийством. Запусти HTTP-сервис через pm2 start app.js -i 4, прокачай запросы, затем убей одного воркера kill -9. Убедись, что PM2 поднимает замену и сервис не теряет запросы (проверь коды ответов). Сравни с однопроцессным вариантом.
  2. Хешер в пуле. Endpoint /hash, который хеширует файл через пул Piscina (4 потока), и endpoint /hash-sync с тем же кодом синхронно. Прогони autocannon -c 10 по обоим и сравни p99 latency и пропускную способность. Убедись, что /hash не блокирует /healthz.
  3. Передача гигабайта. Напиши воркер, принимающий ArrayBuffer через transferList, и вариант без transferList. Замерь время передачи и пиковое потребление памяти обоими способами (process.memoryUsage() в мейне до и после).
  4. Безопасный конвертер. Сервис принимает имя файла через query-параметр и вызывает convert (ImageMagick) для ресайза. Реализуй два варианта: с exec (специально уязвимый) и с execFile. Продемонстрируй инъекцию в первом (?file=a.png;id) и её невозможность во втором.
  5. Пул вручную. Напиши минимальный пул worker_threads своими руками: N постоянных воркеров, очередь заданий, рассылка по round-robin, пересоздание упавшего воркера. Потом замени на Piscina и сравни код.