В комментах к посту про Transactional Outbox задали вопрос про этот паттерн, и мне захотелось сделать пост.
Когда я работала в платформенной команде, у нас был сервис, принимающий сообщения из Кафки и сохраняющий их в долговременное хранилище. Сообщения ни в коем случае нельзя было терять.
«Так бизнес-сообщения никогда нельзя терять», — скажете вы. «Никому не понравится, если он оплатит заказ, а сообщение об этом потеряется».
Согласна 😊 Но тут особенность заключалась в том, что даже некорректные сообщения, которые не удалось обработать, тоже нужно было сохранить для последующего анализа.
Поэтому мы использовали подход с retry-топиками и Dead Letter Queue.
Dead Letter Queue (DLQ) — это специальная очередь (топик в Кафке, очередь в RabbitMQ), в которую складывают сообщения, которые консьюмеру не удалось обработать.
Алгоритм работы обычно такой:
1️⃣ Консьюмер получает сообщение и пробует его обработать. Если всё ок — коммитим оффсет.
2️⃣ Если возникает ошибка, то пытаемся вычитать и обработать сообщение ещё N раз (ретраи)
3️⃣ Если попытки исчерпаны, сообщение в исходном виде отправляется в Dead Letter Queue
4️⃣ Поступившие в Dead Letter Queue сообщения часто анализируются вручную
Возникает вопрос: зачем нужны ретраи?
Ведь от повторных попыток обработки "битое" сообщение не исправится.
Всё верно, если проблема заключается в самом сообщении, например, у него некорректный формат, и его не удается распарсить, то ретраи не спасут. Ретраи помогают, когда сообщение нормальное, но проблема не в нём:
➡️ Если при обработке задействованы какие-то еще сервисы инфраструктуры, которые могут отваливаться.
Например, при обработке мы обращаемся в кэш, а тут в моменте он упал. Или БД не ответила. Или при обработке идет REST-запрос в другой сервис, и он не ответил (вообще так делать не стоит, так как REST-запрос может затянуться, и консьюмера выгонят из консьюмер-группы).
То есть если была временная недоступность какого-то сервиса, а через минуту всё восстановилось. Без ретрая в такой ситуации мы бы потеряли корректное сообщение.
➡️ Если обработка сообщения завязана на другие события, которые должны успеть прийти.
Например, нам пришло сообщение о заказе пользователя, а самого пользователя еще нет в БД. При повторной попытке обработки пользователь уже «приедет» и будет на месте.
Конечно, такой сценарий стоит использовать с осторожностью, и возможно тут лучше подойдет реализация не через ретраи, а с отдельным временным сохранением заказов, ожидающих пользователей.
Число ретраев обычно ограничивают (стандартно — 3 попытки). Хорошей практикой является ретраить не сразу, а настроить экспоненциальное время задержки для следующей попытки. Это можно сделать как в коде приложения с консьюмером, так и через отдельные топики.
В случае с отдельными retry-топиками сообщение после ошибки обработки отправляется в отдельный топик, например,
orders-retry-1min, orders-retry-10min.Отдельный консьюмер читает с нужной задержкой сообщение из retry-топика и либо удачно обрабатывает его, либо перекладывает в следующий retry-топик, а когда попытки исчерпаны — в Dead Letter Queue.
📎 Отмечу, что Kafka в чистом виде не имеет никаких механизмов задержки сообщений, поэтому все эти задержки для ретрай-топиков мы должны настраивать в коде сами. Библиотеки для работы с Kafka могут иметь инструменты для этого.
📎 RabbitMQ имеет поддержку отложенных сообщений (через настройку TTL у сообщений или через плагин).
📎 В NATS Core, как в Kafka, нет встроенной поддержки задержек чтения, нужно реализовывать в коде.
Использовали ли вы retry-очереди и Dead Letter Queue ?