Перейти к содержимому
Skein
← Все проекты

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 — это отдельные шаги, каждый запускается и проверяется независимо.
    Самопроверка

    Покажи артефакт этапа 'batch-ingest' и объясни ключевой компромисс или режим отказа который он закрывает. Senior-ревьюер проверяет что deliverable измерен (не просто 'работает') и компромисс указан с числами.

  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, обновление или ошибка) — задокументированный осознанный выбор, привязанный к семантике источника.
    Самопроверка

    Покажи артефакт этапа 'idempotent-load' и объясни ключевой компромисс или режим отказа который он закрывает. Senior-ревьюер проверяет что deliverable измерен (не просто 'работает') и компромисс указан с числами.

  3. 03Перейди на инкремент с watermark

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

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

    Покажи артефакт этапа 'incremental-watermark' и объясни ключевой компромисс или режим отказа который он закрывает. Senior-ревьюер проверяет что deliverable измерен (не просто 'работает') и компромисс указан с числами.

  4. 04Поставь на load проверки качества

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

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

    Покажи артефакт этапа 'data-quality-gates' и объясни ключевой компромисс или режим отказа который он закрывает. Senior-ревьюер проверяет что deliverable измерен (не просто 'работает') и компромисс указан с числами.

  5. 05Бэкфилл и поздно приходящие данные

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

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

    Покажи артефакт этапа 'backfill-and-late-data' и объясни ключевой компромисс или режим отказа который он закрывает. Senior-ревьюер проверяет что deliverable измерен (не просто 'работает') и компромисс указан с числами.

Стартер

fallowlone/skein-projects

projects/idempotent-etl-pipeline

Открыть на GitHub ↗
  • README.md
  • src/etl.ts
  • test/etl.test.ts
Забрать только этот проект npx degit fallowlone/skein-projects/projects/idempotent-etl-pipeline idempotent-etl-pipeline

Реализуй заглушки, затем гоняй тесты, пока не позеленеют: bun test

Форкни репозиторий и запушь свою работу — workflow grade прогонит тесты и статические проверки на твоих раннерах.

Рубрика

Джуниор Миддл Сеньор
Идемпотентный 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

Материалы