← назад к разделу

Один узел RabbitMQ на ноутбуке ведёт себя идеально. В проде ломается другое. Сервер перезагрузили ночью — куда делись сообщения, которые ждали обработки? Один заказ падает на разборе и крутится по кругу — почему за ним встала вся очередь? Обработчик молчит сорок минут — почему брокер сам закрыл ему канал? Пора обновлять версию — как не остановить приём заказов?

Главное решение здесь одно, остальные из него вытекают: очередь лежит на одном узле или реплицируется по нескольким.

записали сообщение в очередь orders — кластер из 3 узлов узел 2 упал Classic без реплик узел 1 узел 2 узел 3 нет копии orders нет копии подтверждение пришло от одного узла очередь недоступна до возврата узла 2 Quorum 3 реплики узел 1 узел 2 узел 3 реплика лидер реплика лидер подтверждение — когда записали 2 из 3 реплик лидером стал узел 1, сообщения целы цена надёжности: подтверждение ждёт записи на диск на 2 узлах из 3,у Classic — на одном, а без persistent диска не ждут вовсе

Classic живёт на одном узле: упал узел — упала очередь, и ждать придётся его возврата. Quorum подтверждает запись только когда её приняло большинство, поэтому переживает потерю одного узла из трёх — но каждое подтверждение ждёт диска на двух узлах вместо одного.

Обязательно

Зачем нужен кластер и почему узлов три

Пока узел один, его перезагрузка останавливает всё. Кластер убирает единственную точку отказа: несколько узлов работают вместе и знают друг о друге. Метаданные — описания обменников, очередей, привязок, пользователей — есть на каждом узле, их синхронизирует встроенное хранилище. Долгие годы им была база Mnesia из Erlang с известными проблемами восстановления после раскола сети. На смену пришло хранилище Khepri на алгоритме согласования Raft — том же, на котором работают надёжные очереди: в 4.0 поддерживаемый вариант при Mnesia по умолчанию, в 4.2 выбор по умолчанию, в 4.3 единственный. Клиент подключается к любому узлу; если очередь живёт на другом, узел сам перенаправит запрос.

Почему узлов минимум три. Когда узел перестаёт отвечать, оставшиеся не могут отличить два случая: сосед упал — или порвалась сеть, а сосед работает и принимает записи. Если каждая половина решит, что живая она, после восстановления сети две версии одной очереди придётся склеивать руками. Поэтому в типовой настройке (cluster_partition_handling = pause_minority) узел, оказавшийся в меньшинстве, останавливается целиком — ни записи, ни чтения. Из двух узлов большинство не набрать никогда: при любом разрыве остановятся оба. Три узла переживают потерю одного, пять — двух. Другие стратегии — autoheal, ignore, pause_if_all_down — тоже не бесплатны: либо простой меньшинства, либо ручное слияние после раскола.

Что происходит в секунды после падения

Если процесс узла умер, TCP-соединения между узлами закрываются, и остальные узнают об этом сразу. Если узел завис или к нему пропала сеть, соединение остаётся открытым, и его смерть обнаруживает таймер net_ticktime — по умолчанию около минуты. Всё это время очереди, чьи лидеры были на пропавшем узле, не подтверждают публикации. Потом оставшиеся реплики надёжной очереди выбирают нового лидера. Отправитель с включёнными подтверждениями (publisher confirms) видит это как паузу: basic.publish ушёл, а confirm пришёл через несколько секунд, уже от нового лидера. Отправитель без подтверждений не видит ничего — и поэтому без них нельзя утверждать, что при падении «ничего не потеряли»: спросить не у кого.

Проверять это надо учением. На стенде под нагрузкой остановите один узел (rabbitmqctl stop_app) и смотрите на три вещи: насколько выросла задержка подтверждений; переподключились ли клиенты к живым узлам — клиент, знающий адрес только упавшего, будет стучаться в него бесконечно, нужен список всех узлов или балансировщик; сколько накопилось в messages_ready за время выборов. Потом rabbitmqctl start_app — и смотрите, как вернувшаяся реплика догоняет журнал.

