MAATRIX / Блог / Предел одной очереди брокера сообщений: когда он перестаёт быть буфером

Предел одной очереди брокера сообщений: когда он перестаёт быть буфером

MAATRIX

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

Глубина очереди — это не то же самое, что нагрузка

Первая путаница, из-за которой инциденты диагностируют неправильно: глубина очереди (backlog, число сообщений, которые лежат и ждут обработки) и скорость поступления сообщений (rate, сообщений в секунду) — два разных числа, и растущая глубина не обязана означать растущую нагрузку.

Формула Литтла даёт строгую связь между ними: среднее число сообщений в системе (L) равно скорости их поступления (λ) умноженной на среднее время, которое сообщение проводит в системе (W). Если скорость поступления постоянна, а глубина очереди растёт — растёт именно W, время ожидания. Каждое новое сообщение в среднем ждёт дольше предыдущего. Это уже не «буфер сглаживания», а очередь, которая копит долг.

Пока очередь остаётся в пределах десятков-сотен тысяч сообщений, для большинства брокеров это просто данные — брокер их хранит и отдаёт по мере готовности консьюмеров. Проблема начинается на других масштабах: когда структура данных внутри брокера (индекс очереди в памяти, метаданные партиции, состояние Raft-лога) перестаёт помещаться туда, где рассчитана работать быстро — в оперативной памяти или в page cache. С этого момента брокер тратит ресурсы не только на приём новых сообщений, но и на обслуживание разросшегося backlog. Базовые роли — продюсер, консьюмер, брокер как посредник — разобраны в статье как устроена очередь сообщений, здесь сосредоточимся на том, что происходит, когда очередь перестаёт быть маленькой.

RabbitMQ: где память и диск упираются первыми

RabbitMQ хранит очередь как процесс: у классической очереди это Erlang-процесс на конкретном узле, у quorum-очереди — реплицированный через Raft лог с лидером и фолловерами. В обоих случаях обработка сообщений в одной очереди принципиально последовательна — параллелизм внутри RabbitMQ достигается количеством очередей и консьюмеров, а не глубиной одной очереди.

Что растёт вместе с глубиной:

  • Индекс и метаданные сообщений в памяти. Даже если тела сообщений уходят на диск, брокеру всё равно нужно держать в RAM структуру, которая знает про каждое неподтверждённое (unacked) сообщение: позицию, флаги доставки, кому оно ушло. На глубоких очередях с большим числом висящих без подтверждения сообщений эта книга учёта сама становится заметным потребителем памяти.
  • Порог памяти узла. RabbitMQ следит за общим потреблением памяти узла и при достижении заданной доли доступной RAM (vm_memory_high_watermark) включает flow control — блокирует публикацию новых сообщений на все очереди узла, пока память не освободится. С точки зрения приложения это выглядит как внезапный стоп на запись, а не плавная деградация.
  • Время переиндексации при рестарте. Чем глубже очередь на момент падения узла или планового рестарта, тем дольше брокер восстанавливает её состояние при старте — для очередей с миллионами сообщений это могут быть минуты, а не секунды простоя.

Проверить, что происходит с конкретной очередью прямо сейчас, можно без графиков — командой rabbitmqctl:

rabbitmqctl list_queues name messages messages_ready \
  messages_unacknowledged memory consumers

Если колонка memory растёт быстрее, чем messages — на очередь давит не столько объём данных, сколько накладные расходы на их учёт (много мелких сообщений, длинные цепочки unacked, приоритетные подочереди — ниже разберём отдельно). Точную цифру «сколько миллионов сообщений — это уже много» дать честно нельзя: она зависит от размера сообщения, числа приоритетов, режима хранения и версии RabbitMQ. Практический ориентир — не общая глубина сама по себе, а тренд потребления памяти узла в rabbitmqctl status: если он идёт к порогу high_watermark быстрее, чем растут сами данные, разбираться нужно раньше, чем сработает flow control.

Нужен сервер под эту задачу?

Разверните VPS MAATRIX за пару минут: NVMe, AMD EPYC, root-доступ, локации UK, США, Франция и РФ. Оплата картой РФ и по СБП.

Арендовать сервер

Kafka: почему топик растёт «бесконечно», но диск и задержка — нет

У Kafka архитектура принципиально другая: партиция — это append-only лог на диске, консьюмер просто хранит свою позицию (offset) и читает последовательно с неё. Из-за этого глубина очереди сама по себе, в отличие от RabbitMQ, почти не увеличивает нагрузку на CPU брокера — чтение и запись остаются последовательными операциями независимо от того, сколько сообщений накопилось.

