Обновить

Комментарии 4

Благодарю, за проделанную работу!
У меня есть небольшой вопрос: Если задача падает с ошибкой, которая не исчезнет (например, невалидный payload), как вы предлагаете её обрабатывать?

Спасибо! Хороший вопрос — в статье он намеренно оставлен за скобками (в ограничениях подхода 4 есть строчка «Нет встроенного DLQ»), но на практике это первое, что приходится добавлять.

Я бы разбирал его двумя вопросами по порядку.

1. Можно ли такую задачу просто выбросить?

Ответ определяет весь дальнейший объём работы, и очень часто он — «да». Если задача best-effort (обновить кеш, пересчитать метрику, необязательное уведомление — всё, что восстанавливается пересчётом), то самое простое решение и есть правильное: пропустить её.

Но «пропустить» — это не «залогировать ошибку и пойти дальше». Нужны две вещи:

Явное терминальное состояние. Иначе в подходах с локами задача вернётся: лок протухнет по locked_until, и FetchEligibleCandidate выдаст её снова — уже другому воркеру. Невалидный payload превратится в вечный цикл.

UPDATE coordinated_tasks
SET status = 'skipped', done_at = CurrentUtcTimestamp(), last_error = $err
WHERE partition_id = $p AND priority = $pr AND id = $id

Видимость. Счётчик в Prometheus и алерт на его рост. Иначе «выбрасываем невалидные задачи» незаметно превращается в «теряем 30% трафика после релиза с багом в сериализации».

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

2. В каком из подходов возникла проблема?

Тут разница принципиальная, потому что вопрос сводится к одному: где хранить состояние ретрая?

Подход 1 — таблица + polling

Проще всего: состояние уже в строке. Достаточно колонки status, чтобы отличать успешное завершение от отброшенной задачи, — и задача просто перестаёт попадать в WHERE done_at IS NULL. Ретраи с backoff добавляются так же тривиально, через scheduled_at.

Подходы 2 и 3 — CDC и Topic API

Здесь ретраить заведомо проблемное сообщени не получится. Топик не даёт ни отложенной доставки, ни подтверждения отдельного сообщения: оффсет — это watermark, «обработано всё до N». Из этого следует всё остальное.

Отложить одно сообщение нельзя в принципе. Есть только два хода: закоммитить оффсет (= выбросить сообщение) или не коммитить (= остановить всю партицию до перезапуска ридера, head-of-line blocking). Промежуточного варианта «верни мне это сообщение через 30 секунд, а остальные давай дальше» не существует.

Единственный способ получить паузу внутри самого топика — заснуть в консьюмере, но это останавливает партицию целиком: пока мы ждём одну задачу, за ней стоят все остальные. На высоком rate это неприемлемо.

Если порядок значим — а ради него мы в подходе 3 и городили партицирование по user_id — этот путь закрыт: остаётся остановка партиции с алертом и ручной разбор.

Поэтому на практике для подходов 2 и 3 ответ обычно такой: состояние ретрая выносится в таблицу. Сообщение с постоянной ошибкой сразу уезжает в DLQ-таблицу, оффсет коммитится, партиция едет дальше. У CDC тут приятный бонус: исходная строка жива в tasks, так что достаточно сохранить id и причину, payload дублировать не нужно. А повторная обработка становится отдельным процессом, который читает эту таблицу по scheduled_at — то есть, по сути, подходом 1 или 4, приделанным сбоку.

Подход 4 — координированные воркеры

Здесь всё состояние и так в таблице, поэтому доступна полная схема:

Разделить ошибки на два класса. Временные (5xx, 429, timeout, недоступна база) — ретраим. Постоянные (битый payload, неизвестная версия схемы, 4xx) — ретраить бессмысленно, они только жгут ресурсы воркера, сразу в терминальный статус.

Для временных — attempts, лимит попыток и backoff. Отдельная инфраструктура не нужна: scheduled_at из статьи прекрасно работает как retry-delay.

ALTER TABLE coordinated_tasks ADD COLUMN attempts Uint8;
ALTER TABLE coordinated_tasks ADD COLUMN last_error Utf8;
UPDATE coordinated_tasks
SET status = 'pending',
    scheduled_at = $next_attempt_at,   -- now + backoff(attempts) + jitter
    lock_value = NULL, locked_until = NULL,
    last_error = $err
WHERE partition_id = $p AND priority = $pr AND id = $id

attempts инкрементировать в момент захвата, а не в момент ошибки. Иначе задача, которая роняет воркер (panic, OOM, зависание), никогда не досчитает попытки: воркер умирает → лок протухает → задача снова доступна → следующий воркер умирает. Такой poison pill может по очереди выкосить весь пул. Если attempts растёт в ClaimTask, то и crash-loop сходится к терминальному статусу за N итераций.

Продумать redrive — способ вернуть задачи в работу после исправления кода:

UPDATE coordinated_tasks
SET status = 'pending', attempts = 0
WHERE partition_id = $p AND status = 'failed' AND ...

DLQ без redrive — это просто медленный способ выбросить задачу.

Итого:

Топик — это транспорт, а не очередь задач. Он умеет доставлять и переживать рестарты, но у него нет ни таймера, ни счётчика попыток, ни состояния отдельной задачи. Как только к обработке ошибок появляются серьёзные требования, ретраи всё равно переезжают в таблицу — вопрос только в том, спроектируете вы это сразу (подход 4) или приделаете сбоку через полгода.

А самое дешёвое лекарство именно от невалидного payload — не пускать его в очередь: валидация и версионирование схемы на стороне продюсера. Всё остальное — это про то, как пережить payload, который всё-таки просочился.

Второй подход же будет обрабатывать update от самого себя как новое сообщение, т.к. нет проверки что done_at пустой. Или есть какая-то магия которая такие обновления не будет в cdc записывать?

Магии нет, вы правы — это баг в примере, а не упрощение ради читаемости. Changefeed пишет запись на каждое изменение строки, кем бы оно ни было сделано; фильтров по типу операции или по столбцам в опциях ADD CHANGEFEED нет.

Причём получается не дубль, а самоподдерживающийся цикл: done_at = CurrentUtcTimestamp() даёт новое значение на каждом проходе, поэтому каждая обработка порождает новое событие, и консьюмер производит ровно столько сообщений, сколько потребляет — очередь не разгребается никогда.

Минимальная починка — фильтр в консьюмере, newImage уже содержит done_at, лишнее чтение не нужно: if row.DoneAt != nil { continue }. Но правильнее вообще не писать состояние обработки в таблицу-источник: в подходе 2 колонка done_at не нужна, прогресс уже хранится в оффсете consumer group, а если отметка нужна для отчётности — ей место в отдельной таблице без changefeed. done_at приехал в подход 2 вместе со схемой из подхода 1, где он и есть механизм выборки (WHERE done_at IS NULL), и там это уместно. Общее правило: таблица с changefeed — источник событий; если консьюмер пишет в неё же, петлю надо рвать явно. Код примера и текст поправлю.

Зарегистрируйтесь на Хабре, чтобы оставить комментарий

Публикации