Какой тип очереди под какие данные

Выбор типа очереди — ответ на три вопроса о данных: можно ли их потерять, нужно ли перечитывать историю, какой поток ожидается.

Classic — для того, что можно потерять

Классическая очередь живёт на одном узле, копий нет. Узел упал — очередь недоступна до его возврата; диск не пережил падение — сообщения пропали. Раньше её «зеркалировали» политикой ha-mode; механизм объявили устаревшим в 3.9 и удалили в 4.0 — старые политики в новом кластере ничего не делают, и очередь, которую считали зеркалированной, оказывается одиночной.

Взамен она дешевле: нет репликации по сети, нет ожидания записи на нескольких дисках. Ей отдают то, что при потере не жалко: ответы RPC, сброс кэша, метрики, временные очереди под конкретное соединение. Заказы, платежи и события, по которым сверяют деньги, в классическую очередь не кладут.

Quorum — по умолчанию для всего остального

Надёжная очередь реплицируется по нескольким узлам через Raft: один узел — лидер, остальные — последователи. Отправитель получает подтверждение только когда большинство реплик записали сообщение на диск. Три узла переживают потерю одного без остановки очереди; после потери двух очередь ждёт возврата хотя бы одного.

@Bean
Queue orders() {
    return QueueBuilder.durable("orders").quorum().build();
}

В amqp091-go тип очереди задают аргументом объявления:

_, err := ch.QueueDeclare("orders", true, false, false, false, amqp.Table{"x-queue-type": "quorum"})

В amqplib то же через arguments:

await ch.assertQueue('orders', { durable: true, arguments: { 'x-queue-type': 'quorum' } });

В pika и aio-pika тоже через arguments:

channel.queue_declare("orders", durable=True, arguments={"x-queue-type": "quorum"})

Тип по умолчанию задают на уровне виртуального хоста: rabbitmqctl add_vhost shop --default-queue-type quorum — очередь, объявленная без аргумента, станет надёжной.

Цена — задержка на каждое подтверждение (записи на диск ждут на двух узлах из трёх) и втрое больше сетевого и дискового трафика. Ограничения: нет эксклюзивных и автоудаляемых очередей; приоритеты появились только в 4.0 и устроены проще, чем у классических, — два уровня, обычный и высокий; строгие 32 уровня добавили в 4.3.

Streams — для истории и нескольких читателей

Поток — не очередь, а журнал: сообщения дописываются в конец и не удаляются после чтения, каждый читатель держит своё смещение и может начать с любого места. Он нужен, когда новому сервису надо прочитать события за прошлую неделю или когда одно событие читают несколько независимых потребителей.

@Bean
Queue eventsLog() {
    return QueueBuilder.durable("events-log").stream()
        .withArgument("x-max-length-bytes", 10_000_000_000L)
        .build();
}
_, err := ch.QueueDeclare("events-log", true, false, false, false, amqp.Table{
	"x-queue-type":       "stream",
	"x-max-length-bytes": int64(10_000_000_000),
})
await ch.assertQueue('events-log', {
  durable: true,
  arguments: { 'x-queue-type': 'stream', 'x-max-length-bytes': 10_000_000_000 },
});
channel.queue_declare(
    "events-log",
    durable=True,
    arguments={"x-queue-type": "stream", "x-max-length-bytes": 10_000_000_000},
)

Раз сообщения не удаляются, размер ограничивают явно: x-max-length-bytes в 10 ГБ — около пяти миллионов событий по 2 КБ, старые сегменты брокер удаляет сам. Через отдельный stream-протокол поток даёт порядок миллиона мелких сообщений в секунду на кластер из трёх узлов; через обычный AMQP — заметно меньше.

Не подходит поток там, где сообщение — задача для одного исполнителя: снять его после обработки нельзя, обменника недоставленных у него нет, срока жизни отдельного сообщения тоже. Если задача чисто журнальная, а инфраструктуры RabbitMQ ещё нет, честнее взять Kafka — сравнение в AMQP против Kafka.

