DLX или его называют очередью повторных попыток, позволяет реализовать логику, при которой, если при обработке сообщения из RabbitMQ произошла ошибка на консьюмере, то мы отправляем сообщение в специальную очередь. В этой очереди оно находится определённое время N, после чего уходит назад в основную очередь, и RabbitMQ, который работает по push модели, проталкивает это сообщение снова в консьюмер.
Почему просто не делать reject и отправлять сообщение назад в очередь?
Потому что оно сразу вернётся в консьюмер, который всё ещё может не оклематься от своих проблем.
Пример из одного проекта на python, реализующий данных подход. Вся декларация происходит на консьюмере:
async def read_messages(
host: str, username: str, password: str, queue_name: str, callback
) -> None:
connection, channel = await rabbitmq_connection(
host=host, username=username, password=password
)
await channel.set_qos(prefetch_count=settings.rabbitmq.prefetch_count)
logger.info("Соединение с rabbitmq установлено.")
dead_letter_exchange = await channel.declare_exchange(
name=settings.rabbitmq.dlx_exchange_name,
type=aio_pika.ExchangeType.DIRECT,
durable=True,
)
dead_letter_queue = await channel.declare_queue(
name=settings.rabbitmq.dlx_queue_name,
durable=True,
arguments={
"x-message-ttl": settings.rabbitmq.dlx_message_ttl,
"x-dead-letter-exchange": "",
"x-dead-letter-routing-key": settings.rabbitmq.queue_name,
},
)
await dead_letter_queue.bind(
dead_letter_exchange, routing_key=settings.rabbitmq.dlx_routing_key
)
queue = await channel.declare_queue(
name=queue_name,
durable=True,
arguments={
"x-dead-letter-exchange": settings.rabbitmq.dlx_exchange_name,
"x-dead-letter-routing-key": settings.rabbitmq.dlx_routing_key,
},
)
await queue.consume(lambda message: callback(message=message, channel=channel))
await asyncio.Future()
await connection.close()
В dead_letter_exchange - декларируем exchange, который будет отправлять сообщения в очередь повторных попыток. Помним, что в RabbitMQ сообщения всегда попадают сначала в exchange, даже если вам кажется, что вы пишете напрямую в очередь.
dead_letter_queue - декларируем очередь повторных попыток.
Параметры:
x-message-ttl - время в миллисекундах, в течение которого сообщение будет находиться в этой очереди перед повторной отправкой в основную очередь. Это обеспечивает задержку между попытками обработки сообщения.
x-dead-letter-exchange - exchange, куда сообщение должно быть отправлено после истечения TTL. В данном случае это пустая строка "", что означает использование default_exchange.
x-dead-letter-routing-key - routing key, по которому понимаем, в какую очередь отправить сообщение при повторной отправке.
После чего делаем bind между dead_letter_exchange и очередью повторных попыток.
А ниже декларируем основную очередь. Проект небольшой, для реализации кое-каких DevOps штук, поэтому использование default_exchange в данном случае вполне оправдано. Это, собственно, тот exchange, куда сообщения отправляются по умолчанию; binding и routing_key RabbitMQ создаст автоматически. Routing key будет равен имени очереди.
И обработка сообщений в callback функции:
async def message_handle(
message: IncomingMessage, channel: AbstractRobustChannel
) -> None:
try:
async with message.process():
data = json.loads(message.body)
<some logic here>
except Exception as err:
await message.reject(requeue=False)
Здесь используется контекстный менеджер async with, если блок кода внутри выполнится без ошибок, значит в RabbitMQ уйдёт ack, что означает, что сообщение обработано, и оно будет удалено из очереди. Если возникнет ошибка, мы вызываем message.reject(requeue=False), что отклоняет сообщение без повторной постановки в ту же очередь, именно этот параметр помечает сообщение как "мертвое".
#python