Потоки Worker
Исходный код: lib/worker_threads.js
Модуль node:worker_threads позволяет использовать потоки, выполняющие JavaScript параллельно. Чтобы получить к нему доступ:
Модули JavaScript
import worker from 'node:worker_threads';
CommonJS
'use strict';
const worker = require('node:worker_threads');Рабочие потоки (потоки) полезны для выполнения ресурсоёмких операций JavaScript, требующих вычислений на CPU. Для задач с интенсивным вводом-выводом они не слишком полезны. Встроенные асинхронные операции ввода-вывода Node.js эффективнее, чем Worker.
В отличие от child_process или cluster, worker_threads могут совместно использовать память. Это достигается передачей экземпляров ArrayBuffer или совместным использованием экземпляров SharedArrayBuffer.
Модули JavaScript
import {
Worker,
isMainThread,
parentPort,
workerData,
} from 'node:worker_threads';
if (!isMainThread) {
const { parse } = await import('some-js-parsing-library');
const script = workerData;
parentPort.postMessage(parse(script));
}
export default function parseJSAsync(script) {
return new Promise((resolve, reject) => {
const worker = new Worker(new URL(import.meta.url), {
workerData: script,
});
worker.on('message', resolve);
worker.on('error', reject);
worker.on('exit', (code) => {
if (code !== 0)
reject(new Error(`Worker stopped with exit code ${code}`));
});
});
};CommonJS
'use strict';
const {
Worker,
isMainThread,
parentPort,
workerData,
} = require('node:worker_threads');
if (isMainThread) {
module.exports = function parseJSAsync(script) {
return new Promise((resolve, reject) => {
const worker = new Worker(__filename, {
workerData: script,
});
worker.on('message', resolve);
worker.on('error', reject);
worker.on('exit', (code) => {
if (code !== 0)
reject(new Error(`Worker stopped with exit code ${code}`));
});
});
};
} else {
const { parse } = require('some-js-parsing-library');
const script = workerData;
parentPort.postMessage(parse(script));
}В приведённом выше примере для каждого вызова parseJSAsync() создаётся поток Worker. На практике для таких задач следует использовать пул Worker. В противном случае затраты на создание Worker, скорее всего, превысят пользу от их использования.
При реализации пула Worker используйте API AsyncResource, чтобы сообщить диагностическим инструментам (например, для формирования асинхронных трассировок стека) о связи между задачами и их результатами. Пример реализации см. в разделе «Использование AsyncResource для пула потоков Worker» документации async_hooks.
Рабочие потоки по умолчанию наследуют параметры, не относящиеся к конкретному процессу. См. раздел Worker constructor options, чтобы узнать, как настраивать параметры рабочих потоков, в частности параметры argv и execArgv.
worker.getEnvironmentData(key)
-
key<any> Произвольное клонируемое значение JavaScript, которое можно использовать в качестве ключа <Map>. - Возвращает: <any>
В рабочем потоке worker.getEnvironmentData() возвращает клон данных, переданных в worker.setEnvironmentData() потока, который его создал. Каждый новый Worker автоматически получает собственную копию данных среды.
Модули JavaScript
import {
Worker,
isMainThread,
setEnvironmentData,
getEnvironmentData,
} from 'node:worker_threads';
if (isMainThread) {
setEnvironmentData('Hello', 'World!');
const worker = new Worker(new URL(import.meta.url));
} else {
console.log(getEnvironmentData('Hello')); // Prints 'World!'.
}CommonJS
'use strict';
const {
Worker,
isMainThread,
setEnvironmentData,
getEnvironmentData,
} = require('node:worker_threads');
if (isMainThread) {
setEnvironmentData('Hello', 'World!');
const worker = new Worker(__filename);
} else {
console.log(getEnvironmentData('Hello')); // Prints 'World!'.
}
worker.isInternalThread
- Тип: <boolean>
Равно true, если этот код выполняется во внутреннем потоке Worker (например, в потоке загрузчика).
node --experimental-loader ./loader.js main.js copy
Модули JavaScript
// loader.js
import { isInternalThread } from 'node:worker_threads';
console.log(isInternalThread); // trueCommonJS
// loader.js
'use strict';
const { isInternalThread } = require('node:worker_threads');
console.log(isInternalThread); // trueМодули JavaScript
// main.js
import { isInternalThread } from 'node:worker_threads';
console.log(isInternalThread); // falseCommonJS
// main.js
'use strict';
const { isInternalThread } = require('node:worker_threads');
console.log(isInternalThread); // false
worker.isMainThread
- Тип: <boolean>
Равно true, если этот код выполняется не в потоке Worker.
Модули JavaScript
import { Worker, isMainThread } from 'node:worker_threads';
if (isMainThread) {
// This re-loads the current file inside a Worker instance.
new Worker(new URL(import.meta.url));
} else {
console.log('Inside Worker!');
console.log(isMainThread); // Prints 'false'.
}CommonJS
'use strict';
const { Worker, isMainThread } = require('node:worker_threads');
if (isMainThread) {
// This re-loads the current file inside a Worker instance.
new Worker(__filename);
} else {
console.log('Inside Worker!');
console.log(isMainThread); // Prints 'false'.
}
worker.markAsUntransferable(object)
-
object<any> Произвольное значение JavaScript.
Помечает объект как непередаваемый. Если object указан в списке передачи вызова port.postMessage(), будет выброшена ошибка. Если object является примитивным значением, операция ничего не делает.
Это особенно полезно для объектов, которые можно клонировать, а не передавать, и которые используются другими объектами на стороне отправителя. Например, Node.js помечает таким образом ArrayBuffer, используемые в пуле Buffer.
Эту операцию нельзя отменить.
Модули JavaScript
import { MessageChannel, markAsUntransferable } from 'node:worker_threads';
const pooledBuffer = new ArrayBuffer(8);
const typedArray1 = new Uint8Array(pooledBuffer);
const typedArray2 = new Float64Array(pooledBuffer);
markAsUntransferable(pooledBuffer);
const { port1 } = new MessageChannel();
try {
// This will throw an error, because pooledBuffer is not transferable.
port1.postMessage(typedArray1, [ typedArray1.buffer ]);
} catch (error) {
// error.name === 'DataCloneError'
}
// The following line prints the contents of typedArray1 -- it still owns
// its memory and has not been transferred. Without
// `markAsUntransferable()`, this would print an empty Uint8Array and the
// postMessage call would have succeeded.
// typedArray2 is intact as well.
console.log(typedArray1);
console.log(typedArray2);CommonJS
'use strict';
const { MessageChannel, markAsUntransferable } = require('node:worker_threads');
const pooledBuffer = new ArrayBuffer(8);
const typedArray1 = new Uint8Array(pooledBuffer);
const typedArray2 = new Float64Array(pooledBuffer);
markAsUntransferable(pooledBuffer);
const { port1 } = new MessageChannel();
try {
// This will throw an error, because pooledBuffer is not transferable.
port1.postMessage(typedArray1, [ typedArray1.buffer ]);
} catch (error) {
// error.name === 'DataCloneError'
}
// The following line prints the contents of typedArray1 -- it still owns
// its memory and has not been transferred. Without
// `markAsUntransferable()`, this would print an empty Uint8Array and the
// postMessage call would have succeeded.
// typedArray2 is intact as well.
console.log(typedArray1);
console.log(typedArray2);В браузерах нет эквивалента этому API.
worker.isMarkedAsUntransferable(object)
Проверяет, помечен ли объект как непередаваемый с помощью markAsUntransferable().
Модули JavaScript
import { markAsUntransferable, isMarkedAsUntransferable } from 'node:worker_threads';
const pooledBuffer = new ArrayBuffer(8);
markAsUntransferable(pooledBuffer);
isMarkedAsUntransferable(pooledBuffer); // Returns true.CommonJS
'use strict';
const { markAsUntransferable, isMarkedAsUntransferable } = require('node:worker_threads');
const pooledBuffer = new ArrayBuffer(8);
markAsUntransferable(pooledBuffer);
isMarkedAsUntransferable(pooledBuffer); // Returns true.В браузерах нет эквивалента этому API.
worker.markAsUncloneable(object)
-
object<any> Произвольное значение JavaScript.
Помечает объект как непригодный для клонирования. Если object используется в качестве message в вызове port.postMessage(), будет выброшена ошибка. Если object является примитивным значением, операция ничего не делает.
Это не влияет на ArrayBuffer и любые объекты типа Buffer.
Эту операцию нельзя отменить.
Модули JavaScript
import { markAsUncloneable } from 'node:worker_threads';
const anyObject = { foo: 'bar' };
markAsUncloneable(anyObject);
const { port1 } = new MessageChannel();
try {
// This will throw an error, because anyObject is not cloneable.
port1.postMessage(anyObject);
} catch (error) {
// error.name === 'DataCloneError'
}CommonJS
'use strict';
const { markAsUncloneable } = require('node:worker_threads');
const anyObject = { foo: 'bar' };
markAsUncloneable(anyObject);
const { port1 } = new MessageChannel();
try {
// This will throw an error, because anyObject is not cloneable.
port1.postMessage(anyObject);
} catch (error) {
// error.name === 'DataCloneError'
}В браузерах нет эквивалента этому API.
worker.moveMessagePortToContext(port, contextifiedSandbox)
-
port<MessagePort> Порт сообщений для передачи. -
contextifiedSandbox<Object> Объект с контекстом, возвращаемый методомvm.createContext(). -
Возвращает: <MessagePort>
Передаёт MessagePort в другой контекст vm. Исходный объект port становится непригодным к использованию, и его заменяет возвращённый экземпляр MessagePort.
Возвращённый MessagePort является объектом целевого контекста и наследует от его глобального класса Object. Объекты, переданные обработчику port.onmessage(), также создаются в целевом контексте и наследуют от его глобального класса Object.
Однако созданный MessagePort больше не наследует от <EventTarget>, и для получения событий через него можно использовать только port.onmessage().
worker.parentPort
- Тип: <null> | <MessagePort>
Если этот поток является Worker, это MessagePort, обеспечивающий связь с родительским потоком. Сообщения, отправленные с помощью parentPort.postMessage(), доступны в родительском потоке через worker.on('message'), а сообщения, отправленные из родительского потока с помощью worker.postMessage(), доступны в этом потоке через parentPort.on('message').
Модули JavaScript
import { Worker, isMainThread, parentPort } from 'node:worker_threads';
if (isMainThread) {
const worker = new Worker(new URL(import.meta.url));
worker.once('message', (message) => {
console.log(message); // Prints 'Hello, world!'.
});
worker.postMessage('Hello, world!');
} else {
// When a message from the parent thread is received, send it back:
parentPort.once('message', (message) => {
parentPort.postMessage(message);
});
}CommonJS
'use strict';
const { Worker, isMainThread, parentPort } = require('node:worker_threads');
if (isMainThread) {
const worker = new Worker(__filename);
worker.once('message', (message) => {
console.log(message); // Prints 'Hello, world!'.
});
worker.postMessage('Hello, world!');
} else {
// When a message from the parent thread is received, send it back:
parentPort.once('message', (message) => {
parentPort.postMessage(message);
});
}
worker.postMessageToThread(threadId, value[, transferList][, timeout])
-
threadId<number> Идентификатор целевого потока. Если идентификатор потока недействителен, будет выброшена ошибкаERR_WORKER_MESSAGING_FAILED. Если идентификатор целевого потока совпадает с идентификатором текущего потока, будет выброшена ошибкаERR_WORKER_MESSAGING_SAME_THREAD. -
value<any> Отправляемое значение. -
transferList<Object[]> Если вvalueпередан один или несколько объектов типаMessagePort, для этих элементов требуетсяtransferList, иначе будет выброшена ошибкаERR_MISSING_MESSAGE_PORT_IN_TRANSFER_LIST. Дополнительные сведения см. в разделеport.postMessage(). -
timeout<number> Время ожидания доставки сообщения в миллисекундах. По умолчанию —undefined, то есть ожидание продолжается бесконечно. Если время ожидания истечёт, будет выброшена ошибкаERR_WORKER_MESSAGING_TIMEOUT. - Возвращает: <Promise> Промис, который выполняется, если целевой поток успешно обработал сообщение.
Отправляет значение другому рабочему потоку, указанному его идентификатором потока.
Если в целевом потоке нет обработчика события workerMessage, операция выбросит ошибку ERR_WORKER_MESSAGING_FAILED.
Если целевой поток выбросил ошибку при обработке события workerMessage, операция выбросит ошибку ERR_WORKER_MESSAGING_ERRORED.
Этот метод следует использовать, если целевой поток не является непосредственным родительским или дочерним потоком текущего. Если потоки связаны отношением родитель — потомок, используйте require('node:worker_threads').parentPort.postMessage() и worker.postMessage() для связи между ними.
В приведённом ниже примере показано использование postMessageToThread: создаются 10 вложенных потоков, и последний пытается связаться с главным потоком.
Модули JavaScript
import process from 'node:process';
import {
postMessageToThread,
threadId,
workerData,
Worker,
} from 'node:worker_threads';
const channel = new BroadcastChannel('sync');
const level = workerData?.level ?? 0;
if (level < 10) {
const worker = new Worker(new URL(import.meta.url), {
workerData: { level: level + 1 },
});
}
if (level === 0) {
process.on('workerMessage', (value, source) => {
console.log(`${source} -> ${threadId}:`, value);
postMessageToThread(source, { message: 'pong' });
});
} else if (level === 10) {
process.on('workerMessage', (value, source) => {
console.log(`${source} -> ${threadId}:`, value);
channel.postMessage('done');
channel.close();
});
await postMessageToThread(0, { message: 'ping' });
}
channel.onmessage = channel.close;CommonJS
'use strict';
const process = require('node:process');
const {
postMessageToThread,
threadId,
workerData,
Worker,
} = require('node:worker_threads');
const channel = new BroadcastChannel('sync');
const level = workerData?.level ?? 0;
if (level < 10) {
const worker = new Worker(__filename, {
workerData: { level: level + 1 },
});
}
if (level === 0) {
process.on('workerMessage', (value, source) => {
console.log(`${source} -> ${threadId}:`, value);
postMessageToThread(source, { message: 'pong' });
});
} else if (level === 10) {
process.on('workerMessage', (value, source) => {
console.log(`${source} -> ${threadId}:`, value);
channel.postMessage('done');
channel.close();
});
postMessageToThread(0, { message: 'ping' });
}
channel.onmessage = channel.close;
worker.receiveMessageOnPort(port)
-
port<MessagePort> | <BroadcastChannel> -
Возвращает: <Object> | <undefined>
Получает одно сообщение из указанного MessagePort. Если сообщений нет, возвращается undefined, в противном случае возвращается объект с единственным свойством message, содержащим данные сообщения, соответствующего самому старому сообщению в очереди MessagePort.
Модули JavaScript
import { MessageChannel, receiveMessageOnPort } from 'node:worker_threads';
const { port1, port2 } = new MessageChannel();
port1.postMessage({ hello: 'world' });
console.log(receiveMessageOnPort(port2));
// Prints: { message: { hello: 'world' } }
console.log(receiveMessageOnPort(port2));
// Prints: undefinedCommonJS
'use strict';
const { MessageChannel, receiveMessageOnPort } = require('node:worker_threads');
const { port1, port2 } = new MessageChannel();
port1.postMessage({ hello: 'world' });
console.log(receiveMessageOnPort(port2));
// Prints: { message: { hello: 'world' } }
console.log(receiveMessageOnPort(port2));
// Prints: undefinedПри использовании этой функции событие 'message' не генерируется, а обработчик onmessage не вызывается.
worker.resourceLimits
- Тип: <Object>
Предоставляет набор ограничений ресурсов движка JS в этом рабочем потоке. Если параметр resourceLimits был передан конструктору Worker, это свойство содержит те же значения.
При использовании в главном потоке это свойство содержит пустой объект.
worker.SHARE_ENV
- Тип: <symbol>
Специальное значение, которое можно передать в качестве параметра env конструктора Worker, чтобы указать, что текущий поток и поток Worker должны совместно использовать доступ на чтение и запись к одному и тому же набору переменных среды.
Модули JavaScript
import process from 'node:process';
import { Worker, SHARE_ENV } from 'node:worker_threads';
new Worker('process.env.SET_IN_WORKER = "foo"', { eval: true, env: SHARE_ENV })
.on('exit', () => {
console.log(process.env.SET_IN_WORKER); // Prints 'foo'.
});CommonJS
'use strict';
const { Worker, SHARE_ENV } = require('node:worker_threads');
new Worker('process.env.SET_IN_WORKER = "foo"', { eval: true, env: SHARE_ENV })
.on('exit', () => {
console.log(process.env.SET_IN_WORKER); // Prints 'foo'.
});
worker.setEnvironmentData(key[, value])
-
key<any> Произвольное клонируемое значение JavaScript, которое можно использовать в качестве ключа <Map>. -
value<any> Произвольное клонируемое значение JavaScript, которое будет клонировано и автоматически передано всем новым экземплярамWorker. Еслиvalueпередано какundefined, ранее установленное значение дляkeyбудет удалено.
API worker.setEnvironmentData() задаёт содержимое worker.getEnvironmentData() в текущем потоке и во всех новых экземплярах Worker, созданных из текущего контекста.
worker.threadId
- Тип: <integer>
Целочисленный идентификатор текущего потока. В соответствующем объекте Worker (если он существует) доступен как worker.threadId. Это значение уникально для каждого экземпляра Worker в рамках одного процесса.
worker.threadName
Строковый идентификатор текущего потока или null, если поток не выполняется. В соответствующем объекте Worker (если он существует) доступен как worker.threadName.
worker.workerData
Произвольное значение JavaScript, содержащее клон данных, переданных конструктору Worker этого потока.
Данные клонируются так же, как при использовании postMessage(), согласно алгоритму структурного клонирования HTML.
Модули JavaScript
import { Worker, isMainThread, workerData } from 'node:worker_threads';
if (isMainThread) {
const worker = new Worker(new URL(import.meta.url), { workerData: 'Hello, world!' });
} else {
console.log(workerData); // Prints 'Hello, world!'.
}CommonJS
'use strict';
const { Worker, isMainThread, workerData } = require('node:worker_threads');
if (isMainThread) {
const worker = new Worker(__filename, { workerData: 'Hello, world!' });
} else {
console.log(workerData); // Prints 'Hello, world!'.
}Класс: BroadcastChannel extends EventTarget
Экземпляры BroadcastChannel обеспечивают асинхронную связь «один ко многим» со всеми остальными экземплярами BroadcastChannel, подключёнными к каналу с тем же именем.
Модули JavaScript
import {
isMainThread,
BroadcastChannel,
Worker,
} from 'node:worker_threads';
const bc = new BroadcastChannel('hello');
if (isMainThread) {
let c = 0;
bc.onmessage = (event) => {
console.log(event.data);
if (++c === 10) bc.close();
};
for (let n = 0; n < 10; n++)
new Worker(new URL(import.meta.url));
} else {
bc.postMessage('hello from every worker');
bc.close();
}CommonJS
'use strict';
const {
isMainThread,
BroadcastChannel,
Worker,
} = require('node:worker_threads');
const bc = new BroadcastChannel('hello');
if (isMainThread) {
let c = 0;
bc.onmessage = (event) => {
console.log(event.data);
if (++c === 10) bc.close();
};
for (let n = 0; n < 10; n++)
new Worker(__filename);
} else {
bc.postMessage('hello from every worker');
bc.close();
}
new BroadcastChannel(name)
-
name<any> Имя канала для подключения. Допускается любое значение JavaScript, которое можно преобразовать в строку с помощью`${name}`.
broadcastChannel.close()
Закрывает соединение BroadcastChannel.
broadcastChannel.onmessage
- Тип: <Function> Вызывается с одним аргументом
MessageEventпри получении сообщения.
broadcastChannel.onmessageerror
- Тип: <Function> Вызывается, если полученное сообщение не удаётся десериализовать.
broadcastChannel.postMessage(message)
-
message<any> Любое клонируемое значение JavaScript.
broadcastChannel.ref()
Противоположность unref(). Вызов ref() для ранее unref()ированного BroadcastChannel не позволяет программе завершиться, если он остался единственным активным дескриптором (поведение по умолчанию). Если порт был ref()ирован, повторный вызов ref() не имеет эффекта.
broadcastChannel.unref()
Вызов unref() для BroadcastChannel позволяет потоку завершиться, если он является единственным активным дескриптором в системе событий. Если BroadcastChannel уже был unref()ирован, повторный вызов unref() не имеет эффекта.
Класс: MessageChannel
Экземпляры класса worker.MessageChannel представляют собой асинхронный двусторонний канал связи. У MessageChannel нет собственных методов. new MessageChannel() возвращает объект со свойствами port1 и port2, которые относятся к связанным экземплярам MessagePort.
Модули JavaScript
import { MessageChannel } from 'node:worker_threads';
const { port1, port2 } = new MessageChannel();
port1.on('message', (message) => console.log('received', message));
port2.postMessage({ foo: 'bar' });
// Prints: received { foo: 'bar' } from the `port1.on('message')` listenerCommonJS
'use strict';
const { MessageChannel } = require('node:worker_threads');
const { port1, port2 } = new MessageChannel();
port1.on('message', (message) => console.log('received', message));
port2.postMessage({ foo: 'bar' });
// Prints: received { foo: 'bar' } from the `port1.on('message')` listenerКласс: MessagePort
- Расширяет: <EventTarget>
Экземпляры класса worker.MessagePort представляют один из концов асинхронного двустороннего канала связи. Его можно использовать для передачи структурированных данных, областей памяти и других MessagePort между разными Worker.
Эта реализация соответствует MessagePort браузера.
Событие: 'close'
Событие 'close' возникает, когда одна из сторон канала отключается.
Модули JavaScript
import { MessageChannel } from 'node:worker_threads';
const { port1, port2 } = new MessageChannel();
// Prints:
// foobar
// closed!
port2.on('message', (message) => console.log(message));
port2.on('close', () => console.log('closed!'));
port1.postMessage('foobar');
port1.close();CommonJS
'use strict';
const { MessageChannel } = require('node:worker_threads');
const { port1, port2 } = new MessageChannel();
// Prints:
// foobar
// closed!
port2.on('message', (message) => console.log(message));
port2.on('close', () => console.log('closed!'));
port1.postMessage('foobar');
port1.close();Событие: 'message'
-
value<any> Переданное значение
Событие 'message' возникает при получении любого сообщения и содержит клонированные входные данные, переданные в port.postMessage().
Обработчики этого события получают копию параметра value, переданного в postMessage(), и никаких дополнительных аргументов.
Событие: 'messageerror'
-
error<Error> Объект Error
Событие 'messageerror' возникает, если при десериализации сообщения произошла ошибка.
В настоящее время это событие возникает при ошибке во время создания отправленного объекта JS на принимающей стороне. Такие ситуации редки, но возможны, например, если в vm.Context получены объекты определённых API Node.js (где API Node.js сейчас недоступны).
port.close()
Отключает дальнейшую отправку сообщений в обе стороны соединения. Этот метод можно вызвать, если через этот MessagePort больше не будет происходить обмен данными.
Событие 'close' возникает у обоих экземпляров MessagePort, входящих в канал.
port.postMessage(value[, transferList])
-
value<any> -
transferList<Object[]>
Отправляет значение JavaScript принимающей стороне этого канала. Передача value выполняется в соответствии с алгоритмом структурного клонирования HTML.
В частности, основные отличия от JSON таковы:
-
valueможет содержать циклические ссылки. -
valueможет содержать экземпляры встроенных типов JS, таких какRegExp,BigInt,Map,Setи т. д. -
valueможет содержать типизированные массивы, использующие какArrayBuffer, так иSharedArrayBuffer. -
valueможет содержать экземплярыWebAssembly.Module. -
valueне может содержать собственные объекты (на основе C++), кроме следующих:
Модули JavaScript
import { MessageChannel } from 'node:worker_threads';
const { port1, port2 } = new MessageChannel();
port1.on('message', (message) => console.log(message));
const circularData = {};
circularData.foo = circularData;
// Prints: { foo: [Circular] }
port2.postMessage(circularData);CommonJS
'use strict';
const { MessageChannel } = require('node:worker_threads');
const { port1, port2 } = new MessageChannel();
port1.on('message', (message) => console.log(message));
const circularData = {};
circularData.foo = circularData;
// Prints: { foo: [Circular] }
port2.postMessage(circularData);transferList может быть списком объектов <ArrayBuffer>, MessagePort и FileHandle. После передачи они становятся недоступны на отправляющей стороне канала (даже если они не входят в value). В отличие от дочерних процессов, передача дескрипторов, таких как сетевые сокеты, в настоящее время не поддерживается.
Если value содержит экземпляры <SharedArrayBuffer>, к ним можно обращаться из любого потока. Их нельзя указывать в transferList.
value может содержать экземпляры ArrayBuffer, не включённые в transferList; в этом случае базовая память копируется, а не перемещается.
Модули JavaScript
import { MessageChannel } from 'node:worker_threads';
const { port1, port2 } = new MessageChannel();
port1.on('message', (message) => console.log(message));
const uint8Array = new Uint8Array([ 1, 2, 3, 4 ]);
// This posts a copy of `uint8Array`:
port2.postMessage(uint8Array);
// This does not copy data, but renders `uint8Array` unusable:
port2.postMessage(uint8Array, [ uint8Array.buffer ]);
// The memory for the `sharedUint8Array` is accessible from both the
// original and the copy received by `.on('message')`:
const sharedUint8Array = new Uint8Array(new SharedArrayBuffer(4));
port2.postMessage(sharedUint8Array);
// This transfers a freshly created message port to the receiver.
// This can be used, for example, to create communication channels between
// multiple `Worker` threads that are children of the same parent thread.
const otherChannel = new MessageChannel();
port2.postMessage({ port: otherChannel.port1 }, [ otherChannel.port1 ]);CommonJS
'use strict';
const { MessageChannel } = require('node:worker_threads');
const { port1, port2 } = new MessageChannel();
port1.on('message', (message) => console.log(message));
const uint8Array = new Uint8Array([ 1, 2, 3, 4 ]);
// This posts a copy of `uint8Array`:
port2.postMessage(uint8Array);
// This does not copy data, but renders `uint8Array` unusable:
port2.postMessage(uint8Array, [ uint8Array.buffer ]);
// The memory for the `sharedUint8Array` is accessible from both the
// original and the copy received by `.on('message')`:
const sharedUint8Array = new Uint8Array(new SharedArrayBuffer(4));
port2.postMessage(sharedUint8Array);
// This transfers a freshly created message port to the receiver.
// This can be used, for example, to create communication channels between
// multiple `Worker` threads that are children of the same parent thread.
const otherChannel = new MessageChannel();
port2.postMessage({ port: otherChannel.port1 }, [ otherChannel.port1 ]);Объект сообщения клонируется немедленно, поэтому после отправки его можно изменять без побочных эффектов.
Дополнительные сведения о механизмах сериализации и десериализации, лежащих в основе этого API, см. в API сериализации модуля node:v8.
Особенности передачи TypedArray и Buffer
Все экземпляры <TypedArray> | <Buffer> являются представлениями базового <ArrayBuffer>. Иными словами, фактические данные хранятся в ArrayBuffer, а объекты TypedArray и Buffer предоставляют способ просмотра данных и управления ими. Часто над одним экземпляром ArrayBuffer создаётся несколько представлений. При использовании списка передачи для передачи ArrayBuffer следует проявлять особую осторожность, поскольку в результате все экземпляры TypedArray и Buffer, использующие тот же ArrayBuffer, становятся недоступны.
const ab = new ArrayBuffer(10); const u1 = new Uint8Array(ab); const u2 = new Uint16Array(ab); console.log(u2.length); // prints 5 port.postMessage(u1, [u1.buffer]); console.log(u2.length); // prints 0 copy
В частности, для экземпляров Buffer возможность передачи или клонирования базового ArrayBuffer полностью зависит от способа создания экземпляров, который часто невозможно надёжно определить.
Для ArrayBuffer можно вызвать markAsUntransferable(), чтобы указать, что его всегда следует клонировать, а не передавать.
В зависимости от способа создания экземпляра Buffer он может владеть базовым ArrayBuffer или не владеть им. ArrayBuffer нельзя передавать, если не установлено, что экземпляр Buffer владеет им. В частности, передача экземпляров Buffer, созданных из внутреннего пула Buffer (например, с помощью Buffer.from() или Buffer.allocUnsafe()), невозможна: они всегда клонируются, что приводит к отправке копии всего пула Buffer. Это может непреднамеренно увеличить расход памяти и создать потенциальные проблемы безопасности.
Дополнительные сведения о пуле Buffer см. в разделе Buffer.allocUnsafe().
Объекты ArrayBuffer экземпляров Buffer, созданных с помощью Buffer.alloc() или Buffer.allocUnsafeSlow(), всегда можно передавать, однако в результате все остальные существующие представления этих ArrayBuffer становятся недоступны.
Особенности клонирования объектов с прототипами, классами и аксессорами
Поскольку при клонировании объектов используется алгоритм структурного клонирования HTML, неперечисляемые свойства, аксессоры свойств и прототипы объектов не сохраняются. В частности, объекты <Buffer> на принимающей стороне будут прочитаны как обычные объекты <Uint8Array>, а экземпляры классов JavaScript будут клонированы как обычные объекты JavaScript.
const b = Symbol('b');
class Foo {
#a = 1;
constructor() {
this[b] = 2;
this.c = 3;
}
get d() { return 4; }
}
const { port1, port2 } = new MessageChannel();
port1.onmessage = ({ data }) => console.log(data);
port2.postMessage(new Foo());
// Prints: { c: 3 } copy Это ограничение распространяется на многие встроенные объекты, например на глобальный объект URL:
const { port1, port2 } = new MessageChannel();
port1.onmessage = ({ data }) => console.log(data);
port2.postMessage(new URL('https://example.org'));
// Prints: { } copy
port.hasRef()
- Возвращает: <boolean>
Если значение равно true, объект MessagePort будет поддерживать активность цикла событий Node.js.
port.ref()
Противоположность unref(). Вызов ref() для порта, ранее подвергнутого unref(), не позволит программе завершиться, если это единственный оставшийся активный дескриптор (поведение по умолчанию). Если порт подвергнут ref(), повторный вызов ref() не даст результата.
Если обработчики добавляются или удаляются с помощью .on('message'), порт автоматически подвергается ref() или unref() в зависимости от наличия обработчиков события.
port.start()
Начинает получать сообщения через этот MessagePort. При использовании этого порта в качестве источника событий метод вызывается автоматически после добавления обработчиков событий 'message'.
Этот метод существует для совместимости с API Web MessagePort. В Node.js он полезен только для игнорирования сообщений, когда обработчик события отсутствует. Node.js также отличается обработкой .onmessage. Установка этого параметра автоматически вызывает .start(), но его сброс приводит к тому, что сообщения накапливаются в очереди, пока не будет назначен новый обработчик или порт не будет удалён.
port.unref()
Вызов unref() для порта позволяет потоку завершить работу, если это единственный активный дескриптор в системе событий. Если порт уже подвергнут unref(), повторный вызов unref() не даст результата.
Если обработчики добавляются или удаляются с помощью .on('message'), порт автоматически подвергается ref() или unref() в зависимости от наличия обработчиков события.
Класс: Worker
- Наследует: <EventEmitter>
Класс Worker представляет собой независимый поток выполнения JavaScript. В нем доступны большинство API Node.js.
Примечательные отличия среды Worker:
- Потоки
process.stdin,process.stdoutиprocess.stderrмогут перенаправляться родительским потоком. - Свойству
require('node:worker_threads').isMainThreadприсваивается значениеfalse. - Доступен порт сообщений
require('node:worker_threads').parentPort. -
process.exit()останавливает не всю программу, а только отдельный поток, аprocess.abort()недоступен. -
Методы
process.chdir()иprocess, устанавливающие идентификаторы группы или пользователя, недоступны. -
process.env— это копия переменных окружения родительского потока, если не указано иное. Изменения одной копии не видны в других потоках и нативных дополнениях (если только в качестве параметраenvконструкторуWorkerне передано значениеworker.SHARE_ENV). В Windows, в отличие от основного потока, копия переменных окружения работает с учетом регистра. -
process.titleизменить нельзя. - Сигналы не передаются через
process.on('...'). - Выполнение может остановиться в любой момент в результате вызова
worker.terminate(). - Каналы IPC родительских процессов недоступны.
- Модуль
trace_eventsне поддерживается. - Нативные дополнения можно загружать из нескольких потоков, только если они удовлетворяют определенным условиям.
Можно создавать экземпляры Worker внутри других Worker.
Как и в случае с Web Workers и модулем node:cluster, двустороннюю связь можно обеспечить передачей сообщений между потоками. Внутри Worker есть встроенная пара объектов MessagePort, уже связанных друг с другом при создании Worker. Хотя объект MessagePort на стороне родительского потока не предоставляется напрямую, его функциональность доступна через worker.postMessage() и событие worker.on('message') объекта Worker родительского потока.
Чтобы создавать пользовательские каналы обмена сообщениями (это рекомендуется вместо использования глобального канала по умолчанию, поскольку такой подход способствует разделению ответственности), можно создать объект MessageChannel в любом из потоков и передать один из объектов MessagePort этого MessageChannel в другой поток через уже существующий канал, например глобальный.
Подробнее о передаче сообщений и о типах значений JavaScript, которые можно передать через границу потока, см. в разделе port.postMessage().
Модули JavaScript
import assert from 'node:assert';
import {
Worker, MessageChannel, MessagePort, isMainThread, parentPort,
} from 'node:worker_threads';
if (isMainThread) {
const worker = new Worker(new URL(import.meta.url));
const subChannel = new MessageChannel();
worker.postMessage({ hereIsYourPort: subChannel.port1 }, [subChannel.port1]);
subChannel.port2.on('message', (value) => {
console.log('received:', value);
});
} else {
parentPort.once('message', (value) => {
assert(value.hereIsYourPort instanceof MessagePort);
value.hereIsYourPort.postMessage('the worker is sending this');
value.hereIsYourPort.close();
});
}CommonJS
'use strict';
const assert = require('node:assert');
const {
Worker, MessageChannel, MessagePort, isMainThread, parentPort,
} = require('node:worker_threads');
if (isMainThread) {
const worker = new Worker(__filename);
const subChannel = new MessageChannel();
worker.postMessage({ hereIsYourPort: subChannel.port1 }, [subChannel.port1]);
subChannel.port2.on('message', (value) => {
console.log('received:', value);
});
} else {
parentPort.once('message', (value) => {
assert(value.hereIsYourPort instanceof MessagePort);
value.hereIsYourPort.postMessage('the worker is sending this');
value.hereIsYourPort.close();
});
}
new Worker(filename[, options])
-
filename<string> | <URL> Путь к основному скрипту или модулю Worker. Должен быть абсолютным или относительным путем (то есть относительно текущего рабочего каталога), начинающимся с./или../, либо объектом WHATWGURLс использованием протоколаfile:илиdata:. При использовании URLdata:данные интерпретируются на основе MIME-типа с использованием загрузчика модулей ECMAScript. Еслиoptions.evalимеет значениеtrue, это строка с кодом JavaScript, а не путь. -
options<Object>-
argv<any[]> Список аргументов, которые будут преобразованы в строки и добавлены кprocess.argvв Worker. Это во многом аналогичноworkerData, но значения доступны в глобальном объектеprocess.argv, как если бы они были переданы скрипту в качестве параметров CLI. -
env<Object> Если задано, указывает начальное значениеprocess.envв потоке Worker. В качестве специального значения можно использоватьworker.SHARE_ENV, чтобы указать, что родительский и дочерний потоки должны использовать общие переменные окружения; в этом случае изменения объектаprocess.envодного потока влияют и на другой поток. По умолчанию:process.env. -
eval<boolean> Если имеет значениеtrueи первый аргумент являетсяstring, первый аргумент конструктора интерпретируется как скрипт, который выполняется после запуска Worker. -
execArgv<string[]> Список параметров CLI Node.js, передаваемых Worker. Параметры V8 (например,--max-old-space-size) и параметры, влияющие на процесс (например,--title), не поддерживаются. Если задано, это значение доступно в Worker какprocess.execArgv. По умолчанию параметры наследуются от родительского потока. -
stdin<boolean> Если этому параметру присвоено значениеtrue, тоworker.stdinпредоставляет записываемый поток, содержимое которого становится доступно в Worker какprocess.stdin. По умолчанию данные не предоставляются. -
stdout<boolean> Если этому параметру присвоено значениеtrue, тоworker.stdoutне перенаправляется автоматически вprocess.stdoutродительского потока. -
stderr<boolean> Если этому параметру присвоено значениеtrue, тоworker.stderrне перенаправляется автоматически вprocess.stderrродительского потока. -
workerData<any> Любое значение JavaScript, которое копируется и становится доступным какrequire('node:worker_threads').workerData. Копирование выполняется согласно алгоритму структурного клонирования HTML; если объект невозможно скопировать (например, если он содержит объектыfunction), возникает ошибка. -
trackUnmanagedFds<boolean> Если этому параметру присвоено значениеtrue, Worker отслеживает дескрипторы файлов низкого уровня, управляемые черезfs.open()иfs.close(), и закрывает их при завершении Worker, как и другие ресурсы, например сетевые сокеты или дескрипторы файлов, управляемые через APIFileHandle. Этот параметр автоматически наследуется всеми вложенными объектамиWorker. По умолчанию:true. -
transferList<Object[]> Если вworkerDataпередан один или несколько объектов типаMessagePort, для этих элементов необходимо указатьtransferList, иначе возникает ошибкаERR_MISSING_MESSAGE_PORT_IN_TRANSFER_LIST. Подробнее см. в разделеport.postMessage(). -
resourceLimits<Object> Необязательный набор ограничений ресурсов для нового экземпляра движка JS. Достижение этих пределов приводит к завершению экземпляраWorker. Эти ограничения влияют только на движок JS и не распространяются на внешние данные, включая объектыArrayBuffer. Даже при заданных ограничениях процесс все равно может аварийно завершиться при глобальной нехватке памяти.-
maxOldGenerationSizeMb<number> Максимальный размер основной кучи в МБ. Если задан аргумент командной строки--max-old-space-size, он переопределяет это значение. -
maxYoungGenerationSizeMb<number> Максимальный размер области кучи для недавно созданных объектов. Если задан аргумент командной строки--max-semi-space-size, он переопределяет это значение. -
codeRangeSizeMb<number> Размер предварительно выделенного диапазона памяти, используемого для сгенерированного кода. -
stackSizeMb<number> Максимальный размер стека по умолчанию для потока. Небольшие значения могут привести к тому, что экземпляры Worker станут непригодными для использования. По умолчанию:4.
-
-
name<string> Необязательное значениеname, добавляемое к заголовку Worker для отладки или идентификации; итоговый заголовок имеет вид[worker ${id}] ${name}. По умолчанию:''.
-
Событие: 'error'
-
err<Error>
Событие 'error' генерируется, если в потоке Worker возникает необработанное исключение. В этом случае Worker завершается.
Событие: 'exit'
-
exitCode<integer>
Событие 'exit' генерируется после остановки Worker. Если Worker завершился вызовом process.exit(), параметр exitCode содержит переданный код завершения. Если Worker был принудительно остановлен, параметр exitCode имеет значение 1.
Это последнее событие, генерируемое любым экземпляром Worker.
Событие: 'message'
-
value<any> Переданное значение
Событие 'message' генерируется, когда поток Worker вызывает require('node:worker_threads').parentPort.postMessage(). Подробнее см. в описании события port.on('message').
Все сообщения, отправленные из потока Worker, генерируются до события 'exit' объекта Worker.
Событие: 'messageerror'
-
error<Error> Объект Error
Событие 'messageerror' генерируется, если при десериализации сообщения произошла ошибка.
Событие: 'online'
Событие 'online' генерируется, когда поток Worker начинает выполнять код JavaScript.
worker.cpuUsage([prev])
- Возвращает: <Promise>
Этот метод возвращает Promise, который разрешается объектом, идентичным результату process.threadCpuUsage(), либо отклоняется с ошибкой ERR_WORKER_NOT_RUNNING, если Worker больше не выполняется. Этот метод позволяет получать статистику извне самого потока.
worker.getHeapSnapshot([options])
Возвращает читаемый поток со снимком текущего состояния Worker в V8. Подробнее см. в разделе v8.getHeapSnapshot().
Если поток Worker больше не выполняется, что может произойти до генерации события 'exit', возвращенный Promise немедленно отклоняется с ошибкой ERR_WORKER_NOT_RUNNING.
worker.getHeapStatistics()
- Возвращает: <Promise>
Этот метод возвращает Promise, который разрешается объектом, идентичным результату v8.getHeapStatistics(), либо отклоняется с ошибкой ERR_WORKER_NOT_RUNNING, если Worker больше не выполняется. Этот метод позволяет получать статистику извне самого потока.
worker.performance
Объект, с помощью которого можно получать данные о производительности экземпляра Worker. Аналогичен perf_hooks.performance.
performance.eventLoopUtilization([utilization1[, utilization2]])
-
utilization1<Object> Результат предыдущего вызоваeventLoopUtilization(). -
utilization2<Object> Результат предыдущего вызоваeventLoopUtilization(), выполненного доutilization1. - Возвращает: <Object>
Этот вызов аналогичен perf_hooks eventLoopUtilization(), но возвращает значения для экземпляра Worker.
Одно из отличий заключается в том, что, в отличие от основного потока, в Worker начальная загрузка выполняется в цикле событий. Поэтому данные об использовании цикла событий доступны сразу после начала выполнения скрипта Worker.
Неизменяющееся время idle не означает, что Worker завис при начальной загрузке. В следующем примере время idle не накапливается на протяжении всего времени работы Worker, но он по-прежнему может обрабатывать сообщения.
Модули JavaScript
import { Worker, isMainThread, parentPort } from 'node:worker_threads';
if (isMainThread) {
const worker = new Worker(new URL(import.meta.url));
setInterval(() => {
worker.postMessage('hi');
console.log(worker.performance.eventLoopUtilization());
}, 100).unref();
} else {
parentPort.on('message', () => console.log('msg')).unref();
(function r(n) {
if (--n < 0) return;
const t = Date.now();
while (Date.now() - t < 300);
setImmediate(r, n);
})(10);
}CommonJS
'use strict';
const { Worker, isMainThread, parentPort } = require('node:worker_threads');
if (isMainThread) {
const worker = new Worker(__filename);
setInterval(() => {
worker.postMessage('hi');
console.log(worker.performance.eventLoopUtilization());
}, 100).unref();
} else {
parentPort.on('message', () => console.log('msg')).unref();
(function r(n) {
if (--n < 0) return;
const t = Date.now();
while (Date.now() - t < 300);
setImmediate(r, n);
})(10);
}Данные об использовании цикла событий Worker доступны только после генерации события 'online'. Если вызвать этот метод до этого момента или после события 'exit', все свойства будут иметь значение 0.
worker.postMessage(value[, transferList])
-
value<any> -
transferList<Object[]>
Отправляет сообщение Worker, которое принимается через require('node:worker_threads').parentPort.on('message'). Подробнее см. в разделе port.postMessage().
worker.ref()
Противоположность unref(): вызов ref() для Worker, ранее вызвавшего unref(), не позволяет программе завершиться, если он остается единственным активным дескриптором (поведение по умолчанию). Если для Worker был вызван ref(), повторный вызов ref() не дает эффекта.
worker.resourceLimits
- Тип: <Object>
Предоставляет набор ограничений ресурсов движка JS для этого потока Worker. Если конструктору Worker был передан параметр resourceLimits, возвращаемое значение соответствует его значениям.
Если Worker остановлен, возвращается пустой объект.
worker.startCpuProfile(name)
Запускает профилирование ЦП с указанным name, а затем возвращает Promise, который выполняется с ошибкой или объектом с методом stop. Вызов метода stop останавливает сбор данных профиля, после чего возвращается Promise, который выполняется с ошибкой или данными профиля.
const { Worker } = require('node:worker_threads');
const worker = new Worker(`
const { parentPort } = require('worker_threads');
parentPort.on('message', () => {});
`, { eval: true });
worker.on('online', async () => {
const handle = await worker.startCpuProfile('demo');
const profile = await handle.stop();
console.log(profile);
worker.terminate();
}); copy
worker.stderr
- Тип: <stream.Readable>
Это читаемый поток, содержащий данные, записанные в process.stderr в потоке Worker. Если конструктору Worker не передан параметр stderr: true, данные перенаправляются в поток process.stderr родительского потока.
worker.stdin
- Тип: <null> | <stream.Writable>
Если конструктору Worker передан параметр stdin: true, это записываемый поток. Данные, записанные в этот поток, становятся доступны в потоке Worker как process.stdin.
worker.stdout
- Тип: <stream.Readable>
Это читаемый поток, содержащий данные, записанные в process.stdout в потоке Worker. Если конструктору Worker не передан параметр stdout: true, данные перенаправляются в поток process.stdout родительского потока.
worker.terminate()
- Возвращает: <Promise>
Как можно скорее останавливает выполнение всего кода JavaScript в потоке Worker. Возвращает Promise с кодом завершения, который выполняется при генерации события 'exit'.
worker.threadId
- Тип: <integer>
Целочисленный идентификатор указанного потока. Внутри потока Worker он доступен как require('node:worker_threads').threadId. Это значение уникально для каждого экземпляра Worker в пределах одного процесса.
worker.threadName
Строковый идентификатор указанного потока или null, если поток не выполняется. Внутри потока Worker он доступен как require('node:worker_threads').threadName.
worker.unref()
Вызов unref() для Worker позволяет потоку завершиться, если он остается единственным активным дескриптором в системе событий. Если для Worker уже был вызван unref(), повторный вызов unref() не дает эффекта.
worker[Symbol.asyncDispose]()
Псевдоним для worker.terminate().
Примечания
Синхронная блокировка стандартных потоков ввода-вывода
Worker используют передачу сообщений через <MessagePort> для взаимодействия с stdio. Это означает, что вывод stdio, поступающий из Worker, может быть заблокирован синхронным кодом на принимающей стороне, если этот код блокирует цикл событий Node.js.
Модули JavaScript
import {
Worker,
isMainThread,
} from 'node:worker_threads';
if (isMainThread) {
new Worker(new URL(import.meta.url));
for (let n = 0; n < 1e10; n++) {
// Looping to simulate work.
}
} else {
// This output will be blocked by the for loop in the main thread.
console.log('foo');
}CommonJS
'use strict';
const {
Worker,
isMainThread,
} = require('node:worker_threads');
if (isMainThread) {
new Worker(__filename);
for (let n = 0; n < 1e10; n++) {
// Looping to simulate work.
}
} else {
// This output will be blocked by the for loop in the main thread.
console.log('foo');
}Запуск потоков Worker из скриптов предварительной загрузки
Будьте осторожны при запуске потоков Worker из скриптов предварительной загрузки (скриптов, загружаемых и запускаемых с помощью флага командной строки -r). Если параметр execArgv явно не задан, новые потоки Worker автоматически наследуют флаги командной строки запущенного процесса и предварительно загружают те же скрипты, что и основной поток. Если скрипт предварительной загрузки безусловно запускает поток Worker, каждый запущенный поток будет порождать новый, пока приложение не завершится сбоем.
© 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-v22.x/docs/api/worker_threads.html