Bull и очереди

Для начала необходимо установить библиотеку Bull, которая представляет собой мощный инструмент для организации очередей задач в Node.js. Установка производится через npm:

npm install bull

Bull работает поверх Redis, поэтому необходимо, чтобы Redis-сервер был запущен и доступен. Подключение к очереди выполняется следующим образом:

const Queue = require(&

const myQueue = new Queue('my-queue', {
  redis: {
    host: '127.0.0.1',
    port: 6379
  }
});

В данном примере создаётся очередь с именем my-queue. Redis-конфигурация задаётся объектом с указанием хоста и порта.


Добавление задач в очередь

Добавление задач в очередь происходит методом add. Задачи могут содержать данные, необходимые для выполнения:

myQueue.add({ userId: 123, action: 'sendEmail' });

Метод add возвращает промис, который разрешается объектом задачи (Job). Каждая задача имеет уникальный идентификатор и состояние, которое можно отслеживать.

Можно также использовать дополнительные опции, например, задержку выполнения:

myQueue.add(
  { userId: 456, action: 'generateReport' },
  { delay: 5000, attempts: 3, backoff: 2000 }
);

Ключевые моменты:

  • delay — задержка перед выполнением задачи в миллисекундах.
  • attempts — количество попыток выполнения при неудаче.
  • backoff — стратегия повторной попытки, может быть числом (миллисекунды) или объектом для экспоненциального увеличения.

Обработка задач

Очередь обрабатывается с помощью метода process. Обработчик получает объект задачи, содержащий данные и методы для управления выполнением:

myQueue.process(async (job) => {
  const { userId, action } = job.data;
  
  if (action === 'sendEmail') {
    await sendEmailToUser(userId);
  } else if (action === 'generateReport') {
    await generateReportForUser(userId);
  }
  
  return { status: 'completed' };
});

Особенности обработки:

  • Обработчик может быть асинхронным (async) и возвращать промис.
  • Если задача выбрасывает ошибку, Bull автоматически помечает её как failed и применяет стратегию повторной попытки.
  • Можно использовать несколько параллельных обработчиков для повышения производительности:
myQueue.process(5, async (job) => {
  // параллельная обработка до 5 задач одновременно
});

Мониторинг и события очереди

Bull предоставляет богатую систему событий, которые позволяют отслеживать состояние задач:

myQueue.on('completed', (job, result) => {
  console.log(`Задача ${job.id} выполнена с результатом`, result);
});

myQueue.on('failed', (job, err) => {
  console.error(`Задача ${job.id} завершилась с ошибкой:`, err);
});

myQueue.on('stalled', (job) => {
  console.warn(`Задача ${job.id} была «зависшей»`);
});

События:

  • completed — задача успешно выполнена.
  • failed — задача завершилась ошибкой.
  • stalled — задача была заблокирована обработчиком и повторно ставится в очередь.
  • active — задача начала выполняться.
  • waiting — задача ожидает обработки.

Очереди с приоритетом и повторением

Bull позволяет задавать приоритет выполнения задач:

myQueue.add(
  { userId: 789, action: 'urgentTask' },
  { priority: 1 } // чем меньше число, тем выше приоритет
);

Для повторяющихся задач используется параметр repeat:

myQueue.add(
  { action: 'cleanupTempFiles' },
  { repeat: { cron: '0 0 * * *' } } // каждый день в полночь
);

Можно использовать как cron-выражения, так и интервалы в миллисекундах (every: 60000).


Управление задачами и очередями

Bull предоставляет методы для управления существующими задачами и очередями:

// Получить все ожидающие задачи
const waitingJobs = await myQueue.getWaiting();

// Получить все завершенные задачи
const completedJobs = await myQueue.getCompleted();

// Очистить очередь от старых задач
await myQueue.clean(1000, 'completed'); // удаляет задачи старше 1 секунды

Также можно удалять конкретные задачи:

const job = await myQueue.getJob(jobId);
if (job) await job.remove();

Масштабирование и производительность

Bull поддерживает горизонтальное масштабирование:

  • Несколько воркеров могут подключаться к одной очереди.
  • Задачи автоматически распределяются между воркерами.
  • Использование Redis обеспечивает устойчивость к сбоям и отказоустойчивость.

Для повышения производительности рекомендуется:

  • Настроить параллельные обработчики (process(concurrency, handler)).
  • Ограничивать количество повторных попыток задач с ошибками.
  • Регулярно чистить старые задачи для освобождения памяти Redis.

Интеграция с другими инструментами

Bull легко интегрируется с любыми Node.js приложениями:

  • Express: задачи создаются в маршрутах API.
  • NestJS: используется модуль @nestjs/bull для декларативного определения очередей и обработчиков.
  • Puppeteer и тестирование: можно ставить задачи на генерацию скриншотов или парсинг страниц в фоне, не блокируя основной поток.

Пример интеграции с Puppeteer:

myQueue.process(async (job) => {
  const browser = await puppeteer.launch();
  const page = await browser.newPage();
  await page.goto(job.data.url);
  await page.screenshot({ path: `screenshot-${job.data.id}.png` });
  await browser.close();
});

Это позволяет выполнять ресурсоёмкие задачи в фоне, обеспечивая стабильность основного приложения и точный контроль над очередью.


Эта структура демонстрирует все основные возможности Bull: от простого добавления задач до сложных сценариев с приоритетами, повторениями и масштабированием, обеспечивая гибкость и надежность фоновых процессов в Node.js приложениях.