Data EngineeringETLидемпотентность

Идемпотентность ETL: повторный запуск без дублей

2026-07-12 11 мин

Идемпотентный пайплайн — это загрузка, которую можно запустить два, три, десять раз подряд, и в витрине останется ровно один корректный результат, как будто прогон был один. Джоба упала на середине — вы перезапускаете её целиком и не получаете задвоенных заказов, поехавшую выручку и сотню вопросов от продакта «почему DAU вырос вдвое за ночь». Весь фокус в том, чтобы писать данные не через слепой INSERT, а через операции, которые сами затирают прошлый результат по ключу: UPSERT, MERGE или DELETE+INSERT по партиции. Ниже — как это устроено на практике аналитика, который каждый день поддерживает витрины.

Что такое идемпотентность пайплайна простыми словами?

Слово пришло из математики: операция идемпотентна, если применить её один раз или сто раз — результат одинаковый. Выключатель света идемпотентен: щёлкнули «выкл» — лампа погасла, щёлкнули ещё раз «выкл» — она так и осталась погашенной. А вот «прибавить единицу» не идемпотентно: каждый вызов меняет число.

Переносим на ETL. Идемпотентная загрузка приводит таблицу в одно и то же состояние независимо от того, сколько раз вы её запустили. Обработали данные за 11 июля один раз — в витрине данные за 11 июля. Запустили её же ещё дважды — данные за 11 июля не изменились, не задвоились, не пропали. Это свойство самой загрузки, а не удача.

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

Почему повторный запуск задваивает данные?

Потому что самый очевидный способ записать данные — INSERT — по своей природе аппендит, то есть добавляет строки, ничего не проверяя.

-- так писать в витрину нельзя, если джоба может перезапуститься
INSERT INTO mart_orders_daily (dt, user_id, orders, revenue)
SELECT dt, user_id, count(*), sum(amount)
FROM stg_orders
WHERE dt = '2026-07-11'
GROUP BY dt, user_id;

Прогнали один раз — в витрине агрегат за 11 июля. Airflow по таймауту сделал retry, прогнали второй раз — теперь на каждого пользователя по две строки, и sum(revenue) в дашборде удвоился. Никакой ошибки в SQL нет, синтаксис валиден, тест «данные появились» проходит. Просто их появилось в два раза больше, чем надо.

Отдельно больно, когда джоба падает на середине. Представьте, что INSERT успел записать 400 тысяч строк из 900 тысяч, и тут упал из-за нехватки памяти. Если вставка шла без транзакции, эти 400 тысяч остаются в таблице. Вы чините причину, запускаете заново с нуля — и получаете 400 тысяч старых плюс 900 тысяч новых. Витрина двоится частично, что даже хуже полного дубля: цифры «почти правильные», и ошибку замечают через неделю.

Вывод простой: механизм записи должен либо сначала убирать прошлый результат, либо обновлять его на месте по ключу. Разберём оба.

At-least-once или exactly-once — что выбрать?

Это две гарантии доставки, о которых спрашивают на собеседованиях по data engineering, и разница между ними определяет, нужна ли вам вообще идемпотентность.

Практический ответ, который я всегда даю: не гонитесь за exactly-once в транспорте — сделайте идемпотентным приёмник. Тогда at-least-once доставка плюс идемпотентная запись = «effectively-once», результат такой, будто всё прошло ровно один раз. Транспорт может присылать один и тот же заказ трижды — если запись в витрину идёт через UPSERT по order_id, все три раза схлопнутся в одну строку.

Это снимает большую часть головной боли. Вместо того чтобы строить сложную инфраструктуру гарантий, вы переносите ответственность в один слой — в SQL-загрузку, которой управляете сами.

Как сделать загрузку идемпотентной через UPSERT и MERGE?

UPSERT — это «обнови, если ключ уже есть, иначе вставь». В PostgreSQL это INSERT ... ON CONFLICT:
INSERT INTO mart_orders (order_id, user_id, amount, status, updated_at)
SELECT order_id, user_id, amount, status, updated_at
FROM stg_orders
WHERE updated_at > :watermark
ON CONFLICT (order_id) DO UPDATE
SET amount     = EXCLUDED.amount,
    status     = EXCLUDED.status,
    updated_at = EXCLUDED.updated_at;

Ключевое условие — на order_id должен стоять UNIQUE или PRIMARY KEY, иначе базе не по чему ловить конфликт. Теперь сколько раз ни запусти загрузку с одним и тем же заказом — в таблице останется одна строка с последними значениями. Заказ поменял статус с pending на paidUPSERT его обновит, а не добавит вторую строку.

В хранилищах (Snowflake, BigQuery, свежий PostgreSQL, Spark SQL) для того же используют MERGE, где логика match/not matched расписана явно:

MERGE INTO mart_orders AS t
USING stg_orders AS s
ON t.order_id = s.order_id
WHEN MATCHED THEN
    UPDATE SET amount = s.amount, status = s.status, updated_at = s.updated_at
WHEN NOT MATCHED THEN
    INSERT (order_id, user_id, amount, status, updated_at)
    VALUES (s.order_id, s.user_id, s.amount, s.status, s.updated_at);
UPSERT и MERGE идеальны для витрин уровня строки-факта с понятным бизнес-ключом: заказы, платежи, профили пользователей. Синтаксис ON CONFLICT и оконные конструкции удобно отработать на живой базе в SQL-тренажёре, а подсмотреть точный синтаксис под свой диалект — в справочнике по SQL.