Память: где лежат тела сообщений

Очередь на пять миллионов сообщений — что будет с памятью узла? В версиях 4.x надёжная очередь тела сообщений в памяти не держит: там живёт только указатель — минимум 32 байта на сообщение, около мегабайта на 30 тысяч штук независимо от их размера; тело читается с диска по запросу потребителя. Классические очереди с 3.12 устроены так же: режим lazy, который когда-то включали ради этого, убрали, аргументы x-max-in-memory-length и x-max-in-memory-bytes ничего не делают.

Значит, память узла от длины очереди почти не зависит, а диск — напрямую. Бережёт его ограничение длины: x-max-length или x-max-length-bytes вместе с x-overflow. По умолчанию при переполнении брокер выбрасывает самые старые сообщения из головы (drop-head) — для метрик нормально, для заказов нет; заказам ставят reject-publish, и отправитель получает отказ вместо тихой потери. Лимит в 1 ГБ при сообщениях по 2 КБ — это полмиллиона заказов, около суток накопления при пяти заказах в секунду: столько времени есть, чтобы починить потребителей.

Отравленное сообщение и очередь недоставленных

Обработчик упал на сообщении с битым полем и вернул его в очередь. Оно приходит снова, обработчик снова падает. У классической очереди круг не кончается никогда, а с одним потребителем и prefetch в единицу за этим сообщением встаёт вся очередь. Это отравленное сообщение.

С RabbitMQ 4.0 у надёжных очередей ограничитель включён по умолчанию: x-delivery-limit равен 20. Брокер считает доставки в заголовке x-delivery-count и после двадцатой снимает сообщение сам. Дальше два исхода: если у очереди настроен обменник недоставленных (Dead Letter Exchange, DLX), сообщение уходит туда; если нет — брокер молча его выбрасывает. Второй исход — потеря данных без единой строки в логе приложения, поэтому после обновления с 3.13 первое, что проверяют: у всех надёжных очередей с бизнес-данными есть DLX.

Настраивают это политикой, а не в коде каждого сервиса, — тогда и новая очередь под шаблон получит защиту:

rabbitmqctl set_policy orders-dlx "^orders" \
  '{"delivery-limit": 5, "dead-letter-exchange": "orders.dlx"}' \
  --apply-to quorum_queues

Пять доставок вместо двадцати — потому что двадцать бесполезных прогонов по три секунды это минута, на которую сообщение занимает потребителя. Сборка DLX, очереди orders.dlq и повторов с паузой — в паттернах и в коде для Java, Go, Node и Python.

На очередь недоставленных ставят одну проверку: messages_ready в orders.dlq больше нуля дольше пяти минут — сигнал дежурному. Заголовок x-death в каждом сообщении говорит, почему оно здесь: rejected — обработчик отказался, delivery_limit — исчерпаны доставки, expired — вышел срок, maxlen — очередь переполнилась. delivery_limit почти всегда означает ошибку в коде, expired — что потребители не успевают.

Без способа вернуть сообщения в работу после починки кода очередью недоставленных не пользуются: все смотрят, никто не разгребает. Штатный способ — временный Shovel из orders.dlq обратно в исходный обменник:

rabbitmqctl set_parameter shovel dlq-replay \
  '{"src-uri": "amqp://", "src-queue": "orders.dlq",
    "dest-uri": "amqp://", "dest-exchange": "orders",
    "dest-exchange-key": "order.created",
    "src-delete-after": "queue-length"}'

src-delete-after: queue-length — перекачать столько, сколько лежит сейчас, и удалиться. Без этого насос останется навсегда и будет возвращать в работу всё, что попадает в DLQ, — включая то же отравленное сообщение по кругу.

Потребитель, который молчит: consumer_timeout

Обработчик берёт сообщение и работает над ним сорок минут — строит отчёт или перекодирует видео. Через полчаса брокер закрывает ему канал с ошибкой PRECONDITION_FAILED - delivery acknowledgement on channel 1 timed out, сообщение возвращается в очередь и уходит другому потребителю, который начинает ту же работу заново; первый доделывает свою и пытается подтвердить — а канала уже нет. В логах это «потребитель отвалился сам», а в очереди недоставленных потом появляется сообщение с причиной delivery_limit: каждый круг засчитан как доставка.

