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 при первой же правке.