Отслеживание асинхронного контекста
Исходный код: lib/async_hooks.js
Введение
Эти классы используются для связывания состояния и его распространения через обратные вызовы и цепочки промисов. Они позволяют хранить данные на протяжении всего жизненного цикла веб-запроса или любого другого асинхронного периода. Аналогично хранилищу потоковых данных в других языках.
Классы AsyncLocalStorage и AsyncResource являются частью модуля node:async_hooks.
Модули MJS
import { AsyncLocalStorage, AsyncResource } from 'node:async_hooks';
Модули CJS
const { AsyncLocalStorage, AsyncResource } = require('node:async_hooks'); Класс: AsyncLocalStorage
Этот класс создает хранилища, которые остаются согласованными во время асинхронных операций.
Хотя вы можете создать собственную реализацию на основе модуля node:async_hooks, рекомендуется использовать AsyncLocalStorage, так как это производительная и безопасная в плане памяти реализация, включающая значительные, неявные для реализации, оптимизации.
В следующем примере используется AsyncLocalStorage, чтобы создать простой логгер, который присваивает идентификаторы входящим HTTP-запросам и включает их в сообщения, записанные в рамках каждого запроса.
Модули MJS
import http from 'node:http';
import { AsyncLocalStorage } from 'node:async_hooks';
const asyncLocalStorage = new AsyncLocalStorage();
function logWithId(msg) {
const id = asyncLocalStorage.getStore();
console.log(`${id !== undefined ? id : '-'}:`, msg);
}
let idSeq = 0;
http.createServer((req, res) => {
asyncLocalStorage.run(idSeq++, () => {
logWithId('start');
// Imagine any chain of async operations here
setImmediate(() => {
logWithId('finish');
res.end();
});
});
}).listen(8080);
http.get('http://localhost:8080');
http.get('http://localhost:8080');
// Prints:
// 0: start
// 1: start
// 0: finish
// 1: finish
Модули CJS
const http = require('node:http');
const { AsyncLocalStorage } = require('node:async_hooks');
const asyncLocalStorage = new AsyncLocalStorage();
function logWithId(msg) {
const id = asyncLocalStorage.getStore();
console.log(`${id !== undefined ? id : '-'}:`, msg);
}
let idSeq = 0;
http.createServer((req, res) => {
asyncLocalStorage.run(idSeq++, () => {
logWithId('start');
// Imagine any chain of async operations here
setImmediate(() => {
logWithId('finish');
res.end();
});
});
}).listen(8080);
http.get('http://localhost:8080');
http.get('http://localhost:8080');
// Prints:
// 0: start
// 1: start
// 0: finish
// 1: finish Каждый экземпляр AsyncLocalStorage поддерживает независимый контекст хранения. Несколько экземпляров могут безопасно существовать одновременно, не рискуя вмешиваться в данные друг друга.
new AsyncLocalStorage()
Создает новый экземпляр AsyncLocalStorage. Хранилище предоставляется только в рамках вызова run() или после вызова enterWith().
Статический метод: AsyncLocalStorage.bind(fn)
-
fn<Функция> Функция, которую нужно привязать к текущему контексту выполнения. - Возвращает: <Функция> Новая функция, которая вызывает
fnв захваченном контексте выполнения.
Привязывает заданную функцию к текущему контексту выполнения.
Статический метод: AsyncLocalStorage.snapshot()
- Возвращает: <Функция> Новая функция с сигнатурой
(fn: (...args) : R, ...args) : R.
Захватывает текущий контекст выполнения и возвращает функцию, принимающую функцию в качестве аргумента. При каждом вызове возвращаемой функции она вызывает переданную функцию в захваченном контексте.
const asyncLocalStorage = new AsyncLocalStorage(); const runInAsyncScope = asyncLocalStorage.run(123, () => AsyncLocalStorage.snapshot()); const result = asyncLocalStorage.run(321, () => runInAsyncScope(() => asyncLocalStorage.getStore())); console.log(result); // returns 123 copy
AsyncLocalStorage.snapshot() может заменить использование AsyncResource для простых целей отслеживания асинхронного контекста, например:
class Foo {
#runInAsyncScope = AsyncLocalStorage.snapshot();
get() { return this.#runInAsyncScope(() => asyncLocalStorage.getStore()); }
}
const foo = asyncLocalStorage.run(123, () => new Foo());
console.log(asyncLocalStorage.run(321, () => foo.get())); // returns 123 copy
asyncLocalStorage.disable()
Отключает экземпляр AsyncLocalStorage. Все последующие вызовы asyncLocalStorage.getStore() вернут undefined до тех пор, пока снова не будет вызван asyncLocalStorage.run() или asyncLocalStorage.enterWith().
При вызове asyncLocalStorage.disable(), все текущие связанные с экземпляром контексты будут завершены.
Вызов asyncLocalStorage.disable() необходим перед тем, как asyncLocalStorage может быть удален сборщиком мусора. Это не относится к хранилищам, предоставляемым asyncLocalStorage, так как эти объекты удаляются сборщиком мусора вместе с соответствующими асинхронными ресурсами.
Используйте этот метод, когда asyncLocalStorage больше не используется в текущем процессе.
asyncLocalStorage.getStore()
- Возвращает: <любой тип>
Возвращает текущее хранилище. Если вызвана вне асинхронного контекста, инициализированного вызовом asyncLocalStorage.run() или asyncLocalStorage.enterWith(), возвращает undefined.
asyncLocalStorage.enterWith(store)
-
store<любой тип>
Переходит в контекст на оставшуюся часть текущего синхронного выполнения и сохраняет хранилище через последующие асинхронные вызовы.
Пример:
const store = { id: 1 };
// Replaces previous store with the given store object
asyncLocalStorage.enterWith(store);
asyncLocalStorage.getStore(); // Returns the store object
someAsyncOperation(() => {
asyncLocalStorage.getStore(); // Returns the same object
}); copy Это переключение будет продолжаться в течение всего синхронного выполнения. Это означает, что если, например, контекст вводится в обработчик событий, последующие обработчики событий также будут выполняться в этом контексте, если не будут явно привязаны к другому контексту с помощью AsyncResource. Вот почему run() предпочтительнее enterWith(), если нет веских причин использовать последний метод.
const store = { id: 1 };
emitter.on('my-event', () => {
asyncLocalStorage.enterWith(store);
});
emitter.on('my-event', () => {
asyncLocalStorage.getStore(); // Returns the same object
});
asyncLocalStorage.getStore(); // Returns undefined
emitter.emit('my-event');
asyncLocalStorage.getStore(); // Returns the same object copy
asyncLocalStorage.run(store, callback[, ...args])
-
store<любой тип> -
callback<Функция> -
...args<любой тип>
Выполняет функцию синхронно в контексте и возвращает ее результат. Хранилище недоступно вне функции обратного вызова. Хранилище доступно для любых асинхронных операций, созданных внутри обратного вызова.
Дополнительные параметры args передаются в функцию обратного вызова.
Если функция обратного вызова генерирует ошибку, ошибка генерируется и run() также. Стек-трейс не затрагивается этим вызовом, и контекст завершается.
Пример:
const store = { id: 2 };
try {
asyncLocalStorage.run(store, () => {
asyncLocalStorage.getStore(); // Returns the store object
setTimeout(() => {
asyncLocalStorage.getStore(); // Returns the store object
}, 200);
throw new Error();
});
} catch (e) {
asyncLocalStorage.getStore(); // Returns undefined
// The error will be caught here
} copy
asyncLocalStorage.exit(callback[, ...args])
-
callback<Функция> -
...args<любой тип>
Выполняет функцию синхронно вне контекста и возвращает ее результат. Хранилище недоступно в функции обратного вызова или асинхронных операциях, созданных внутри обратного вызова. Любой вызов getStore() в функции обратного вызова всегда вернет undefined.
Дополнительные параметры args передаются в функцию обратного вызова.
Если функция обратного вызова генерирует ошибку, ошибка генерируется и exit() также. Стек-трейс не затрагивается этим вызовом, и контекст повторно вводится.
Пример:
// Within a call to run
try {
asyncLocalStorage.getStore(); // Returns the store object or value
asyncLocalStorage.exit(() => {
asyncLocalStorage.getStore(); // Returns undefined
throw new Error();
});
} catch (e) {
asyncLocalStorage.getStore(); // Returns the same object or value
// The error will be caught here
} copy Использование с async/await
Если в асинхронной функции должен быть выполнен только один вызов await в контексте, используйте следующий шаблон:
async function fn() {
await asyncLocalStorage.run(new Map(), () => {
asyncLocalStorage.getStore().set('key', value);
return foo(); // The return value of foo will be awaited
});
} copy В этом примере хранилище доступно только в функции обратного вызова и функциях, вызываемых foo. Вне run, вызов getStore вернет undefined.
Устранение неполадок: потеря контекста
В большинстве случаев AsyncLocalStorage работает без проблем. В редких случаях текущее хранилище теряется в одной из асинхронных операций.
Если ваш код основан на обратных вызовах, достаточно преобразовать его в промисы с помощью util.promisify(), чтобы он начал работать с нативными промисами.
Если вам необходимо использовать API, основанный на обратных вызовах, или ваш код предполагает реализацию кастомных thenable, используйте класс AsyncResource для связывания асинхронной операции с правильным контекстом выполнения. Найдите вызов функции, ответственный за потерю контекста, залогировав содержимое asyncLocalStorage.getStore() после вызовов, которые, по вашему мнению, ответственны за потерю. Когда код выводит undefined, последний вызванный обратный вызов, вероятно, ответственен за потерю контекста.
Класс: AsyncResource
Класс AsyncResource предназначен для расширения асинхронными ресурсами встраивающего компонента. С его помощью пользователи могут легко вызывать жизненные циклы собственных ресурсов.
Хук init будет вызван при создании объекта AsyncResource.
Ниже приведён обзор API AsyncResource.
Модули MJS
import { AsyncResource, executionAsyncId } from 'node:async_hooks';
// AsyncResource() is meant to be extended. Instantiating a
// new AsyncResource() also triggers init. If triggerAsyncId is omitted then
// async_hook.executionAsyncId() is used.
const asyncResource = new AsyncResource(
type, { triggerAsyncId: executionAsyncId(), requireManualDestroy: false },
);
// Run a function in the execution context of the resource. This will
// * establish the context of the resource
// * trigger the AsyncHooks before callbacks
// * call the provided function `fn` with the supplied arguments
// * trigger the AsyncHooks after callbacks
// * restore the original execution context
asyncResource.runInAsyncScope(fn, thisArg, ...args);
// Call AsyncHooks destroy callbacks.
asyncResource.emitDestroy();
// Return the unique ID assigned to the AsyncResource instance.
asyncResource.asyncId();
// Return the trigger ID for the AsyncResource instance.
asyncResource.triggerAsyncId();
Модули CJS
const { AsyncResource, executionAsyncId } = require('node:async_hooks');
// AsyncResource() is meant to be extended. Instantiating a
// new AsyncResource() also triggers init. If triggerAsyncId is omitted then
// async_hook.executionAsyncId() is used.
const asyncResource = new AsyncResource(
type, { triggerAsyncId: executionAsyncId(), requireManualDestroy: false },
);
// Run a function in the execution context of the resource. This will
// * establish the context of the resource
// * trigger the AsyncHooks before callbacks
// * call the provided function `fn` with the supplied arguments
// * trigger the AsyncHooks after callbacks
// * restore the original execution context
asyncResource.runInAsyncScope(fn, thisArg, ...args);
// Call AsyncHooks destroy callbacks.
asyncResource.emitDestroy();
// Return the unique ID assigned to the AsyncResource instance.
asyncResource.asyncId();
// Return the trigger ID for the AsyncResource instance.
asyncResource.triggerAsyncId();
new AsyncResource(type[, options])
-
type<строка> Тип асинхронного события. -
options<Объект>-
triggerAsyncId<число> Идентификатор контекста выполнения, создавшего это асинхронное событие. По умолчанию:executionAsyncId(). -
requireManualDestroy<логическое значение> Если установлено вtrue, отключаетemitDestroyпри сборе мусора объекта. Обычно не нужно устанавливать (даже еслиemitDestroyвызывается вручную), если только не извлекаетсяasyncIdресурса и не вызываетсяemitDestroyчувствительного API с ним. При установке вfalse, вызовemitDestroyпри сборе мусора выполнится только при наличии хотя бы одного активного хукаdestroy. По умолчанию:false.
-
Пример использования:
class DBQuery extends AsyncResource {
constructor(db) {
super('DBQuery');
this.db = db;
}
getInfo(query, callback) {
this.db.get(query, (err, data) => {
this.runInAsyncScope(callback, null, err, data);
});
}
close() {
this.db = null;
this.emitDestroy();
}
} copy Статический метод: AsyncResource.bind(fn[, type[, thisArg]])
-
fn<Функция> Функция для привязки к текущему контексту выполнения. -
type<строка> Необязательное имя для ассоциации с базовымAsyncResource. -
thisArg<любой>
Привязывает заданную функцию к текущему контексту выполнения.
Возвращаемая функция будет иметь свойство asyncResource, ссылающееся на AsyncResource, к которому привязана функция.
asyncResource.bind(fn[, thisArg])
Привязывает заданную функцию для выполнения к области видимости этого AsyncResource.
Возвращаемая функция будет иметь свойство asyncResource, ссылающееся на AsyncResource, к которому привязана функция.
asyncResource.runInAsyncScope(fn[, thisArg, ...args])
-
fn<Функция> Функция для вызова в контексте выполнения этого асинхронного ресурса. -
thisArg<любой> Приёмник, используемый для вызова функции. -
...args<любой> Необязательные аргументы для передачи функции.
Вызывает предоставленную функцию с предоставленными аргументами в контексте выполнения асинхронного ресурса. Это установит контекст, вызовет хуки AsyncHooks до обратных вызовов, вызовет функцию, вызовет хуки AsyncHooks после обратных вызовов и восстановит исходный контекст выполнения.
asyncResource.emitDestroy()
- Возвращает: <AsyncResource> Ссылка на
asyncResource.
Вызывает все хуки destroy. Это должно вызываться только один раз. Будет выброшено исключение, если оно вызывается более одного раза. Это обязательно вызывать вручную. Если ресурс оставлен для сбора сборщиком мусора, то хуки destroy никогда не будут вызваны.
asyncResource.asyncId()
- Возвращает: <число> Уникальный
asyncId, назначенный ресурсу.
asyncResource.triggerAsyncId()
- Возвращает: <число> Тот же
triggerAsyncId, что передаётся конструкторуAsyncResource.
Использование AsyncResource для пула потоков Worker
Следующий пример демонстрирует, как использовать класс AsyncResource для правильного отслеживания асинхронных операций в пуле Worker. Другие пулы ресурсов, например, пулы подключений к базе данных, могут следовать аналогичной модели.
Предположим, задача заключается в сложении двух чисел, используя файл под названием task_processor.js со следующим содержимым:
Модули MJS
import { parentPort } from 'node:worker_threads';
parentPort.on('message', (task) => {
parentPort.postMessage(task.a + task.b);
});
Модули CJS
const { parentPort } = require('node:worker_threads');
parentPort.on('message', (task) => {
parentPort.postMessage(task.a + task.b);
}); Пул потоков Worker вокруг него может использовать следующую структуру:
Модули MJS
import { AsyncResource } from 'node:async_hooks';
import { EventEmitter } from 'node:events';
import path from 'node:path';
import { Worker } from 'node:worker_threads';
const kTaskInfo = Symbol('kTaskInfo');
const kWorkerFreedEvent = Symbol('kWorkerFreedEvent');
class WorkerPoolTaskInfo extends AsyncResource {
constructor(callback) {
super('WorkerPoolTaskInfo');
this.callback = callback;
}
done(err, result) {
this.runInAsyncScope(this.callback, null, err, result);
this.emitDestroy(); // `TaskInfo`s are used only once.
}
}
export default class WorkerPool extends EventEmitter {
constructor(numThreads) {
super();
this.numThreads = numThreads;
this.workers = [];
this.freeWorkers = [];
this.tasks = [];
for (let i = 0; i < numThreads; i++)
this.addNewWorker();
// Any time the kWorkerFreedEvent is emitted, dispatch
// the next task pending in the queue, if any.
this.on(kWorkerFreedEvent, () => {
if (this.tasks.length > 0) {
const { task, callback } = this.tasks.shift();
this.runTask(task, callback);
}
});
}
addNewWorker() {
const worker = new Worker(new URL('task_processor.js', import.meta.url));
worker.on('message', (result) => {
// In case of success: Call the callback that was passed to `runTask`,
// remove the `TaskInfo` associated with the Worker, and mark it as free
// again.
worker[kTaskInfo].done(null, result);
worker[kTaskInfo] = null;
this.freeWorkers.push(worker);
this.emit(kWorkerFreedEvent);
});
worker.on('error', (err) => {
// In case of an uncaught exception: Call the callback that was passed to
// `runTask` with the error.
if (worker[kTaskInfo])
worker[kTaskInfo].done(err, null);
else
this.emit('error', err);
// Remove the worker from the list and start a new Worker to replace the
// current one.
this.workers.splice(this.workers.indexOf(worker), 1);
this.addNewWorker();
});
this.workers.push(worker);
this.freeWorkers.push(worker);
this.emit(kWorkerFreedEvent);
}
runTask(task, callback) {
if (this.freeWorkers.length === 0) {
// No free threads, wait until a worker thread becomes free.
this.tasks.push({ task, callback });
return;
}
const worker = this.freeWorkers.pop();
worker[kTaskInfo] = new WorkerPoolTaskInfo(callback);
worker.postMessage(task);
}
close() {
for (const worker of this.workers) worker.terminate();
}
}
Модули CJS
const { AsyncResource } = require('node:async_hooks');
const { EventEmitter } = require('node:events');
const path = require('node:path');
const { Worker } = require('node:worker_threads');
const kTaskInfo = Symbol('kTaskInfo');
const kWorkerFreedEvent = Symbol('kWorkerFreedEvent');
class WorkerPoolTaskInfo extends AsyncResource {
constructor(callback) {
super('WorkerPoolTaskInfo');
this.callback = callback;
}
done(err, result) {
this.runInAsyncScope(this.callback, null, err, result);
this.emitDestroy(); // `TaskInfo`s are used only once.
}
}
class WorkerPool extends EventEmitter {
constructor(numThreads) {
super();
this.numThreads = numThreads;
this.workers = [];
this.freeWorkers = [];
this.tasks = [];
for (let i = 0; i < numThreads; i++)
this.addNewWorker();
// Any time the kWorkerFreedEvent is emitted, dispatch
// the next task pending in the queue, if any.
this.on(kWorkerFreedEvent, () => {
if (this.tasks.length > 0) {
const { task, callback } = this.tasks.shift();
this.runTask(task, callback);
}
});
}
addNewWorker() {
const worker = new Worker(path.resolve(__dirname, 'task_processor.js'));
worker.on('message', (result) => {
// In case of success: Call the callback that was passed to `runTask`,
// remove the `TaskInfo` associated with the Worker, and mark it as free
// again.
worker[kTaskInfo].done(null, result);
worker[kTaskInfo] = null;
this.freeWorkers.push(worker);
this.emit(kWorkerFreedEvent);
});
worker.on('error', (err) => {
// In case of an uncaught exception: Call the callback that was passed to
// `runTask` with the error.
if (worker[kTaskInfo])
worker[kTaskInfo].done(err, null);
else
this.emit('error', err);
// Remove the worker from the list and start a new Worker to replace the
// current one.
this.workers.splice(this.workers.indexOf(worker), 1);
this.addNewWorker();
});
this.workers.push(worker);
this.freeWorkers.push(worker);
this.emit(kWorkerFreedEvent);
}
runTask(task, callback) {
if (this.freeWorkers.length === 0) {
// No free threads, wait until a worker thread becomes free.
this.tasks.push({ task, callback });
return;
}
const worker = this.freeWorkers.pop();
worker[kTaskInfo] = new WorkerPoolTaskInfo(callback);
worker.postMessage(task);
}
close() {
for (const worker of this.workers) worker.terminate();
}
}
module.exports = WorkerPool; Без явного отслеживания, добавленного объектами WorkerPoolTaskInfo, может показаться, что обратные вызовы связаны с отдельными объектами Worker. Однако создание Worker не связано с созданием задач и не предоставляет информации о том, когда задачи были запланированы.
Этот пул можно использовать следующим образом:
Модули MJS
import WorkerPool from './worker_pool.js';
import os from 'node:os';
const pool = new WorkerPool(os.availableParallelism());
let finished = 0;
for (let i = 0; i < 10; i++) {
pool.runTask({ a: 42, b: 100 }, (err, result) => {
console.log(i, err, result);
if (++finished === 10)
pool.close();
});
}
Модули CJS
const WorkerPool = require('./worker_pool.js');
const os = require('node:os');
const pool = new WorkerPool(os.availableParallelism());
let finished = 0;
for (let i = 0; i < 10; i++) {
pool.runTask({ a: 42, b: 100 }, (err, result) => {
console.log(i, err, result);
if (++finished === 10)
pool.close();
});
} Интеграция AsyncResource с EventEmitter
Обработчики событий, инициированные EventEmitter, могут выполняться в другом контексте выполнения, отличном от контекста, который был активен при вызове eventEmitter.on().
Следующий пример демонстрирует, как использовать класс AsyncResource для правильной ассоциации обработчика событий с правильным контекстом выполнения. Тот же подход может быть применён к Stream или аналогичному классу, управляемому событиями.
Модули MJS
import { createServer } from 'node:http';
import { AsyncResource, executionAsyncId } from 'node:async_hooks';
const server = createServer((req, res) => {
req.on('close', AsyncResource.bind(() => {
// Execution context is bound to the current outer scope.
}));
req.on('close', () => {
// Execution context is bound to the scope that caused 'close' to emit.
});
res.end();
}).listen(3000);
Модули CJS
const { createServer } = require('node:http');
const { AsyncResource, executionAsyncId } = require('node:async_hooks');
const server = createServer((req, res) => {
req.on('close', AsyncResource.bind(() => {
// Execution context is bound to the current outer scope.
}));
req.on('close', () => {
// Execution context is bound to the scope that caused 'close' to emit.
});
res.end();
}).listen(3000);
© 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/async_context.html