DLQ на PostgreSQL: а что, если сломается сама DLQ?

В прошлый раз мы разобрали как одной сообщение кладёт весь прод. Но, допустим, мы наконец победили poison pill. Consumer получил сообщение, не смог его обработать несколько раз и отправил в DLQ.

Красота.

Теперь хотя бы основной поток не стоит из-за одного битого сообщения.

Но есть один неприятный вопрос: а что произойдёт, если PostgreSQL в этот момент тоже недоступна?

Вот здесь многие реализации DLQ начинают выглядеть уже не так красиво.

Представим обычный сценарий:

Kafka ↓ consumer ↓ обработка ↓ ERROR ↓ PostgreSQL

onsumer поймал ошибку и хочет записать сообщение в dead_letter_queue. А PostgreSQL лежит. Что теперь делать?

Если просто бросить ошибку наверх — сообщение останется в Kafka и consumer попробует обработать его снова. И снова получит ту же ошибку.

Получается примерно так:

Kafka ↓ consumer ↓ ERROR ↓ PostgreSQL DOWN ↓ retry ↓ ERROR ↓ PostgreSQL DOWN ↓ retry ↓ ...

Мы вроде бы вынесли проблемные сообщения в DLQ, но сами же сделали DLQ зависимой от ещё одного сервиса. И это только первая проблема.

Есть ещё replay.

Допустим, через два дня мы нашли причину ошибки и починили consumer. В DLQ лежит 50 тысяч сообщений. Нужно вернуть их обратно в Kafka.

На первый взгляд всё просто:

PostgreSQL ↓ replay ↓ Kafka ↓ consumer

Но что произойдёт, если consumer успешно обработал сообщение, а PostgreSQL не успела отметить его как replayed?

При следующем запуске мы можем отправить его ещё раз. А если Kafka приняла сообщение, но соединение оборвалось до того, как мы получили подтверждение?

Отправим повторно.

И внезапно наша задача «просто сделать replay» превращается в задачу про идемпотентность. Поэтому в DLQ я бы обязательно хранил не только само сообщение и текст ошибки.

Например:

id topic partition offset message_key message_value error_message retry_count status created_at replayed_at

А status мог бы быть примерно таким:

NEW PROCESSING REPLAYED FAILED

Тогда у нас появляется нормальная модель работы с сообщением.

Мы можем найти все новые ошибки:

SELECT * FROM dead_letter_queue WHERE status = 'NEW';

Посмотреть, сколько сообщений упало по конкретной причине. Найти сообщения, которые уже отправляли на replay. И главное — не пытаться управлять всем этим через один флаг is_replayed = true. Потому что между «мы отправили сообщение в Kafka» и «мы точно знаем, что его можно больше не трогать» есть целая куча неприятных сценариев.

И ещё один момент, о котором часто вспоминают слишком поздно.

DLQ тоже нужно чистить.

Если просто складывать туда всё подряд, через несколько месяцев PostgreSQL превратится в кладбище сообщений. И тут уже начинаются обычные задачи эксплуатации: партиционирование таблицы, шардирование, retention, индексы, архивирование, мониторинг размера, алерты.

То есть DLQ, которую мы сначала воспринимали как маленькую таблицу «на случай ошибки», постепенно превращается в полноценную подсистему.

И это, пожалуй, главный вывод: DLQ — это не место, куда мы складываем то, что не смогли обработать.

Это часть системы обработки ошибок. У неё есть свои SLA, свои failure-сценарии, свой lifecycle и свои проблемы с консистентностью.

И чем дороже ошибка в бизнесе, тем внимательнее к этому стоит относиться.

В следующей части как раз соберём всё это руками: Go + Kafka + PostgreSQL, poison pill, DLQ и replay. Без магии — чтобы было понятно, где именно всё ломается.

Селькин Андрей
TGM / Go-разработчик из Fintech. Инженерные заметки: https://t.me/andrei_selkin_outbox
1