Комментарии 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 — источник событий; если консьюмер пишет в неё же, петлю надо рвать явно. Код примера и текст поправлю.

Обработка отложенных задач c YDB: от таблицы до распределённого координатора