Пошук уроків, статей та іншого контенту
Налаштуєте producer і worker, пріоритети, паралельність та контроль станів завдань.
BullMQ — бібліотека для роботи з чергами завдань поверх Redis. Вона дає змогу розділити:
producer — додає завдання до черги;
worker — отримує завдання та виконує їх;
queue — зберігає завдання й керує їхнім станом.
Наприклад, HTTP-запит може лише поставити надсилання листа в чергу, а фактичне надсилання виконає worker. Користувач не чекатиме завершення повільної операції.
У NestJS інтеграція з BullMQ доступна через пакет @nestjs/bullmq.
BullMQ потребує запущеного Redis.
npm install @nestjs/bullmq bullmqДля локального запуску Redis можна використати Docker:
docker run --name nest-redis -p 6379:6379 -d redis:7Якщо контейнер із такою назвою вже існує, його можна запустити повторно:
docker start nest-redisСтворимо модуль для черги листів.
// src/jobs/jobs.module.ts
import { Module } from '@nestjs/common';
import { BullModule } from '@nestjs/bullmq';
import { EmailProcessor } from './email.processor';
import { EmailsService } from './emails.service';
export const EMAIL_QUEUE = 'emails';
@Module({
imports: [
BullModule.forRoot({
connection: {
host: 'localhost',
port: 6379,
},
}),
BullModule.registerQueue({
name: EMAIL_QUEUE,
}),
],
providers: [EmailsService, EmailProcessor],
exports: [EmailsService],
})
export class JobsModule {}BullModule.forRoot() налаштовує підключення до Redis, а registerQueue() реєструє конкретну чергу.
Модуль потрібно імпортувати в кореневий модуль застосунку:
// src/app.module.ts
import { Module } from '@nestjs/common';
import { JobsModule } from './jobs/jobs.module';
@Module({
imports: [JobsModule],
})
export class AppModule {}Для додавання завдань у чергу використовується Queue. У NestJS черга впроваджується через @InjectQueue().
// src/jobs/emails.service.ts
import { InjectQueue } from '@nestjs/bullmq';
import { Injectable } from '@nestjs/common';
import { Queue } from 'bullmq';
import { EMAIL_QUEUE } from './jobs.module';
export interface EmailJob {
to: string;
subject: string;
body: string;
}
@Injectable()
export class EmailsService {
constructor(
@InjectQueue(EMAIL_QUEUE)
private readonly emailQueue: Queue<EmailJob>,
) {}
async enqueueEmail(data: EmailJob, urgent = false) {
const job = await this.emailQueue.add('send-email', data, {
priority: urgent ? 1 : 10,
attempts: 3,
backoff: {
type: 'exponential',
delay: 1000,
},
removeOnComplete: {
age: 3600,
count: 1000,
},
removeOnFail: {
age: 24 * 3600,
},
});
return {
id: job.id,
name: job.name,
};
}
async getEmailState(id: string) {
const job = await this.emailQueue.getJob(id);
if (!job) {
return null;
}
return {
id: job.id,
name: job.name,
state: await job.getState(),
progress: job.progress,
attemptsMade: job.attemptsMade,
failedReason: job.failedReason,
returnValue: job.returnvalue,
};
}
}Метод queue.add() приймає:
ім’я завдання;
дані завдання;
параметри виконання.
У прикладі завдання має ім’я send-email, а його дані містять адресу отримувача, тему та текст листа.
Параметр attempts визначає максимальну кількість спроб виконання завдання.
{
attempts: 3,
backoff: {
type: 'exponential',
delay: 1000,
},
}Для експоненційної затримки спроби відбуватимуться приблизно через:
1 секунду;
2 секунди;
4 секунди.
Якщо worker викине помилку, BullMQ виконає наступну спробу відповідно до цієї конфігурації.
Redis зберігає інформацію про завдання, тому завершені завдання потрібно очищати:
{
removeOnComplete: {
age: 3600,
count: 1000,
},
}У цьому прикладі BullMQ зберігатиме не більше 1000 останніх завершених завдань і видалятиме завдання, старші за одну годину.
Producer можна викликати з контролера.
// src/jobs/emails.controller.ts
import {
Body,
Controller,
Get,
NotFoundException,
Param,
Post,
Query,
} from '@nestjs/common';
import { EmailsService } from './emails.service';
interface CreateEmailRequest {
to: string;
subject: string;
body: string;
urgent?: boolean;
}
@Controller('emails')
export class EmailsController {
constructor(private readonly emailsService: EmailsService) {}
@Post()
async createEmail(@Body() input: CreateEmailRequest) {
return this.emailsService.enqueueEmail(
{
to: input.to,
subject: input.subject,
body: input.body,
},
input.urgent ?? false,
);
}
@Get(':id')
async getEmail(@Param('id') id: string) {
const result = await this.emailsService.getEmailState(id);
if (!result) {
throw new NotFoundException('Завдання не знайдено');
}
return result;
}
}Тепер запит:
POST /emails
Content-Type: application/json
{
"to": "user@example.com",
"subject": "Підтвердження реєстрації",
"body": "Ваш обліковий запис створено",
"urgent": true
}не надсилає лист безпосередньо. Він лише створює завдання в Redis і повертає його ідентифікатор.
У NestJS worker створюється за допомогою WorkerHost і декоратора @Processor().
// src/jobs/email.processor.ts
import { Processor, WorkerHost } from '@nestjs/bullmq';
import { Job } from 'bullmq';
import { EMAIL_QUEUE } from './jobs.module';
import { EmailJob } from './emails.service';
@Processor({
name: EMAIL_QUEUE,
concurrency: 5,
})
export class EmailProcessor extends WorkerHost {
async process(job: Job<EmailJob>) {
if (job.name !== 'send-email') {
throw new Error(`Невідомий тип завдання: ${job.name}`);
}
await job.updateProgress(25);
// Тут має бути виклик реального поштового сервісу.
await this.sendEmail(job.data);
await job.updateProgress(100);
return {
sentTo: job.data.to,
sentAt: new Date().toISOString(),
};
}
private async sendEmail(data: EmailJob): Promise<void> {
// Імітуємо повільну зовнішню операцію.
await new Promise((resolve) => setTimeout(resolve, 500));
console.log(`Лист надіслано: ${data.to} — ${data.subject}`);
}
}Метод process() отримує об’єкт Job. Через нього можна прочитати:
job.id — ідентифікатор;
job.name — ім’я;
job.data — дані;
job.attemptsMade — кількість уже виконаних спроб;
job.progress — поточний прогрес.
Якщо process() успішно завершиться, завдання перейде до стану completed. Якщо метод викине помилку, завдання стане failed або буде повторно заплановане, якщо залишилися спроби.
Пріоритет задається під час додавання завдання:
await queue.add('send-email', data, {
priority: 1,
});У BullMQ менше числове значення означає вищий пріоритет:
await queue.add('send-email', urgentEmail, {
priority: 1,
});
await queue.add('send-email', regularEmail, {
priority: 10,
});Завдання з пріоритетом 1 має оброблятися раніше за завдання з пріоритетом 10, якщо вони очікують у черзі.
Зазвичай зручно зафіксувати рівні пріоритетів у константах:
const PRIORITY = {
urgent: 1,
normal: 10,
low: 20,
} as const;
await this.emailQueue.add('send-email', data, {
priority: urgent ? PRIORITY.urgent : PRIORITY.normal,
});Пріоритет впливає на порядок отримання завдань, але не перериває завдання, яке worker уже виконує.
Параметр concurrency визначає, скільки завдань один worker може обробляти одночасно:
@Processor({
name: EMAIL_QUEUE,
concurrency: 5,
})
export class EmailProcessor extends WorkerHost {
// ...
}За concurrency: 5 один екземпляр worker може мати до п’яти активних завдань одночасно.
Паралельність корисна для завдань, які переважно очікують на:
HTTP-відповідь;
відповідь бази даних;
завершення файлової операції;
зовнішній поштовий сервіс.
Значення потрібно підбирати з урахуванням обмежень зовнішніх сервісів. Надто велика кількість паралельних операцій може призвести до перевищення rate limit або навантаження на Redis.
Кілька екземплярів NestJS із тим самим worker також збільшують загальну пропускну здатність:
Екземпляр 1: concurrency = 5
Екземпляр 2: concurrency = 5
Загалом: до 10 активних завданьПорядок завершення завдань за паралельної обробки не гарантується. Завдання, яке було поставлене першим, може завершитися пізніше за наступне.
Протягом життєвого циклу завдання може мати такі стани:
waiting — очікує в черзі;
active — зараз виконується worker;
completed — успішно завершене;
failed — завершилося помилкою і більше не має доступних спроб;
delayed — відкладене до певного моменту;
paused — черга призупинена.
Поточний стан можна отримати через job.getState():
const job = await queue.getJob(jobId);
if (!job) {
throw new Error('Завдання не знайдено');
}
const state = await job.getState();
console.log({
id: job.id,
state,
progress: job.progress,
});Результат виконання повертається з process() і доступний у job.returnvalue після переходу завдання до completed.
Причина останньої помилки доступна в job.failedReason.
Worker може повідомляти прогрес:
await job.updateProgress(25);
// Виконується перший етап
await job.updateProgress(75);
// Виконується другий етап
await job.updateProgress(100);Прогрес може бути числом або довільним JSON-значенням. Наприклад:
await job.updateProgress({
current: 25,
total: 100,
});Значення прогресу можна прочитати з job.progress під час запиту статусу.
Оновлення прогресу не змінює основний стан завдання. Завдання залишається active, доки process() не завершиться або не викине помилку.
Нижче наведено основні файли застосунку, у якому producer додає листи, а worker обробляє їх паралельно.
// src/jobs/jobs.module.ts
import { Module } from '@nestjs/common';
import { BullModule } from '@nestjs/bullmq';
import { EmailProcessor } from './email.processor';
import { EmailsController } from './emails.controller';
import { EmailsService } from './emails.service';
export const EMAIL_QUEUE = 'emails';
@Module({
imports: [
BullModule.forRoot({
connection: {
host: 'localhost',
port: 6379,
},
}),
BullModule.registerQueue({
name: EMAIL_QUEUE,
}),
],
controllers: [EmailsController],
providers: [EmailsService, EmailProcessor],
})
export class JobsModule {}// src/jobs/emails.service.ts
import { InjectQueue } from '@nestjs/bullmq';
import { Injectable } from '@nestjs/common';
import { Queue } from 'bullmq';
import { EMAIL_QUEUE } from './jobs.module';
export interface EmailJob {
to: string;
subject: string;
body: string;
}
@Injectable()
export class EmailsService {
constructor(
@InjectQueue(EMAIL_QUEUE)
private readonly queue: Queue<EmailJob>,
) {}
async add(data: EmailJob, urgent = false) {
const job = await this.queue.add('send-email', data, {
priority: urgent ? 1 : 10,
attempts: 3,
backoff: {
type: 'exponential',
delay: 1000,
},
removeOnComplete: 100,
removeOnFail: 1000,
});
return { id: job.id };
}
async state(id: string) {
const job = await this.queue.getJob(id);
if (!job) {
return null;
}
return {
id: job.id,
state: await job.getState(),
progress: job.progress,
attemptsMade: job.attemptsMade,
failedReason: job.failedReason,
returnValue: job.returnvalue,
};
}
}// src/jobs/email.processor.ts
import { Processor, WorkerHost } from '@nestjs/bullmq';
import { Job } from 'bullmq';
import { EMAIL_QUEUE } from './jobs.module';
import { EmailJob } from './emails.service';
@Processor({
name: EMAIL_QUEUE,
concurrency: 5,
})
export class EmailProcessor extends WorkerHost {
async process(job: Job<EmailJob>) {
await job.updateProgress(50);
// Тут має бути реальна інтеграція з поштовим сервісом.
await new Promise((resolve) => setTimeout(resolve, 500));
console.log(`Оброблено лист для ${job.data.to}`);
await job.updateProgress(100);
return {
ok: true,
recipient: job.data.to,
};
}
}// src/jobs/emails.controller.ts
import { Body, Controller, Get, Param, Post } from '@nestjs/common';
import { EmailsService } from './emails.service';
@Controller('emails')
export class EmailsController {
constructor(private readonly emailsService: EmailsService) {}
@Post()
add(
@Body()
body: {
to: string;
subject: string;
body: string;
urgent?: boolean;
},
) {
return this.emailsService.add(
{
to: body.to,
subject: body.subject,
body: body.body,
},
body.urgent ?? false,
);
}
@Get(':id')
async state(@Param('id') id: string) {
return this.emailsService.state(id);
}
}Після запуску застосунку можна:
виконати POST /emails;
отримати ідентифікатор завдання;
виконати GET /emails/:id;
спостерігати за зміною стану від waiting до active, а потім до completed.
Producer і worker не обов’язково мають працювати в одному процесі.
У невеликому застосунку їх можна зареєструвати в одному NestJS-модулі. Для масштабування worker часто запускають окремим процесом або окремим сервісом:
API NestJS → Redis → Worker NestJSОбидва процеси повинні:
використовувати однакове підключення до Redis;
використовувати однакову назву черги;
мати сумісну структуру даних завдань.
API додає завдання через Queue, а worker читає їх через @Processor().
Клас processor має бути вказаний у providers модуля:
@Module({
providers: [EmailProcessor],
})
export class JobsModule {}Якщо processor не зареєстрований, завдання залишатимуться у стані waiting.
Producer і worker повинні використовувати точно однакову назву:
const EMAIL_QUEUE = 'emails';Черги emails і email є різними чергами.
Якщо Redis не запущений або вказано неправильний порт, producer і worker не зможуть підключитися до черги.
Високе значення concurrency не завжди пришвидшує систему. Воно може перевантажити зовнішній сервіс або базу даних.
За паралельної обробки завдання завершуються в різний час. Не слід покладатися на те, що порядок завершення збігатиметься з порядком додавання.
Тимчасові помилки мережі або зовнішнього сервісу можуть зробити завдання failed після першої невдалої спроби. Для таких операцій варто налаштовувати attempts і backoff.
Якщо не налаштувати removeOnComplete і removeOnFail, Redis поступово зберігатиме дедалі більше старих завдань.
BullMQ зберігає завдання в Redis і розділяє додавання та виконання роботи.
Producer додає завдання через Queue.
Worker у NestJS створюється через @Processor() і WorkerHost.
Пріоритет задається параметром priority; менше число означає вищий пріоритет.
concurrency визначає кількість завдань, які один worker обробляє одночасно.
attempts і backoff налаштовують повторні спроби.
Стан завдання можна отримати через job.getState().
Прогрес оновлюється методом job.updateProgress().
Завершені та невдалі завдання варто очищати за допомогою removeOnComplete і removeOnFail.