NNODE LOOP LABruntime observatoryNNEON · Статьи на 90 языках
CONNECTING
v24.18.0linux/x64
Microtasks → check → Redis jobs
07

Promises, setImmediate и BullMQ

Закрепите Promise-цепочки, async/await, комбинаторы и setImmediate, затем проследите реальный lifecycle BullMQ job через Redis.

PROCESS IDтекущий сервер
UPTIMEпосле запуска
LOOP DELAY P95perf_hooks
UTILIZATIONevent loop
HTTP ROUNDTRIPbrowser → server
LIVE TRACE

Временная шкала

ГОТОВ
0 ms
События появятся здесьЗапустите выбранный сценарий
#ВРЕМЯИСТОЧНИКСОБЫТИЕ
Ожидаю запуск эксперимента…
ГЛАВА 07
ПОДРОБНЫЙ РАЗБОР · ОТ БАЗЫ К КОДУ

Разбираем: Promises, setImmediate и BullMQ

Этот раздел можно читать до запуска опыта. После теории вернитесь к live trace и сопоставьте каждый шаг с реальным событием.

СНАЧАЛА ПРОСТЫМИ СЛОВАМИ

Promise — это не «фоновый поток», а коробка для одного будущего результата. Коробка сначала pending, затем навсегда становится fulfilled со значением либо rejected с ошибкой. `then`, `catch` и `finally` не меняют исходную коробку: каждый вызов создаёт следующий Promise в цепочке. BullMQ решает уже другую задачу — сохраняет задания в Redis, чтобы отдельные Workers могли взять их сейчас, позже или после перезапуска процесса.

ТЕХНИЧЕСКАЯ ОСНОВА

Executor внутри `new Promise(executor)` вызывается синхронно, а реакции `then/catch/finally` планируются как microtasks. `async`-функция всегда возвращает Promise; `await` приостанавливает только эту функцию и продолжает её через microtask. `setImmediate` регистрирует callback check-фазы Node и не означает «выполнить немедленно». BullMQ находится уровнем выше Event Loop: Queue записывает job в Redis, Worker получает и обрабатывает его, а QueueEvents наблюдает глобальные события через Redis.

Зачем это знатьЭти уровни часто смешивают: Promise координирует получение результата внутри программы, setImmediate переносит продолжение на следующую check-фазу, а BullMQ переживает HTTP-запросы и связывает несколько процессов или машин. Правильный выбор делает ошибки предсказуемыми и не теряет тяжёлую работу при рестарте.
ГДЕ ВЫПОЛНЯЕТСЯ РАБОТА
01ВАШ JS-КОДfunctions · callbacks
02NODE APIsfs · crypto · timers
03V8 + LIBUVheap · loop · pool
04ОПЕРАЦИОННАЯ СИСТЕМАI/O · threads · memory
01 · СЛОВАРЬ

Термины этого эксперимента

Сначала поймите слова — затем порядок выполнения.

01

Promise

Объект одного будущего outcome: fulfilled со значением или rejected с причиной. После settlement состояние уже не меняется.

02

Executor

Функция, переданная в `new Promise`. Она вызывается синхронно и получает `resolve` и `reject`; помещать туда код без реальной callback-обёртки обычно не нужно.

03

Pending / fulfilled / rejected

Три состояния Promise. Fulfilled и rejected вместе называют settled; переход из settled обратно в pending невозможен.

04

then chain

Цепочка новых Promise. Возвращённое значение становится результатом следующего звена, возвращённый Promise ожидается, а throw превращается в rejection.

05

Microtask

Высокоприоритетное продолжение Promise. Microtasks очищаются после текущего JavaScript callback до перехода к следующей фазе Event Loop.

06

async / await

`async` оборачивает результат в Promise. `await` подписывается на Promise и приостанавливает только текущую async-функцию, не поток Node.

07

Promise.all

Ждёт успеха всех входов и сохраняет порядок результатов. Первый rejection отклоняет общий Promise, но не отменяет оставшуюся работу.

08

Promise.allSettled

Ждёт все входы и возвращает объекты со статусами fulfilled/rejected; удобно, когда ошибка одного задания не должна скрывать остальные результаты.

09

Promise.race / any

