open atlas
← Все проекты

infra · advanced · 9d

Идемпотентный ETL-конвейер

Конвейеры не падают аккуратно — они падают в три часа ночи, на середине загрузки, и кто-то их перезапускает. Этот проект учит одному свойству, которое отличает любительский скрипт от настоящей дата-инженерии: прогон, который можно повторить сколько угодно раз и всё равно получить ровно одну копию каждой строки. Ты построишь пакетную загрузку, идемпотентный load, watermark для инкрементальных выгрузок и проверки качества данных, которые останавливают мусор, прежде чем он отравит всё ниже по течению.

Почти любая платформа данных — это башня из конвейеров, и те, что выживают при встрече с продакшеном, держатся на одном хребте: их можно перезапускать без страха. Идемпотентность, устойчивый watermark и проверки качества — это не продвинутые излишества, а минимальная планка для конвейера, которому можно доверять в три часа ночи, когда планировщик запускает retry, а ты спишь. Собери это однажды руками — и поймёшь, что на самом деле автоматизирует любой оркестратор (Airflow, dbt, Dagster): дисциплину вычислять один и тот же верный ответ, сколько бы раз работа ни выполнялась. Момент, когда число дубликатов остаётся нулём через перезапуск, сбой и бэкфилл, — это момент, когда ты перестаёшь писать скрипты и начинаешь инженерить данные.

Результат

Запускаемый ETL-джоб, который загружает источник в таблицу хранилища, может быть перезапущен на тех же данных без создания дубликатов (это проверяется подсчётом строк), продвигает сохранённый watermark, чтобы на следующем прогоне забирать только новые строки, и аварийно завершается с понятным отчётом, когда проверка качества данных не проходит.

Этапы

