Кластер
Исходный код: lib/cluster.js
Кластеры процессов Node.js могут использоваться для запуска нескольких экземпляров Node.js, способных распределять рабочую нагрузку между своими потоками приложения. Если изоляция процесса не требуется, используйте модуль worker_threads вместо этого, который позволяет запускать несколько потоков приложения в рамках одного экземпляра Node.js.
Модуль cluster позволяет легко создавать дочерние процессы, которые все совместно используют порты сервера.
MJS модули
import cluster from 'node:cluster';
import http from 'node:http';
import { availableParallelism } from 'node:os';
import process from 'node:process';
const numCPUs = availableParallelism();
if (cluster.isPrimary) {
console.log(`Primary ${process.pid} is running`);
// Fork workers.
for (let i = 0; i < numCPUs; i++) {
cluster.fork();
}
cluster.on('exit', (worker, code, signal) => {
console.log(`worker ${worker.process.pid} died`);
});
} else {
// Workers can share any TCP connection
// In this case it is an HTTP server
http.createServer((req, res) => {
res.writeHead(200);
res.end('hello world\n');
}).listen(8000);
console.log(`Worker ${process.pid} started`);
}
CJS модули
const cluster = require('node:cluster');
const http = require('node:http');
const numCPUs = require('node:os').availableParallelism();
const process = require('node:process');
if (cluster.isPrimary) {
console.log(`Primary ${process.pid} is running`);
// Fork workers.
for (let i = 0; i < numCPUs; i++) {
cluster.fork();
}
cluster.on('exit', (worker, code, signal) => {
console.log(`worker ${worker.process.pid} died`);
});
} else {
// Workers can share any TCP connection
// In this case it is an HTTP server
http.createServer((req, res) => {
res.writeHead(200);
res.end('hello world\n');
}).listen(8000);
console.log(`Worker ${process.pid} started`);
} Запуск Node.js теперь будет использовать порт 8000 для обмена между рабочими процессами:
$ node server.js Primary 3596 is running Worker 4324 started Worker 4520 started Worker 6056 started Worker 5644 started copy
В Windows пока невозможно настроить сервер именованных каналов в рабочем процессе.
Как это работает
Рабочие процессы создаются с помощью метода child_process.fork(), чтобы они могли взаимодействовать с родительским процессом через IPC и передавать друг другу дескрипторы сервера.
Модуль cluster поддерживает два метода распределения входящих соединений.
Первый (и по умолчанию на всех платформах, кроме Windows) — подход round-robin, где основной процесс прослушивает порт, принимает новые подключения и распределяет их между рабочими процессами по кругу, с некоторыми встроенными средствами для предотвращения перегрузки рабочих процессов.
Второй подход — основной процесс создает сокет прослушивания и отправляет его заинтересованным рабочим процессам. Рабочие процессы затем принимают входящие подключения напрямую.
Второй подход, теоретически, должен обеспечить наилучшую производительность. Однако на практике распределение часто оказывается очень несбалансированным из-за особенностей планировщика операционной системы. Были замечены случаи, когда более 70% всех подключений приходилось всего на два процесса из восьми.
Поскольку server.listen() передаёт большую часть работы основному процессу, существует три случая, когда поведение процесса Node.js и рабочего процесса кластера отличается:
-
server.listen({fd: 7})Поскольку сообщение передаётся главному процессу, за прослушивание дескриптора 7 в родительском процессе отвечает родительский процесс, и дескриптор передаётся рабочему процессу, а не рабочему процессу, который ссылается на дескриптор 7. -
server.listen(handle)Прослушивание по предоставленным дескрипторам заставляет рабочий процесс использовать предоставленный дескриптор, а не взаимодействовать с основным процессом. -
server.listen(0)Обычно серверы прослушивают случайный порт. Однако в кластере каждый рабочий процесс получает тот же «случайный» порт каждый раз, когда онlisten(0). По сути, порт случайный в первый раз, но предсказуемый в последующие. Чтобы прослушивать уникальный порт, генерируйте номер порта на основе идентификатора рабочего процесса кластера.
Node.js не предоставляет логику маршрутизации. Поэтому важно проектировать приложение таким образом, чтобы оно не слишком сильно полагалось на объекты данных в памяти для таких вещей, как сессии и авторизация.
Так как рабочие процессы — это отдельные процессы, их можно завершать или перезапускать в зависимости от потребностей программы, не затрагивая другие рабочие процессы. Пока некоторые рабочие процессы всё ещё активны, сервер будет продолжать принимать подключения. Если ни один рабочий процесс не активен, существующие подключения будут прерваны, а новые подключения отклонены. Node.js не управляет количеством рабочих процессов автоматически. Ответственность за управление пулом рабочих процессов лежит на приложении в зависимости от его потребностей.
Хотя основное применение модуля node:cluster — это сетевое программирование, его также можно использовать для других задач, требующих рабочих процессов.
Класс: Worker
- Расширяет: <EventEmitter>
Объект Worker содержит всю публичную информацию и методы о работнике. В основном процессе он может быть получен с помощью cluster.workers. В рабочем процессе он может быть получен с помощью cluster.worker.
Событие: 'disconnect'
Аналогично событию cluster.on('disconnect'), но специфично для этого работника.
cluster.fork().on('disconnect', () => {
// Worker has disconnected
}); copy Событие: 'error'
Это событие такое же, как и предоставляемое child_process.fork().
В рабочем процессе также может использоваться process.on('error').
Событие: 'exit'
-
code<число> Код выхода, если выход нормальный. -
signal<строка> Название сигнала (например,'SIGHUP'), который привел к завершению процесса.
Аналогично событию cluster.on('exit'), но специфично для этого работника.
MJS модули
import cluster from 'node:cluster';
if (cluster.isPrimary) {
const worker = cluster.fork();
worker.on('exit', (code, signal) => {
if (signal) {
console.log(`worker was killed by signal: ${signal}`);
} else if (code !== 0) {
console.log(`worker exited with error code: ${code}`);
} else {
console.log('worker success!');
}
});
}
CJS модули
const cluster = require('node:cluster');
if (cluster.isPrimary) {
const worker = cluster.fork();
worker.on('exit', (code, signal) => {
if (signal) {
console.log(`worker was killed by signal: ${signal}`);
} else if (code !== 0) {
console.log(`worker exited with error code: ${code}`);
} else {
console.log('worker success!');
}
});
} Событие: 'listening'
-
address<Объект>
Аналогично событию cluster.on('listening'), но специфично для этого работника.
MJS модули
cluster.fork().on('listening', (address) => {
// Worker is listening
});
CJS модули
cluster.fork().on('listening', (address) => {
// Worker is listening
}); Не испускается в рабочем процессе.
Событие: 'message'
-
message<Объект> -
handle<undefined> | <Объект>
Аналогично событию 'message' модуля cluster, но специфично для этого работника.
В рабочем процессе также может использоваться process.on('message').
См. process событие: 'message'.
Вот пример использования системы сообщений. Он сохраняет счет в основном процессе количества полученных HTTP запросов рабочими процессами:
MJS модули
import cluster from 'node:cluster';
import http from 'node:http';
import { availableParallelism } from 'node:os';
import process from 'node:process';
if (cluster.isPrimary) {
// Keep track of http requests
let numReqs = 0;
setInterval(() => {
console.log(`numReqs = ${numReqs}`);
}, 1000);
// Count requests
function messageHandler(msg) {
if (msg.cmd && msg.cmd === 'notifyRequest') {
numReqs += 1;
}
}
// Start workers and listen for messages containing notifyRequest
const numCPUs = availableParallelism();
for (let i = 0; i < numCPUs; i++) {
cluster.fork();
}
for (const id in cluster.workers) {
cluster.workers[id].on('message', messageHandler);
}
} else {
// Worker processes have a http server.
http.Server((req, res) => {
res.writeHead(200);
res.end('hello world\n');
// Notify primary about the request
process.send({ cmd: 'notifyRequest' });
}).listen(8000);
}
CJS модули
const cluster = require('node:cluster');
const http = require('node:http');
const process = require('node:process');
if (cluster.isPrimary) {
// Keep track of http requests
let numReqs = 0;
setInterval(() => {
console.log(`numReqs = ${numReqs}`);
}, 1000);
// Count requests
function messageHandler(msg) {
if (msg.cmd && msg.cmd === 'notifyRequest') {
numReqs += 1;
}
}
// Start workers and listen for messages containing notifyRequest
const numCPUs = require('node:os').availableParallelism();
for (let i = 0; i < numCPUs; i++) {
cluster.fork();
}
for (const id in cluster.workers) {
cluster.workers[id].on('message', messageHandler);
}
} else {
// Worker processes have a http server.
http.Server((req, res) => {
res.writeHead(200);
res.end('hello world\n');
// Notify primary about the request
process.send({ cmd: 'notifyRequest' });
}).listen(8000);
} Событие: 'online'
Аналогично событию cluster.on('online'), но специфично для этого работника.
cluster.fork().on('online', () => {
// Worker is online
}); copy Не испускается в рабочем процессе.
worker.disconnect()
- Возвращает: <cluster.Worker> Ссылка на
worker.
В рабочем процессе эта функция закроет все серверы, подождёт события 'close' на этих серверах, а затем отключит канал IPC.
В основном процессе отправляется внутреннее сообщение работнику, заставляя его вызвать .disconnect() на себе.
Вызывает установку .exitedAfterDisconnect.
После закрытия сервера он больше не будет принимать новые подключения, но подключения могут быть приняты любым другим слушающим работником. Существующие подключения будут закрыты обычным образом. Когда больше нет подключений, см. server.close(), канал IPC с работником закроется, что позволит ему завершиться корректно.
Вышесказанное относится только к серверным подключениям, клиентские подключения не закрываются автоматически рабочими процессами, и отключение не ожидает их закрытия перед выходом.
В рабочем процессе process.disconnect существует, но это не эта функция; это disconnect().
Поскольку долгоживущие серверные подключения могут блокировать выход из рабочих процессов, может быть полезно отправить сообщение, чтобы были предприняты действия, специфичные для приложения, по их закрытию. Также может быть полезно реализовать таймаут, убивая работника, если событие 'disconnect' не было испущено через некоторое время.
if (cluster.isPrimary) {
const worker = cluster.fork();
let timeout;
worker.on('listening', (address) => {
worker.send('shutdown');
worker.disconnect();
timeout = setTimeout(() => {
worker.kill();
}, 2000);
});
worker.on('disconnect', () => {
clearTimeout(timeout);
});
} else if (cluster.isWorker) {
const net = require('node:net');
const server = net.createServer((socket) => {
// Connections never end
});
server.listen(8000);
process.on('message', (msg) => {
if (msg === 'shutdown') {
// Initiate graceful close of any connections to server
}
});
} copy
worker.exitedAfterDisconnect
Это свойство true , если работник завершился из-за .disconnect(). Если работник завершился другим способом, оно false. Если работник еще не завершился, оно undefined.
Логическое значение worker.exitedAfterDisconnect позволяет различать добровольный и случайный выход, основной процесс может выбрать не перезапускать работника на основе этого значения.
cluster.on('exit', (worker, code, signal) => {
if (worker.exitedAfterDisconnect === true) {
console.log('Oh, it was just voluntary – no need to worry');
}
});
// kill worker
worker.kill(); copy
worker.id
Каждый новый работник получает свой уникальный идентификатор, этот идентификатор хранится в id.
Пока работник жив, это ключ, который индексирует его в cluster.workers.
worker.isConnected()
Эта функция возвращает true , если работник подключен к своему основному процессу через канал IPC, false в противном случае. Работник подключается к своему основному процессу после создания. Он отключается после испускания события 'disconnect'.
worker.isDead()
Эта функция возвращает true , если процесс работника завершился (из-за выхода или сигнала). В противном случае она возвращает false.
MJS модули
import cluster from 'node:cluster';
import http from 'node:http';
import { availableParallelism } from 'node:os';
import process from 'node:process';
const numCPUs = availableParallelism();
if (cluster.isPrimary) {
console.log(`Primary ${process.pid} is running`);
// Fork workers.
for (let i = 0; i < numCPUs; i++) {
cluster.fork();
}
cluster.on('fork', (worker) => {
console.log('worker is dead:', worker.isDead());
});
cluster.on('exit', (worker, code, signal) => {
console.log('worker is dead:', worker.isDead());
});
} else {
// Workers can share any TCP connection. In this case, it is an HTTP server.
http.createServer((req, res) => {
res.writeHead(200);
res.end(`Current process\n ${process.pid}`);
process.kill(process.pid);
}).listen(8000);
}
CJS модули
const cluster = require('node:cluster');
const http = require('node:http');
const numCPUs = require('node:os').availableParallelism();
const process = require('node:process');
if (cluster.isPrimary) {
console.log(`Primary ${process.pid} is running`);
// Fork workers.
for (let i = 0; i < numCPUs; i++) {
cluster.fork();
}
cluster.on('fork', (worker) => {
console.log('worker is dead:', worker.isDead());
});
cluster.on('exit', (worker, code, signal) => {
console.log('worker is dead:', worker.isDead());
});
} else {
// Workers can share any TCP connection. In this case, it is an HTTP server.
http.createServer((req, res) => {
res.writeHead(200);
res.end(`Current process\n ${process.pid}`);
process.kill(process.pid);
}).listen(8000);
}
worker.kill([signal])
-
signal<строка> Имя сигнала завершения для отправки процессу работника. По умолчанию:'SIGTERM'
Эта функция завершит работу работника. В основном процессе это делается путем отключения worker.process, а после отключения - завершением с помощью signal. В рабочем процессе это делается путем завершения процесса с помощью signal.
Функция kill() завершает процесс работника без ожидания корректного отключения, она имеет такое же поведение, как и worker.process.kill().
Этот метод является псевдонимом worker.destroy() для обратной совместимости.
В рабочем процессе process.kill() существует, но это не эта функция; это kill().
worker.process
Все рабочие процессы создаются с помощью child_process.fork(), возвращаемый объект из этой функции хранится как .process. В рабочем процессе глобальный process хранится.
См.: Модуль процесса дочернего процесса.
Рабочие процессы вызовут process.exit(0) , если событие 'disconnect' произойдёт на process и .exitedAfterDisconnect не true. Это защищает от случайного отключения.
worker.send(message[, sendHandle[, options]][, callback])
-
message<Объект> -
sendHandle<Обработчик> -
options<Объект> Аргументoptions, если он присутствует, является объектом, используемым для параметризации отправки определенных типов обработчиков.optionsподдерживает следующие свойства:-
keepOpen<логическое> Значение, которое может быть использовано при передаче экземпляровnet.Socket. Приtrue, сокет остается открытым в процессе отправки. По умолчанию:false.
-
-
callback<Функция> - Возвращает: <логическое>
Отправить сообщение работнику или основному процессу, необязательно с обработчиком.
В основном процессе это отправляет сообщение определенному работнику. Это идентично ChildProcess.send().
В рабочем процессе это отправляет сообщение основному процессу. Это идентично process.send().
Этот пример вернёт все сообщения от основного процесса:
if (cluster.isPrimary) {
const worker = cluster.fork();
worker.send('hi there');
} else if (cluster.isWorker) {
process.on('message', (msg) => {
process.send(msg);
});
} copy Событие: 'disconnect'
-
worker<cluster.Worker>
Выдаётся после того, как канал IPC работника отключился. Это может произойти, когда работник завершился нормально, был убит или отключен вручную (например, с помощью worker.disconnect()).
Возможна задержка между событиями 'disconnect' и 'exit'. Эти события можно использовать для определения, застрял ли процесс в очистке или есть долгоживущие подключения.
cluster.on('disconnect', (worker) => {
console.log(`The worker #${worker.id} has disconnected`);
}); copy Событие: 'exit'
-
worker<cluster.Worker> -
code<number> Код завершения, если завершение было нормальным. -
signal<string> Название сигнала (например,'SIGHUP'), по которому процесс был убит.
Когда любой из работников умирает, модуль cluster выпустит событие 'exit'.
Это можно использовать для перезапуска работника, вызвав .fork() снова.
cluster.on('exit', (worker, code, signal) => {
console.log('worker %d died (%s). restarting...',
worker.process.pid, signal || code);
cluster.fork();
}); copy Событие: 'fork'
-
worker<cluster.Worker>
При создании нового работника модуль cluster выпустит событие 'fork'. Это можно использовать для регистрации активности работников и создания пользовательской задержки.
const timeouts = [];
function errorMsg() {
console.error('Something must be wrong with the connection ...');
}
cluster.on('fork', (worker) => {
timeouts[worker.id] = setTimeout(errorMsg, 2000);
});
cluster.on('listening', (worker, address) => {
clearTimeout(timeouts[worker.id]);
});
cluster.on('exit', (worker, code, signal) => {
clearTimeout(timeouts[worker.id]);
errorMsg();
}); copy Событие: 'listening'
-
worker<cluster.Worker> -
address<Object>
После вызова listen() из работника, когда событие 'listening' генерируется на сервере, событие 'listening' также будет сгенерировано в основном процессе.
Обработчик события выполняется с двумя аргументами: worker содержит объект работника, а объект address содержит следующие свойства подключения: address, port, и addressType. Это очень полезно, если работник прослушивает более одного адреса.
cluster.on('listening', (worker, address) => {
console.log(
`A worker is now connected to ${address.address}:${address.port}`);
}); copy Событие addressType может быть одним из:
-
4(TCPv4) -
6(TCPv6) -
-1(сокет доменной области) -
'udp4'или'udp6'(UDPv4 или UDPv6)
Событие: 'message'
-
worker<cluster.Worker> -
message<Object> -
handle<undefined> | <Object>
Выдаётся, когда основной процесс cluster получает сообщение от любого работника.
Событие: 'online'
-
worker<cluster.Worker>
После создания нового работника, работник должен ответить сообщением о запуске. Когда основной процесс получает сообщение о запуске, он выпустит это событие. Разница между 'fork' и 'online' в том, что 'online' генерируется при создании работника основным процессом, а 'online' - при запуске работника.
cluster.on('online', (worker) => {
console.log('Yay, the worker responded after it was forked');
}); copy Событие: 'setup'
-
settings<Object>
Выдаётся каждый раз, когда вызывается .setupPrimary().
Объект settings — это объект cluster.settings в момент вызова .setupPrimary() и носит лишь рекомендательный характер, так как вызовы .setupPrimary() могут происходить в одном такте.
Если точность важна, используйте cluster.settings.
cluster.disconnect([callback])
-
callback<Function> Вызывается, когда все работники отключены, и закрыты дескрипторы.
Вызывает .disconnect() для каждого работника в cluster.workers.
При отключении закрываются все внутренние дескрипторы, что позволяет главному процессу завершиться корректно, если нет других ожидающих событий.
Метод принимает необязательный аргумент-обработчик, который будет вызван по завершении.
Вызывать можно только из основного процесса.
cluster.fork([env])
-
env<Object> Пара ключ-значение для добавления в среду процесса работника. - Возвращает: <cluster.Worker>
Запустить новый процесс работника.
Вызывать можно только из основного процесса.
cluster.isMaster
Устаревшее псевдоним для cluster.isPrimary.
cluster.isPrimary
Истинно, если процесс является основным. Это определяется process.env.NODE_UNIQUE_ID. Если process.env.NODE_UNIQUE_ID неопределённо, то isPrimary равно true.
cluster.isWorker
Истинно, если процесс не является основным (это отрицание cluster.isPrimary).
cluster.schedulingPolicy
Политика планирования, либо cluster.SCHED_RR для циклического распределения, либо cluster.SCHED_NONE для переопределения операционной системой. Это глобальный параметр и фактически замораживается, как только создаётся первый работник или вызывается .setupPrimary() (что произойдёт раньше).
SCHED_RR является значением по умолчанию на всех операционных системах, кроме Windows. В Windows будет использоваться SCHED_RR как только libuv сможет эффективно распределять дескрипторы IOCP без значительных потерь производительности.
cluster.schedulingPolicy также может быть установлено через переменную среды NODE_CLUSTER_SCHED_POLICY. Допустимые значения 'rr' и 'none'.
cluster.settings
-
<Объект>
-
execArgv<строка[]> Список строковых аргументов, переданных исполняемому файлу Node.js. По умолчанию:process.execArgv. -
exec<строка> Путь к файлу worker. По умолчанию:process.argv[1]. -
args<строка[]> Строковые аргументы, переданные worker. По умолчанию:process.argv.slice(2). -
cwd<строка> Текущий рабочий каталог процесса worker. По умолчанию:undefined(унаследован от родительского процесса). -
serialization<строка> Указывает тип сериализации, используемой для отправки сообщений между процессами. Возможные значения:'json'и'advanced'. Более подробную информацию см. в разделе Расширенная сериализация дляchild_process. По умолчанию:false. -
silent<логическое> Выводить ли вывод в stdio родительского процесса. По умолчанию:false. -
stdio<Массив> Настраивает stdio разветвлённых процессов. Поскольку модуль cluster полагается на IPC для работы, эта настройка должна содержать запись'ipc'. При предоставлении этой опции она переопределяетsilent. См.child_process.spawn()иstdio. -
uid<число> Устанавливает идентификатор пользователя процесса. (См.setuid(2).) -
gid<число> Устанавливает идентификатор группы процесса. (См.setgid(2).) -
inspectPort<число> | <Функция> Устанавливает порт инспектора worker. Это может быть число или функция, которая не принимает аргументов и возвращает число. По умолчанию каждый worker получает свой порт, инкрементированный от порта основного процессаprocess.debugPort. -
windowsHide<логическое> Скрыть окно консоли разветвлённых процессов, которое обычно создаётся в системах Windows. По умолчанию:false.
-
После вызова .setupPrimary() (или .fork()) этот объект настроек будет содержать настройки, включая значения по умолчанию.
Этот объект не предназначен для изменения или ручного задания.
cluster.setupMaster([settings])
Устаревшее псевдоним для .setupPrimary().
cluster.setupPrimary([settings])
-
settings<Объект> См.cluster.settings.
setupPrimary используется для изменения поведения по умолчанию 'fork'. После вызова настройки будут присутствовать в cluster.settings.
Любые изменения настроек влияют только на последующие вызовы .fork() и не оказывают влияния на уже запущенные worker.
Единственный атрибут worker, который нельзя установить через .setupPrimary() — это env переданный в .fork().
Указанные выше значения по умолчанию применяются только к первому вызову; значения по умолчанию для последующих вызовов — текущие значения на момент вызова cluster.setupPrimary().
Модули MJS
import cluster from 'node:cluster';
cluster.setupPrimary({
exec: 'worker.js',
args: ['--use', 'https'],
silent: true,
});
cluster.fork(); // https worker
cluster.setupPrimary({
exec: 'worker.js',
args: ['--use', 'http'],
});
cluster.fork(); // http worker
Модули CJS
const cluster = require('node:cluster');
cluster.setupPrimary({
exec: 'worker.js',
args: ['--use', 'https'],
silent: true,
});
cluster.fork(); // https worker
cluster.setupPrimary({
exec: 'worker.js',
args: ['--use', 'http'],
});
cluster.fork(); // http worker Этот метод может быть вызван только из основного процесса.
cluster.worker
Ссылка на текущий объект worker. Не доступен в основном процессе.
Модули MJS
import cluster from 'node:cluster';
if (cluster.isPrimary) {
console.log('I am primary');
cluster.fork();
cluster.fork();
} else if (cluster.isWorker) {
console.log(`I am worker #${cluster.worker.id}`);
}
Модули CJS
const cluster = require('node:cluster');
if (cluster.isPrimary) {
console.log('I am primary');
cluster.fork();
cluster.fork();
} else if (cluster.isWorker) {
console.log(`I am worker #${cluster.worker.id}`);
}
cluster.workers
Хеш, хранящий активные объекты worker, ключи — поле id. Это упрощает циклический проход по всем worker. Доступен только в основном процессе.
Worker удаляется из cluster.workers после того, как worker отключился и завершил работу. Порядок между этими двумя событиями не может быть определён заранее. Однако гарантируется, что удаление из списка cluster.workers происходит до отправки последнего события 'disconnect' или 'exit'.
Модули MJS
import cluster from 'node:cluster';
for (const worker of Object.values(cluster.workers)) {
worker.send('big announcement to all workers');
}
Модули CJS
const cluster = require('node:cluster');
for (const worker of Object.values(cluster.workers)) {
worker.send('big announcement to all workers');
}
© Joyent, Inc. and other Node contributors
Licensed under the MIT License.
Node.js is a trademark of Joyent, Inc. and is used with its permission.
We are not endorsed by or affiliated with Joyent.
https://nodejs.org/dist/latest-v18.x/docs/api/cluster.html