TTL, приоритеты и лимиты очередей

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

TTL (time to live) — срок жизни сообщения или очереди. Просроченное сообщение брокер удаляет сам: либо совсем, либо в DLX, если он настроен. Лимит очереди — потолок по числу сообщений или байтам, при достижении которого брокер решает, кем пожертвовать: старыми сообщениями или новыми.

Зачем ограничивать очередь

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

Есть и вторая причина, менее драматичная, но более частая. Огромная часть сообщений имеет смысл только «сейчас»: код из СМС, пуш «водитель подъезжает», инвалидация кеша. Доставить их через два часа не просто бесполезно — вредно. Такому сообщению нужен срок годности.

TTL: два уровня

Срок жизни задаётся либо на всю очередь, либо на конкретное сообщение.

ЧтоГде задаётсяТипКогда применять
x-message-ttlаргумент очереди или policyчисло, миллисекундыВсе сообщения потока протухают одинаково. Основной рабочий вариант.
expirationсвойство сообщениястрока! с числом мсСрок зависит от самого сообщения: у кода из СМС минута, у отчёта — час.
x-expiresаргумент очередичисло, миллисекундыУдалить саму очередь, если ей столько времени не пользовались.
import pika

conn = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
ch = conn.channel()

ch.queue_declare(
    queue='notifications',
    durable=True,
    arguments={
        'x-message-ttl': 300000,                 # 5 минут: старше уже не нужно
        'x-max-length': 50000,                   # потолок: 50 тыс. сообщений
        'x-overflow': 'reject-publish',          # переполнение → отказ издателю
        'x-max-priority': 5,                     # 5 уровней приоритета
        'x-dead-letter-exchange': 'notif.dlx',   # протухшее не пропадёт бесследно
    },
)

ch.basic_publish(
    exchange='',
    routing_key='notifications',
    body=b'{"user": 42, "text": "Код подтверждения"}',
    properties=pika.BasicProperties(
        delivery_mode=2,      # сообщение переживёт перезапуск брокера
        priority=5,           # срочное: вперёд обычных
        expiration='60000',   # 60 секунд. СТРОКА, не число!
    ),
)

Разбор по строкам:

  • expiration='60000' — да, это строка. Так устроен протокол AMQP, и передача числа вернёт ошибку. Ошибаются на этом все, включая опытных.
  • Если у очереди есть и x-message-ttl, и у сообщения expiration, побеждает меньшее значение.
  • x-dead-letter-exchange рядом с TTL — важная пара. Без DLX просроченное сообщение исчезает бесследно. С DLX оно попадает в карантин, и вы хотя бы узнаете, сколько уведомлений не доехало.

Гоча: TTL проверяется только в голове очереди

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

Приоритеты

Приоритетная очередь пропускает срочное вперёд обычного. Включается аргументом x-max-priority при объявлении, дальше каждое сообщение получает свойство priority.

Здесь важно снять типичное ожидание. Приоритет — не про скорость доставки, а про порядок разбора накопившегося. Если потребители справляются и очередь пуста, сообщение уходит первому свободному потребителю мгновенно — приоритет ни на что не влияет, сортировать нечего. Приоритет начинает работать ровно тогда, когда есть затор.

Второй подвох — prefetch. Если потребитель забрал по basic_qos(prefetch_count=1000) тысячу сообщений в локальный буфер, срочное сообщение приедет в брокер и встанет за этой тысячей: они уже отданы, переставить их брокер не может. Приоритеты и большой prefetch несовместимы — держите prefetch_count в пределах 1–20, если порядок важен.

И ещё три ограничения, которые лучше знать заранее:

  • Разумный максимум — 5 уровней. Технически можно до 255, но каждый уровень — это отдельная внутренняя подочередь, и брокер начинает есть память на ровном месте. Пяти хватает всем.
  • Сообщение без свойства priority считается приоритетом 0 — самым низким.
  • Приоритеты не спасают от голодания: если поток срочных сообщений не иссякает, обычные не будут обработаны никогда. Если это неприемлемо — разводите потоки по разным очередям с разными пулами потребителей. Часто это вообще более честное решение, чем приоритеты.

Лимит длины и режимы переполнения

Потолок задаётся двумя аргументами: x-max-length (число сообщений) и x-max-length-bytes (суммарный объём тел). Работает тот, который сработает раньше. А вот что произойдёт при переполнении, решает x-overflow:

РежимЧто делаетКогда выбирать
drop-head (по умолчанию)Выбрасывает самое старое сообщение из головы очереди. Если настроен DLX — отправляет его туда с причиной maxlen.Поток телеметрии, котировок, координат: свежие данные ценнее старых.
reject-publishОтказывает издателю: с publisher confirms тот получает basic.nack. Очередь не растёт, сообщение не принято.Заказы, платежи: терять нельзя, лучше честно сказать источнику «не могу принять».
reject-publish-dlxТо же, что выше, но отвергнутое сообщение ещё и уходит в DLX.Когда нужен и отказ издателю, и сохранённая копия для разбора.