`race` берёт первый settled outcome, включая ошибку. `any` берёт первый fulfilled, а если успехов нет — отклоняется с AggregateError.

10

setImmediate

Node API для callback check-фазы. Это не microtask и не гарантия «раньше setTimeout» без знания контекста регистрации.

11

BullMQ

Redis-backed очередь задач Node.js. Она хранит jobs вне памяти HTTP-процесса и координирует producers, workers, retries, delays и события.

12

Redis

Внешнее хранилище состояний BullMQ: waiting, active, delayed, completed и failed. Без доступного Redis BullMQ не является долговечной очередью.

13

Job

Именованное задание с сериализуемыми data и options: attempts, backoff, delay, priority и правила удаления результата.

14

BullMQ Worker

Потребитель jobs. Processor может быть async, но обычный Worker не превращает CPU-тяжёлый JavaScript в безопасную фоновую работу автоматически.

15

Idempotency

Свойство обработчика, при котором повторный запуск с тем же jobId не создаёт лишний платёж, письмо или запись. Важно для retries и восстановления stalled jobs.

02 · МЕХАНИКА

Что происходит по шагам

Каждый шаг соответствует наблюдаемому состоянию runtime.

  1. 01
    Executor запускается сейчас

    `new Promise(executor)` синхронно вызывает executor. Он может вызвать resolve, но код после resolve внутри executor всё равно дойдёт до конца.

  2. 02
    Outcome фиксируется один раз

    Первый resolve/reject побеждает. Если resolve получает другой Promise, новый Promise принимает его eventual outcome.

  3. 03
    Реакции становятся microtasks

    `then`, `catch` и `finally` не входят в текущий стек. Они получат шанс после завершения синхронного callback.

  4. 04
    Каждое звено возвращает новый Promise

    Следующий `then` ждёт именно return предыдущего. Забытый return часто запускает работу, но рвёт цепочку ожидания и обработки ошибок.

  5. 05
    await разворачивает то же правило

    Код до await выполняется сейчас; продолжение после settlement планируется microtask. try/catch ловит rejection так же, как `.catch()`.

  6. 06
    Комбинатор задаёт политику

    all, allSettled, race и any не запускают переданные операции — они получают уже созданные значения/Promise и по-разному агрегируют outcomes.

  7. 07
    setImmediate переносит фазу

    Callback попадает в check. Promise из текущего callback обычно выполнится раньше, потому что microtasks очищаются до продолжения фаз.

  8. 08
    Producer добавляет BullMQ job

    `queue.add()` возвращает Promise после записи данных job в Redis. Это подтверждение постановки в очередь, а не завершения бизнес-работы.

  9. 09
    Worker переводит waiting → active

    Свободный Worker резервирует job, вызывает processor и завершает его return/throw состоянием completed/failed.

  10. 10
    События, retries и cleanup

    QueueEvents сообщает глобальный lifecycle; attempts/backoff возвращают временно упавшие jobs; removeOnComplete/removeOnFail не дают Redis расти бесконечно.

03 · КОНТЕКСТ

Где результат требует оговорки

Эти детали объясняют, почему похожий код иногда даёт другой trace.

01

Promise не делает синхронный код асинхронным

`new Promise(() => heavyCpu())` немедленно выполнит heavyCpu в текущем потоке. Для CPU-bound работы нужны Worker Threads или отдельный процесс.

02

resolve — не мгновенный вызов then

resolve фиксирует outcome или начинает следовать другому Promise. Зарегистрированный then всё равно выполнится microtask после текущего стека.

03

finally прозрачен не всегда

Обычный return из finally не заменяет исходное значение, но throw или rejected Promise из finally заменит outcome новой ошибкой.

04

catch ловит ошибки выше по цепочке

Он получает rejection исходной операции, throw из then и rejected Promise, который вернул then. Ошибки в коде после отдельного, не возвращённого Promise могут уйти мимо.

05

Promise.all fail-fast, но не cancel-fast

Общий Promise отклоняется при первой ошибке, однако сетевой запрос, таймер или job продолжаются, пока их API не получит собственный сигнал отмены.

06

Timeout через race сам ничего не отменяет

Проигравший Promise продолжает работу. Для fetch нужен AbortController; для BullMQ — отдельная прикладная стратегия отмены job.