Когда лучше DELETE+INSERT по партиции?

Когда вы пересобираете витрину целыми кусками — обычно по дню — а не точечно обновляете строки. Тогда проще снести всю партицию за дату и записать заново:

BEGIN;

DELETE FROM mart_orders_daily
WHERE dt = '2026-07-11';

INSERT INTO mart_orders_daily (dt, user_id, orders, revenue)
SELECT dt, user_id, count(*), sum(amount)
FROM stg_orders
WHERE dt = '2026-07-11'
GROUP BY dt, user_id;

COMMIT;

Два момента, без которых это не работает. Первое — обе операции обязаны быть в одной транзакции. Если DELETE прошёл, а INSERT упал, транзакция откатит всё разом, и старые данные за день останутся на месте; без транзакции вы получили бы пустую партицию. Второе — DELETE и INSERT должны бить строго по одному и тому же фильтру партиции (dt = '2026-07-11'), иначе снесёте не то, что пересоберёте.

Этот приём хорош для агрегатов и daily-витрин, где ключ строки не всегда стабилен (пользователь мог отвалиться, группа схлопнуться), а вот «состояние за день» пересчитывается детерминированно. В Spark и Hive аналог — INSERT OVERWRITE PARTITION, он делает ровно то же: перезаписывает партицию, а не дописывает.

Как выбрать между UPSERT и DELETE+INSERT? Грубое правило: обновляете отдельные записи по стабильному ключу — UPSERT или MERGE; пересобираете целую партицию с нуля — DELETE+INSERT в транзакции.

Как дедуплицировать по ключу и зачем нужен watermark?

Даже с правильной записью в источнике встречаются дубли — очередь прислала событие дважды, джоба-предок отработала с наложением. Поэтому перед загрузкой я почти всегда оставляю по ключу одну свежую версию строки через оконную функцию:

WITH ranked AS (
    SELECT *,
           row_number() OVER (
               PARTITION BY order_id
               ORDER BY updated_at DESC
           ) AS rn
    FROM stg_orders
)
SELECT order_id, user_id, amount, status, updated_at
FROM ranked
WHERE rn = 1;
PARTITION BY order_id группирует все версии заказа, ORDER BY updated_at DESC ставит самую свежую первой, rn = 1 оставляет только её. В pandas та же логика — в две строки:
df = (df.sort_values("updated_at")
        .drop_duplicates(subset="order_id", keep="last"))

Теперь про watermark — это метка «до какого момента я уже загрузил данные». Обычно это максимальный updated_at в витрине:

SELECT max(updated_at) AS watermark FROM mart_orders;

На следующем прогоне вы берёте из источника только строки новее watermark — это инкрементальная загрузка, она в разы дешевле полного пересчёта. Но есть подвох: данные приходят с опозданием, и запись с updated_at чуть меньше watermark может появиться уже после того, как вы прогнали джобу. Поэтому watermark берут с запасом — перечитывают небольшой хвост (например, минус несколько часов) и полагаются на UPSERT или дедуп, чтобы повторно прочитанные строки не задвоились. Именно связка «watermark с перехлёстом плюс идемпотентная запись» даёт и дёшево, и без потерь. Отработать оконные функции и дедуп на реальных данных можно в Python-тренажёре и на практических задачах.

Как проверить, что пайплайн идемпотентен?

Самый честный тест — запустить загрузку дважды подряд на одних входных данных и убедиться, что витрина не изменилась. Я обычно смотрю на два числа:

SELECT
    count(*)                    AS rows_total,
    count(DISTINCT order_id)    AS keys_unique
FROM mart_orders;

Если rows_total больше keys_unique — по ключу есть дубли, идемпотентность сломана. На идемпотентной витрине эти числа равны, и повторный прогон их не двигает. Для агрегатов сравнивают контрольные суммы: sum(revenue) и count(*) до второго запуска и после должны совпасть до копейки.

Хорошая привычка — зашить такую проверку прямо в пайплайн отдельным шагом после загрузки: посчитать дубли по ключу и уронить джобу с понятной ошибкой, если они есть. Так вы ловите задвоение в тот же день, а не когда его замечает продакт в дашборде DAU.

Частые грабли, из-за которых витрина двоится

Собрал то, на чём спотыкаются чаще всего:

Идемпотентность — это не отдельная модная технология, а дисциплина записи: любой шаг пишет данные так, что повтор не портит результат. Освоите её — и ночные алерты «витрина двоится» перестанут будить вас в 4 утра, а на собеседовании по data engineering вы спокойно объясните разницу между at-least-once и exactly-once. Кстати, это один из навыков, за которые заметно доплачивают: посмотрите вилки в разборе зарплат аналитиков и инженеров данных и типовые вопросы с собеседований.

Если хотите системно прокачать SQL под такие задачи — от оконных функций до MERGE — в Pro открыты все 425 SQL-задач с настоящим PostgreSQL прямо в браузере, разбор кейсов и AI-собес без лимитов. Начать можно бесплатно, а дальше — по мере того, как втянетесь.

Связанная тема — Бэкфилл данных: что это, идемпотентность, SQL/Airflow.

Отработай SQL на практике
545 SQL-задач с автопроверкой — первые открыты без регистрации.
SQL-тренажёр →