Отслеживание асинхронного контекста
Исходный код: 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не извлекается, и к нему не применяется чувствительный APIemitDestroy. Если установлено в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.bind(fn[, thisArg])
-
fn<Функция> Функция, которая должна быть привязана к текущему контекстуAsyncResource. -
thisArg<любой тип>
Привязывает данную функцию к текущему контексту 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-v20.x/docs/api/async_context.html