Promises, setImmediate и BullMQ
Закрепите Promise-цепочки, async/await, комбинаторы и setImmediate, затем проследите реальный lifecycle BullMQ job через Redis.
Временная шкала
Разбираем: 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
Объект одного будущего outcome: fulfilled со значением или rejected с причиной. После settlement состояние уже не меняется.
Executor
Функция, переданная в `new Promise`. Она вызывается синхронно и получает `resolve` и `reject`; помещать туда код без реальной callback-обёртки обычно не нужно.
Pending / fulfilled / rejected
Три состояния Promise. Fulfilled и rejected вместе называют settled; переход из settled обратно в pending невозможен.
then chain
Цепочка новых Promise. Возвращённое значение становится результатом следующего звена, возвращённый Promise ожидается, а throw превращается в rejection.
Microtask
Высокоприоритетное продолжение Promise. Microtasks очищаются после текущего JavaScript callback до перехода к следующей фазе Event Loop.
async / await
`async` оборачивает результат в Promise. `await` подписывается на Promise и приостанавливает только текущую async-функцию, не поток Node.
Promise.all
Ждёт успеха всех входов и сохраняет порядок результатов. Первый rejection отклоняет общий Promise, но не отменяет оставшуюся работу.
Promise.allSettled
Ждёт все входы и возвращает объекты со статусами fulfilled/rejected; удобно, когда ошибка одного задания не должна скрывать остальные результаты.
Promise.race / any
`race` берёт первый settled outcome, включая ошибку. `any` берёт первый fulfilled, а если успехов нет — отклоняется с AggregateError.
setImmediate
Node API для callback check-фазы. Это не microtask и не гарантия «раньше setTimeout» без знания контекста регистрации.
BullMQ
Redis-backed очередь задач Node.js. Она хранит jobs вне памяти HTTP-процесса и координирует producers, workers, retries, delays и события.
Redis
Внешнее хранилище состояний BullMQ: waiting, active, delayed, completed и failed. Без доступного Redis BullMQ не является долговечной очередью.
Job
Именованное задание с сериализуемыми data и options: attempts, backoff, delay, priority и правила удаления результата.
BullMQ Worker
Потребитель jobs. Processor может быть async, но обычный Worker не превращает CPU-тяжёлый JavaScript в безопасную фоновую работу автоматически.
Idempotency
Свойство обработчика, при котором повторный запуск с тем же jobId не создаёт лишний платёж, письмо или запись. Важно для retries и восстановления stalled jobs.
Что происходит по шагам
Каждый шаг соответствует наблюдаемому состоянию runtime.
- 01Executor запускается сейчас
`new Promise(executor)` синхронно вызывает executor. Он может вызвать resolve, но код после resolve внутри executor всё равно дойдёт до конца.
- 02Outcome фиксируется один раз
Первый resolve/reject побеждает. Если resolve получает другой Promise, новый Promise принимает его eventual outcome.
- 03Реакции становятся microtasks
`then`, `catch` и `finally` не входят в текущий стек. Они получат шанс после завершения синхронного callback.
- 04Каждое звено возвращает новый Promise
Следующий `then` ждёт именно return предыдущего. Забытый return часто запускает работу, но рвёт цепочку ожидания и обработки ошибок.
- 05await разворачивает то же правило
Код до await выполняется сейчас; продолжение после settlement планируется microtask. try/catch ловит rejection так же, как `.catch()`.
- 06Комбинатор задаёт политику
all, allSettled, race и any не запускают переданные операции — они получают уже созданные значения/Promise и по-разному агрегируют outcomes.
- 07setImmediate переносит фазу
Callback попадает в check. Promise из текущего callback обычно выполнится раньше, потому что microtasks очищаются до продолжения фаз.
- 08Producer добавляет BullMQ job
`queue.add()` возвращает Promise после записи данных job в Redis. Это подтверждение постановки в очередь, а не завершения бизнес-работы.
- 09Worker переводит waiting → active
Свободный Worker резервирует job, вызывает processor и завершает его return/throw состоянием completed/failed.
- 10События, retries и cleanup
QueueEvents сообщает глобальный lifecycle; attempts/backoff возвращают временно упавшие jobs; removeOnComplete/removeOnFail не дают Redis расти бесконечно.
Где результат требует оговорки
Эти детали объясняют, почему похожий код иногда даёт другой trace.
Promise не делает синхронный код асинхронным
`new Promise(() => heavyCpu())` немедленно выполнит heavyCpu в текущем потоке. Для CPU-bound работы нужны Worker Threads или отдельный процесс.
resolve — не мгновенный вызов then
resolve фиксирует outcome или начинает следовать другому Promise. Зарегистрированный then всё равно выполнится microtask после текущего стека.
finally прозрачен не всегда
Обычный return из finally не заменяет исходное значение, но throw или rejected Promise из finally заменит outcome новой ошибкой.
catch ловит ошибки выше по цепочке
Он получает rejection исходной операции, throw из then и rejected Promise, который вернул then. Ошибки в коде после отдельного, не возвращённого Promise могут уйти мимо.
Promise.all fail-fast, но не cancel-fast
Общий Promise отклоняется при первой ошибке, однако сетевой запрос, таймер или job продолжаются, пока их API не получит собственный сигнал отмены.
Timeout через race сам ничего не отменяет
Проигравший Promise продолжает работу. Для fetch нужен AbortController; для BullMQ — отдельная прикладная стратегия отмены job.
setImmediate зависит от места
Внутри одного I/O callback check выполняется перед следующим timers-проходом, поэтому immediate раньше нулевого timer. В main-модуле универсального порядка нет.
BullMQ Worker — роль, а не Worker Thread
BullMQ Worker может жить в том же или отдельном процессе/на другой машине. Его async processor хорош для I/O, но CPU-цикл всё равно блокирует Event Loop своего процесса.
Состояние находится в Redis
Promise исчезает вместе с процессом. BullMQ job остаётся доступным Workers после завершения HTTP-запроса и может пережить рестарт producer-а.
Доставка не равна уникальному эффекту
Retries и восстановление требуют идемпотентности: используйте стабильный jobId, уникальные ключи БД или таблицу обработанных операций.
Сначала разберитесь, какие части Node участвуют в выполнении.
Затем уберите служебные детали и рассмотрите только главную идею.
После этого сопоставьте модель с кодом, который создаёт live trace.
Минимальная модель без служебного кода
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));Полный код, который выполняет сценарий
Это не альтернативный пример: ниже показаны функции и файлы, используемые кнопкой запуска.
Код сформирован из реальной серверной функции. Для сценариев с отдельным процессом или Worker показаны все участвующие файлы.
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-поток открытым до завершения сценария.
Практические шаблоны, которые можно подсмотреть
Сравнивайте цель, код и оговорки — не запоминайте синтаксис без модели.
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 · Цепочка 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 · Тот же 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 · Выбираем 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 · 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 · 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 · 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 · 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.
Популярные заблуждения
Миф слева, корректная модель справа.
`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-событие.
Ответьте своими словами
Если ответ получается объяснить без терминов из документации, ментальная модель уже начала складываться.
- Почему console.log внутри executor появляется раньше console.log из then?
- Что получит следующий then, если предыдущий ничего не вернул?
- Чем Promise.allSettled полезнее Promise.all для пакетной обработки независимых файлов?
- Почему setImmediate и BullMQ нельзя называть одной и той же очередью?
- Как сделать отправку письма безопасной при повторной попытке BullMQ job?