Dead Letter Queue и повторы

Потребитель уронил сообщение. Что с ним будет дальше — и почему «просто вернуть в очередь» однажды положит вам весь сервис.

Dead Letter Exchange (DLX) — обменник, в который RabbitMQ сам переправляет сообщения, отвергнутые потребителем, просроченные по TTL или вытесненные из переполненной очереди. Очередь, привязанная к этому обменнику, называется DLQ (Dead Letter Queue) — карантин, куда складывают всё, что не удалось обработать.

Зачем: история одного бесконечного цикла

Классическая авария. Сервис оплат читает очередь orders. Кто-то опубликовал заказ, в котором вместо суммы пришла строка "—". Потребитель падает на разборе, ловит исключение и делает то, что кажется безопасным: возвращает сообщение обратно в очередь.

except Exception:
    ch.basic_nack(delivery_tag=method.delivery_tag, requeue=True)  # мина замедленного действия

Сообщение мгновенно возвращается в голову очереди, брокер тут же отдаёт его тому же (или соседнему) потребителю, тот снова падает. Получается плотный цикл на скорости в тысячи итераций в секунду: процессор в потолок, лог растёт гигабайтами, живые заказы стоят за этим одним битым. Такое сообщение называют poison message — «отравленное». Одна кривая запись останавливает всю обработку.

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

Три двери в DLX

Сообщение попадает в dead letter exchange в трёх случаях — и это стоит запомнить, потому что чаще всего в DLQ прилетает не то, что ждёшь:

  • Отказ потребителяbasic_nack или basic_reject с requeue=False. Явное решение: «я не смогу это обработать никогда».
  • Истёк TTL — срок жизни сообщения вышел, пока оно лежало в очереди. Именно на этом свойстве строится отложенный повтор, до которого мы дойдём ниже.
  • Переполнение — очередь упёрлась в x-max-length и вытеснила самое старое сообщение (об этом подробно в следующем уроке).

У quorum queues есть и четвёртая дверь — превышен x-delivery-limit — но о ней в уроке про кластер.

Объявляем карантин

DLX — это обычный обменник, а DLQ — обычная очередь. Никакой магии, только два аргумента при объявлении рабочей очереди.

import pika

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

# 1. Обменник-карантин и очередь под него
ch.exchange_declare(exchange='orders.dlx', exchange_type='direct', durable=True)
ch.queue_declare(queue='orders.dead', durable=True)
ch.queue_bind(queue='orders.dead', exchange='orders.dlx', routing_key='dead')

# 2. Рабочая очередь: всё отвергнутое уезжает в orders.dlx
ch.queue_declare(
    queue='orders',
    durable=True,
    arguments={
        'x-dead-letter-exchange': 'orders.dlx',
        'x-dead-letter-routing-key': 'dead',
    },
)

conn.close()

Построчно:

  • x-dead-letter-exchange — единственный обязательный аргумент. Указывает, куда брокер отправит «мёртвое» сообщение.
  • x-dead-letter-routing-key — необязательный, но почти всегда нужный. Без него сообщение поедет в DLX с исходным ключом маршрутизации. Если ключ был order.paid, а в DLX такой привязки нет, сообщение просто исчезнет: обменник, которому некуда доставить, молча выбрасывает сообщение. Явный ключ dead убирает этот класс ошибок целиком.
  • durable=True — и у очереди-карантина тоже. DLQ, которая исчезает при рестарте брокера, бессмысленна.

Важный нюанс: аргументы очереди неизменяемы. Если очередь orders уже существует без DLX, повторное объявление с новыми аргументами вернёт ошибку PRECONDITION_FAILED - inequivalent arg (код 406) и убьёт канал. Придётся либо создавать очередь заново, либо задавать настройки через policy — этот путь разберём в следующем уроке.

Потребитель, который умеет сдаваться

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

import json
import pika


class PermanentError(Exception):
    pass  # данные битые: повторять бессмысленно


def handle(order):
    if 'amount' not in order:
        raise PermanentError('нет поля amount')
    print('Оплачиваем заказ', order['id'])


def on_message(ch, method, properties, body):
    try:
        handle(json.loads(body))
    except (PermanentError, json.JSONDecodeError):
        # requeue=False → сообщение уходит в DLX, а не по кругу
        ch.basic_nack(delivery_tag=method.delivery_tag, requeue=False)
        return
    ch.basic_ack(delivery_tag=method.delivery_tag)


conn = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
ch = conn.channel()
ch.basic_qos(prefetch_count=20)
ch.basic_consume(queue='orders', on_message_callback=on_message)
ch.start_consuming()

Ключевая строка одна: requeue=False. Она означает «в очередь не возвращай, поступай по инструкции» — а инструкция записана в аргументах очереди, то есть отправь в DLX. Поставьте здесь True — и вернётесь к бесконечному циклу из начала урока.

Счётчик попыток и отложенный повтор

Но не каждая ошибка — приговор. Если база данных на секунду ушла в дедлок, сообщение стоит повторить. Только не сразу и не бесконечно: нужна пауза (чтобы дать сервису очухаться) и потолок попыток (чтобы не крутить вечно).

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