Но у роста backlog в Kafka есть три других предела, которые никуда не деваются:

  • Диск. Пока консьюмер отстаёт, брокер обязан хранить непрочитанные сегменты — а retention (log.retention.hours, log.retention.bytes) настроен по времени или объёму, а не по факту прочтения. Растущий лаг напрямую превращается в растущее использование диска, и если консьюмер не догоняет неделями, брокер рано или поздно упирается в физическую ёмкость раздела. Реальный случай, когда сегменты не удалялись и диск забился за выходные, разобран в статье Kafka не удаляла старые сегменты и забила диск.
  • Page cache и промах мимо памяти. Пока консьюмер читает данные, недавно попавшие в page cache, чтение практически бесплатно для диска. Как только глубина backlog превышает объём RAM под page cache, догоняющий консьюмер начинает читать старые сегменты, которых в кэше уже нет — последовательное чтение превращается в конкуренцию за реальный дисковый I/O с текущей записью, и задержка чтения резко растёт именно при попытке разгрести backlog. Как считать память под этот сценарий — в статье сколько RAM нужно для Kafka.
  • Число партиций. Каждая партиция — это файловые дескрипторы, поток репликации, запись в метаданных контроллера. Проблему глубокого backlog часто пытаются лечить увеличением числа партиций, а на больших значениях (счёт на тысячи партиций на брокер, точный порог зависит от версии и железа) начинает расти время ребалансировки и восстановления кластера после сбоя — уже независимо от глубины самой очереди.

SQS и managed-очереди: предел не исчезает, он прячется

У managed-брокеров вроде SQS физическая инфраструктура очереди скрыта от вас — Amazon декларирует практически неограниченную глубину стандартной очереди. Это правда лишь отчасти: предел перестаёт быть вопросом «хватит ли диска брокеру», но появляется на уровне механики самого сервиса.

  • Visibility timeout и повторная доставка. Забранное консьюмером сообщение становится невидимым для других на время visibility timeout. Если консьюмер не успел удалить его за это время — сообщение возвращается в очередь видимым снова. На глубокой очереди с медленными или падающими консьюмерами это создаёт эффект снежного кома: одни и те же сообщения обрабатываются повторно, а реальная нагрузка на downstream-сервисы оказывается в разы выше номинальной скорости поступления.
  • Dead-letter queue как обязательный предохранитель. Без настроенной DLQ и maxReceiveCount «ядовитое» сообщение, которое всегда падает с ошибкой, будет возвращаться в очередь до бесконечности — и часть видимого backlog окажется не новой работой, а одним и тем же сообщением, зацикленным по кругу.
  • Возраст самого старого сообщения — лучший индикатор, чем счётчик глубины. Метрика ApproximateAgeOfOldestMessage в CloudWatch показывает не «сколько сообщений», а «насколько давно застряло самое старое» — она первой сигнализирует о структурной проблеме, тогда как ApproximateNumberOfMessagesVisible может расти и от рутинного дневного пика.
  • Стоимость. У managed-очереди предел смещается из инженерной плоскости в финансовую: оплата идёт за количество запросов, и очередь, копящая сообщения неделями, а потом долго вычерпываемая, генерирует кратно больше API-вызовов на poll, чем та же нагрузка без задержки.

Партиционирование и приоритеты: рычаги масштабирования, а не решение проблемы

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

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

  • Горячий ключ партиционирования. Если сообщения распределяются по партициям не равномерно (например, по ID клиента, а у вас один клиент генерирует половину трафика), то одна партиция всё равно накопит глубокий backlog, пока остальные будут пустовать. Увеличение общего числа партиций такую проблему не решает — решает пересмотр ключа шардирования.
  • Общий downstream-бутылочное горлышко. Если все консьюмеры всех партиций пишут в одну и ту же базу данных или дёргают один и тот же внешний API с его собственным лимитом, увеличение числа партиций и консьюмеров упирается не в очередь, а в этот общий ресурс — глубина очереди в таком случае сигнализирует о пределе не брокера, а того, что стоит за консьюмерами.

Приоритеты очередей (несколько уровней приоритета внутри одной очереди RabbitMQ через x-max-priority, либо несколько отдельных очередей с потреблением по весу) дают срочным сообщениям возможность обгонять backlog — но это перераспределение порядка обработки, а не рост пропускной способности. Если суммарная скорость поступления стабильно выше скорости обработки, приоритетная схема не предотвращает рост очереди — она определяет, кто будет расти первым: очередь низкого приоритета начнёт не просто отставать, а голодать полностью, пока не иссякнет поток высокоприоритетных сообщений. Есть и обратная сторона: каждый уровень приоритета в RabbitMQ — отдельная внутренняя подочередь со своими накладными расходами на память, поэтому десятки уровней сами по себе заметно увеличивают потребление RAM узла — практичнее держаться единиц уровней, а не тонкой шкалы важности.

МеханизмЧто решаетЧто не решаетТипичный побочный эффект на масштабе
ПартиционированиеПараллельная скорость потребленияОбщий downstream-лимит, перекос по ключуГорячая партиция копит backlog даже при высоком параллелизме
ПриоритетыПорядок обработки при перегрузкеНедостаточную суммарную мощностьНизкоприоритетная очередь растёт быстрее и голодает
Больше консьюмеровПропускную способность до предела downstreamПределы самого downstream-ресурсаРост ошибок и таймаутов на стороне БД/API вместо ускорения