07

setImmediate зависит от места

Внутри одного I/O callback check выполняется перед следующим timers-проходом, поэтому immediate раньше нулевого timer. В main-модуле универсального порядка нет.

08

BullMQ Worker — роль, а не Worker Thread

BullMQ Worker может жить в том же или отдельном процессе/на другой машине. Его async processor хорош для I/O, но CPU-цикл всё равно блокирует Event Loop своего процесса.

09

Состояние находится в Redis

Promise исчезает вместе с процессом. BullMQ job остаётся доступным Workers после завершения HTTP-запроса и может пережить рестарт producer-а.

10

Доставка не равна уникальному эффекту

Retries и восстановление требуют идемпотентности: используйте стабильный jobId, уникальные ключи БД или таблицу обработанных операций.

01
Теория

Сначала разберитесь, какие части Node участвуют в выполнении.

02
Упрощённый код

Затем уберите служебные детали и рассмотрите только главную идею.

03
Runtime-код

После этого сопоставьте модель с кодом, который создаёт live trace.

04 · Упрощённый код

Минимальная модель без служебного кода

src/demos.js · учебный фрагментJavaScript
import { setImmediate as waitImmediate } from 'node:timers/promises';
import { Queue, QueueEvents, Worker } from 'bullmq';

const value = await Promise.resolve(2)
  .then(number => number * 3)
  .then(async number => {
    await waitImmediate(); // check phase
    return number + 1;
  })
  .catch(error => {
    console.error(error);
    return 0;
  })
  .finally(() => console.log('cleanup'));

const connection = { host: 'redis', port: 6379 };
const queue = new Queue('examples', { connection });
const events = new QueueEvents('examples', { connection });
const worker = new Worker('examples', async job => {
  return { doubled: job.data.value * 2 };
}, { connection });

await events.waitUntilReady();
const job = await queue.add('double', { value });
console.log(await job.waitUntilFinished(events, 5000));
05 · Runtime-код

Полный код, который выполняет сценарий

Это не альтернативный пример: ниже показаны функции и файлы, используемые кнопкой запуска.

ФАКТИЧЕСКИЙ SOURCE

Код сформирован из реальной серверной функции. Для сценариев с отдельным процессом или Worker показаны все участвующие файлы.

src/demos.js
сценарий241 строк
import { randomUUID } from 'node:crypto';
import { setImmediate as waitForImmediate } from 'node:timers/promises';
import {
  Queue,
  QueueEvents,
  Worker as BullWorker,
} from 'bullmq';

const sleep = (ms) =>
  new Promise((resolve) => setTimeout(resolve, ms));

function redisConnectionOptions(redisUrl) {
  const parsed = new URL(redisUrl);
  if (parsed.protocol !== 'redis:' && parsed.protocol !== 'rediss:') {
    throw new Error('REDIS_URL должен использовать redis:// или rediss://');
  }

  const database = Number.parseInt(parsed.pathname.slice(1) || '0', 10);
  return {
    host: parsed.hostname,
    port: Number.parseInt(parsed.port || '6379', 10),
    username: parsed.username || undefined,
    password: parsed.password || undefined,
    db: Number.isFinite(database) ? database : 0,
    tls: parsed.protocol === 'rediss:' ? {} : undefined,
    connectTimeout: 1200,
    maxRetriesPerRequest: null,
    enableOfflineQueue: false,
    retryStrategy: () => null,
  };
}