Таймер consumer_timeout, по умолчанию 30 минут, защищает от обработчика, который взял сообщение и завис. Если работа честно длится дольше, его поднимают — глобально в rabbitmq.conf (consumer_timeout = 7200000, в миллисекундах — два часа) или для одной очереди аргументом x-consumer-timeout. Но сначала стоит спросить, почему одно сообщение требует сорока минут: обычно работу режут на шаги и подтверждают каждый.

Обновление кластера без остановки

Вышла новая версия, а останавливать приём заказов нельзя. Обновление по узлам (rolling upgrade) — единственный штатный способ: узлы останавливают и обновляют по одному, остальные держат нагрузку, а кластер несколько минут живёт в смешанном режиме версий.

Подготовка занимает больше времени, чем само обновление. Первое — прочитать заметки к выпуску всех версий между текущей и целевой: там бывают обязательные промежуточные шаги, и перескакивать через них нельзя — с 3.13 сначала поднимаются до 4.2 и только потом на 4.3. Второе — требования к Erlang: новая версия RabbitMQ может требовать новый Erlang, а это второй перезапуск на том же узле. Третье — включить все стабильные переключатели возможностей: rabbitmqctl enable_feature_flag all.

Переключатели (feature flags) и делают смешанный режим возможным. Новое поведение — другой формат хранения, новый протокол между узлами — включается не установкой версии, а явным переключением, один раз и навсегда. Пока переключатель выключен, обновлённый узел ведёт себя по-старому и умеет разговаривать со старыми соседями. Узел новой версии откажется запускаться в кластере, где не включён переключатель, ставший для неё обязательным, — отсюда правило «сначала включить всё, потом обновлять». Когда последний узел обновлён, enable_feature_flag all запускают ещё раз — включить то, что привезла новая версия.

Порядок на каждом узле: убедиться, что в кластере нет тревог по памяти и диску; проверить, что узел можно выключить — rabbitmq-diagnostics check_if_node_is_quorum_critical завершится с ошибкой, если хоть одна надёжная очередь после остановки потеряет большинство, например потому что одна её реплика уже недоступна; перевести узел в режим обслуживания — rabbitmq-upgrade drain передаёт лидерство его очередей другим и закрывает клиентские соединения; остановить, обновить пакет, запустить, rabbitmq-upgrade revive; дождаться, пока реплики на этом узле догонят журналы, — и только тогда идти к следующему. Если начать второй узел, пока первый не догнал, надёжная очередь останется с одной живой репликой из трёх и остановится.

Другой способ — поднять рядом новый кластер и перелить трафик через Federation или Shovel: дороже, зато обновляет сразу и Erlang, и операционную систему.

Регионы: Federation и Shovel

Хочется растянуть кластер на два региона: три узла тут, три там, и потеря региона ничего не сломает. Работать это не будет. Raft подтверждает каждую запись большинством, и при задержке в десятки миллисекунд между регионами каждое подтверждение ждёт полёта туда и обратно, а таймеры обнаружения отказа срабатывают от обычных сетевых колебаний: узлы теряют друг друга, меньшинство останавливается, очереди подвисают на секунды. Точного порога в документации нет, это накопленная практика.

Правильная схема — независимый кластер в каждом регионе и связь между ними на уровне сообщений.

заказы принимаются в регионе 1, копия уходит в регион 2 регион 1 упал: сервисы переключились на регион 2 регион 1 вернулся: связь ожила, отставшие сообщения доехали сервисы кластер A — регион 1 узел 1 узел 2 узел 3 обменник orders очередь orders кластер B — регион 2 узел 1 узел 2 узел 3 обменник orders очередь orders federation link копия с задержкой связь оборвана подтверждение приходит от своего кластера — копию не ждём потерялось не больше, чем связь не успела скопироватьэти сообщения лежат в очередях A и приедут после возврата то, что было в полёте при обрыве, придёт второй разпотребитель узнаёт повтор по id заказа — идемпотентность два независимых кластера: у каждого свой Raft и свои лидеры

