Идемпотентный пайплайн — это загрузка, которую можно запустить два, три, десять раз подряд, и в витрине останется ровно один корректный результат, как будто прогон был один. Джоба упала на середине — вы перезапускаете её целиком и не получаете задвоенных заказов, поехавшую выручку и сотню вопросов от продакта «почему 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, и разница между ними определяет, нужна ли вам вообще идемпотентность.
- At-least-once («хотя бы один раз») — сообщение или батч гарантированно обработается, но может обработаться и дважды. Так работают почти все очереди и оркестраторы: если консьюмер не подтвердил успех, система шлёт данные повторно. Дубли — норма, с ними живут.
- Exactly-once («ровно один раз») — каждая запись обрабатывается строго один раз. Звучит идеально, но честный end-to-end exactly-once дорог и хрупок: нужна распределённая транзакция между источником, обработкой и приёмником, а на стыке разных систем она почти всегда протекает.
Практический ответ, который я всегда даю: не гонитесь за 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 на paid — UPSERT его обновит, а не добавит вторую строку.
В хранилищах (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.
Частые грабли, из-за которых витрина двоится
Собрал то, на чём спотыкаются чаще всего:
INSERTбез ключа и без транзакции — база разработки. Любой ретрай оркестратора аппендит дубли.DELETE+INSERTвне одной транзакции — падение между шагами оставляет пустую или частичную партицию.- Недетерминированный ключ — если в ключ строки залезает
now(), случайный id или порядок чтения файлов, каждый прогон генерит «новую» строку, и никакойUPSERTне спасёт. - Watermark без перехлёста — опоздавшие данные, пришедшие впритык к границе, теряются навсегда.
MERGEилиUPSERTпо неуникальному ключу — если на колонке нетUNIQUE, база либо ругнётся, либо (в некоторых движках) обновит случайную из совпавших строк. Всегда проверяйте, что ключ и правда уникален.- Дедуп по неправильному
ORDER BY— оставили «первую» версию вместо последней и заморозили устаревший статус заказа.
Идемпотентность — это не отдельная модная технология, а дисциплина записи: любой шаг пишет данные так, что повтор не портит результат. Освоите её — и ночные алерты «витрина двоится» перестанут будить вас в 4 утра, а на собеседовании по data engineering вы спокойно объясните разницу между at-least-once и exactly-once. Кстати, это один из навыков, за которые заметно доплачивают: посмотрите вилки в разборе зарплат аналитиков и инженеров данных и типовые вопросы с собеседований.
Если хотите системно прокачать SQL под такие задачи — от оконных функций до MERGE — в Pro открыты все 425 SQL-задач с настоящим PostgreSQL прямо в браузере, разбор кейсов и AI-собес без лимитов. Начать можно бесплатно, а дальше — по мере того, как втянетесь.
Связанная тема — Бэкфилл данных: что это, идемпотентность, SQL/Airflow.