async function runBullMqRoundtrip(emit) {
  if (!process.env.REDIS_URL) {
    emit(
      'bullmq',
      'warning',
      'BullMQ-пример пропущен: задайте REDIS_URL или запустите проект через Docker Compose',
    );
    return;
  }

  const queueName = `node-loop-lab-${process.pid}-${randomUUID().slice(0, 8)}`;
  const connection = redisConnectionOptions(process.env.REDIS_URL);
  let queue;
  let worker;
  let queueEvents;
  const reportedConnectionErrors = new Set();
  const reportConnectionError = (source, error) => {
    const key = `${source}:${error.message}`;
    if (reportedConnectionErrors.has(key)) return;
    reportedConnectionErrors.add(key);
    emit(source, 'error', `${source} error: ${error.message}`);
  };

  try {
    queue = new Queue(queueName, {
      connection,
      defaultJobOptions: {
        attempts: 3,
        backoff: { type: 'exponential', delay: 100 },
        removeOnComplete: false,
        removeOnFail: false,
      },
    });
    queueEvents = new QueueEvents(queueName, { connection });
    worker = new BullWorker(
      queueName,
      async (job) => {
        emit(
          'bullmq-worker',
          'callback',
          `BullMQ Worker взял job ${job.id}: ${job.name}`,
        );
        await sleep(35);
        return {
          thumbnail: `${job.data.file}.webp`,
          processedBy: process.pid,
        };
      },
      { connection, concurrency: 1 },
    );

    queue.on('error', (error) => reportConnectionError('bullmq-queue', error));
    queueEvents.on('error', (error) =>
      reportConnectionError('bullmq-events', error),
    );
    worker.on('error', (error) =>
      reportConnectionError('bullmq-worker', error),
    );

    await Promise.all([
      queue.waitUntilReady(),
      queueEvents.waitUntilReady(),
      worker.waitUntilReady(),
    ]);

    emit(
      'bullmq-producer',
      'schedule',
      'Queue.add сохраняет job в Redis; HTTP-запрос не выполняет job сам',
    );
    const job = await queue.add('make-thumbnail', {
      file: 'avatar.png',
    });
    emit('redis', 'info', `Redis хранит job ${job.id} в состоянии waiting`);

    const result = await job.waitUntilFinished(queueEvents, 5000);
    emit(
      'bullmq-events',
      'result',
      `QueueEvents получил completed для job ${job.id}: ${result.thumbnail}`,
    );
  } catch (error) {
    emit(
      'bullmq',
      'warning',
      `BullMQ недоступен: ${error.message}. Promise-часть сценария уже выполнена.`,
    );
  } finally {
    await Promise.allSettled([worker?.close(), queueEvents?.close()]);
    if (queue) {
      try {
        await queue.obliterate({ force: true });
      } catch {
        // Redis мог стать недоступен во время очистки учебной очереди.
      }
      await queue.close().catch(() => {});
    }
  }
}

async function promisesImmediateBullMq(emit) {
  emit(
    'call-stack',
    'sync',
    'Создаём Promise: executor выполняется синхронно прямо сейчас',
  );

  const basePromise = new Promise((resolve) => {
    emit('promise-executor', 'sync', 'executor вызвал resolve(2)');
    resolve(2);
    emit(
      'promise-executor',
      'sync',
      'Код после resolve ещё выполняется, но повторно изменить outcome уже нельзя',
    );
  });

  const chainedPromise = basePromise
    .then((value) => {
      emit('microtasks', 'callback', `Первый then получил ${value} и вернул ${value * 3}`);
      return value * 3;
    })
    .then(async (value) => {
      await sleep(25);
      emit(
        'timers',
        'callback',
        `Второй then дождался Promise и получил ${value}`,
      );
      return value + 1;
    })
    .finally(() => {
      emit(
        'microtasks',
        'callback',
        'finally выполнился без подмены успешного результата',
      );
    });

  const recoveredPromise = Promise.reject(new Error('учебная ошибка'))
    .catch((error) => {
      emit('microtasks', 'warning', `catch обработал: ${error.message}`);
      return 'fallback';
    });

  const immediatePromise = new Promise((resolve) => {
    setImmediate(() => {
      emit(
        'check',
        'callback',
        'setImmediate callback выполняется в check-фазе',
      );
      resolve('immediate');
    });
  });

  emit(
    'call-stack',
    'schedule',
    'then/catch/setImmediate зарегистрированы; текущий стек освобождается',
  );
  const [chainResult, recovered] = await Promise.all([
    chainedPromise,
    recoveredPromise,
    immediatePromise,
  ]);
  emit(
    'result',
    'result',
    `Цепочка завершилась значением ${chainResult}; catch вернул ${recovered}`,
  );

  const aggregateTasks = [
    sleep(45).then(() => 'slow'),
    sleep(10).then(() => 'fast'),
    sleep(25).then(() => 'middle'),
  ];
  const allResult = await Promise.all(aggregateTasks);
  emit(
    'promise-all',
    'result',
    `Promise.all сохранил входной порядок: ${allResult.join(', ')}`,
  );

  const settled = await Promise.allSettled([
    Promise.resolve('ok'),
    Promise.reject(new Error('expected failure')),
  ]);
  emit(
    'promise-all-settled',
    'result',
    `Promise.allSettled вернул статусы: ${settled.map((item) => item.status).join(', ')}`,
  );

  const raceWinner = await Promise.race([
    sleep(30).then(() => 'timer 30ms'),
    waitForImmediate().then(() => 'setImmediate'),
  ]);
  emit('check', 'result', `Promise.race: первым завершился ${raceWinner}`);

  await waitForImmediate();
  emit(
    'check',
    'callback',
    'await setImmediate() из node:timers/promises продолжил функцию в check-фазе',
  );

  await runBullMqRoundtrip(emit);
}

