RabbitMQ: обмены, очереди, маршрутизация и гарантии доставки
Кратко о теме
Заголовок раздела «Кратко о теме»RabbitMQ — это брокер с «умным» сервером и «глупым» потребителем, в противоположность Kafka, где логика в основном на стороне клиента. Модель, которую надо держать в голове, задаётся протоколом AMQP 0-9-1: producer никогда не публикует «в очередь», он публикует в exchange (обмен) с некоторым routing key; exchange по своим bindings (привязкам) решает, в какие queues положить копию сообщения; consumer подписывается на конкретную очередь и получает сообщения оттуда. Очередь — это буфер с сообщениями, exchange — правило маршрутизации без собственного хранилища. Если ни одна привязка не совпала, сообщение молча уничтожается (или возвращается публикующему, если он выставил флаг mandatory, либо уходит в alternate exchange).
Типов exchange четыре. direct — точное совпадение routing key с binding key. fanout — копия в каждую привязанную очередь, routing key игнорируется. topic — сопоставление по шаблону с разделителем . и подстановками * (ровно одно слово) и # (ноль и более слов). headers — маршрутизация по заголовкам вместо ключа, применяется редко. Плюс есть безымянный exchange по умолчанию ("") типа direct, к которому каждая очередь автоматически привязана по своему имени — именно поэтому в туториалах публикуют «прямо в очередь», указывая её имя как routing key.
Ключевое отличие от Kafka по семантике потребления: очередь RabbitMQ — деструктивная. Сообщение доставляется одному потребителю, тот подтверждает его (basic.ack), и брокер его удаляет. Реплея как в Kafka нет — «перечитать топик с начала» нельзя, потому что нет лога с оффсетами (исключение — тип очереди stream, появившийся в 3.9, который как раз является append-only логом с оффсетами и неразрушающим чтением). Отсюда следует главный архитектурный вывод: чтобы одно событие обработали N разных сервисов, нужно N очередей, привязанных к одному exchange, — по очереди на сервис. А несколько инстансов одного сервиса читают из одной общей очереди и делят нагрузку по схеме competing consumers.
Гарантии доставки — at-least-once с обеих сторон и только при явной настройке. Со стороны publisher: publisher confirms (канал переводится в confirm mode, брокер асинхронно подтверждает принятие каждого сообщения), persistent delivery mode и durable-очередь — без этого набора сообщение может исчезнуть при рестарте брокера. Со стороны consumer: autoAck=false и ручной Ack/Nack после успешной обработки; если соединение потребителя оборвалось до ack, сообщение вернётся в очередь и уйдёт другому. Exactly-once брокер не даёт ни в каком режиме, поэтому обработчик обязан быть идемпотентным (дедуп по message_id или по бизнес-ключу).
Про типы очередей в современных версиях: классические зеркалируемые очереди (classic mirrored queues) объявлены устаревшими в 3.8/3.9 и удалены в RabbitMQ 4.0; штатный отказоустойчивый вариант сегодня — quorum queues (репликация через Raft, нечётное число реплик, гарантированный консенсус вместо eventual-консистентности зеркал). Для потоковых сценариев — streams. Метаданные кластера исторически хранились в Mnesia, в ветке 3.13/4.x её вытесняет Raft-хранилище Khepri.
Вопросы и ответы
Заголовок раздела «Вопросы и ответы»Как в RabbitMQ организовали обработку одного события несколькими сервисами, каждый из которых запущен в нескольких экземплярах? Сейчас выбрали бы тот же RabbitMQ или что-то другое? Почему?
Заголовок раздела «Как в RabbitMQ организовали обработку одного события несколькими сервисами, каждый из которых запущен в нескольких экземплярах? Сейчас выбрали бы тот же RabbitMQ или что-то другое? Почему?»Коротко. Один exchange (fanout или topic) и отдельная durable-очередь на каждый сервис-подписчик, привязанная к этому exchange. Все инстансы одного сервиса подключаются к своей общей очереди как competing consumers — брокер раздаёт им сообщения по кругу с учётом prefetch, так что каждое событие обрабатывается ровно одним инстансом сервиса, но всеми сервисами независимо. Очередь принадлежит подписчику, а не издателю: издатель про подписчиков не знает вообще.
Глубже. Разберём по частям, потому что на собеседовании обычно докапываются именно до деталей.
Fan-out между сервисами. Очередь в RabbitMQ деструктивна, поэтому «одна очередь на всех» дала бы конкуренцию между разными сервисами, а не размножение события. Правильная схема — очередь на каждого потребителя: orders (topic) → billing.order.created, notify.order.created, analytics.order.created. Это ровно тот же эффект, что consumer groups в Kafka, только группы материализованы в виде физических очередей, и их создаёт сам подписчик при старте (declare идемпотентен). Очереди обязательно durable, сообщения persistent, иначе рестарт брокера съест буфер. Топология обычно объявляется самим сервисом на старте, чтобы новый подписчик не требовал ручных действий в UI.
Балансировка между инстансами. Все инстансы делают basic.consume на одну очередь. Без настройки брокер раздаёт сообщения round-robin, не глядя на занятость потребителя, — быстрый инстанс простаивает, медленный копит невыполненное. Лечится basic.qos(prefetch_count=N) на канал: брокер не отдаёт потребителю больше N неподтверждённых сообщений. Практический ориентир — небольшое значение (10–50) для длинной обработки и большее для быстрой; prefetch=1 даёт идеальный fair dispatch ценой лишних round-trip.
Что ломается и как чинится. Порядок сообщений сохраняется только внутри очереди при одном потребителе; как только инстансов больше одного или включён requeue, глобального порядка нет. Если порядок нужен per-entity, ставят плагин rabbitmq_consistent_hash_exchange и хешируют по ключу сущности в M очередей — получается аналог партиций. Ядовитые сообщения (постоянно падающие) заворачиваются в DLX: Nack(requeue=false) отправляет сообщение в dead-letter exchange, а не в бесконечный цикл; у quorum-очередей для этого есть встроенный x-delivery-limit. Ретраи с задержкой делаются либо плагином rabbitmq_delayed_message_exchange, либо связкой «очередь с x-message-ttl + DLX обратно в рабочую очередь». Из-за at-least-once обработчики обязаны быть идемпотентны.
Стоит ли выбрать RabbitMQ снова. Это вопрос на инженерное суждение, отвечать надо через критерии, а не «Kafka лучше». RabbitMQ уместен, когда нужна сложная маршрутизация (topic/headers), per-message ack и retry, приоритеты, delayed-сообщения, low-latency RPC, тысячи очередей и относительно скромный поток (десятки тысяч сообщений в секунду). Kafka уместна, когда событие — это факт, который может понадобиться перечитать: нужен replay и хранение истории, нужен строгий порядок по ключу, нужны сотни тысяч–миллионы сообщений в секунду, нужны потоковые агрегации и подключение новых потребителей задним числом без участия издателя. Честный ответ звучит примерно так: «для событийной интеграции с несколькими независимыми потребителями и требованием переиграть историю сейчас взял бы Kafka — новый подписчик там не требует заранее созданной очереди и может прочитать прошлое; для task queue с ретраями, приоритетами и отложенным выполнением остался бы на RabbitMQ. Промежуточный вариант — RabbitMQ Streams, если не хочется тащить в инфраструктуру ещё один кластер ради replay». Отдельно стоит упомянуть операционную сторону: RabbitMQ проще в эксплуатации на малых объёмах, но плохо переносит длинные очереди — при разрастании включается flow control и падает throughput, тогда как для Kafka «лежащие» данные это норма.
package main
import ( "context" "log" "time"
amqp "github.com/rabbitmq/amqp091-go")
// Подписчик сервиса billing: своя очередь на общий exchange.func consume(ctx context.Context, ch *amqp.Channel) error { if err := ch.ExchangeDeclare("orders", "topic", true, false, false, false, nil); err != nil { return err } q, err := ch.QueueDeclare("billing.order.created", true, false, false, false, amqp.Table{ "x-queue-type": "quorum", "x-dead-letter-exchange": "orders.dlx", "x-delivery-limit": int32(5), }) if err != nil { return err } if err := ch.QueueBind(q.Name, "order.created", "orders", false, nil); err != nil { return err } // Не более 20 неподтверждённых сообщений на этот канал. if err := ch.Qos(20, 0, false); err != nil { return err }
deliveries, err := ch.ConsumeWithContext(ctx, q.Name, "", false, false, false, false, nil) if err != nil { return err } for d := range deliveries { if err := handle(d.Body); err != nil { // requeue=false — сообщение уходит в DLX, а не в бесконечный цикл. _ = d.Nack(false, false) continue } _ = d.Ack(false) } return nil}
func handle(body []byte) error { _ = body; return nil }
// Издатель: публикация с подтверждением от брокера.func publish(ctx context.Context, ch *amqp.Channel, body []byte) error { if err := ch.Confirm(false); err != nil { return err } ctx, cancel := context.WithTimeout(ctx, 5*time.Second) defer cancel()
dc, err := ch.PublishWithDeferredConfirmWithContext(ctx, "orders", "order.created", true, // mandatory: вернуть сообщение, если никуда не смаршрутизировалось false, // immediate: устарел, всегда false amqp.Publishing{ ContentType: "application/json", DeliveryMode: amqp.Persistent, MessageId: "01J8Z9...", // для дедупликации на стороне потребителя Timestamp: time.Now(), Body: body, }) if err != nil { return err } acked, err := dc.WaitContext(ctx) if err != nil { return err } if !acked { log.Println("broker nack: сообщение не принято, нужен повтор") } return nil}В RabbitMQ куда отправляем сообщения и откуда забираем?
Заголовок раздела «В RabbitMQ куда отправляем сообщения и откуда забираем?»Коротко. Отправляем в exchange с routing key, забираем из очереди. Producer никогда не пишет в очередь напрямую: exchange по привязкам (binding key + тип обмена) кладёт копии сообщения в подходящие очереди, а consumer подписывается на конкретную очередь.
Глубже. Иллюзия «публикуем в очередь» возникает из-за default exchange — безымянного обмена "" типа direct, к которому брокер автоматически привязывает каждую очередь по её имени. Поэтому publish(exchange: "", routingKey: "my-queue") попадает именно в my-queue; в проде так делать не стоит, потому что теряется вся гибкость маршрутизации и издатель жёстко привязывается к именам очередей подписчиков.
Полезно помнить крайние случаи. Если сообщение не подошло ни под одну привязку, оно бесследно удаляется — это самая частая причина «сообщения пропадают»; чтобы это увидеть, публикуют с mandatory=true и слушают канал возвратов (Channel.NotifyReturn), либо настраивают на exchange параметр alternate-exchange, куда сваливается всё несмаршрутизированное. Exchange можно привязать к другому exchange (exchange-to-exchange binding) — так строят многоуровневую маршрутизацию. Забирать из очереди можно двумя способами: push-модель basic.consume (брокер сам шлёт сообщения по мере появления — это правильный способ) и pull-модель basic.get (разовый опрос, дорогой round-trip на каждое сообщение, годится для отладки и редких заданий). Само сообщение состоит из свойств (content-type, delivery-mode, message-id, correlation-id, reply-to, priority, expiration, headers) и непрозрачного для брокера тела []byte.
Работал ли с mq rabbit?
Заголовок раздела «Работал ли с mq rabbit?»Коротко. Вопрос на личный опыт: интервьюер проверяет не факт «да/нет», а глубину — какие задачи решали, какую топологию строили, какие грабли ловили. Отвечать надо конкретно: контекст → топология → настройки надёжности → инцидент и как чинили.
Глубже. Каркас хорошего ответа на 60–90 секунд.
- Контекст и объёмы. Что за поток, зачем брокер (асинхронная обработка, интеграция сервисов, отложенные задачи), порядок величин: сколько сообщений в секунду, какой размер, сколько потребителей. Цифры сразу отличают реальный опыт от «поднимал в docker-compose».
- Топология. Какие exchange и почему именно такого типа, как назывались очереди, как разделяли подписчиков, был ли DLX, были ли ретраи с задержкой, использовали ли quorum-очереди или классические.
- Надёжность. Publisher confirms, persistent + durable, ручной ack, prefetch и его подобранное значение, идемпотентность обработчика, transactional outbox если писали в БД и в брокер одновременно.
- Эксплуатация. Мониторинг (глубина очереди,
messages_unacknowledged,publish/deliverrate, memory/disk alarms), поведение при разрастании очереди, реконнекты клиента, как выкатывали изменения топологии. - Конкретный инцидент. Самая ценная часть: «потребитель падал на битом JSON и бесконечно возвращал сообщение в очередь — добавили DLX и
x-delivery-limit», «выставилиprefetch=1000, один инстанс забирал всю очередь себе — снизили до 30», «забылиdurableна очереди, рестарт брокера потерял задачи», «очередь выросла до миллионов сообщений, брокер упёрся в память и включил flow control — вынесли лишний трафик и включили lazy-режим». - Границы честности. Если опыта мало — так и сказать, но сразу перевести в модель: «в проде вёл RabbitMQ только как потребитель, топологию не проектировал; но модель exchange → binding → queue и разницу с Kafka понимаю, могу расписать». Это лучше, чем плавать в деталях.
Типичные ошибки: отвечать одним словом «да, работал»; путать exchange и очередь; рассказывать документацию вместо своего опыта; заявлять «настраивал кластер», не умея объяснить разницу между зеркалированием и quorum-очередями.
Частые ошибки на собесе
Заголовок раздела «Частые ошибки на собесе»- Говорят «публикуем в очередь» — публикуют всегда в exchange, а очередь получает сообщение по привязке. Через default exchange это выглядит так же, но модель другая.
- Считают, что несколько потребителей на одной очереди получат по копии сообщения. Нет: очередь деструктивна, копию получает один. Копия каждому сервису — это отдельная очередь на сервис.
- Забывают, что несмаршрутизированное сообщение молча удаляется, и потом ищут «потерю сообщений» в приложении вместо
mandatory/alternate exchange. - Путают durable-очередь и persistent-сообщение: нужны обе настройки одновременно, иначе рестарт брокера всё равно потеряет данные.
- Обещают exactly-once. RabbitMQ даёт at-least-once; дедупликацию делает потребитель.
- Утверждают, что порядок сообщений гарантирован. Он держится только в пределах одной очереди с одним потребителем и ломается при requeue и нескольких инстансах.
- Рассказывают про mirrored queues как про актуальный способ отказоустойчивости — они удалены в RabbitMQ 4.0, современный ответ это quorum queues.
- Отвечают «RabbitMQ — это как Kafka» и не могут назвать отличие: у RabbitMQ нет реплея из лога (кроме streams), а маршрутизация на стороне брокера, а не потребителя.
- Ставят огромный
prefetch«для скорости» и получают перекос нагрузки и раздутую задержку при падении инстанса.
Что почитать
Заголовок раздела «Что почитать»- RabbitMQ Tutorials — канонические 7 туториалов (work queues, pub/sub, routing, topics, RPC), в том числе на Go.
- AMQP 0-9-1 Model Explained — точное описание exchange, bindings, queues и свойств сообщения.
- Quorum Queues и Streams — современные типы очередей и их ограничения.
- Publisher Confirms and Consumer Acknowledgements — как именно устроены гарантии доставки и prefetch.
- rabbitmq/amqp091-go — официальный Go-клиент (форк streadway/amqp), примеры реконнекта и confirm-режима.