Когда сервисы начинают общаться через брокер, возникает вопрос: как правильно организовать очереди, обменники и подписки? Большинство задач укладываются в несколько типовых схем. Разберём каждую на примерах amqplib.
Раздать задачи нескольким обработчикам — Work Queue
Представьте: пользователи загружают фотографии, и каждую надо сжать и нарезать в несколько размеров. Это долго. Делать это прямо в HTTP-запросе нельзя — пользователь будет ждать минуты.
Решение — work queue (очередь задач): сохранить задание в очередь и отдать её одному из пула обработчиков.
Очередь одна, потребителей несколько: брокер раздаёт сообщения по одному, и каждое достаётся ровно кому-то одному. Больше потребителей — быстрее разбирается очередь, но общий порядок между ними уже не сохраняется.
Брокер сам распределяет задачи между обработчиками. Каждое сообщение получает ровно один из них.
const channel = await connection.createChannel();
await channel.assertQueue('images.to-process', {
durable: true,
arguments: { 'x-queue-type': 'quorum' },
});
await channel.prefetch(20);
await channel.consume('images.to-process', async (msg) => {
if (!msg) return;
const job: ImageJob = JSON.parse(msg.content.toString());
// обработка изображения
channel.ack(msg);
});
prefetch(20) означает: обработчик держит до 20 неподтверждённых сообщений одновременно. Если запустить 5 копий сервиса — получится до 100 параллельных обработчиков.
Когда брать: фоновые задачи — обработка файлов, отправка писем, генерация отчётов, любые «положили в очередь — кто-то возьмёт».
Отправить событие всем сервисам сразу — Publish/Subscribe
Другая задача: произошло событие «конфигурация обновлена» — каждый сервис должен обновить свой кеш. Нельзя знать заранее, кто именно подписан и сколько сервисов запущено.
Это publish/subscribe (рассылка всем): издатель отправляет одно сообщение, а копию получают все подписчики одновременно.
Здесь нужен fanout exchange — он копирует каждое сообщение во все привязанные очереди. Каждый сервис объявляет свою очередь и привязывает её к общему обменнику.
await channel.assertExchange('cache.invalidation', 'fanout', { durable: true });
const { queue } = await channel.assertQueue('', {
exclusive: true,
autoDelete: true,
durable: false,
});
await channel.bindQueue(queue, 'cache.invalidation', '');
await channel.consume(queue, (msg) => {
if (!msg) return;
const event: CacheInvalidationEvent = JSON.parse(msg.content.toString());
cache.evict(event.key);
channel.ack(msg);
});
exclusive + autoDelete — очередь принадлежит одному соединению и удаляется при отключении. При перезапуске сервиса не накапливается мусор в брокере.
Когда брать: инвалидация кешей, broadcast-уведомления всему кластеру, обновления конфигурации.
Направить событие в нужный обработчик — Routing
Иногда нужно не «всем», а «именно тому, кому надо». Например: событие order.created должно идти в сервис выполнения и в аудит, а order.payment-failed — только в алерты.
Это routing (точечная маршрутизация): direct exchange смотрит на ключ маршрутизации сообщения и отправляет его только в очереди с совпадающим binding key.
await channel.assertExchange('orders', 'direct', { durable: true });
const quorum = { durable: true, arguments: { 'x-queue-type': 'quorum' } };
await channel.assertQueue('orders.fulfillment', quorum);
await channel.assertQueue('orders.audit', quorum);
await channel.assertQueue('orders.alerts', quorum);
await channel.bindQueue('orders.fulfillment', 'orders', 'order.created');
await channel.bindQueue('orders.audit', 'orders', 'order.created');
await channel.bindQueue('orders.audit', 'orders', 'order.cancelled');
await channel.bindQueue('orders.alerts', 'orders', 'order.payment-failed');
order.created→ fulfillment + audit.order.cancelled→ только audit.order.payment-failed→ только alerts.
Когда брать: явное разделение потоков — алерты отдельно от аудита, основной обработчик отдельно от мониторинга.
Подписаться по маске — Topic
Routing хорош для жёстких правил. Но что если сервис хочет подписаться на «все события по заказам»? Или «всё из EU-региона»?
Topic exchange позволяет задавать подписки через шаблоны. Ключи сообщений строятся через точку (order.created.eu), а в подписке можно использовать:
*— ровно одно слово,#— ноль и более слов.
await channel.assertExchange('events', 'topic', { durable: true });
const quorum = { durable: true, arguments: { 'x-queue-type': 'quorum' } };
await channel.assertQueue('audit.orders', quorum);
await channel.assertQueue('dashboard.eu', quorum);
await channel.assertQueue('alerts.critical', quorum);
await channel.bindQueue('audit.orders', 'events', 'order.#');
await channel.bindQueue('dashboard.eu', 'events', '*.*.eu');
await channel.bindQueue('alerts.critical', 'events', 'payment.failed.#');
Сообщение с ключом order.cancelled.eu попадёт в audit.orders (по order.#) и в dashboard.eu (по *.*.eu).
Когда брать: события с иерархической структурой, когда нужно гибко подписываться без переделки топологии при добавлении новых типов событий.
Порядок по ключу при нескольких обработчиках
Work Queue выше раздаёт сообщения нескольким обработчикам, и это ровно то, что нужно, пока сообщения независимы. Как только они связаны (два обновления одного заказа, создание и отмена одной брони), параллельная обработка ломает порядок: второе сообщение попадает к свободному обработчику раньше, чем первый закончил.
В Kafka порядок держит ключ партиции. У RabbitMQ порядок гарантирован только внутри одной очереди с одним потребителем, и из этого растут два приёма.
Single active consumer. Аргумент очереди x-single-active-consumer: true означает: подписчиков может быть много, но сообщения получает один, остальные ждут и мгновенно подхватывают работу, если активный отвалился. Порядок сохранён, отказоустойчивость есть, параллелизма нет. Это правильный выбор, когда очередь и так не перегружена, а порядок важен целиком.
await channel.assertQueue('bookings', {
durable: true,
arguments: { 'x-queue-type': 'quorum', 'x-single-active-consumer': true },
});
Consistent hash exchange. Плагин rabbitmq_consistent_hash_exchange добавляет обменник, который раскладывает сообщения по нескольким очередям по хешу ключа маршрутизации или заголовка. Сообщения одного заказа всегда попадают в одну очередь, у каждой очереди свой единственный потребитель, и получается то же, что партиции в Kafka: порядок внутри ключа плюс параллелизм между ключами. Цена ручное управление: очереди создаёте вы, число очередей меняется только с остановкой, а перекос по ключам («один крупный клиент») складывается в одну очередь.
Что не работает: «сделаем один потребитель, но с prefetch в десять и параллельной обработкой». Обработчики разберут сообщения из общей выдачи и выполнят их вперемешку, порядок теряется так же, как при нескольких потребителях. Если нужна и скорость, и порядок, ключ обязан определять, кто обрабатывает; других вариантов нет.
Запрос-ответ через очередь — RPC
Иногда нужен синхронный ответ, но HTTP не подходит: сервис за NAT, нет публичного адреса, или хочется балансировки по пулу обработчиков.
RPC через очередь: клиент отправляет запрос и ждёт ответа. Брокер доставляет запрос одному из обработчиков, тот отвечает в отдельную очередь-ответ. Для сопоставления запроса и ответа используется correlation-id.
В amqplib это собирается вручную; RabbitMQ даёт быструю псевдоочередь amq.rabbitmq.reply-to (direct reply-to):
// Клиент
const pending = new Map<string, (quote: PriceQuote) => void>();
await channel.consume('amq.rabbitmq.reply-to', (msg) => {
if (!msg) return;
pending.get(msg.properties.correlationId)?.(JSON.parse(msg.content.toString()));
pending.delete(msg.properties.correlationId);
}, { noAck: true });
function quote(request: QuoteRequest): Promise<PriceQuote> {
const correlationId = randomUUID();
return new Promise((resolve) => {
pending.set(correlationId, resolve);
channel.publish('pricing.exchange', 'pricing.quote',
Buffer.from(JSON.stringify(request)),
{ replyTo: 'amq.rabbitmq.reply-to', correlationId });
});
}
// Сервер
await channel.consume('pricing.quote', (msg) => {
if (!msg) return;
const quote = computeQuote(JSON.parse(msg.content.toString()));
channel.sendToQueue(msg.properties.replyTo, Buffer.from(JSON.stringify(quote)), {
correlationId: msg.properties.correlationId,
});
channel.ack(msg);
});
Клиент подписывается на amq.rabbitmq.reply-to, проставляет reply-to и correlation-id, сопоставляет ответы с запросами. Сервер публикует ответ в очередь из свойства replyTo запроса.
Когда брать: нужен синхронный вызов, но HTTP не работает (NAT, firewall, нет публичного адреса); нужна балансировка запросов по пулу обработчиков.
Когда не брать: если HTTP/gRPC просто работает — RPC через брокер сложнее в отладке и дороже.
Что делать с повторной доставкой — Idempotent Consumer
AMQP гарантирует at-least-once доставку: одно и то же сообщение может прийти дважды. Это случается, когда брокер не получил подтверждение обработки (например, из-за сетевой проблемы) и переотправляет сообщение.
Обработчик обязан уметь работать с повторами, не ломая бизнес-логику.
Ключ идемпотентности в базе данных
Самый надёжный способ — запоминать уже обработанные сообщения:
await channel.consume('payments', async (msg) => {
if (!msg) return;
const event: PaymentEvent = JSON.parse(msg.content.toString());
await dataSource.transaction(async (tx) => {
if (await tx.processedEvents.existsByIdempotencyKey(event.idempotencyKey)) {
return; // уже обработано — просто подтверждаем получение
}
await tx.processedEvents.save({ idempotencyKey: event.idempotencyKey });
await tx.accounts.debit(event.accountId, event.amount);
});
channel.ack(msg);
});
Таблица processed_events с уникальным индексом на idempotency_key. Если два одинаковых сообщения придут одновременно — база данных поймает дубль через нарушение уникального ограничения.
Проверка состояния объекта
Если событие переводит объект в новое состояние, достаточно проверить текущее:
async function onOrderConfirmed(event: OrderConfirmedEvent): Promise<void> {
await dataSource.transaction(async (tx) => {
const order = await tx.orders.findByIdOrFail(event.orderId);
if (order.status === OrderStatus.CONFIRMED) {
return; // уже в нужном состоянии
}
order.confirm();
await tx.orders.save(order);
});
}
Не требует отдельной таблицы — состояние уже хранится в бизнес-объекте.
Retry с задержкой и Dead Letter Queue
Что если обработчик упал не из-за бага, а из-за временной недоступности внешнего сервиса? Нужно попробовать снова, но не немедленно.
На amqplib delayed retry собирается через x-message-ttl и Dead Letter Exchange:
await channel.assertQueue('orders.retry', {
durable: true,
arguments: {
'x-queue-type': 'quorum',
'x-message-ttl': 30_000, // ждём 30 секунд
'x-dead-letter-exchange': 'orders',
'x-dead-letter-routing-key': 'order.created',
},
});
Поток: обработчик отклонил сообщение → оно попадает в retry очередь → через 30 секунд по истечении TTL уходит через DLX обратно в основную очередь → новая попытка.
Количество попыток считается через заголовок x-death.count — его нужно проверять вручную, встроенного ограничителя нет.
Ещё одна ловушка: очередь с TTL отдаёт сообщения строго по порядку, и сообщение с задержкой 5 секунд, вставшее за десятиминутным, дождётся его. Поэтому под каждую задержку заводят свою очередь повторов (orders.retry.5s, orders.retry.30s, orders.retry.10m) либо ставят плагин rabbitmq_delayed_message_exchange. Он включается на брокере (rabbitmq-plugins enable rabbitmq_delayed_message_exchange) и добавляет особый тип обменника, у которого задержка задаётся у каждого сообщения заголовком:
await channel.assertExchange('orders.delayed', 'x-delayed-message', {
durable: true,
arguments: { 'x-delayed-type': 'direct' },
});
channel.publish('orders.delayed', 'order.created', Buffer.from(JSON.stringify(event)), {
persistent: true,
headers: { 'x-delay': 30_000 }, // задержка этого сообщения, мс
});
Одна очередь, произвольные задержки, никакой головной блокировки. Цена тоже есть: плагин держит отложенные сообщения в своей встроенной базе на том узле, куда они опубликованы, это не очередь, и она не реплицируется. Падение узла до срока доставки теряет отложенные сообщения, а большое их число заметно ест память. Для повторов и коротких задержек плагин удобен, а для «напомнить через месяц» берут не брокер, а таблицу в базе и задачу по расписанию.
Сообщения, которые не получилось обработать после всех попыток, уходят в Dead Letter Queue (DLQ) — отдельную очередь для разбора вручную или через алерты.
Гарантированная публикация — Outbox
Бывает задача: сохранить заказ в базу данных и опубликовать событие — атомарно. Если сначала сохранить, потом опубликовать, то сервис может упасть между двумя операциями. Событие потеряется.
Outbox pattern: событие сохраняется в ту же транзакцию, что и бизнес-данные. Отдельный процесс читает таблицу и публикует в AMQP.
@Injectable()
export class OrderService {
async confirm(orderId: OrderId): Promise<void> {
await this.dataSource.transaction(async (tx) => {
const order = await tx.orders.findByIdOrFail(orderId);
order.confirm();
await tx.orders.save(order);
await tx.outbox.save({
id: randomUUID(),
routingKey: 'order.confirmed',
exchange: 'orders',
payload: JSON.stringify(new OrderConfirmedEvent(orderId)),
});
});
}
@Interval(500)
async publishOutbox(): Promise<void> {
const batch = await this.outboxRepo.fetchUnpublished(100);
for (const event of batch) {
this.channel.publish(event.exchange, event.routingKey, Buffer.from(event.payload));
await this.outboxRepo.markPublished(event.id);
}
}
}
Либо оба изменения зафиксированы, либо ни одного. Дубли возможны (публикация прошла, но пометить как отправленное не успело) — поэтому получатель всё равно должен быть идемпотентным.
Шпаргалка по выбору
| Задача | Паттерн | Тип обменника |
|---|---|---|
| Распределить нагрузку между обработчиками | Work Queue | direct (default) |
| Broadcast события всем сервисам | Publish/Subscribe | fanout |
| Разные события в разные очереди | Routing | direct |
| Подписка по маске на иерархические события | Topic | topic |
| Синхронный вызов через очередь | RPC | direct + reply-to |
| Защита от повторной доставки | Idempotent Consumer | любой |
| Retry с задержкой | Delayed Retry | direct + DLX |
| Атомарная публикация вместе с записью в БД | Outbox | direct |
Коротко
- Work Queue — одна очередь, несколько обработчиков, каждое сообщение получает ровно один. Для фоновых задач.
- Publish/Subscribe — fanout exchange копирует сообщение во все привязанные очереди. Для broadcast-событий.
- Routing — direct exchange смотрит на ключ маршрутизации. Для точного разделения потоков.
- Topic — как routing, но с шаблонами
*и#. Для иерархических событий с гибкой подпиской. - RPC через очередь — запрос-ответ через брокер с
reply-toиcorrelation-id. Для вызовов без HTTP. - Idempotent Consumer — at-least-once означает возможные дубли. Защита: ключ идемпотентности в БД или проверка состояния объекта.
- Delayed Retry — TTL + DLX: сообщение «паркуется» на время, потом возвращается.
- Outbox — событие сохраняется в той же транзакции, что и данные. Атомарность без двухфазного коммита.
Что почитать дальше
- Протокол AMQP — модель exchange/binding/queue изнутри.
- Spring AMQP — конфигурация, RabbitTemplate, аннотации.
- RabbitMQ в production — Quorum Queues, кластеризация, мониторинг.
- AMQP vs Kafka — какой брокер и когда брать.