# Отстойник: у очереди нет потребителей, она просто выдерживает паузу
ch.queue_declare(
    queue='orders.retry.30s',
    durable=True,
    arguments={
        'x-message-ttl': 30000,                 # 30 секунд выдержки
        'x-dead-letter-exchange': '',           # обменник по умолчанию
        'x-dead-letter-routing-key': 'orders',  # и обратно в рабочую очередь
    },
)

MAX_ATTEMPTS = 3


def on_message(ch, method, properties, body):
    headers = properties.headers or {}
    attempt = headers.get('x-attempt', 0)
    try:
        handle(json.loads(body))
        ch.basic_ack(delivery_tag=method.delivery_tag)
    except TransientError:                      # сеть, таймаут, дедлок
        if attempt + 1 >= MAX_ATTEMPTS:
            ch.basic_nack(delivery_tag=method.delivery_tag, requeue=False)
            return
        headers['x-attempt'] = attempt + 1
        ch.basic_publish(
            exchange='',
            routing_key='orders.retry.30s',
            body=body,
            properties=pika.BasicProperties(delivery_mode=2, headers=headers),
        )
        ch.basic_ack(delivery_tag=method.delivery_tag)
    except PermanentError:
        ch.basic_nack(delivery_tag=method.delivery_tag, requeue=False)

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

               ошибка, попытка 1-2
  [orders] ─────────────────────────▶ [orders.retry.30s]
     ▲                                         │
     └────────── TTL 30 с, DLX по умолчанию ───┘
     │
     │ попытки исчерпаны: nack(requeue=False)
     ▼
  [orders.dlx] ──▶ [orders.dead]   ← сюда смотрит дежурный

Хотите экспоненциальную задержку — заведите несколько отстойников: retry.10s, retry.1m, retry.10m, и выбирайте очередь по номеру попытки. Это ровно та схема, которую в больших системах называют «backoff-лестницей».

Как это работает

Когда брокер отправляет сообщение в DLX, он не выбрасывает его молча, а дописывает в заголовки массив x-death — историю смертей:

x-death:
  - count:        3
    reason:       rejected
    queue:        orders
    exchange:     shop
    routing-keys: [order.paid]
    time:         2026-07-14 11:02:57

Поле reason отвечает на главный вопрос разбора: rejected (потребитель отверг), expired (истёк TTL), maxlen (вытеснено из переполненной очереди). Поле count — сколько раз сообщение умирало по этому маршруту.

И вот тонкость, на которой спотыкаются все: x-death.count нельзя использовать как счётчик попыток в схеме с отстойником. Когда потребитель делает basic_publish, он публикует новое сообщение — история заново начинается с нуля. Именно поэтому в примере выше мы ведём собственный заголовок x-attempt. На x-death можно опираться, только если сообщение ходит по кругу «очередь → DLX → очередь» без переиздания.

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

  • requeue=True в обработчике исключений. Главный источник ночных инцидентов. Возвращать в очередь можно только осознанно — например, при штатном завершении процесса, когда сообщение точно не начато.
  • DLQ без потребителя и без алерта. Карантин, в который никто не смотрит, — просто медленная утечка данных. Заведите алерт «в orders.dead появилось хотя бы одно сообщение» и дашборд с телом последнего.
  • Забыли x-dead-letter-routing-key. Сообщения уходят в DLX с исходным ключом, привязок нет, всё тихо исчезает. Хуже, чем потеря с ошибкой: потери нет в логах.
  • DLX без TTL на самой DLQ. Карантин тоже растёт. Дайте ему x-message-ttl в неделю-две, иначе однажды он съест диск.
  • x-death.count как счётчик попыток при переиздании. Всегда 1 — и цикл повторов становится вечным.
  • Одна DLQ на все очереди сразу. Работает, но разбирать её невозможно: непонятно, откуда что прилетело. Заводите карантин на каждый значимый поток.

Итоги

  • DLX — штатный механизм брокера: сообщение уходит туда при nack(requeue=False), истечении TTL или переполнении очереди.
  • requeue=True на исключении — прямой путь к бесконечному циклу и остановке обработки.
  • Различайте временные и постоянные ошибки: первые — в отстойник на повтор, вторые — сразу в карантин.
  • Отложенный повтор собирается из очереди с x-message-ttl и DLX, указывающим обратно на рабочую очередь.
  • Счётчик попыток ведите сами (свой заголовок), а не через x-death, — при переиздании история обнуляется.
  • Публикуйте в отстойник до ack, а не после: дубль лучше потери.
Проверьте себя
1. Потребитель ловит исключение на битом JSON и вызывает basic_nack(requeue=True). Что произойдёт?
AСообщение уйдёт в DLX и будет разобрано вручную
BСообщение вернётся в очередь, снова будет доставлено, снова упадёт — и так по кругу
CRabbitMQ автоматически удалит сообщение после третьей неудачи
DСообщение будет доставлено с задержкой в 30 секунд
2. Зачем в схеме отложенного повтора вести собственный заголовок x-attempt вместо x-death.count?
Ax-death.count доступен только в quorum queues
Bx-death.count считает только истёкшие по TTL сообщения
CПри повторной публикации сообщения через basic_publish история x-death начинается заново, и счётчик всегда будет равен 1
Dx-death вырезается санитайзером RabbitMQ при проходе через DLX