Кластеры не знают друг о друге — их связывает обычный AMQP-клиент с подтверждениями, поэтому копия отстаёт на секунды и после обрыва досылается заново. Переключение на регион 2 теряет только то, что связь не успела скопировать, а возврат региона 1 приносит опоздавшие сообщения и повторы: приём обязан узнавать уже обработанное.

Federation связывает обменник с обменником или очередь с очередью в другом кластере. Обменник-к-обменнику копирует публикации: всё, что опубликовано в orders региона 1 и ушло бы в очередь, асинхронно доезжает в orders региона 2. Очередь-к-очереди не копирует, а перемещает сообщения туда, где есть свободные потребители, — это разделение нагрузки между кластерами, а не резервная копия. Shovel — насос: читает из очереди и публикует в обменник, безусловно и в одну сторону; годится для разового переноса и для возврата из очереди недоставленных выше.

Оба — обычные AMQP-клиенты с подтверждениями: связь читает сообщение у источника, публикует у получателя и подтверждает источнику только после подтверждения получателя. Если связь оборвалась между публикацией и подтверждением, после переподключения источник отдаст сообщение снова — получатель увидит его дважды. Значит, потребители в принимающем регионе обязаны быть идемпотентными — узнавать сообщение, которое уже обрабатывали; как — в паттернах.

Второе следствие — про переключение. Когда регион 1 упал целиком, сервисы переключаются на регион 2 и теряют максимум то, что связь не успела скопировать — обычно секунды. Но эти сообщения не пропали: они лежат в очередях региона 1 и приедут, когда он вернётся, — с опозданием, вперемешку с новыми и, возможно, уже обработанные по обходному пути. У копирующей схемы есть и третье: подтверждения не реплицируются, поэтому сообщение, разобранное в регионе 1, в регионе 2 всё ещё лежит в очереди. Либо его обрабатывают в обоих регионах, и идемпотентность обязательна вдвойне, либо потребители резервного региона стоят выключенными до переключения.

Обратное давление: где встаёт публикация

Симптом: публикация, которая занимала миллисекунды, стала занимать секунды, потом пошли таймауты, а брокер жив и потребители работают. Это не сбой, а защита — брокер перестал принимать, потому что ему некуда складывать.

сервис-отправитель basic.publish ушёл ждёт confirm blocked сокет не читается: кадры лежат в TCP-буфере через 5 с — таймаут, наружу HTTP 503 узел RabbitMQ очередь растёт на диске память ≥ watermark диск < disk_free_limit тревога на весь кластер blocked на всех узлах потребители читают и подтверждают их соединения не трогают: очередь разгружается

Тревога поднимается по любому из двух порогов на любом узле и действует на весь кластер. Брокер не отвечает отправителю отказом — он перестаёт читать его сокет, поэтому в сервисе публикация выглядит как зависание, а не как ошибка. Соединения потребителей не трогают: именно им и нужно разгрузить очередь.

Порога два: память узла выше доли, заданной vm_memory_high_watermark.relative, или свободное место на диске ниже disk_free_limit. Дисковый порог по умолчанию — 50 МБ, это меньше минуты записи при тысяче сообщений по килобайту в секунду; на проде его поднимают до гигабайт, чтобы тревога срабатывала раньше, чем кончится место под журналы Raft. Брокер не отвечает отказом — он перестаёт читать сокеты соединений, через которые публикуют: кадры basic.publish остаются в TCP-буфере, подтверждение не приходит, сервис упирается в свой таймаут. Клиенту приходит уведомление connection.blocked; в Spring AMQP это события ConnectionBlockedEvent и ConnectionUnblockedEvent, в amqp091-go канал из NotifyBlocked, в amqplib события blocked и unblocked у соединения, в pika add_on_connection_blocked_callback; по ним приложение переводят в режим «новые задания не принимаем».