Именно вызовы emit(...) превращаются в строки live trace. await и Promise удерживают HTTP-поток открытым до завершения сценария.

06 · РЕЦЕПТЫ

Практические шаблоны, которые можно подсмотреть

Сравнивайте цель, код и оговорки — не запоминайте синтаксис без модели.

01

01 · Оборачиваем callback в Promise

Используйте constructor на границе старого callback API; для уже promise-based API новый constructor не нужен.

import { readFile } from 'node:fs';

function readText(path) {
  return new Promise((resolve, reject) => {
    readFile(path, 'utf8', (error, text) => {
      if (error) return reject(error);
      resolve(text);
    });
  });
}
  • Executor запускается синхронно.
  • resolve/reject вызываются позже из callback.
  • Для fs уже существует fs/promises — это учебная обёртка.
02

02 · Цепочка then и обязательный return

Каждый return задаёт input следующего звена и связывает ошибки в одну цепочку.

getUser(42)
  .then(user => {
    return getOrders(user.id);
  })
  .then(orders => orders.filter(order => order.paid))
  .then(paid => console.log(paid))
  .catch(error => console.error(error))
  .finally(() => console.log('finished'));
  • Можно сократить первый then до `.then(user => getOrders(user.id))`.
  • throw внутри любого then попадёт в catch.
  • finally не получает paid или error аргументом.
03

03 · Тот же flow через async/await

async/await меняет синтаксис, но остаётся цепочкой Promise и microtasks.

async function printPaidOrders(userId) {
  try {
    const user = await getUser(userId);
    const orders = await getOrders(user.id);
    const paid = orders.filter(order => order.paid);
    console.log(paid);
    return paid;
  } catch (error) {
    console.error(error);
    throw error;
  } finally {
    console.log('finished');
  }
}
  • Возвращаемое paid станет fulfilled value.
  • Повторный throw не скрывает ошибку от вызывающего кода.
  • Последовательные await нужны только при реальной зависимости.
04

04 · Выбираем Promise-комбинатор

Сначала запускаем операции, затем выбираем политику ожидания.

const tasks = urls.map(url => fetch(url));

const everyResponse = await Promise.all(tasks);
const everyOutcome = await Promise.allSettled(tasks);
const firstSettled = await Promise.race(tasks);
const firstSuccess = await Promise.any(tasks);
  • all: все успехи или первая ошибка.
  • allSettled: полный отчёт без раннего reject.
  • race: первый success или error; any: первый success.
05

05 · Timeout с настоящей отменой fetch

Promise.race недостаточно: AbortController сообщает проигравшей операции, что её надо остановить.

async function fetchWithTimeout(url, timeoutMs = 2000) {
  const controller = new AbortController();
  const timer = setTimeout(() => controller.abort(), timeoutMs);

  try {
    return await fetch(url, { signal: controller.signal });
  } finally {
    clearTimeout(timer);
  }
}
  • AbortError всё равно надо обработать вызывающему коду.
  • Не каждое Promise API поддерживает отмену.
  • finally гарантирует очистку timer.
06

06 · setImmediate: callback и Promise API

Передаём управление check-фазе, не обещая задержку в миллисекундах.

import { setImmediate as waitImmediate } from 'node:timers/promises';

setImmediate(() => {
  console.log('check phase callback');
});

await waitImmediate();
console.log('continued in check phase');
  • Promise.then текущего callback обычно выполнится раньше.
  • Это Node API, не браузерный стандарт.
  • Для разбивки CPU-цикла одного immediate недостаточно — лучше Worker.
