Пошук уроків, статей та іншого контенту
Підключите BullMQ і створите черги для надійного виконання фонових завдань через Redis.
BullMQ — це бібліотека для створення черг фонових завдань у Node.js. Вона використовує Redis для зберігання:
завдань, які очікують виконання;
завдань, що виконуються;
успішно завершених і невдалих завдань;
кількості спроб і часу повторного запуску.
У NestJS BullMQ зазвичай використовується у трьох частинах:
Producer — додає завдання в чергу.
Queue — логічна черга, зареєстрована в NestJS.
Worker або Processor — отримує та виконує завдання.
Наприклад, HTTP-запит може швидко додати завдання на надсилання листа в Redis, а окремий worker виконає це завдання у фоновому режимі.
Для інтеграції з NestJS встановіть @nestjs/bullmq і bullmq:
npm install @nestjs/bullmq bullmqBullMQ потребує запущеного Redis. Локальний Redis має бути доступний за адресою:
redis://localhost:6379Перевірити доступність Redis можна командою:
redis-cli pingОчікувана відповідь:
PONGПідключення BullMQ виконується через BullModule.forRoot().
Створимо модуль із чергою emails:
// app.module.ts
import { Module } from '@nestjs/common';
import { BullModule } from '@nestjs/bullmq';
import { EmailProducer } from './email.producer';
import { EmailProcessor } from './email.processor';
@Module({
imports: [
BullModule.forRoot({
connection: {
host: 'localhost',
port: 6379,
},
}),
BullModule.registerQueue({
name: 'emails',
}),
],
providers: [EmailProducer, EmailProcessor],
})
export class AppModule {}BullModule.forRoot() налаштовує спільне підключення до Redis.
BullModule.registerQueue() реєструє конкретну чергу. Назва черги має збігатися в усіх частинах застосунку, де ця черга використовується.
Щоб додати завдання в чергу, потрібно отримати її через @InjectQueue().
// email.producer.ts
import { Injectable } from '@nestjs/common';
import { InjectQueue } from '@nestjs/bullmq';
import { Queue } from 'bullmq';
export interface SendEmailJob {
recipient: string;
subject: string;
body: string;
}
@Injectable()
export class EmailProducer {
constructor(
@InjectQueue('emails')
private readonly emailQueue: Queue<SendEmailJob>,
) {}
async addEmailJob(data: SendEmailJob): Promise<string> {
const job = await this.emailQueue.add('send-email', data, {
attempts: 3,
backoff: {
type: 'exponential',
delay: 1000,
},
removeOnComplete: true,
removeOnFail: false,
});
return job.id as string;
}
}Метод queue.add() приймає:
ім’я завдання;
дані завдання;
параметри виконання.
У прикладі завдання називається send-email, а його дані містять отримувача, тему та текст листа.
Параметр attempts: 3 означає, що BullMQ спробує виконати завдання максимум три рази.
Параметр backoff визначає затримку між спробами:
backoff: {
type: 'exponential',
delay: 1000,
}Для експоненційної затримки інтервал збільшується після кожної невдалої спроби. Це корисно, якщо тимчасово недоступний зовнішній сервіс.
Повторна спроба відбувається лише тоді, коли processor завершує виконання помилкою. Якщо помилку перехопити й не викинути повторно, BullMQ вважатиме завдання успішним.
Обробник черги створюється за допомогою @Processor() і має наслідуватися від WorkerHost.
// email.processor.ts
import { Injectable, Logger } from '@nestjs/common';
import { Processor, WorkerHost } from '@nestjs/bullmq';
import { Job } from 'bullmq';
import { SendEmailJob } from './email.producer';
@Processor('emails')
@Injectable()
export class EmailProcessor extends WorkerHost {
private readonly logger = new Logger(EmailProcessor.name);
async process(job: Job<SendEmailJob>): Promise<void> {
if (job.name !== 'send-email') {
throw new Error(`Невідомий тип завдання: ${job.name}`);
}
this.logger.log(
`Підготовка листа для ${job.data.recipient}`,
);
// Тут може бути виклик сервісу надсилання електронних листів.
await this.sendEmail(job.data);
this.logger.log(`Лист для ${job.data.recipient} успішно надіслано`);
}
private async sendEmail(data: SendEmailJob): Promise<void> {
// Імітація асинхронної операції надсилання листа.
await new Promise((resolve) => setTimeout(resolve, 1000));
if (!data.recipient.includes('@')) {
throw new Error('Некоректна адреса отримувача');
}
}
}Метод process() викликається для кожного завдання з черги emails.
Якщо метод завершується успішно, завдання вважається виконаним. Якщо метод викидає помилку, BullMQ:
позначає поточну спробу як невдалу;
застосовує налаштовану затримку;
запускає завдання повторно, якщо залишилися спроби;
залишає завдання невдалим після вичерпання спроб.
Нижче наведено мінімальний приклад застосунку NestJS, який додає та обробляє завдання для надсилання листів.
// app.module.ts
import { Module } from '@nestjs/common';
import { BullModule } from '@nestjs/bullmq';
import { EmailProducer } from './email.producer';
import { EmailProcessor } from './email.processor';
@Module({
imports: [
BullModule.forRoot({
connection: {
host: 'localhost',
port: 6379,
},
}),
BullModule.registerQueue({
name: 'emails',
}),
],
providers: [EmailProducer, EmailProcessor],
exports: [EmailProducer],
})
export class AppModule {}// email.producer.ts
import { Injectable } from '@nestjs/common';
import { InjectQueue } from '@nestjs/bullmq';
import { Queue } from 'bullmq';
export interface SendEmailJob {
recipient: string;
subject: string;
body: string;
}
@Injectable()
export class EmailProducer {
constructor(
@InjectQueue('emails')
private readonly emailQueue: Queue<SendEmailJob>,
) {}
async addEmailJob(data: SendEmailJob): Promise<string> {
const job = await this.emailQueue.add('send-email', data, {
attempts: 3,
backoff: {
type: 'exponential',
delay: 1000,
},
removeOnComplete: true,
removeOnFail: false,
});
return job.id as string;
}
}// email.processor.ts
import { Injectable, Logger } from '@nestjs/common';
import { Processor, WorkerHost } from '@nestjs/bullmq';
import { Job } from 'bullmq';
import { SendEmailJob } from './email.producer';
@Processor('emails')
@Injectable()
export class EmailProcessor extends WorkerHost {
private readonly logger = new Logger(EmailProcessor.name);
async process(job: Job<SendEmailJob>): Promise<void> {
this.logger.log(
`Отримано завдання ${job.id}: ${job.name}`,
);
if (job.name !== 'send-email') {
throw new Error(`Невідомий тип завдання: ${job.name}`);
}
await new Promise((resolve) => setTimeout(resolve, 1000));
if (!job.data.recipient.includes('@')) {
throw new Error('Некоректна адреса отримувача');
}
this.logger.log(
`Лист для ${job.data.recipient} оброблено`,
);
}
}// email.controller.ts
import { Body, Controller, Post } from '@nestjs/common';
import { EmailProducer, SendEmailJob } from './email.producer';
@Controller('emails')
export class EmailController {
constructor(
private readonly emailProducer: EmailProducer,
) {}
@Post()
async createEmailJob(
@Body() data: SendEmailJob,
): Promise<{ jobId: string }> {
const jobId = await this.emailProducer.addEmailJob(data);
return { jobId };
}
}Не забудьте додати контролер у модуль:
// app.module.ts
import { Module } from '@nestjs/common';
import { BullModule } from '@nestjs/bullmq';
import { EmailProducer } from './email.producer';
import { EmailProcessor } from './email.processor';
import { EmailController } from './email.controller';
@Module({
imports: [
BullModule.forRoot({
connection: {
host: 'localhost',
port: 6379,
},
}),
BullModule.registerQueue({
name: 'emails',
}),
],
controllers: [EmailController],
providers: [EmailProducer, EmailProcessor],
})
export class AppModule {}Після запуску застосунку можна створити завдання HTTP-запитом:
curl -X POST http://localhost:3000/emails \
-H "Content-Type: application/json" \
-d '{"recipient":"user@example.com","subject":"Вітаємо","body":"Ваше завдання виконано"}'Контролер одразу поверне ідентифікатор завдання, а його обробка відбудеться окремо:
{
"jobId": "1"
}У консолі processor з’являться повідомлення про отримання та обробку завдання.
Черга корисна тоді, коли фонове завдання не повинно загубитися через помилку HTTP-запиту або тимчасову недоступність сервісу.
Під час додавання завдання BullMQ записує його в Redis. Worker може:
бути тимчасово вимкненим;
перезапуститися;
обробити завдання після відновлення роботи;
повторити завдання після помилки.
Це відокремлює приймання запитів від тривалих операцій. Наприклад, сервер може швидко відповісти клієнту, не очікуючи завершення надсилання листа.
Надійність залежить від правильної конфігурації Redis і самої операції. Якщо завдання виконує зовнішній запит, цей запит також має бути безпечним для повторного виконання. Наприклад, повторна спроба не повинна створювати дубльований платіж або дубльований запис без додаткової перевірки.
Producer та processor можуть бути частинами одного NestJS-застосунку, як у прикладі вище. У такому випадку один процес:
приймає HTTP-запити;
додає завдання;
обробляє завдання.
Але їх також можна запускати окремо:
API-застосунок лише додає завдання в Redis;
окремий worker-застосунок обробляє чергу.
Обидва застосунки мають підключатися до того самого Redis і використовувати однакову назву черги.
Якщо Redis недоступний, NestJS не зможе підключити BullMQ, а додавання завдань завершуватиметься помилкою.
Перевірте:
чи запущений Redis;
чи правильні host і port;
чи не використовується інший порт;
чи однакові налаштування підключення для producer і processor.
Ці назви мають збігатися:
BullModule.registerQueue({
name: 'emails',
});@Processor('emails')Назви emails і email — це різні черги.
providersКлас із @Processor() має бути доступним NestJS через providers модуля:
@Module({
providers: [EmailProcessor],
})
export class AppModule {}Якщо цього не зробити, NestJS не створить worker для черги.
Неправильний варіант:
async process(job: Job): Promise<void> {
try {
await this.runTask(job);
} catch (error) {
console.error(error);
}
}У цьому випадку метод завершується успішно після catch, тому BullMQ не запустить повторну спробу.
Якщо помилку потрібно залогувати, її слід викинути повторно:
async process(job: Job): Promise<void> {
try {
await this.runTask(job);
} catch (error) {
console.error(error);
throw error;
}
}У Redis зберігаються і дані завдання, тому не варто передавати в чергу великі файли або великі об’єкти. Зазвичай у завданні передають ідентифікатор ресурсу, а сам ресурс processor отримує під час виконання.
Через повторні спроби одна операція може бути запущена більше одного разу. Обробник має враховувати це, особливо під час:
надсилання повідомлень;
створення замовлень;
зміни балансу;
виклику зовнішніх API.
BullMQ використовує Redis для зберігання та виконання фонових завдань.
BullModule.forRoot() налаштовує підключення до Redis.
BullModule.registerQueue() реєструє чергу в NestJS.
@InjectQueue() дає доступ до черги для додавання завдань.
@Processor() і WorkerHost використовуються для обробки завдань.
attempts і backoff налаштовують повторні спроби.
Processor має викидати помилку, якщо завдання не виконано.
Producer і processor можуть працювати в одному або різних застосунках.
Надійність черги не скасовує потребу робити фонові операції безпечними для повторного виконання.