Кроме тревог есть более мягкий механизм — кредиты между процессами внутри узла. Отправитель, который публикует быстрее, чем очередь пишет на диск, получает состояние flow: соединение притормаживают, не блокируя. Симптом — пропускная способность упёрлась в потолок, а не встала.

Приложению заблокированное соединение — сигнал «сейчас нельзя», а не ошибка сети: не копить сообщения в памяти без предела, а возвращать вызывающему честный отказ или складывать в буфер с явным лимитом. Причину лечат на другой стороне: добавляют потребителей, ставят очередям лимит длины с reject-publish.

Мониторинг: по каким числам видно, что не справляемся

Плагин rabbitmq_prometheus отдаёт метрики на порту 15692. Смотреть надо не на абсолютные значения, а на соотношения — почти каждая беда видна как нарушение равенства, которое в норме держится.

МетрикаЧто значитПризнак беды
rabbitmq_queue_messages_readyСообщений ждёт обработкиРастёт при том же числе потребителей — они не успевают
rabbitmq_queue_messages_unackedВыдано потребителям без подтвержденияДержится на prefetch × потребители и не падает — все окна забиты
rabbitmq_queue_consumersАктивных потребителейНоль там, где должен быть хотя бы один
rabbitmq_connectionsОткрытых соединенийОбнулились разом — сеть или перезапуск; растут без остановки — утечка соединений
rabbitmq_alarms_memory_used_watermark, rabbitmq_alarms_free_disk_space_watermarkТревога по памяти и дискуЕдиница — отправители уже заблокированы

Абсолютные пороги у каждого свои. Память сравнивают не с размером машины, а с порогом тревоги: rabbitmq_process_resident_memory_bytes выше 80 % от rabbitmq_resident_memory_limit_bytes — до блокировки отправителей осталось немного. Диск — rabbitmq_disk_space_available_bytes — сравнивают со своим disk_free_limit плюс час записи в текущем темпе, а не с круглым числом. И messages_ready в очереди недоставленных больше нуля — уже повод смотреть, что там.

Главное соотношение — рост messages_ready против скорости потребителей, аналог consumer lag в Kafka: очередь, которая растёт третий час, означает, что система не справляется, даже если все узлы зелёные.

Ориентиры по производительности

Цифры пропускной способности из чужих замеров на свой стенд не переносятся: они зависят от размера сообщения, подтверждений, записи на диск, числа очередей и задержки диска.

У надёжной очереди подтверждение ждёт записи на диск на большинстве реплик, поэтому её потолок — задержка fsync и то, сколько сообщений брокер группирует в одну запись. У классической без записи на диск ограничение — сеть и процессор, отсюда разница в разы. Размер сообщения бьёт по всему сразу: тело надёжной очереди летит на три узла и пишется на три диска, поэтому 1–10 КБ — комфортный размер, от 100 КБ пропускная способность заметно проседает, а больше мегабайта кладут в объектное хранилище и передают ссылку. Число очередей — отдельный ресурс: каждая надёжная очередь — свой журнал и своя группа Raft со своими сердцебиениями; тысячи — норма, десятки тысяч — нагрузка на узлы и на выборы лидеров, сотни тысяч — рискованно.

Порядок величин в опубликованных замерах разработчиков — десятки тысяч сообщений в секунду на надёжную очередь и сотни тысяч на классическую без записи на диск, на мелких сообщениях. Свой потолок измеряют штатным rabbitmq-perf-test с боевыми параметрами:

perf-test --uri amqp://localhost --quorum-queue --queue orders-bench \
  --producers 4 --consumers 4 --size 2000 --confirm 100 --flag persistent

Замер без подтверждений и без записи на диск даст красивую цифру, которой в проде не будет.

Дополнительно: при первом чтении можно пропустить

Глубже: отставание два часа: расширить, отбросить, перенаправитьрасширенное