0/5 · 0%
  1. 01Загрузи источник пакетом

    Начни со скучной, честной версии: прочитай каждую строку из источника — таблицы, CSV или Parquet-файла — и вставь её в staging-таблицу в хранилище. Пока без инкрементальной логики, без дедупликации, просто полная загрузка, которую можно прочитать от начала до конца. Смысл начать отсюда — сделать форму данных и контракт источника явными до любой оптимизации: какой у строки естественный ключ, какая колонка говорит, когда она изменилась, и источник — это снимок или append-only лог? Эти три ответа определяют каждое последующее решение, и ошибиться в них — самая дорогая ошибка во всём конвейере. Держи extract, transform и load отдельными шагами, чтобы рассуждать о каждом; не поддавайся искушению свернуть их в один хитрый запрос.

    Критерии готовности
    • Одна команда читает весь источник и кладёт каждую строку в staging-таблицу, и ты можешь назвать естественный ключ строки и её колонку с временем изменения.
    • Extract, transform и load — это отдельные шаги, каждый запускается и проверяется независимо.
  2. 02Сделай load идемпотентным

    Теперь запусти полную загрузку дважды и посмотри, как растут дубликаты — именно этот баг призван убить данный этап. Идемпотентная загрузка даёт одну и ту же итоговую таблицу, запустилась она раз или пять, а значит load должен быть по ключу, а не append. Замени слепой INSERT на upsert: INSERT ... ON CONFLICT (natural_key) DO UPDATE в Postgres или MERGE. Глубокая мысль в том, что конвейеры живут в мире at-least-once — планировщики делают retry, операторы перезапускают, частичный сбой бросает тебя на середине загрузки, — поэтому корректность не может зависеть от того, что джоб выполнится ровно один раз. Осознанно реши, что значит «конфликт»: повторно увиденный ключ — это no-op, last-writer-wins или ошибка? Каждый выбор кодирует свою правду об источнике, и неверный тихо портит историю.

    Критерии готовности
    • Двойной запуск одной и той же загрузки на идентичном входе оставляет в целевой таблице одинаковое число строк оба раза — это доказано count-запросом.
    • Поведение upsert при конфликте (no-op, обновление или ошибка) — задокументированный осознанный выбор, привязанный к семантике источника.
  3. 03Перейди на инкремент с watermark

    Перечитывать весь источник на каждом прогоне нормально на 10 тысячах строк и катастрофа на 10 миллиардах. Решение — watermark: маленький кусочек устойчивого состояния (максимальный change-timestamp или sequence id, который ты успешно обработал), позволяющий следующему прогону извлекать только строки новее этой отметки. Храни его транзакционно, чтобы watermark продвигался только если покрываемая им загрузка действительно закоммитилась; иначе сбой между «load прошёл» и «watermark сохранён» навсегда пропустит данные. Здесь идемпотентность и инкрементальность обязаны сотрудничать: поскольку перезапуски всё равно могут повторно забрать граничную строку, именно upsert из прошлого этапа делает это перекрытие безвредным. Выбирай границу аккуратно — бери замкнутый интервал и смирись с переобработкой края, потому что пропуск строки — это тихая потеря данных, а переобработка под идемпотентной загрузкой бесплатна.

    Критерии готовности
    • После успешного прогона следующий извлекает только строки новее сохранённого watermark, а сбой на середине загрузки оставляет watermark непродвинутым.
    • Повторный прогон через границу watermark не создаёт дубликатов, потому что загрузка идемпотентна по естественному ключу.
  4. 04Поставь на load проверки качества

    Конвейер, который грузит что угодно, хуже отсутствия конвейера, потому что он отмывает плохие данные в доверенные таблицы, на которых строят другие. Добавь проверки качества данных, которые запускаются между transform и load и имеют право прервать загрузку: число строк в ожидаемом коридоре, естественный ключ в батче уникален и не null, обязательные колонки не null, числовая колонка в разумных границах. Сложный вопрос дизайна — что должен делать «провал»: остановить всю загрузку или отправить проблемные строки в карантин и загрузить остальное? Остановка защищает корректность, но тормозит бизнес; карантин держит поток данных, но требует пути для переобработки отложенного. В любом случае проверки должны быть громкими: тихий провал качества — самый опасный баг в дата-инженерии, ведь его находят недели спустя в дашборде, которому уже никто не верит.

    Критерии готовности
    • Минимум три проверки (коридор числа строк, уникальность/не-null ключа, диапазон значений) выполняются перед load и могут прервать его с читаемым отчётом.
    • Намеренно испорченный батч не проходит проверку и не попадает в целевую таблицу; провал виден, а не молчалив.
  5. 05Бэкфилл и поздно приходящие данные

    Две сеньорные реальности ломают аккуратную forward-only модель. Первая — бэкфилл: кому-то нужно перезалить прошлый март, потому что трансформация была неверной, поэтому твой джоб должен принимать явный диапазон дат и переобрабатывать его, не задваивая сегодняшние строки и не ломая живой watermark. Вторая — поздно приходящие данные: строка со вчерашним временем появляется после того, как сегодняшний watermark её уже миновал, и наивная инкрементальная выгрузка пропустит её навсегда. Обе проблемы — на самом деле один урок: watermark это оптимизация, а не источник истины, и именно идемпотентность загрузки позволяет безопасно переобработать любое окно. Реши их, сделав бэкфилл ограниченным идемпотентным перезапуском по выбранному интервалу и расширив инкрементальное окно небольшим допуском на опоздание (перечитывай хвостовой буфер), чтобы строки не по порядку получили второй шанс лечь ровно один раз.

    Критерии готовности
    • Прогон бэкфилла по явному диапазону дат переобрабатывает только этот диапазон, не создаёт дубликатов и оставляет forward-watermark нетронутым.
    • Строка, вставленная с временем старше текущего watermark, всё равно ложится ровно один раз — благодаря буферу опоздания и идемпотентному upsert.

Рубрика

