А теперь специально всё сломаем: DLQ руками на Go + PostgreSQL
Теперь давайте не будем больше говорить про архитектуру. Лучше один раз всё сломаем руками.
Kafka для этого не понадобится. Возьмём REST API вместо producer, Go-приложение вместо consumer и PostgreSQL для хранения проблемных сообщений.
В итоге хотим получить вот такой сценарий:
Структура проекта будет примерно такой:
Сначала поднимем PostgreSQL
Для демо достаточно обычного docker-compose:
Запускаем:
Теперь создадим таблицу для сообщений, которые не удалось обработать.
Здесь я специально сохраняю и исходное сообщение, и ошибку. Потом, когда в DLQ окажется несколько тысяч записей, очень не хочется смотреть на них и гадать, почему каждая из них туда попала. DLQ без причины ошибки — это просто склад JSON-ов.
Теперь немного Go
Для события хватит простой структуры:
А вот для DLQ нам уже нужна дополнительная информация:
Здесь есть важное разделение:
- Event — это бизнес-сообщение.
- DeadLetterMessage — это уже сообщение плюс технический контекст, который понадобится для расследования.
Ломаем обработчик
Сам consumer для демо можно сделать совсем примитивным:
Никакой реальной бизнес-логики здесь нет. Нам нужно только одно: иметь сообщение, которое гарантированно падает.
Поэтому договариваемся, что:
всегда вызывает ошибку.
Теперь проблему можно воспроизводить сколько угодно раз.
Отправляем первое сообщение
Обычное событие:
Оно успешно проходит обработку.
Теперь отправляем то самое:
Consumer получает сообщение. Обработка падает. Вместо того чтобы бесконечно пытаться обработать его снова, сохраняем сообщение в PostgreSQL. Теперь оно лежит в DLQ и больше не блокирует основной поток.
Записываем сообщение в DLQ
Сам repository здесь выглядит довольно скучно:
А consumer связывает обработку и DLQ:
То есть любое сообщение, которое мы не смогли обработать, уходит в отдельный контур.
Но положить сообщение в DLQ — только половина работы
Допустим, через несколько часов мы нашли баг. Исправили код. Теперь нужно вернуть сообщение обратно в обработку.
Добавим:
И вот здесь есть место, где очень легко сделать ошибку.
Не удаляйте сообщение из DLQ перед обработкой.
Например, такой код выглядит вполне логично:
Но представьте, что Process() упал.
Получаем:
Вот вам и потеря данных.
Правильнее сначала обработать сообщение и только после успешного результата удалить его:
Теперь если обработка снова упала — запись никуда не делась. Можно исправить проблему и попробовать ещё раз.
Проверяем
Смотрим, что лежит в DLQ:
Берём id сообщения:
Пока poison_pill считается ошибкой, replay будет падать, а сообщение останется в DLQ.
Теперь исправляем обработчик.
Убираем условие:
И запускаем replay ещё раз. На этот раз сообщение успешно обрабатывается. После этого его можно удалить из DLQ.
Получается довольно простой цикл:
И где здесь Kafka?
Её всё ещё нет. И это нормально.
Сейчас у нас:
В реальном сервисе HTTP можно заменить Kafka:
При этом сама логика обработки проблемных сообщений практически не меняется.
В этом и была идея демо. Не написать «свою Kafka», а отделить получение сообщения от его обработки. Тогда источник сообщений можно заменить, не переписывая весь механизм работы с ошибками.
И да, такой стенд сегодня вполне можно попросить собрать AI-кодером. docker-compose, миграцию, модели и HTTP-ручки он напишет довольно быстро.Но я бы всё равно внимательно посмотрел на три места: миграцию, транзакции и replay.
Потому что создать таблицу dead_letter_queue — действительно пять минут. А вот сделать так, чтобы она через полгода не стала причиной нового факапа, — уже инженерная задача.