Создание надежных систем фоновых задач с использованием BullMQ и Redis
Изучите, как создавать надежные фоновые задачи в Node.js с использованием BullMQ и Redis, рассмотрев повторные попытки, конкурентность, идемпотентность и мониторинг.
Отправка подтверждающего электронного письма, генерация отчета, обработка платежа — множество задач на стороне сервера не обязательно должны быть завершены до того, как вы ответите пользователю. В этом руководстве рассматривается создание надежных систем фоновых задач с использованием BullMQ в сочетании с Redis.
Представьте себе серверную часть, где практически каждая задача выполняется непосредственно в рамках цикла HTTP-запроса. Нужно отправить электронное письмо? Сделайте это прямо здесь. Нужно сгенерировать PDF? То же самое. Нужно обработать какие-то данные в фоновом режиме? Также делайте это прямо в процессе.
Сначала такой подход работает неплохо. Но затем он перестает справляться со своей задачей.
API начинает работать медленнее. Запросы начинают тайм-аутить. А если какой-то внешний сервис выходит из строя, весь запрос может потерпеть неудачу вместе с ним.
Именно в этот момент фоновые задачи становятся полезными.
Вместо того чтобы заставлять API выполнять каждый шаг перед ответом, вы можете поместить работу в очередь и позволить специальному процессу обработать её отдельно.
Хорошим вариантом в экосистеме Node.js является BullMQ, использующий Redis для хранения данных. Вот как все это связано между собой.
1. Что такое фоновая задача?
Фоновая задача — это любой элемент работы, который не обязан выполняться синхронно в рамках HTTP-запроса.
Рассмотрим типичный процесс регистрации. Когда кто-то создаёт учётную запись, API может потребоваться:
- Создать запись пользователя
- Отправить письмо с приветствием
- Сгенерировать PDF-документ с приветствием
- Отправить уведомление
- Обновить какую-либо другую вспомогательную систему
Вы могли бы попытаться выполнить всё это непосредственно перед ответом:
Client
↓
API
↓
Create User
↓
Send Email
↓
Generate PDF
↓
Send Notification
↓
Response
Но это заставляет пользователя ждать завершения каждого из этих шагов.
Лучший подход выглядит следующим образом:
Client
↓
API
↓
Create User
↓
Add Job to Queue
↓
Response
И отдельно:
Queue
↓
Worker
↓
Send Email
↓
Done
Поскольку API больше не обязано выполнять все задачи перед ответом, оно отвечает гораздо быстрее.
2. Зачем нам очередь?
Предположим, отправка электронного письма занимает 1 секунду, генерация PDF — 2 секунды, а вызов другого API — ещё 1 секунду. Ваш конечный пункт может застревать на несколько секунд перед тем, как отправить ответ — что создаёт плохой опыт для пользователя.
Что ещё хуже, что произойдёт, если поставщик электронной почты недоступен? Запрос может завершиться с ошибкой, даже если создание пользователя прошло успешно. Это ненужная зависимость между двумя не связанными между собой аспектами.
Очередь разрывает эту связь:
┌──────────────┐
│ Node API │
└──────┬───────┘
↓
Add Job
↓
┌──────────────┐
│ Redis │
│ Queue │
└──────┬───────┘
↓
┌──────────────┐
│ Worker │
└──────┬───────┘
↓
Email / PDF / API / etc.
Благодаря такому разделению API и фоновая задача имеют по отдельной ответственности.
3. Что такое BullMQ?
BullMQ — это библиотека очередей для Node.js, которая использует Redis для хранения и координации задач. Её архитектура в общих чертах выглядит следующим образом:
Producer
↓
Queue
↓
Worker
↓
Job Processing
Производитель — это компонент, создающий задачи. Очередь хранит их. Рабочий процесс — это тот, кто фактически их обрабатывает.
Например:
await emailQueue.add("welcome-email", {
userId: user.id,
email: user.email
});
По сути, API формулирует запрос так:
«Вот некоторая работа, которую необходимо выполнить».
Он сам не обязан выполнять эту работу.
4. Создание очереди
Вот как выглядит минимальная настройка очереди BullMQ:
import { Queue } from "bullmq";
const connection = {
host: "localhost",
port: 6379
};const emailQueue = new Queue("email", {
connection
});
После этого вы можете добавлять задачи в неё:
await emailQueue.add("welcome-email", {
userId: "123",
email: "user@example.com"
});
Redis скрытно обрабатывает хранение всего состояния, связанного с очередью. Концептуально это можно представить так:
email queue
Job 1
Job 2
Job 3
Job 4
Job 5
Затем рабочий процесс берет эти задачи и обрабатывает их.
5. Создание рабочего процесса
Рабочий процесс — это тот компонент, который фактически выполняет работу:
import { Worker } from "bullmq";
const worker = new Worker(
"email",
async (job) => {
console.log("Processing:", job.name); await sendWelcomeEmail(
job.data.email
);
},
{
connection
}
);
Если объединить все элементы, поток действий будет выглядеть следующим образом:
API
↓
emailQueue.add()
↓
Redis
↓
Worker
↓
sendWelcomeEmail()
API не должен ждать завершения отправки электронного письма — и в этом заключается основное преимущество фоновых задач.
6. Что происходит, если задача проваливается?
Именно здесь очередь начинает демонстрировать свое реальное преимущество по сравнению с обычным вызовом сервиса.
Представьте такую схему:
API
↓
Email Service
↓
ERROR
При прямом вызове API вам приходится немедленно принимать решение относительно возникшей ошибки.
Очередь предоставляет еще один вариант: задачу можно просто попробовать выполнить снова.
Вот пример:
await emailQueue.add(
"welcome-email",
{
email: "user@example.com"
},
{
attempts: 3
}
);
При такой настройке заданию разрешается несколько попыток перед тем, как сдаться.
Визуально процесс выглядит следующим образом:
Attempt 1
↓
Failed
↓
Attempt 2
↓
Failed
↓
Attempt 3
↓
Success
Этот подход крайне полезен при работе с ненадежными сторонними сервисами.
Тем не менее, повторные попытки не должны быть бесконечными или осуществляться бездумно.
Избегайте ситуаций, когда задание продолжает пытаться выполниться бесконечно без каких-либо ограничений.
7. Повторные попытки с задержкой
Предположим, что внешний сервис временно выходит из строя.
Необходимо избегать ситуаций, подобных этой:
FAIL
RETRY IMMEDIATELY
FAIL
RETRY IMMEDIATELY
FAIL
RETRY IMMEDIATELY
Мгновенные повторные попытки могут усугубить ситуацию с нестабильным сервисом.
Решение заключается в введении задержки между попытками.
Например:
await emailQueue.add(
"welcome-email",
{
email: "user@example.com"
},
{
attempts: 5,
backoff: {
type: "exponential",
delay: 5000
}
}
);
Вот как это выглядит в концептуальном плане:
Attempt 1 → Fail
↓
5 sec
↓
Attempt 2 → Fail
↓
10 sec
↓
Attempt 3 → Fail
↓
20 sec
↓
Attempt 4 → Success
Точное время выполнения зависит от настроек стратегии повторных попыток и задержек.
Но основная идея остаётся прежней:
Дайте временным сбоям возможность восстановиться перед повторной попыткой.
8. Задания с отложенной выполнением
Не каждое задание нужно запускать сразу после создания.
Например:
Отправить напоминание через 24 часа после регистрации.
BullMQ позволяет запланировать выполнение задания на более позднее время.
await emailQueue.add(
"reminder",
{
userId: "123"
},
{
delay: 24 * 60 * 60 * 1000
}
);
Концептуально:
Create Job
↓
Wait 24 hours
↓
Worker processes job
Этот паттерн встречается в таких ситуациях, как:
- Напоминательные электронные письма
- Запланированные уведомления
- Истечение пробного периода
- Напоминания о оплате
- Последующие сообщения
9. Несколько рабочих процессов
Теперь представьте систему, принимающую тысячи заданий в минуту.
Один рабочий процесс может не справиться.
Вы можете расширить масштабы, запуская несколько рабочих процессов одновременно:
Redis Queue
↓
┌──────────┼──────────┐
↓ ↓ ↓
Worker 1 Worker 2 Worker 3
↓ ↓ ↓
Jobs Jobs Jobs
Каждый из них самостоятельно берет задачи из очереди.
Например, если у нас есть:
1000 email jobs
Вы можете увидеть что-то вроде:
Worker 1 → Job 1, 4, 7...
Worker 2 → Job 2, 5, 8...
Worker 3 → Job 3, 6, 9...
Добавление большего количества рабочих процессов — один из способов повышения пропускной способности.
Но будьте осторожны:
Простое добавление большего количества рабочих процессов не гарантирует улучшения.
У вашей базы данных, поставщика электронной почты, процессора, памяти и любых последующих сервисов есть свои собственные пределы пропускной способности.
10. Конкурентность
Помимо запуска нескольких рабочих процессов, BullMQ также позволяет настроить количество задач, которые один рабочий процесс может обрабатывать одновременно.
Например:
const worker = new Worker(
"email",
async (job) => {
await sendEmail(job.data.email);
},
{
connection,
concurrency: 5
}
);
Это позволяет одному рабочему процессу обрабатывать несколько задач параллельно.
Концептуально:
Worker
├── Job 1
├── Job 2
├── Job 3
├── Job 4
└── Job 5
Более высокая конкурентность может повысить пропускную способность.
Но не стоит просто увеличивать уровень конкурентности до 100 без тщательного обдумывания.
Если каждая задача обращается к вашей базе данных, высокая конкурентность может легко привести к её перегрузке.
Параметры конкурентности следует настраивать так, чтобы они соответствовали тому, что действительно может обработать ваша нагрузка.
11. Ограничение скорости
Иногда узким местом вовсе не является ваш собственный система — это сторонняя служба, от которой вы зависите.
Допустим, ваш поставщик электронной почты ограничивает количество запросов в секунду фиксированным числом.
Если у вас вдруг появится:
10,000 jobs
вы не захотите отправлять их все сразу.
Очередь может снизить скорость обработки задач.
Результативная архитектура выглядит так:
10,000 Jobs
↓
Queue
↓
Rate Limit
↓
Worker
↓
External API
Это гораздо безопаснее, чем отправка тысяч одновременных запросов поставщику.
12. Идемпотентность задач имеет значение
Эта следующая концепция является одной из наиболее важных идей в обработке задач в фоновом режиме.
Рассмотрим задачу обработки платежей:
Process Payment
Рабочий процесс её запускает.
Платеж пройдёт успешно.
Но непосредственно перед тем, как рабочий процесс отметит её как завершённую, процесс падает.
Очередь, выполняя именно то, для чего она предназначена, пытается выполнить задачу снова.
Без защитных мер можно оказаться в ситуации, когда клиенту будет взиматься плата второй раз.
Это реальная и дорогостоящая проблема.
Чтобы этого избежать, задачи следует делать идемпотентными, где это возможно.
На практике это означает, что двукратное выполнение одной и той же задачи не должно приводить к нежелательным дублирующим эффектам.
Один из распространённых подходов — использование уникального идентификатора платежа в качестве ключа:
payment:order_123
Затем, прежде чем начинать какую-либо работу, необходимо проверить:
Has this payment already been completed?
↓
Yes → Don't charge again
↓
No → Process payment
У самого BullMQ нет встроенного механизма для этого.
Обеспечение идемпотентности лежит на коде вашего приложения.
13. Для неудачных заданий нужна стратегия
Не все сбои одинаковы, и не все из них стоят повторной попытки.
Рассмотрим несколько примеров:
Invalid email
Invalid user ID
Missing database record
Invalid payment information
Повторная выполнение этих заданий пять раз ничего не решит.
Полезно разделить сбои на две категории:
Временные сбои
К ним относятся такие ситуации:
- Время ожидания ответа сети
- Зависимость, которая временно недоступна
- Прерванное соединение с базой данных
Именно в таких случаях имеет смысл попробовать снова позже.
Постоянные сбои
К ним относятся такие ситуации:
- Некорректные входные данные
- Ресурс, на который делается ссылка, больше не существует
В таких случаях повторная попытка бесполезна — задачу необходимо сразу направить в механизм обработки ошибок.
Хорошо спроектированная система очередей не следует простому универсальному правилу:
Retry everything
Вместо этого она использует более продуманный алгоритм работы:
Understand why it failed
↓
Temporary?
/ \
YES NO
↓ ↓
Retry Handle failure
14. Обработка неработающих/неудачных задач
Как бы вы ни старались, некоторые задачи могут сорваться таким образом, что их невозможно исправить путем повторной попытки. Вам нужен доступ к этим задачам, чтобы они не просто исчезали.
Например, может возникнуть ситуация вроде:
Failed Jobs
──────────────
Job 101 → Email invalid
Job 102 → Payment failed
Job 103 → API timeout
Как только вы сможете увидеть эти ошибки, у вас появятся варианты действий:
- Записать информацию об ошибке для последующего анализа
- Уведомить вашу команду
- Дать кому-то возможность вручную попробовать снова
- Исправить проблемные данные, вызвавшие ошибку
- Перенаправьте задачу в специализированный рабочий процесс обработки сбоев
Конкретный способ реализации зависит от потребностей вашей системы. Самое важное — это один принцип:
Неудавшиеся операции никогда не должны исчезать бесследно.
15. Очередь против задачи Cron
Легко перепутать эти два инструмента, но они решают разные задачи.
Задача cron-задачи — указать:
"Выполнить эту задачу в определённое время."
Задача очереди — указать:
"Обработать этот элемент работы."
На практике эти два инструмента часто хорошо работают вместе. Например:
Cron
↓
Find users whose trial expires today
↓
Create jobs
↓
Queue
↓
Workers
↓
Send emails
Это позволяет разделить логику планирования и логику обработки. Обычно это более чистый подход, чем когда один процесс Cron пытается выполнить всю работу самостоятельно.
16. События в очереди и мониторинг
Как только вы начнете использовать это в производственных условиях, вам потребуется возможность отслеживать то, что на самом деле происходит внутри очереди.
Показатели, которые стоит отслеживать, включают:
- Задания, ожидающие обработки
- Задания, которые в настоящее время обрабатываются
- Задания, которые были успешно завершены
- Задания, которые потерпели неудачу
- Время, затрачиваемое на обработку
- Количество попыток повторной обработки
- Общий размер очереди
Представьте панель управления, которая вдруг покажет что-то вроде этого:
Waiting Jobs
Normal: 50
Current: 25,000
Такой резкий скачок — это тревожный сигнал. Это может означать:
- Ваши процессы прекратили работу
- Внешний API замедлил свою работу
- Ваша база данных находится под сильной нагрузкой
- Нагрузка на систему резко возросла
- В результате недавней развертки появилась ошибка
Если вы не отслеживаете свою очередь, эти проблемы могут накапливаться незаметно, пока пользователи не начнут замечать, что что-то не так.
17. Не ставьте всё в очередь
Наличие BullMQ не означает, что каждая операция должна выполняться в фоновом режиме.
Возьмём, к примеру:
GET /profile
Здесь пользователь ожидает данные своего профиля немедленно. Откладывание этой задачи в фоновую очередь не имеет смысла — это лишь приведёт к ненужной задержке.
Очередь имеет смысл тогда, когда:
- Выполнение задачи занимает время
- Задачу можно выполнить асинхронно
- Возможна необходимость повторной попытки выполнения
- Задача требует больших ресурсов
- Задача зависит от внешних сервисов, которые не являются полностью надёжными
- Результат не обязателен для немедленного ответа
Полезный вопрос, который стоит задать:
Действительно ли пользователю нужен этот результат до того, как вы отправите HTTP-ответ?
Если нет, стоит рассмотреть возможность выполнения этой работы в фоновом режиме.
18. Архитектура в стиле производства
Собирая всё это вместе, типичная конфигурация выглядит следующим образом:
Client
↓
Node.js API
↓
┌──────┴──────┐
↓ ↓
PostgreSQL Redis
↓
Queue
↓
┌──────────┼──────────┐
↓ ↓ ↓
Worker 1 Worker 2 Worker 3
↓ ↓ ↓
Email PDF Notifications
Слой API обрабатывает всё, что должно произойти немедленно. PostgreSQL (или любая другая выбранная вами база данных) хранит стабильные бизнес-данные. Redis обеспечивает инфраструктуру очередей и другие краткосрочные задачи, для которых он подходит. Рабочие процессы занимаются всем, что может выполняться асинхронно.
Такое разделение обязанностей значительно упрощает масштабирование всей системы.
19. Ошибки, которых стоит избегать
Ошибка 1: Выполнение всего внутри HTTP-запроса
Это приводит к API, которые одновременно медленные и нестабильные.
Ошибка 2: Повторные попытки без ограничений
Некоторые сбои просто не исчезнут самостоятельно, независимо от количества попыток.
Ошибка 3: Игнорирование идемпотентности
Если задача выполняется дважды, это может привести к дублированию побочных эффектов, которых вы не ожидали.
Ошибка 4: Разрешение неограниченной конкурентности
Без ограничений существует риск перегрузки систем, от которых зависят ваши задачи.
Ошибка 5: Игнорирование мониторинга
Очередь, которая продолжает расти без контроля, представляет собой операционную проблему, готовую проявиться.
Ошибка 6: Использование Redis в качестве системы хранения данных
Состояние очереди и основные бизнес-данные служат разным целям и не должны смешиваться.
Ошибка 7: Применение асинхронности ко всему
Некоторые операции действительно требуют завершения до того, как можно будет отправить ответ.
20. Лучшая когнитивная модель
Прежде чем понять концепцию очередей, естественным импульсом является мышление о обработке запросов следующим образом:
Request
↓
Do everything
↓
Response
Более полезная модель выглядит так:
Request
↓
Do what must happen immediately
↓
Queue what can happen later
↓
Response
Затем следует:
Queue
↓
Worker
↓
Process
↓
Retry if appropriate
↓
Complete / Fail
Именно это изменение — разделение того, что должно произойти сейчас, и того, что может произойти позже, — является основной идеей всего этого.
Итог
BullMQ полезен не просто потому, что это широко используемая библиотека для Node.js. Он полезен потому, что обработка фоновых задач решает реальную архитектурную проблему.
Если какая-либо задача:
- медленная
- может быть выполнена заново
- не должна блокировать ответ
- зависит от внешнего сервиса
- требует больших ресурсов
то её, скорее всего, не следует включать в цикл обработки HTTP-запросов.
Очередь обеспечивает место для выполнения задач. Redis предоставляет базовую инфраструктуру. BullMQ отвечает за управление заданиями. Рабочие процессы выполняют саму обработку данных. Механизмы повторных попыток устраняют временные сбои. Настройки конкурентности позволяют контролировать пропускную способность. Мониторинг сообщает о возникновении проблем. А тщательно продуманная архитектура на уровне приложения гарантирует, что задания могут безопасно выполняться несколько раз при необходимости.
Основной урок заключается в следующем:
Не всё должно решаться в рамках цикла запрос-ответ.
Иногда правильным ответом будет просто:
«Я принял задание. Мы займёмся остальным».
Связанная литература
- Проектирование бэкендов для реального времени: комнаты, сохранение данных и масштабирование — Узнайте, как спроектировать бэкенд для реального времени с использованием Socket.IO, PostgreSQL и Redis, с упором на комнаты, порядок сохранения сообщений, информацию о присутствии пользователей и масштабирование на несколько серверов.
- Основы кэширования в Redis: шаблоны, подводные камни и вопросы для собеседований — Ознакомьтесь с тем, как работает кэширование в Redis в приложениях на Node.js, начиная с методов кэширования и времени действия кэша, заканчивая защитой от перегрузки, политиками удаления данных и типичными вопросами на собеседованиях.