Джуниор Миддл Сеньор
Идемпотентный load Load — это слепой INSERT; повторный запуск на тех же данных удваивает каждую строку в целевой таблице. INSERT ... ON CONFLICT (natural_key) DO UPDATE заменяет слепую вставку; повторный запуск на идентичных данных не меняет число строк, а поведение при конфликте (обновление или no-op) — задокументированный осознанный выбор. Ты рассуждаешь о том, что кодирует «конфликт»: no-op означает, что источник — снапшот (первая запись выигрывает); update означает last-writer-wins (источник отражает текущее состояние); ошибка — что повторных ключей не ожидается. Неверный выбор молча портит историю; ты называешь семантику источника, на которой основан твой.
Долговечность watermark Watermark записывается в файл или память; сбой между load и сохранением оставляет следующий прогон повторно забирающим уже загруженные строки или, хуже, пропускающим строки, которые упавший load не закоммитил. Watermark хранится транзакционно в базе данных и продвигается в том же коммите, что и load, поэтому сбой оставляет его непродвинутым и следующий прогон безопасно перебирает то же окно. Ты обрабатываешь граничные строки: замкнутый интервал намеренно перебирает граничную строку, а идемпотентный upsert поглощает перекрытие без дубликатов. Ты объясняешь, почему пропуск строки — тихая потеря данных, а переобработка бесплатна — и показываешь, что буфер опоздания для поздних данных расширяет тот же принцип.
Проверки качества данных Конвейер грузит всё, что приходит; плохие строки попадают в целевую таблицу и обнаруживаются позже в дашборде, которому уже никто не верит. Коридор числа строк, уникальность ключа, ненулевые обязательные колонки и проверка диапазона значений выполняются между transform и load; намеренно испорченный батч не проходит проверку и не попадает в целевую таблицу — провал виден, а не молчалив. У тебя задокументирована политика «провала»: остановка или карантин. Ты можешь сформулировать, почему остановка защищает корректность, но тормозит бизнес, а карантин держит поток, но требует пути переобработки — и показываешь, что строки в карантине доступны для запросов и имеют понятный механизм повтора.
Бэкфилл и поздние данные Нет способа переобработать прошлый диапазон дат; строка с устаревшей меткой времени молча пропускается, когда watermark её уже миновал. Бэкфилл принимает явный диапазон дат и идемпотентно перезапускается по нему, не трогая live watermark; буфер опоздания перечитывает хвостовое окно, давая строкам не по порядку второй шанс. Ты можешь сформулировать компромисс: расширение буфера опоздания увеличивает стоимость переобработки за прогон; ты замерил границу — например, 2-часовой буфер перечитывает N строк за X секунд — и выбрал значение, которое можешь обосновать фактическим распределением задержек источника.
Эталонный разбор (спойлер)

Почему конвейеры живут в мире at-least-once: планировщики делают retry, операторы перезапускают после сбоев, а частичная загрузка оставляет целевую таблицу в промежуточном состоянии. Идемпотентность превращает перезапуск из риска целостности данных в безопасный no-op. INSERT ... ON CONFLICT — это Postgres-примитив, кодирующий это на уровне строк; ключ — выбрать правильное разрешение конфликта для семантики источника.

Watermark — это оптимизация, а не источник истины: он нужен, чтобы избежать повторного сканирования всего источника на каждом прогоне, а не определять, какие строки существуют. Идемпотентность load делает переобработку окна безопасной. Коммит watermark в той же транзакции, что и load, означает, что единственные согласованные состояния — «окно загружено и watermark продвинут» или «ни то ни другое».

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

Сделай по-сеньорски

  • Замени единый watermark на пер-партиционные (например, по одному на каждый шард источника), чтобы медленная или бэкфилящаяся партиция не тормозила и не пропускала данные в остальных.
  • Сделай конвейер наблюдаемым: публикуй метрики на каждый прогон (строк на входе, апсертнуто, отклонено, watermark до/после, длительность) и алёрти, когда срабатывает проверка качества или число строк уходит за ожидаемый коридор.
  • Добавь историю медленно меняющегося измерения (SCD Type 2) для одной таблицы, чтобы обновления создавали новые версионные строки вместо перезаписи, и докажи, что upsert сохраняет эту историю идемпотентной при перезапусках.

Навыки

designing idempotent loads with upsert/merge keysincremental extraction with a durable watermarkwriting data-quality assertions that gate a loadreasoning about at-least-once delivery and re-runsbackfilling a date range without double-countinghandling late-arriving and out-of-order data

Рекомендуемый стек

pythonpostgresSQL (INSERT ... ON CONFLICT / MERGE)a source table or CSV/Parquet drop