07

07 · BullMQ producer и Worker

Producer быстро сохраняет job в Redis, а Worker выполняет processor независимо от HTTP-контроллера.

import { Queue, Worker } from 'bullmq';

const connection = { host: '127.0.0.1', port: 6379 };
const queue = new Queue('email', { connection });

await queue.add('welcome', { userId: 42 }, {
  jobId: 'welcome:42',
  attempts: 3,
  backoff: { type: 'exponential', delay: 1000 },
  removeOnComplete: 1000,
  removeOnFail: 5000,
});

const worker = new Worker('email', async job => {
  return sendWelcomeEmail(job.data.userId);
}, { connection, concurrency: 10 });

worker.on('error', error => console.error(error));
  • Queue и Worker могут находиться в разных процессах или машинах.
  • jobId помогает дедупликации постановки, но не заменяет идемпотентность эффекта.
  • concurrency полезна для I/O; CPU-heavy processor выносите отдельно.
08

08 · QueueEvents и ожидание конкретного job

Различаем локальные Worker events и глобальные события очереди.

import { QueueEvents } from 'bullmq';

const queueEvents = new QueueEvents('email', { connection });
await queueEvents.waitUntilReady();

queueEvents.on('completed', ({ jobId, returnvalue }) => {
  console.log(jobId, returnvalue);
});

const job = await queue.add('welcome', { userId: 42 });
const result = await job.waitUntilFinished(queueEvents, 10_000);

await queueEvents.close();
  • QueueEvents использует отдельное Redis-соединение.
  • Timeout ожидания не обязан отменять сам job.
  • В долгоживущем сервисе QueueEvents обычно создают один раз, а не на каждый job.
07 · НЕ ПЕРЕПУТАЙТЕ

Популярные заблуждения

Миф слева, корректная модель справа.

МИФ

`new Promise(async (resolve) => ...)` — обычный способ использовать await.

НА САМОМ ДЕЛЕ

Promise-constructor ожидает синхронный executor; async executor создаёт второй Promise, ошибку которого внешний constructor не связывает автоматически. Обычно вынесите async-функцию отдельно.

МИФ

then изменяет исходный Promise.

НА САМОМ ДЕЛЕ

then всегда возвращает новый Promise. Исходный outcome остаётся прежним.

МИФ

Если не написать return внутри then, следующий then всё равно подождёт.

НА САМОМ ДЕЛЕ

Без return цепочка получает undefined и продолжает раньше запущенной внутри операции.

МИФ

await блокирует Event Loop до ответа.

НА САМОМ ДЕЛЕ

Он приостанавливает одну async-функцию. Event Loop продолжает обслуживать другие callbacks, пока awaited Promise pending.

МИФ

Promise.all запускает функции параллельно.

НА САМОМ ДЕЛЕ

Операции запускаются при вызове функций. Promise.all только агрегирует уже полученные Promise/значения.

МИФ

setImmediate означает выполнить прямо сейчас.

НА САМОМ ДЕЛЕ

Он только регистрирует callback check-фазы, который ждёт свободного стека и подходящего прохода Event Loop.

МИФ

BullMQ — это большая очередь callbacks Event Loop.

НА САМОМ ДЕЛЕ

BullMQ — распределённая прикладная очередь в Redis. Callback/microtask очереди принадлежат runtime конкретного Node-процесса.

МИФ

Успешный queue.add означает, что письмо уже отправлено.

НА САМОМ ДЕЛЕ

Он означает, что job сохранён. Завершение подтверждает Worker/состояние completed или соответствующее QueueEvents-событие.

08 · САМОПРОВЕРКА

Ответьте своими словами

Если ответ получается объяснить без терминов из документации, ментальная модель уже начала складываться.

  1. Почему console.log внутри executor появляется раньше console.log из then?
  2. Что получит следующий then, если предыдущий ничего не вернул?
  3. Чем Promise.allSettled полезнее Promise.all для пакетной обработки независимых файлов?
  4. Почему setImmediate и BullMQ нельзя называть одной и той же очередью?
  5. Как сделать отправку письма безопасной при повторной попытке BullMQ job?