Пошук уроків, статей та іншого контенту
Налаштуєте RabbitMQ, черги та обмінники для надійної доставки повідомлень.
RabbitMQ — брокер повідомлень, який приймає повідомлення від producer і доставляє їх consumer через черги.
Типовий потік має такий вигляд:
Producer публікує повідомлення.
Exchange визначає, у які черги його передати.
Queue зберігає повідомлення до моменту обробки.
Consumer отримує повідомлення та підтверджує його обробку.
У NestJS RabbitMQ використовується як транспорт для мікросервісів:
ClientProxy надсилає повідомлення;
@MessagePattern() визначає обробник;
RmqContext дає доступ до оригінального RabbitMQ-повідомлення;
ручне підтвердження через ack() дозволяє не втрачати повідомлення під час помилок.
Для локальної розробки зручно запустити RabbitMQ у Docker:
docker run -d \
--name rabbitmq \
-p 5672:5672 \
-p 15672:15672 \
rabbitmq:3-managementПорти:
5672 — AMQP-з'єднання для застосунків;
15672 — вебінтерфейс RabbitMQ Management.
За замовчуванням RabbitMQ створює користувача:
логін: guest;
пароль: guest.
Підключення до брокера матиме вигляд:
amqp://guest:guest@localhost:5672Для NestJS-мікросервісів потрібні пакети:
npm install @nestjs/microservices amqplib
npm install -D @types/amqplib@nestjs/microservices містить транспорт RabbitMQ, а amqplib знадобиться для явного створення exchange та публікації повідомлень із підтвердженням брокера.
Створимо мікросервіс, який обробляє повідомлення з черги orders.created.
main.tsimport { NestFactory } from '@nestjs/core';
import {
MicroserviceOptions,
Transport,
} from '@nestjs/microservices';
import { AppModule } from './app.module';
async function bootstrap() {
const app = await NestFactory.createMicroservice<MicroserviceOptions>(
AppModule,
{
transport: Transport.RMQ,
options: {
urls: ['amqp://guest:guest@localhost:5672'],
queue: 'orders.created',
queueOptions: {
durable: true,
},
noAck: false,
prefetchCount: 10,
},
},
);
await app.listen();
}
bootstrap();Основні параметри:
urls — адреси RabbitMQ;
queue — черга, з якої читає мікросервіс;
durable: true — черга переживає перезапуск RabbitMQ;
noAck: false — повідомлення потрібно підтверджувати вручну;
prefetchCount: 10 — consumer отримує не більше 10 непідтверджених повідомлень одночасно.
prefetchCount допомагає не перевантажувати один екземпляр мікросервісу. RabbitMQ не передаватиме йому нові повідомлення, доки частина попередніх не буде підтверджена.
import {
Controller,
Logger,
} from '@nestjs/common';
import {
Ctx,
MessagePattern,
Payload,
RmqContext,
} from '@nestjs/microservices';
@Controller()
export class OrdersController {
private readonly logger = new Logger(OrdersController.name);
@MessagePattern('order.created')
async handleOrderCreated(
@Payload() payload: { orderId: string; userId: string },
@Ctx() context: RmqContext,
) {
const channel = context.getChannelRef();
const message = context.getMessage();
try {
this.logger.log(`Обробка замовлення ${payload.orderId}`);
await this.processOrder(payload);
// Підтверджуємо повідомлення лише після успішної обробки.
channel.ack(message);
} catch (error) {
this.logger.error(
`Помилка обробки замовлення ${payload.orderId}`,
error instanceof Error ? error.stack : undefined,
);
// Повідомлення не повертається в чергу,
// щоб нескінченно не повторювати невдалу обробку.
channel.nack(message, false, false);
}
}
private async processOrder(payload: {
orderId: string;
userId: string;
}) {
// Тут може бути запис у базу даних або виклик іншого сервісу.
await new Promise((resolve) => setTimeout(resolve, 100));
}
}Поки повідомлення не підтверджене через ack(), RabbitMQ вважає його незавершеним.
Якщо consumer аварійно завершиться до виклику ack(), RabbitMQ зможе доставити це повідомлення повторно після перепідключення іншого consumer.
ClientProxyNestJS має вбудований клієнт RabbitMQ. Його можна зареєструвати через ClientsModule.
app.module.tsimport { Module } from '@nestjs/common';
import {
ClientsModule,
Transport,
} from '@nestjs/microservices';
import { OrdersService } from './orders.service';
@Module({
imports: [
ClientsModule.register([
{
name: 'ORDERS_SERVICE',
transport: Transport.RMQ,
options: {
urls: ['amqp://guest:guest@localhost:5672'],
queue: 'orders.created',
queueOptions: {
durable: true,
},
},
},
]),
],
providers: [OrdersService],
exports: [OrdersService],
})
export class AppModule {}orders.service.tsimport {
Inject,
Injectable,
OnModuleInit,
} from '@nestjs/common';
import { ClientProxy } from '@nestjs/microservices';
@Injectable()
export class OrdersService implements OnModuleInit {
constructor(
@Inject('ORDERS_SERVICE')
private readonly ordersClient: ClientProxy,
) {}
async onModuleInit() {
await this.ordersClient.connect();
}
publishOrderCreated(orderId: string, userId: string) {
return this.ordersClient.emit('order.created', {
orderId,
userId,
});
}
}Метод emit() призначений для подій, коли producer не очікує відповіді від consumer.
Для шаблону запит-відповідь використовується send(), але для фонової обробки подій зазвичай підходить саме emit().
RabbitMQ використовує три важливі поняття:
exchange — приймає повідомлення та маршрутизує їх;
routing key — ключ маршрутизації;
binding — правило зв'язку між exchange і queue.
Популярні типи exchange:
direct — точний збіг routing key;
topic — маршрутизація за шаблоном, наприклад order.*;
fanout — надсилання в усі прив'язані черги;
headers — маршрутизація за заголовками.
У стандартній конфігурації Transport.RMQ NestJS працює безпосередньо з чергою. Для явної роботи з іменованим exchange потрібно створити його через AMQP API.
import amqp from 'amqplib';
async function setupRabbitTopology() {
const connection = await amqp.connect(
'amqp://guest:guest@localhost:5672',
);
const channel = await connection.createChannel();
await channel.assertExchange('orders', 'direct', {
durable: true,
});
await channel.assertQueue('orders.created', {
durable: true,
});
await channel.bindQueue(
'orders.created',
'orders',
'order.created',
);
console.log('RabbitMQ topology створено');
await channel.close();
await connection.close();
}
setupRabbitTopology().catch((error) => {
console.error('Не вдалося створити topology', error);
process.exitCode = 1;
});Цей код створює:
durable exchange orders;
durable queue orders.created;
binding із routing key order.created.
Його можна запустити окремо перед запуском мікросервісів:
npx ts-node setup-rabbit.tsКоли повідомлення потрібно публікувати саме в іменований exchange, використовуйте amqplib.
import amqp from 'amqplib';
type OrderCreatedEvent = {
orderId: string;
userId: string;
};
async function publishOrderCreated(event: OrderCreatedEvent) {
const connection = await amqp.connect(
'amqp://guest:guest@localhost:5672',
);
const channel = await connection.createConfirmChannel();
await channel.assertExchange('orders', 'direct', {
durable: true,
});
const message = Buffer.from(
JSON.stringify({
pattern: 'order.created',
data: event,
}),
);
channel.publish(
'orders',
'order.created',
message,
{
persistent: true,
contentType: 'application/json',
},
);
// Очікуємо підтвердження від RabbitMQ.
await channel.waitForConfirms();
console.log('Повідомлення підтверджено брокером');
await channel.close();
await connection.close();
}
publishOrderCreated({
orderId: 'order-1001',
userId: 'user-42',
}).catch((error) => {
console.error('Не вдалося опублікувати повідомлення', error);
process.exitCode = 1;
});Об'єкт повідомлення містить pattern і data, тому NestJS зможе зіставити його з декоратором:
@MessagePattern('order.created')Параметр persistent: true позначає повідомлення як стійке. Разом із durable exchange і durable queue це зменшує ризик втрати повідомлення під час перезапуску RabbitMQ.
Важливо: стійкість повідомлення не замінює підтвердження producer. createConfirmChannel() і waitForConfirms() дають змогу дочекатися підтвердження, що брокер прийняв повідомлення.
RabbitMQ має два основні режими обробки:
{
noAck: true,
}RabbitMQ вважає повідомлення обробленим одразу після доставки consumer.
Перевага — простота. Недолік — повідомлення може бути втрачено, якщо процес завершиться під час обробки.
{
noAck: false,
}Consumer сам вирішує, коли підтвердити повідомлення:
channel.ack(message);Якщо обробка не вдалася:
channel.nack(message, false, false);Аргументи nack() означають:
повідомлення, яке потрібно відхилити;
чи потрібно відхилити також інші повідомлення;
чи повернути повідомлення назад у чергу.
Параметр true у третьому аргументі спричинить повторну доставку:
channel.nack(message, false, true);Це може бути корисно для тимчасових помилок, але небезпечно для постійних помилок: одне повідомлення може нескінченно оброблятися повторно.
Для базової надійної доставки потрібно узгодити кілька параметрів:
Черга має бути durable.
Exchange має бути durable.
Повідомлення має бути persistent.
Producer має використовувати publisher confirms.
Consumer має використовувати ручний ack.
Потрібно визначити стратегію для невдалих повідомлень.
Приклад конфігурації consumer:
{
transport: Transport.RMQ,
options: {
urls: ['amqp://guest:guest@localhost:5672'],
queue: 'orders.created',
queueOptions: {
durable: true,
},
noAck: false,
prefetchCount: 10,
},
}Однак RabbitMQ не гарантує, що бізнес-операція не буде виконана двічі. Якщо consumer обробив замовлення, але завершився до ack(), повідомлення може бути доставлене повторно.
Тому обробники подій мають бути ідемпотентними. Наприклад, перед створенням платежу можна перевірити, чи не був уже оброблений eventId.
Producer і consumer мають використовувати однакову чергу або узгоджену exchange-топологію. Помилка в назві призводить до того, що consumer не отримує повідомлення.
noAck: true для важливих подійАвтоматичне підтвердження може призвести до втрати повідомлення під час падіння процесу. Для важливих подій використовуйте noAck: false.
ack() до завершення операціїНе підтверджуйте повідомлення перед записом у базу або завершенням зовнішнього виклику. Інакше збій після ack() призведе до втрати події.
Durable queue не зберігає недовговічні повідомлення після перезапуску брокера. Для надійної доставки потрібні і durable-об'єкти, і persistent-повідомлення.
Безумовний виклик:
channel.nack(message, false, true);може зациклити проблемне повідомлення. Для постійних помилок його потрібно відхилити без повернення в основну чергу або передати в окрему чергу помилок.
prefetchCountОдин consumer може отримати надто багато повідомлень до завершення обробки. Обмеження prefetch допомагає рівномірніше розподіляти навантаження між екземплярами мікросервісу.
Порядок запуску:
Запустіть RabbitMQ у Docker.
Створіть exchange, queue і binding.
Запустіть NestJS consumer.
Опублікуйте подію через ClientProxy або amqplib.
Перевірте логи consumer.
Зупиніть consumer до виклику ack() і переконайтеся, що повідомлення доставляється повторно після запуску.
У вебінтерфейсі RabbitMQ можна перевірити:
кількість повідомлень у черзі;
кількість consumer;
кількість непідтверджених повідомлень;
bindings між exchange і queue.
RabbitMQ доставляє повідомлення через exchange та queue.
NestJS підтримує RabbitMQ через Transport.RMQ.
@MessagePattern() обробляє події з відповідним pattern.
noAck: false дає змогу підтверджувати повідомлення вручну.
ack() потрібно викликати після успішної бізнес-операції.
Durable exchange і queue захищають топологію від перезапуску брокера.
persistent: true потрібен для стійких повідомлень.
Publisher confirms дозволяють перевірити прийняття повідомлення брокером.
Для іменованих exchange можна використовувати amqplib.
Обробники повідомлень мають бути готовими до повторної доставки.