Мониторинг выше показывает, что messages_ready растёт быстрее, чем потребители забирают, и к утру очередь уведомлений отстаёт на два часа. Разгребать это нужно в определённом порядке, и порядок одинаков для RabbitMQ и Kafka.

Сначала остановить рост. Понять, откуда поток: всплеск от бизнеса, повторная отправка после сбоя, цикл повторов из-за отравленного сообщения (тогда лечится не расширением, а DLQ). Если это внешний поток, включают обратное давление на отправителе: x-max-length с overflow: reject-publish на очереди, лимит частоты на API, пауза продьюсера. Наращивать потребление при неограниченном притоке бессмысленно.

Расширить обработку. У RabbitMQ очередь масштабируется числом потребителей напрямую: поднять ещё экземпляры или потоки слушателя (concurrency у контейнера), поднять prefetch, если обработчик быстрый. У Kafka потолок это число партиций: потребителей больше партиций стоят без дела, и это тот случай, когда расчёт партиций из статьи про основы был неверен. Обработчик ускоряют, убирая необязательное: не слать письмо на каждое событие, а собирать пачку; пропускать дубли по отметке идемпотентности до обращения к внешним сервисам.

Отбросить протухшее. Уведомление «ваш заказ собран» через два часа после «заказ доставлен» вредно, а не бесполезно. Потребитель проверяет отметку времени события и молча подтверждает всё старше порога; у RabbitMQ порог можно поставить и снаружи политикой message-ttl на очередь, что протухнет и уйдёт в DLX без обработки; в крайнем случае очередь чистят целиком (purge), если её содержимое можно восстановить из источника правды. У Kafka аналог это сдвиг оффсета группы к концу, kafka-consumer-groups --reset-offsets --to-latest при остановленной группе, для топиков, где старое не нужно.

Перенаправить. Если старое всё же нужно, но свежее важнее: свежий поток пускают в новую очередь (у RabbitMQ перепривязкой, у Kafka новой группой с чтением с конца), её обрабатывают в приоритетном режиме, а старую очередь разгребает отдельный, медленный потребитель в фоне; в RabbitMQ Shovel умеет перелить накопленное в отдельную очередь. Так пользователи получают свежие уведомления сейчас, а история доезжает потом.

После. Каждое такое утро это ошибка размера: сколько партиций, потребителей и запаса заложили; по каким числам звать дежурного. Метрику «возраст самого старого сообщения» ставят на тревогу с порогом, за которым ещё можно расширить, а не отбросить.

Коротко

  • Узлов минимум три: из двух большинство не набрать, при любом разрыве остановятся оба.
  • Тип очереди — ответ про данные: терять нельзя — Quorum; не жалко — Classic; нужна перечитка — Streams. Зеркалирование классических удалено в 4.0.
  • С 4.0 надёжная очередь снимает сообщение после 20 доставок и без DLX молча выбрасывает; DLX ставят политикой до первого случая, вместе со способом вернуть из DLQ в работу.
  • Обработчик дольше 30 минут без подтверждения теряет канал — consumer_timeout.
  • Обновление — по узлам: feature flags включены заранее, перед каждым узлом check_if_node_is_quorum_critical.
  • Кластер не тянут через регионы; Federation и Shovel доставляют «хотя бы раз» — приём идемпотентен.
  • Тревога по памяти или диску блокирует отправителей на всём кластере; приложение обязано уметь отказывать.
  • Беду ищут по соотношениям: рост messages_ready при тех же потребителях, unacked = prefetch × потребители.
  • Отставание разгребают по порядку: остановить приток, расширить обработку (у Kafka не выше числа партиций), отбросить протухшее по отметке времени или TTL, перенаправить свежий поток в отдельную очередь; потом пересчитать размер.

Что почитать дальше

  • Протокол AMQP — подтверждения, prefetch и consumer_timeout изнутри: откуда берётся unacked.
  • Код повторов и DLX, то, что здесь задано политикой: Java, Go, Node, Python.
  • Паттерны через AMQP — идемпотентный потребитель, обязательный для приёма из Federation.
  • AMQP против Kafka — когда журнал нужен настолько, что Streams мало.