Режим reject-publish работает по-настоящему только вместе с подтверждениями публикации — иначе издатель просто не узнает об отказе и будет думать, что всё доставлено.

ch.confirm_delivery()   # включаем publisher confirms

try:
    ch.basic_publish(exchange='', routing_key='notifications', body=payload)
except pika.exceptions.NackError:
    # брокер отверг публикацию: очередь упёрлась в x-max-length
    metrics.inc('queue_full')
    raise
except pika.exceptions.UnroutableError:
    # некому доставить: нет подходящей привязки
    log.warning('сообщение некуда маршрутизировать')

Как это работает: аргументы против policy

Всё описанное можно задать двумя способами — аргументами при объявлении очереди (как в примерах) и policy на стороне брокера. Разница принципиальная.

Аргументы вшиваются в очередь навсегда. Захотите поменять TTL с 5 минут на 10 — повторное объявление вернёт PRECONDITION_FAILED - inequivalent arg 'x-message-ttl' (код 406) и разорвёт канал. Придётся выкачивать сообщения, удалять очередь, создавать заново — в продакшене посреди рабочего дня это так себе развлечение.

Policy применяется к очередям по регулярному выражению и меняется на живой системе:

rabbitmqctl set_policy notif-limits '^notifications' \
  '{"message-ttl": 300000, "max-length": 50000, "overflow": "reject-publish"}' \
  --apply-to queues --priority 1

Обратите внимание на имена ключей: в policy они без префикса x-message-ttl, а не x-message-ttl. Практическое правило: то, что может измениться (TTL, лимиты, overflow), задавайте через policy; то, что определяет саму природу очереди (x-max-priority, тип очереди), — аргументами при объявлении, потому что через policy это и не выставить.

Частые ошибки

  • expiration=60000 числом. Ждёт строку. Классика, которая ловится только в рантайме.
  • TTL без DLX. Сообщения тихо испаряются, метрика «доставлено» падает, причина не находится неделями.
  • Попытка использовать TTL как отложенную доставку. «Положу в очередь с TTL, потребитель прочитает через час» — не работает: потребитель заберёт сообщение сразу же. Отложенная доставка — только через связку «очередь-отстойник без потребителей + DLX» из прошлого урока.
  • Приоритеты при большом prefetch. Срочные сообщения встают в хвост уже выданной пачки, и вся конструкция превращается в декорацию.
  • x-max-priority: 255 «на всякий случай». Каждый уровень — отдельная структура в памяти. Берите 3–5.
  • Смена аргументов у живой очереди. Ошибка 406 и оборванный канал. Меняемое — через policy.
  • drop-head для платежей. Режим по умолчанию молча выбрасывает самые старые сообщения. Для денег нужен reject-publish.

Итоги

  • Неограниченная очередь однажды упрётся в память, включит flow control и заблокирует издателей — то есть уронит и синхронную часть приложения.
  • x-message-ttl — на очередь, expiration (строкой!) — на сообщение, x-expires — на саму очередь. Из двух TTL побеждает меньший.
  • TTL проверяется в голове очереди, поэтому просроченные сообщения могут какое-то время числиться в messages_ready.
  • Приоритеты имеют смысл только при заторе и умирают от большого prefetch; часто честнее развести потоки по разным очередям.
  • x-overflow: drop-head — жертвуем старым (телеметрия), reject-publish — отказываем издателю (деньги).
  • Меняемые настройки задавайте через policy, а не аргументами, иначе получите 406 при первой же правке.
Проверьте себя
1. Очередь объявлена с x-max-priority: 5, потребитель работает с prefetch_count=1000. Очередь забита. Почему срочные сообщения всё равно обрабатываются с большой задержкой?
AПриоритеты работают только в quorum queues
BПотребитель уже забрал в свой буфер 1000 сообщений, и брокер не может переставить их местами — срочное встанет за ними
CЗначение priority должно быть строкой, иначе оно игнорируется
DПриоритет учитывается только при включённых publisher confirms
2. Очередь с заказами упёрлась в x-max-length. Какой режим x-overflow выбрать, чтобы не потерять ни одного заказа?
Adrop-head — он используется по умолчанию и удаляет самые старые
Breject-publish — брокер откажет издателю (basic.nack), и тот сможет повторить попытку
Cx-expires — очередь удалится и создастся заново
DНикакой: при переполнении RabbitMQ всегда расширяет очередь автоматически