Отслеживание асинхронного контекста
Исходный код: lib/async_hooks.js
Введение
Эти классы используются для ассоциации состояния и его распространения через обратные вызовы и цепочки промисов. Они позволяют хранить данные на протяжении всего жизненного цикла веб-запроса или любой другой асинхронной операции. Это аналогично хранилищу с локальной областью видимости в других языках.
Классы AsyncLocalStorage и AsyncResource являются частью модуля async_hooks.
MJS модули
import { AsyncLocalStorage, AsyncResource } from 'async_hooks';
CJS модули
const { AsyncLocalStorage, AsyncResource } = require('async_hooks'); Класс: AsyncLocalStorage
Этот класс создает хранилища, которые остаются согласованными во время асинхронных операций.
Хотя вы можете создать собственную реализацию на основе модуля async_hooks, рекомендуется использовать AsyncLocalStorage, поскольку это производительная и безопасная для памяти реализация, включающая значительные оптимизации, не очевидные для реализации.
Следующий пример использует AsyncLocalStorage, чтобы создать простой логгер, который присваивает идентификаторы входящим HTTP-запросам и включает их в сообщения, записанные в рамках каждого запроса.
MJS модули
import http from 'http';
import { AsyncLocalStorage } from '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('http');
const { AsyncLocalStorage } = require('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.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
}); Это переключение будет продолжаться на всю синхронную операцию. Это означает, что, например, если контекст был введен в обработчике события, последующие обработчики событий также будут работать в этом контексте, если не связаны с другим контекстом с помощью 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
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
}
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
} Использование с 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
});
} В этом примере хранилище доступно только в функции-обработчике и функциях, вызываемых foo. Вне run, вызов getStore вернет undefined.
Устранение неполадок: потеря контекста
В большинстве случаев код вашего приложения или библиотеки не должен испытывать проблем с AsyncLocalStorage. Но в редких случаях вы можете столкнуться с ситуациями, когда текущее хранилище теряется в одной из асинхронных операций. В таких случаях рассмотрите следующие варианты.
Если ваш код основан на обратных вызовах, достаточно преобразовать его в промисы с помощью util.promisify(), чтобы он начал работать с нативными промисами.
Если вам нужно продолжить использование API на основе обратных вызовов или ваш код предполагает реализацию пользовательского thenable, используйте класс AsyncResource для сопоставления асинхронной операции с правильным контекстом выполнения. Для этого вам нужно будет определить вызов функции, ответственный за потерю контекста. Вы можете сделать это, отобразив содержимое asyncLocalStorage.getStore() после вызовов, которые, по вашему мнению, ответственны за потерю. Когда код отображает undefined, последний вызванный обратный вызов, вероятно, ответственен за потерю контекста.
Класс: AsyncResource
Класс AsyncResource предназначен для расширения разработчиками асинхронных ресурсов. Используя его, пользователи могут легко запускать жизненный цикл своих собственных ресурсов.
Обработчик init будет срабатывать при создании экземпляра AsyncResource.
Ниже представлен обзор API AsyncResource.
Модули MJS
import { AsyncResource, executionAsyncId } from '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('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();
}
} Статический метод: 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 'worker_threads';
parentPort.on('message', (task) => {
parentPort.postMessage(task.a + task.b);
});
Модули CJS
const { parentPort } = require('worker_threads');
parentPort.on('message', (task) => {
parentPort.postMessage(task.a + task.b);
}); Пул потоков Worker вокруг него может использовать следующую структуру:
Модули MJS
import { AsyncResource } from 'async_hooks';
import { EventEmitter } from 'events';
import path from 'path';
import { Worker } from '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_processer.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('async_hooks');
const { EventEmitter } = require('events');
const path = require('path');
const { Worker } = require('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 'os';
const pool = new WorkerPool(os.cpus().length);
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('os');
const pool = new WorkerPool(os.cpus().length);
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 'http';
import { AsyncResource, executionAsyncId } from '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('http');
const { AsyncResource, executionAsyncId } = require('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-v16.x/docs/api/async_context.html