сеньорчикОткрыть в Telegram
← вся теориятеория к собесу · ETL и оркестрация

Паттерны пайплайнов данных

Паттерны пайплайнов

Как забирать данные - полностью или инкрементом, батчем или потоком, по расписанию или по событию это набор осознанных выборов, а не привычка. Собес проверяет, выбираешь ли ты паттерн по требованиям, а не по моде.

Стержень: инкремент дёшев, но требует watermark и ловли поздних правок; мониторинг пайплайна отвечает на три вопроса, а не на один.

// Формулировки: «full load или инкремент?», «что такое CDC?», «когда стриминг оправдан?».

Full, инкремент, CDC

Full load - полная перегрузка: проста и самочинится (каждый раз свежая копия), но дорога на объёме. Инкремент дёшев, но требует watermark (докуда обработано), ловли поздних правок и периодической сверки, что легко недоделать.

CDC (change data capture) - поток изменений прямо из WAL (write-ahead log)/binlog источника (Debezium-класс): забирает вставки, update и delete почти в реальном времени, не нагружая источник запросами; цена - сложнее в сопровождении.

// Ключевая ловушка CDC - забыть delete: если ловить только вставки и апдейты, удалённые в источнике строки будут жить в DWH (data warehouse) вечно. А инкремент по updated_at, который источник обновляет не всегда, молча теряет изменения.

CDC
change data capture - поток изменений из лога БД
watermark
граница «докуда обработано» у инкремента

Снапшот, лог, батч, стриминг

Снапшот против лога изменений: снапшот отвечает на «как сейчас», лог - на «что и когда менялось». DWH обычно нужны оба - текущее состояние для витрин и история для аудита и SCD (slowly changing dimension).

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

// Стриминговый пайплайн ради отчёта, который смотрят раз в день, классическая переинженерия: сложность реального времени без единого потребителя этой скорости.

снапшот / лог
состояние «как сейчас» / история «что когда менялось»
батч / стриминг
выбор по требованию к свежести данных

Оркестрация, мониторинг, восстановление

Оркестрация по данным, а не только по времени: сенсоры и триггеры на готовность источника вместо «в 3 ночи, авось долетело». Event-driven там, где расписание врёт. Fan-out/fan-in распараллеливает независимые ветки (по источникам, партициям) и сводит их на барьере - критический путь короче.

Мониторинг пайплайна - три вопроса: успел ли (SLA - service level agreement), сколько привёз (объёмы против истории), сходится ли (сверка с источником). Зелёный статус отвечает только на первый.

// Восстановление проектируют заранее: что перезапускаем при падении слоя, в каком порядке, откуда бэкфилим. Runbook пишут до инцидента, иначе падение в 4 утра превращается в раскопки зависимостей вслепую.

event-driven
запуск по готовности данных, не по часам
SLA пайплайна
к какому времени данные обязаны быть готовы

Как отвечать: «Full load или инкремент - как выбираешь, и что такое CDC?»

Выбираю по объёму и по тому, могу ли надёжно определить, что изменилось. Full load беру, когда таблица небольшая или нет надёжного признака изменений: он прост и самочинится - каждый раз свежая полная копия, ничего не накапливается криво. Как только объём делает полную перегрузку дорогой, перехожу на инкремент, но тогда обязан завести watermark - докуда обработано, ловить поздние правки скользящим окном и периодически сверяться с источником, иначе тихо потеряю строки. Самый честный инкремент это CDC: изменения читаются прямо из WAL или binlog источника, инструментом класса Debezium, и приходят вставки, апдейты и удаления почти в реальном времени, не нагружая источник запросами. Тут критично не забыть про delete, иначе удалённые в источнике строки останутся в DWH навсегда. То есть маленькое и без признака изменений - full, большое с надёжным CDC - поток изменений.

Почему это сильный ответ: критерий выбора (объём + наличие признака изменений), названы обязательные части инкремента (watermark, поздние, сверка), определён CDC и назван его главный подводный камень - delete.

На чём валят

  • Инкремент по updated_at, который источник не всегда обновляет - тихие потери.
  • CDC без обработки delete - удалённые в источнике строки живут в DWH вечно.
  • Стриминг-пайплайн для отчёта, который смотрят раз в день.
  • Запуск по расписанию без проверки готовности источника - регулярные полупустые загрузки.
  • Нет runbook'а восстановления: инцидент в 4 утра превращается в архео-раскопки зависимостей.

Проверьте себя

