Идемпотентность в пайплайнах: скучно, но спасает
Любой регулярный джоб рано или поздно перезапустят: сеть моргнула и сработал ретрай, источник опоздал и нужен бэкфилл, кто-то нажал «rerun» руками. Вопрос не «если», а «когда».
Если таск не идемпотентен, перезапуск задваивает данные. Причём тихо — ошибок нет, пайплайн зелёный, а через месяц в отчёте выручка выше реальной, и никто не помнит, что тогда был ретрай.
Как выглядит проблема
# так делать не стоит
INSERT INTO orders_daily
SELECT * FROM staging_orders WHERE dt = '{{ ds }}';
Второй запуск за ту же дату — и строки просто добавляются ещё раз.
Как чинится
Самое простое и надёжное — сначала удалить партицию, потом записать:
DELETE FROM orders_daily WHERE dt = '{{ ds }}';
INSERT INTO orders_daily
SELECT * FROM staging_orders WHERE dt = '{{ ds }}';
Обе операции — в одной транзакции. Тогда сколько раз ни запусти, результат один и тот же. Это стоит лишних пары секунд на удаление и экономит неделю разбирательств «откуда взялись лишние продажи в марте».
Варианты для разных движков:
- PostgreSQL —
INSERT ... ON CONFLICT DO UPDATEпо ключу - ClickHouse —
ALTER TABLE ... DROP PARTITIONперед вставкой - Hive / Spark —
INSERT OVERWRITEвместоINSERT INTO
Ещё пара граблей
Не берите текущее время внутри таска. now() делает результат
зависимым от момента запуска — при бэкфилле получите не то, что было в тот день.
Используйте дату интервала, которую даёт планировщик ({{ ds }} в Airflow).
Внешние вызовы тоже стоит защищать. Отправка письма или запись в очередь при ретрае продублируется. Если операция не идемпотентна по своей природе — выносите её в отдельный таск и передавайте ключ идемпотентности.
Хорошая проверка: запустить таск дважды подряд и сравнить результат. Если данные изменились — идемпотентности нет, даже если кажется, что есть.