Когда два сервиса хотят обменяться данными, простейшее решение — вызов напрямую: один делает HTTP-запрос к другому. Это работает, пока второй сервис всегда доступен, успевает принять все запросы и не падает под нагрузкой.
Как только один из этих пунктов перестаёт выполняться, появляется брокер сообщений. Первый сервис кладёт сообщение в брокер и занимается своими делами; второй забирает его, когда готов. Но чтобы брокер знал, куда класть каждое сообщение и кому его отдавать, нужен протокол — набор правил маршрутизации.
AMQP 0.9.1 (Advanced Message Queuing Protocol) — такой протокол. Его главный носитель — RabbitMQ; поддерживает его и Apache Qpid. Важно не путать версии: существует ещё AMQP 1.0, и это по сути другой протокол, а не следующая редакция. Его понимают Azure Service Bus, ActiveMQ и современный RabbitMQ, но модель обменников и привязок, о которой пойдёт речь дальше, — из версии 0.9.1. Именно AMQP 0.9.1 обычно имеют в виду, когда говорят «AMQP» без уточнения версии.
В RabbitMQ 4.x поддержка AMQP 1.0 перестала быть надстройкой и стала равноправной частью ядра: включена по умолчанию, работает на том же порту, и клиент на 1.0 может читать и писать в те же очереди, что клиент на 0.9.1. Для читателя это значит две практические вещи. Если в команде есть сервис на платформе, где нормальная библиотека только для 1.0 (часто так у .NET и у интеграций с облачными шинами), его можно подключить к тому же брокеру, не ставя мост. И терминология в документации теперь двойная: то, что в 0.9.1 называется каналом и подтверждением, в 1.0 называется сессией, ссылкой и расчётом, а модель обменников для 1.0 предоставляется отдельным уровнем совместимости. Обычные библиотеки, Spring AMQP, amqp091-go, amqplib и pika, говорят на 0.9.1, и менять в них ничего не нужно.
Ниже — вся модель на одном примере: отправитель называет не очередь, а точку входа и метку на сообщении, а дальше копии раскладывает брокер.
Продьюсер называет только exchange и ключ payment-failed — куда попадёт сообщение, решают bindings: две привязки из трёх совпали, значит alerts и audit-log получат по копии, а fulfillment — ничего. Понадобился четвёртый подписчик — добавили ещё один binding, код продьюсера не менялся.
Соединение и канал: почему в коде везде channel
Любой пример работы с брокером начинается с двух объектов, о которых обычно не говорят вовсе: соединение и канал.
Соединение (connection) — это обычное TCP-соединение с брокером. Оно дорогое: рукопожатие, аутентификация, память на стороне брокера. Открывать его на каждое сообщение — та же ошибка, что открывать соединение с базой на каждый запрос, поэтому соединение держат одно на приложение, иногда два: отдельно на публикацию и на чтение.
Канал (channel) — логическое соединение внутри одного TCP. Все команды протокола идут по каналу: объявить очередь, опубликовать, подписаться, подтвердить. Именно поэтому в примерах везде channel.basicPublish(...), а не connection.basicPublish(...). Каналов внутри соединения делают много, они дешёвые.
Три следствия здесь важнее самих определений.
Канал нельзя использовать из нескольких потоков. Он не потокобезопасен: два потока, публикующие в один канал, перемешают кадры протокола, и выглядеть это будет как случайный разрыв соединения, а не как ошибка в коде. Правило: канал на поток, в Go на горутину. Spring AMQP держит для этого пул каналов и выдаёт каждому потоку свой, поэтому там об этом не думают; в amqp091-go пула нет, и канал, положенный в поле структуры и разделённый между горутинами, ломается ровно так; в Node процесс однопоточный, и одного канала на публикацию хватает. Именно поэтому нельзя хранить канал в поле и переиспользовать его откуда угодно.
Настройки живут на канале, а не на соединении. Подтверждения публикации и ограничение выдачи (basic.qos) включаются для конкретного канала. Открыли новый — включайте заново; забытое включение выглядит как «подтверждения не приходят», хотя код их запрашивает.
Падение канала — не падение соединения. Ошибка протокола (публикация в несуществующий exchange, повторное подтверждение одного и того же сообщения) закрывает канал, а соединение продолжает жить. Библиотека обычно открывает новый канал сама, но исключение вы увидите, и текст будет про закрытый канал — причину надо искать в предыдущей операции, а не в текущей.
И одна практическая деталь: число каналов на соединение ограничено (channel_max, в новых версиях брокера 2047 по умолчанию). Утечка каналов — когда канал открывают и не закрывают — типичная причина того, что приложение перестаёт публиковать через сутки работы.
Как сообщение попадает от отправителя к получателю
В обычной почте вы пишете адрес на конверте, а сортировочный центр решает, на какой склад его везти. В AMQP устроено похоже:
- Producer (продьюсер) — отправляет сообщение.
- Exchange — «сортировочный центр». Получает сообщение и решает, в какие очереди его положить.
- Queue (очередь) — «склад». Хранит сообщения до тех пор, пока потребитель не заберёт их.
- Consumer (потребитель) — забирает сообщения из очереди и обрабатывает.
- Binding — правило, по которому exchange знает, в какую очередь направить сообщение.
- Routing key — «адрес» на конверте, метка-строка, которую продьюсер ставит при отправке.
Ключевая идея: продьюсер не выбирает очередь — он указывает exchange и routing key. Куда именно попадёт сообщение, решают bindings.
Producer не знает, кому пишет: он публикует в обменник, а связи решают, в какие очереди сообщение попадёт. Появился новый потребитель — добавили связь, отправителя трогать не пришлось.
Четыре типа exchange
Direct exchange
Алерты о неудачной оплате должны идти в очередь алертов и в аудит, а события о создании заказа — в сборку: маршрут определяет точное имя события. Так работает direct exchange: сообщение попадает в очередь, если её binding key точно совпадает с routing key.
Публикуем в exchange orders сообщение с routing key payment-failed:
| Binding key очереди | Очередь | Придёт ли сообщение |
|---|---|---|
payment-failed | alerts | да — ключ совпал |
payment-failed | audit-log | да — ключ совпал |
order-created | fulfillment | нет — ключ другой |
Частный случай — exchange с именем "" (пустая строка, default direct). Если указать routing key = имя очереди, сообщение уйдёт прямо туда. Это самый простой способ отправить что-то в конкретную очередь.
Topic exchange
Routing key здесь — строка из слов через точку: order.created.eu, metric.cpu.high. В bindings можно использовать маски:
*— ровно одно слово#— ноль или больше слов
Публикуем в exchange events сообщение с routing key order.created.eu:
| Маска в binding | Очередь | Подходит ли |
|---|---|---|
order.created.* | fulfillment-orders | да — * закрывает eu |
order.# | audit-all-orders | да — # закрывает весь остаток ключа |
*.created.eu | eu-monitoring | да — * закрывает order |
metric.cpu.* | alerts | нет — первое слово другое |
Topic exchange подходит, когда потребители хотят подписаться не на всё, а на определённое подмножество событий.
Fanout exchange
Routing key не учитывается совсем. Сообщение уходит во все очереди, привязанные к этому exchange.
Публикуем в exchange broadcast, routing key не указываем:
Очередь, привязанная к broadcast | Что получит |
|---|---|
service-a-cache | копию сообщения |
service-b-cache | копию сообщения |
service-c-cache | копию сообщения |
Классическая задача — инвалидация кеша или рассылка системного события всем сервисам одновременно.
Headers exchange
На практике headers exchange нужен редко — обычно ту же задачу решают через topic exchange со структурированным routing key. Но иногда признаков для маршрутизации несколько и они независимы (формат файла, приоритет), и втискивать их в одну строку ключа неудобно. Тогда маршрутизируют по заголовкам сообщения, а не по routing key: в binding указывают набор пар ключ: значение и условие совпадения: x-match: all (все должны совпасть) или x-match: any (хотя бы один).
Публикуем в exchange files сообщение с заголовками format: pdf, priority: high:
| Условие binding | Очередь | Подходит ли |
|---|---|---|
format: pdf, priority: high, x-match: all | vip-pdf | да — совпали оба заголовка |
format: pdf, x-match: any | pdf-all | да — хватает одного совпадения |
priority: high, x-match: any | high-prio | да — хватает одного совпадения |
Что такое очередь и как она живёт
Очередь хранит сообщения до тех пор, пока потребитель их не заберёт. При объявлении отвечают на три вопроса.
Должна ли очередь пережить перезапуск брокера? Для бизнес-событий — да, поэтому её объявляют durable. Выбор тут скорее формальный: недолговечные очереди RabbitMQ давно помечает как отмирающую возможность и советует объявлять durable все.
Нужна ли очередь кому-то, кроме одного соединения? Очередь для ответов RPC создаёт сам клиент, и после его отключения она никому не нужна — такую объявляют exclusive: она доступна только одному соединению и удаляется при его разрыве.
Что делать с очередью, когда от неё отключился последний потребитель? Если это временная очередь подписчика — например, для живых уведомлений в одну вкладку браузера, — пусть исчезнет сама: auto-delete. Для рабочих очередей это опасно: перезапустили все экземпляры сервиса — очередь удалилась вместе с накопленными сообщениями.
Каждое сообщение тоже бывает persistent или transient. Разница не в том, где сообщение лежит: современный RabbitMQ пишет на диск и то, и другое. Разница в том, что переживёт перезапуск брокера: persistent он поднимет с диска, transient выбросит.
Стандарт для важных данных: durable очередь + persistent сообщение.
Binding: кто и когда привязывает очередь
Binding создаётся потребителем при запуске — он объявляет, какую очередь создать и к какому exchange привязать:
channel.queue_declare("payment-failed-alerts", durable=True)
channel.queue_bind(
queue="payment-failed-alerts",
exchange="orders",
routing_key="payment-failed"
)
channel.basic_consume(queue="payment-failed-alerts", on_message_callback=handler)
Virtual host: несколько окружений на одном брокере
На одном брокере хотят держать и prod, и staging, и очереди соседней команды — так, чтобы чужой сервис не мог подписаться на вашу очередь. Для этого есть virtual host (vhost): пространство имён со своими exchange, очередями и правами доступа, изолированное от других vhost. Обычно их несколько: / (по умолчанию), /prod, /staging. Это удобный способ разделить окружения без запуска отдельных серверов.
Ack, nack, reject — как потребитель подтверждает обработку
Когда брокер отдаёт сообщение потребителю, оно ещё не удаляется из очереди — только «резервируется». Оно будет удалено только после подтверждения.
У потребителя три варианта:
basic.ack— обработано успешно, удалить из очереди.basic.nack(requeue=true)— не обработано, вернуть в очередь. Брокер кладёт сообщение обратно на прежнее место, поэтому оно прилетает почти сразу: если обработка падает всегда, получается горячий круг из одного и того же сообщения.basic.nack(requeue=false)илиbasic.reject(requeue=false)— не обработано, удалить (или отправить в Dead Letter Exchange, если настроен).
Три исхода одного выданного сообщения: посмотрите, что каждая ветка делает с копией, которая до подтверждения остаётся в очереди.
Если потребитель упал, не отправив ack, брокер видит обрыв соединения и автоматически возвращает сообщение в очередь. Это гарантия at-least-once: каждое сообщение будет обработано хотя бы один раз, даже если потребитель падал в процессе.
Обрыв — не единственный случай. Потребитель может держать соединение и молчать: взял сообщение и завис на нём. Для этого у брокера есть таймер consumer_timeout (по умолчанию 30 минут): если ack не пришёл за это время, брокер закрывает канал с ошибкой, а неподтверждённые сообщения возвращает в очередь. Именно отсюда берётся частая жалоба «канал закрылся сам».
Автоподтверждение
Если использовать basic.consume(autoAck=true), брокер сразу считает сообщение доставленным — без ожидания ack. Если потребитель упадёт до конца обработки, сообщение будет потеряно. Подходит только для метрик и логов, где потеря допустима.
Dead Letter Exchange: куда уходят отвергнутые сообщения
Выше DLX упоминался дважды в скобках, и настало время сказать, что это и почему без него теряют данные.
Сообщение уходит из очереди не только после успешной обработки. Ещё оно уходит, когда потребитель отверг его без возврата (nack с requeue=false), когда истёк срок жизни (x-message-ttl на очереди или на самом сообщении) и когда очередь переполнилась и вытесняет старое (x-max-length). Во всех трёх случаях по умолчанию сообщение просто исчезает: ни тела, ни следа.
Dead Letter Exchange — это exchange, в который брокер отправляет такие сообщения вместо удаления. Настраивается он как свойство очереди:
x-dead-letter-exchange: orders.dlx
x-dead-letter-routing-key: orders.failed
Дальше это обычный exchange: к нему привязывают очередь недоставленных, и оттуда сообщения читают, разбирают и, устранив причину, отправляют обратно.
Что стоит знать, чтобы потом не удивляться:
- Ключ маршрутизации сохраняется, если не задан
x-dead-letter-routing-key. Это удобно, когда DLX повторяет структуру основного обменника, и мешает, когда все отказы нужно собрать в одну очередь — тогда ключ задают явно. - Причина записана в заголовках. Брокер добавляет
x-death— массив с причиной (rejected,expired,maxlen), именем очереди, временем и числом попаданий. По нему разбирают очередь недоставленных, а число попаданий используют, чтобы сообщение не ходило по кругу вечно. - Текста ошибки там нет. Брокер не знает, какое исключение случилось у потребителя, он видит только отказ. Стектрейс в заголовки добавляет клиентская библиотека, если её попросить: в Spring AMQP это
RepublishMessageRecoverer, в amqp091-go, amqplib и pika заголовок с текстом ошибки ставят руками при переотправке в очередь мёртвых писем. - DLX — это не повтор. Сообщение в очереди недоставленных лежит и ждёт человека. Автоматический повтор с задержкой строят отдельно, через очередь с TTL и обратным DLX, и разбирает это статья про шаблоны обмена.
Практическое правило: очередь без DLX означает «отказ равен потере». Для очередей с бизнес-событиями DLX настраивают сразу, вместе с самой очередью, а не после первого инцидента.
Prefetch — сколько сообщений отдавать за раз
Если потребитель подтверждает сообщения автоматически или лимит на предварительную выдачу не задан, брокер отправляет ему столько сообщений, сколько может. Если потребитель медленный — все сообщения скопятся у него в буфере, а другие потребители будут простаивать без работы.
basic.qos ограничивает количество неподтверждённых сообщений у одного потребителя:
basic.qos(prefetch_count=10)
После этого брокер не пришлёт 11-е сообщение, пока хотя бы одно из первых десяти не получит ack.
Практические ориентиры: 1–10 для медленных обработчиков (секунды на сообщение), 100–1000 для быстрых (миллисекунды). Без явной настройки поведение непредсказуемо.
У basic.qos есть второй, редко замечаемый параметр — флаг global. Без него (значение по умолчанию) лимит считается на каждого потребителя: если на одном канале подписаны два потребителя с prefetch_count=10, брокер может держать у них двадцать неподтверждённых сообщений. С global=true те же десять делятся на весь канал целиком. Обычно нужен первый вариант, а знать про второй стоит потому, что при нескольких потребителях на канале фактический предел оказывается выше ожидаемого — и потребитель, который «должен был брать по десять», забирает всю очередь.
Отсюда и связь с числом потоков обработчика, которую легко упустить. Лимит выдачи должен быть не меньше числа потоков, иначе часть их будет простаивать: пять потоков при prefetch_count=1 означают, что четыре ждут, пока первый подтвердит. Разумная отправная точка — число потоков плюс небольшой запас на время обработки. И наоборот, большой лимит на многопоточном потребителе означает, что при его падении в повторную доставку уйдёт много сообщений сразу.
Publisher confirms — гарантия на стороне отправителя
По умолчанию basic.publish возвращает управление немедленно, не дожидаясь, пока сообщение окажется в очереди. Если брокер упадёт в этот момент — сообщение пропадёт, а продьюсер об этом не узнает.
Publisher confirms — расширение, при котором брокер отправляет basic.ack только после того, как принял сообщение. Что значит «принял», зависит от сообщения: persistent-сообщение в durable-очереди подтверждается после записи на диск, transient — сразу после попадания в очередь, диска брокер ждать не будет.
channel.confirm_select()
channel.basic_publish(exchange="orders", routing_key="order.created", body=...)
# ждём ack от брокера перед продолжением
Для важных бизнес-событий publisher confirms нужно включать всегда. Для метрик и логов — по желанию.
Сообщение, которое никуда не подошло: mandatory и alternate exchange
Обиднее всего первая ошибка такого вида: опубликовали, брокер не возразил, а в очереди пусто. Причина почти всегда одна — ключ маршрутизации не совпал ни с одной привязкой, и обменник отправил сообщение в никуда. Это штатное поведение протокола, а не сбой: обменник не обязан знать, нужны ли кому-то его сообщения.
Заметить это можно двумя способами.
Флаг mandatory при публикации. С ним брокер, не найдя ни одной подходящей очереди, вернёт сообщение отправителю служебным кадром basic.return, и клиент вызовет обработчик возврата. В Spring AMQP это setMandatory(true) вместе с setReturnsCallback(...), в amqp091-go флаг mandatory у PublishWithContext и канал из NotifyReturn, в amqplib опция mandatory: true и событие return у канала, в pika mandatory=True и add_on_return_callback. Без флага возврата не будет, и ошибки тоже не будет.
Alternate exchange. У обменника есть свойство alternate-exchange: всё, что не подошло ни к одной привязке, уходит в указанный обменник. К нему привязывают очередь «непонятные сообщения» и смотрят её так же, как очередь недоставленных. Это ловушка на стороне брокера, и работает она независимо от того, что настроил отправитель.
Важно не путать оба механизма с подтверждениями публикации. Подтверждение означает «брокер принял сообщение» — и он действительно его принял, а потом выбросил, потому что доставлять было некуда. Поэтому подтверждения и возвраты включают вместе: первое отвечает на вопрос «дошло ли до брокера», второе — «нашлась ли очередь».
И вопрос, который стоит решить заранее: что делать с вернувшимся сообщением. Записать в лог и забыть — значит потерять его так же, только с записью. Обычное решение — сохранить в свою таблицу неотправленных вместе с ключом маршрутизации, потому что почти всегда это ошибка конфигурации: очередь не объявлена, привязка не создана, опечатка в ключе.
TTL и ограничения на размер очереди
При объявлении очереди можно задать дополнительные параметры:
x-message-ttl— время жизни сообщения в очереди. Просроченное сообщение удаляется или уходит в Dead Letter Exchange.x-expires— через сколько миллисекунд неиспользуемая очередь удалится сама. Удобно для временных очередей.x-max-length/x-max-length-bytes— максимальное число сообщений или байт. Что случится при превышении, задаёт аргументx-overflow:drop-head(по умолчанию) выбрасывает самые старые сообщения из головы очереди,reject-publishотказывает отправителю,reject-publish-dlx— то же, но отказанное уходит в Dead Letter Exchange.
Практические случаи: не дать очереди вырасти до гигабайт, пока потребитель не работает; автоматически отбрасывать устаревшие команды (например, «обновить кеш» старше 30 секунд уже неактуален).
Что выбрать для каждого случая
| Задача | Тип exchange |
|---|---|
| Одно сообщение — один обработчик из пула | "" (default direct), routing key = имя очереди |
| Событие → несколько конкретных очередей | direct |
Иерархические события (order.*.eu, metric.cpu.#) | topic |
| Один сигнал — все слышат | fanout |
| Маршрутизация по нескольким атрибутам | headers (редко) |
Глубже: приоритеты, отложенная доставка и сроки сообщенийрасширенное
Очередь по умолчанию честная: первым пришёл, первым вышел. Три задачи требуют иного, и у каждой в RabbitMQ своя механика с оговорками.
Приоритеты. Классическая очередь с аргументом x-max-priority (обычно до 10, хотя допустимо 255) отдаёт сообщения с бо́льшим приоритетом раньше; приоритет ставит отправитель в свойстве priority. Надёжные очереди с версии 4.0 знают два уровня: приоритет до 4 обычный, от 5 высокий, и этого хватает для «срочные уведомления вперёд». Оговорки: приоритет работает только среди сообщений, которые лежат в очереди; уже отданное потребителю по prefetch не вернуть, поэтому у приоритетных очередей prefetch держат маленьким; и низкий приоритет при постоянном потоке высокого не дождётся очереди никогда, отдельная очередь для «медленных» честнее.
Отложенная доставка. «Отправить напоминание через час» в самом AMQP не предусмотрено. Штатный обходной путь это очередь ожидания с x-message-ttl и Dead Letter Exchange: сообщение лежит час, протухает и переезжает в рабочую очередь, о чём раздел про повторы в статье про паттерны. Работает для фиксированных задержек, по очереди на каждую. Произвольная задержка на сообщение это плагин rabbitmq_delayed_message_exchange: exchange типа x-delayed-message, заголовок x-delay в миллисекундах, и сообщение появляется в очереди через указанное время. Плагин хранит отложенные сообщения на одном узле, не реплицирует их и не рассчитан на миллионы ожидающих; для напоминаний это нормально, для «отложить платёж на сутки» нет, там берут таблицу заданий в базе с временем запуска.
Сроки. TTL бывает на очередь (x-message-ttl, все сообщения в ней живут не дольше) и на сообщение (свойство expiration от отправителя). У второго есть неочевидное свойство: срок проверяется только у головы очереди. Сообщение с коротким сроком, стоящее за сообщением с длинным, не протухнет, пока то не уйдёт; поэтому «задержка через TTL на сообщение» в одной очереди работает неправильно, и очередь ожидания делают со сроком на очередь, а не на сообщения. TTL на очередь целиком (x-expires) удаляет саму очередь после простоя, и это для временных очередей ответов. Дедлайн бизнес-операции («не позднее 18:00») ни один из TTL не выражает: его кладут в тело или заголовок, и потребитель сам решает, что делать с опоздавшим.
Kafka всего этого не умеет по устройству: журнал не переупорядочивают, и приоритеты там это разные топики, задержка это отдельный топик с потребителем, который ждёт, а срок это проверка времени записи потребителем.
Коротко
- Producer публикует в exchange, не в очередь напрямую. Exchange раскладывает по очередям по правилам bindings.
- Direct — точное совпадение routing key. Topic — маски
*и#. Fanout — всем без разбора. Headers — по заголовкам сообщения. - durable + persistent — стандарт для данных, которые нельзя терять. TTL и max-length защищают от бесконтрольного роста очереди.
- Ack подтверждает успешную обработку. Без ack — сообщение вернётся в очередь при обрыве соединения. Это at-least-once.
- Prefetch ограничивает, сколько неподтверждённых сообщений у потребителя. Без него — нагрузка распределится неравномерно.
- Publisher confirms дают гарантию на стороне отправителя: брокер подтверждает сохранение перед тем, как вернуть управление.
- Приоритеты через
x-max-priorityу классических и два уровня у надёжных очередей, отложенная доставка через очередь ожидания с TTL и DLX или плагин отложенных сообщений, срок на сообщение проверяется только в голове очереди. - Соединение одно на приложение, а команды идут через каналы: канал на поток,
basic.qosи подтверждения включаются на канале, а утечка каналов упирается вchannel_max. - Очередь без DLX означает «отказ равен потере»: настраивается свойством очереди, причина попадания лежит в заголовке
x-death, а повтор с задержкой — отдельная конструкция. - Сообщение с ключом, который ни с чем не совпал, исчезает молча: ловят его флагом
mandatoryс обработчиком возврата илиalternate-exchangeна стороне брокера.
Что почитать дальше
- RabbitMQ в production — кластеризация, Quorum Queues, мониторинг.
- Практический код: Spring AMQP для Java, паттерны на Go, Node и Python.
- Messaging-паттерны через AMQP — work queue, RPC, pub/sub на практике.
- AMQP vs Kafka — когда очереди, когда лог-модель.