Пять вопросов из банка по этой подтеме. Всего их 8, остальные разбираются в тренажёре.

  1. #pipeline_patterns1 / 5
    Чем оркестрация пайплайнов (Airflow) отличается от хореографии (событийной)?
    A)Оркестрация означает децентрализованную реакцию сервисов на события, а хореография — центрального дирижёра
    B)Это одно и то же: оба термина описывают запуск задач строго по фиксированному расписанию времени
    C)Хореография лучше оркестрации во всех сценариях, поэтому оркестраторы вроде Airflow устарели
    D)Оркестрация — центральный дирижёр задаёт порядок и зависимости; хореография — сервисы реагируют на события друг друга без центра
    показать ответ и разбор
    +D)Оркестрация — центральный дирижёр задаёт порядок и зависимости; хореография — сервисы реагируют на события друг друга без центра

    // разбор: Оркестрация: центральный компонент (Airflow) явно знает граф зависимостей и дирижирует — запускает задачи в нужном порядке, следит за статусами, ретраит. Плюсы — прозрачность, единая точка контроля и мониторинга; минус — центральная зависимость и связанность. Хореография: нет дирижёра, каждый сервис реагирует на события других (шаг завершился → событие → следующий подхватил), связь через шину/события. Плюсы — слабая связанность, масштаб; минус — сложнее отследить сквозной поток и отладить. Батч-DWH тяготеет к оркестрации, событийные/микросервисные потоки — к хореографии.

  2. #pipeline_patterns2 / 5
    Почему пайплайн дробят на мелкие идемпотентные шаги, а не пишут одним большим скриптом?
    A)Ради рестартуемости, наблюдаемости, параллелизма и гранулярных ретраев: упал шаг — перезапускают его, а не всё с нуля
    B)Мелкие шаги работают медленнее одного большого скрипта, зато их проще читать глазами
    C)Дробление нужно чтобы обойти ограничение на максимальную длину SQL-запроса в СУБД
    D)Один большой шаг предпочтительнее, потому что его не получится случайно перезапустить не полностью
    показать ответ и разбор
    +A)Ради рестартуемости, наблюдаемости, параллелизма и гранулярных ретраев: упал шаг — перезапускают его, а не всё с нуля

    // разбор: Один монолитный шаг при сбое на середине приходится перезапускать целиком (дорого, и он мог оставить полусостояние). Разбиение на мелкие идемпотентные задачи даёт: рестартуемость (перезапустить только упавший шаг с его входа), наблюдаемость (видно, где именно встало, метрики по шагам), параллелизм (независимые шаги идут одновременно), точечные ретраи и понятные зависимости. Каждый шаг идемпотентен, поэтому повтор безопасен. Цена — больше оркестрации и накладных расходов на передачу данных между шагами, но на проде это окупается управляемостью.

  3. #pipeline_patterns3 / 5
    Зачем разделять extract, load и transform на отдельные шаги, а не делать всё в одном?
    A)Разделение нужно лишь для красоты архитектурной схемы, на практике всё можно делать в одном шаге
    B)Сохранённое сырьё → переиграть трансформацию без повторного extract
    C)Extract и transform обязаны идти вместе, иначе данные потеряют типы при сохранении в сыром виде
    D)Load должен идти строго перед extract, иначе трансформация не найдёт исходных данных
    показать ответ и разбор
    +B)Сохранённое сырьё → переиграть трансформацию без повторного extract

    // разбор: Слив всего в один шаг «выгрузил-преобразовал-записал» означает, что при любой правке логики или баге надо снова дёргать источник (нагрузка, лимиты, а иногда данные уже недоступны). Разделение на E, L, T сохраняет сырые данные (raw/bronze) отдельно, поэтому трансформацию (T) можно сколько угодно раз переигрывать на уже загруженном сырье, чинить и развивать логику без повторного extract. Это же даёт аудит (что реально пришло), разные витрины из одного сырья и рестарт с середины. Современный ELT именно поэтому грузит сырьё раньше трансформаций.

  4. #pipeline_patterns4 / 5
    Пайплайн должен стартовать, только когда апстрим доложил данные. Как выразить это надёжно?
    A)Просто поставить downstream-пайплайн по расписанию на пару часов после ожидаемого финиша апстрима
    B)Запускать downstream раньше апстрима и надеяться, что нужные данные подъедут по ходу выполнения
    C)Через готовность данных (sensor/датасет-триггер/зависимость), а не фиксированный сдвиг по времени от апстрима
    D)Объединить оба пайплайна в один гигантский шаг, чтобы вопрос готовности данных вообще не возникал
    показать ответ и разбор
    +C)Через готовность данных (sensor/датасет-триггер/зависимость), а не фиксированный сдвиг по времени от апстрима

    // разбор: Ставить свой пайплайн на «через 2 часа после апстрима» хрупко: апстрим задержался — считаешь по неполным/старым данным. Надёжно — гейт по факту готовности: сенсор ждёт появления файла/партиции/маркера, dataset-триггер запускает потребителя при обновлении датасета, либо явная межзадачная/меж-DAG зависимость. Так downstream стартует ровно когда данные реально готовы, а не когда «обычно» готовы. Дополняют freshness/SLA-проверкой, чтобы не ждать вечно и поднять алерт при опоздании апстрима.

  5. #pipeline_patterns5 / 5
    Зачем каждый прогон пайплайна логируют с run id, метриками и происхождением данных (lineage)?
    A)Логи прогонов нужны для юридического аудита и в ежедневной эксплуатации не применяются
    B)Достаточно логировать факт успеха или падения, метрики и происхождение данных избыточны
    C)Lineage данных можно восстановить вручную, поэтому фиксировать его при прогоне не нужно
    D)Для наблюдаемости и отладки: видно, что и когда обработано, откуда данные, где сбой — и можно точечно переиграть
    показать ответ и разбор
    +D)Для наблюдаемости и отладки: видно, что и когда обработано, откуда данные, где сбой — и можно точечно переиграть

    // разбор: Без следов прогон — чёрный ящик: непонятно, обработалось ли, сколько строк, откуда данные и почему цифры разъехались. Логирование run id, метрик (объём, длительность, DQ-результаты) и lineage (какие входы породили какой выход) даёт наблюдаемость: быстро найти сбойный прогон, понять радиус поражения бага (какие витрины затронуты), провести аудит и точечно переналить только пострадавшие партиции. Это фундамент эксплуатации: алерты, разбор инцидентов, backfill-трекинг. Особенно ценно, когда данные из пайплайна используют для решений и отчётности.

дальше

Теорию прочитали. Навык ставится повторением

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