Как разделять очереди задач в Node.js приложении?

«Как разделять очереди задач в Node.js приложении?» — вопрос из категории Брокеры сообщений, который задают на 26% собеседований Node.js Разработчик. Ниже — развёрнутый ответ с разбором ключевых моментов.

Ответ

В Node.js-приложениях я разделяю очереди задач по нескольким критериям, используя библиотеки вроде bull или bee-queue.

1. Разделение по типу задач

const Queue = require('bull');

// Создаем отдельные очереди для разных типов задач
const emailQueue = new Queue('email', 'redis://127.0.0.1:6379');
const imageProcessingQueue = new Queue('image-processing', 'redis://127.0.0.1:6379');
const reportGenerationQueue = new Queue('reports', 'redis://127.0.0.1:6379');

// Добавление задач
emailQueue.add('welcome-email', { userId: 123, email: 'user@example.com' });
imageProcessingQueue.add('resize', { imageId: 456, sizes: ['thumb', 'medium'] });

2. Разделение по приоритету

const priorityQueue = new Queue('tasks', {
  redis: { port: 6379, host: '127.0.0.1' },
  defaultJobOptions: {
    attempts: 3,
    backoff: { type: 'exponential', delay: 1000 }
  }
});

// Задачи с разным приоритетом
priorityQueue.add('high-priority', { task: 'urgent' }, { priority: 1 }); // Высокий
priorityQueue.add('low-priority', { task: 'background' }, { priority: 100 }); // Низкий

3. Разделение по воркерам/процессам

// worker-процесс для CPU-intensive задач
const { Worker } = require('worker_threads');
const cpuIntensiveQueue = new Queue('cpu-tasks');

cpuIntensiveQueue.process(2, (job) => { // 2 параллельных воркера
  return new Promise((resolve, reject) => {
    const worker = new Worker('./image-processor.js', {
      workerData: job.data
    });

    worker.on('message', resolve);
    worker.on('error', reject);
    worker.on('exit', (code) => {
      if (code !== 0) reject(new Error(`Worker stopped with exit code ${code}`));
    });
  });
});

// I/O задачи можно обрабатывать больше воркеров
const ioQueue = new Queue('io-tasks');
ioQueue.process(10); // 10 параллельных воркеров для I/O

4. Разделение по задержке выполнения

const scheduledQueue = new Queue('scheduled');

// Задача выполнится через 1 час
scheduledQueue.add('cleanup', { type: 'temp-files' }, {
  delay: 60 * 60 * 1000, // 1 час в миллисекундах
  removeOnComplete: true
});

// Периодические задачи (cron)
scheduledQueue.add('daily-report', {}, {
  repeat: { cron: '0 9 * * *' } // Каждый день в 9:00
});

5. Разделение по обработчикам ошибок

const mainQueue = new Queue('main');
const deadLetterQueue = new Queue('dead-letters');

mainQueue.process(async (job) => {
  try {
    await processJob(job);
  } catch (error) {
    // После 3 неудачных попыток перемещаем в dead letter queue
    if (job.attemptsMade >= 3) {
      await deadLetterQueue.add('failed-job', {
        originalJob: job.data,
        error: error.message,
        failedAt: new Date()
      });
    }
    throw error;
  }
});

В production-окружении я также разделяю очереди по окружениям (dev/staging/prod) через префиксы в Redis и настраиваю мониторинг через Bull Board или аналоги.