Spec-Zone.ru › Node.js 16 LTS

Отслеживание асинхронного контекста

Устойчивость: 2 - Стабильно

Исходный код: lib/async_hooks.js

Введение

Эти классы используются для ассоциации состояния и его распространения через обратные вызовы и цепочки промисов. Они позволяют хранить данные на протяжении всего жизненного цикла веб-запроса или любой другой асинхронной операции. Это аналогично хранилищу с локальной областью видимости в других языках.

Классы AsyncLocalStorage и AsyncResource являются частью модуля async_hooks.

MJS модули

import { AsyncLocalStorage, AsyncResource } from 'async_hooks';

CJS модули

const { AsyncLocalStorage, AsyncResource } = require('async_hooks');

Класс: AsyncLocalStorage

История
Версия Изменения
v16.4.0

AsyncLocalStorage теперь стабилен. Раньше он был экспериментальным.

v13.10.0, v12.17.0

Добавлен в: v13.10.0, v12.17.0

Этот класс создает хранилища, которые остаются согласованными во время асинхронных операций.

Хотя вы можете создать собственную реализацию на основе модуля 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()

Добавлен в: v13.10.0, v12.17.0

Создает новый экземпляр AsyncLocalStorage. Хранилище доступно только внутри вызова run() или после вызова enterWith().

asyncLocalStorage.disable()

Добавлен в: v13.10.0, v12.17.0
Устойчивость: 1 - Экспериментальная

Отключает экземпляр AsyncLocalStorage. Все последующие вызовы asyncLocalStorage.getStore() вернут undefined до тех пор, пока снова не будут вызваны asyncLocalStorage.run() или asyncLocalStorage.enterWith().

При вызове asyncLocalStorage.disable(), все текущие контексты, связанные с экземпляром, будут завершены.

Вызов asyncLocalStorage.disable() необходим перед тем, как asyncLocalStorage может быть удален из памяти. Это не относится к хранилищам, предоставляемым asyncLocalStorage, поскольку эти объекты удаляются из памяти вместе с соответствующими асинхронными ресурсами.

Используйте этот метод, когда asyncLocalStorage больше не используется в текущем процессе.

asyncLocalStorage.getStore()

Добавлен в: v13.10.0, v12.17.0
  • Возвращает: <любой тип>

Возвращает текущее хранилище. Если вызов осуществляется вне асинхронного контекста, инициированного вызовом asyncLocalStorage.run() или asyncLocalStorage.enterWith(), он возвращает undefined.

asyncLocalStorage.enterWith(store)

Добавлен в: v13.11.0, v12.17.0
Устойчивость: 1 - Экспериментальная
  • 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])

Добавлен в: v13.10.0, v12.17.0
  • 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])

Добавлен в: v13.10.0, v12.17.0
Устойчивость: 1 - Экспериментальная
  • 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

История
Версия Изменения
v16.4.0

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 извлекается, и чувствительный API emitDestroy вызывается с ним. При значении 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]])

История
Версия Изменения
v16.0.0

Добавлен необязательный thisArg.

v14.8.0, v12.19.0

Добавлен в: v14.8.0, v12.19.0

  • fn <Функция> Функция, которая должна быть привязана к текущему контексту выполнения.
  • type <строка> Необязательное имя, которое нужно сопоставить с базовым AsyncResource.
  • thisArg <любое>

Привязывает заданную функцию к текущему контексту выполнения.

Возвращаемая функция будет иметь свойство asyncResource, ссылающееся на AsyncResource, к которому привязана функция.

asyncResource.bind(fn[, thisArg])

История
Версия Изменения
v16.0.0

Добавлен необязательный thisArg.

v14.8.0, v12.19.0

Добавлен в: v14.8.0, v12.19.0

  • fn <Функция> Функция, которая должна быть привязана к текущему AsyncResource.
  • thisArg <любое>

Привязывает указанную функцию к области видимости этого AsyncResource.

Возвращаемая функция будет иметь свойство asyncResource, ссылающееся на AsyncResource, к которому привязана функция.

asyncResource.runInAsyncScope(fn[, thisArg, ...args])

Добавлен в: v9.6.0
  • 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

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API