Когда рост очереди — это шторм, а когда — структурная проблема

Главный диагностический вопрос — не «насколько глубока очередь сейчас», а «стабилизируется ли рост, если ничего не менять». Практически это проверяется не по одной метрике глубины, а по двум отдельным графикам рядом: скорость поступления сообщений и скорость их обработки.

  • Временный всплеск. Скорость поступления кратковременно превысила скорость обработки (маркетинговая рассылка, батч-джоб, ретрай после сбоя внешнего сервиса), но сама скорость обработки осталась прежней и после всплеска превышает текущую скорость поступления — глубина идёт вниз без вмешательства. Ровно тот сценарий, для которого очередь и покупали.
  • Структурная проблема. Скорость обработки уткнулась в потолок и не растёт вместе с числом консьюмеров — потому что упирается не в число воркеров, а в общий ресурс за ними (соединение к БД, rate limit внешнего API, узкое место в коде обработчика). Глубина растёт даже в часы низкой нагрузки, тренд не выходит на плато. Здесь ожидание бессмысленно: сколько ни жди, скорость обработки выше не станет сама по себе. Формальный расчёт, с какой глубины воркеры физически не догонят очередь, и как посчитать их нужное число через формулу Литтла — в статье потолок очереди задач: с какой глубины не догнать.

Отдельно стоит смотреть не только на «нормальные» сообщения, но и на застрявшие. Один зависший или упавший в бесконечный ретрай консьюмер способен годами держать очередь в статусе «растёт», даже если остальная система обрабатывает трафик штатно — на графике это выглядит как рост при стабильной, не аномальной скорости поступления. Отдельная категория структурной проблемы — очередь без TTL и без ограничения максимальной длины: разбор инцидента, где нехватка TTL тихо съела всю память брокера за несколько месяцев, — в статье сообщения без TTL съели память брокера. TTL и x-max-length (или аналог вашего брокера) — не оптимизация, а предохранитель: явное решение, что делать с сообщением, если система физически не успевает его обработать, лучше, чем узнать ответ по факту падения узла.

Нужен сервер под эту задачу?

Разверните VPS MAATRIX за пару минут: NVMe, AMD EPYC, root-доступ, локации UK, США, Франция и РФ. Оплата картой РФ и по СБП.

Арендовать сервер

Нужны сами нейросети для контента?

Генерируйте изображения, видео и озвучку нейросетями на falapi.io — десятки моделей в одном окне. Оплата картой РФ и по СБП.

Частые вопросы

С какой глубины очередь RabbitMQ или Kafka считать «слишком большой»?

Универсального числа нет — предел зависит от размера сообщений, доступной RAM/диска, числа приоритетов и версии брокера. Ориентируйтесь не на абсолютную цифру, а на тренд: растёт ли потребление памяти узла (RabbitMQ) или использование диска относительно retention (Kafka) быстрее, чем можно объяснить объёмом самих данных.

Поможет ли увеличение числа партиций Kafka, если консьюмеры не успевают за растущим лагом?

Только если узкое место — параллелизм чтения, а не общий downstream-ресурс. Если все консьюмеры пишут в одну БД с ограниченным пулом соединений, дополнительные партиции дадут больше параллельных читателей, которые упрутся в тот же самый лимит записи.

Что делать в первую очередь, если очередь выросла из-за одного зависшего консьюмера?

Найти и перезапустить именно его — глубина сама начнёт падать, как только остальные консьюмеры снова получат доступ к сообщениям, которые раньше держал зависший воркер (unacked/in-flight). Массовое добавление новых консьюмеров без диагностики зависшего часто не решает проблему, а только увеличивает конкуренцию за общий downstream-ресурс.

Нужен ли TTL или ограничение максимальной длины на всех очередях без исключения?

Не обязательно на всех, но на любой очереди, где отсутствие обработки сообщения дольше определённого времени делает его бесполезным (устаревшее уведомление, просроченная задача), TTL стоит поставить явно — это дешевле, чем аварийный разбор переполненной очереди постфактум.

Как понять, что причина не в брокере, а в downstream-сервисе за консьюмерами?

Добавьте консьюмеров или партиций и посмотрите на скорость обработки отдельно от глубины очереди. Если она не растёт пропорционально числу воркеров — узкое место не в очереди, а в ресурсе, который эти воркеры делят между собой (соединения к БД, лимиты внешнего API, блокировки в коде).

Обсудить статью, задать вопрос или начать новую тему

Есть вопрос по этой статье, идея для обсуждения или просто хотите поделиться опытом? Сообщество MAATRIX ждёт. Для общения, пожалуйста, зарегистрируйтесь в нашем личном кабинете.

Перейти в сообщество →