Паттерны пайплайнов данных
Как забирать данные - полностью или инкрементом, батчем или потоком, по расписанию или по событию это набор осознанных выборов, а не привычка. Собес проверяет, выбираешь ли ты паттерн по требованиям, а не по моде.
Стержень: инкремент дёшев, но требует 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, остальные разбираются в тренажёре.
- Чем оркестрация пайплайнов (Airflow) отличается от хореографии (событийной)?A)Оркестрация означает децентрализованную реакцию сервисов на события, а хореография — центрального дирижёраB)Это одно и то же: оба термина описывают запуск задач строго по фиксированному расписанию времениC)Хореография лучше оркестрации во всех сценариях, поэтому оркестраторы вроде Airflow устарелиD)Оркестрация — центральный дирижёр задаёт порядок и зависимости; хореография — сервисы реагируют на события друг друга без центра
показать ответ и разбор
+D)Оркестрация — центральный дирижёр задаёт порядок и зависимости; хореография — сервисы реагируют на события друг друга без центра// разбор: Оркестрация: центральный компонент (Airflow) явно знает граф зависимостей и дирижирует — запускает задачи в нужном порядке, следит за статусами, ретраит. Плюсы — прозрачность, единая точка контроля и мониторинга; минус — центральная зависимость и связанность. Хореография: нет дирижёра, каждый сервис реагирует на события других (шаг завершился → событие → следующий подхватил), связь через шину/события. Плюсы — слабая связанность, масштаб; минус — сложнее отследить сквозной поток и отладить. Батч-DWH тяготеет к оркестрации, событийные/микросервисные потоки — к хореографии.
- Почему пайплайн дробят на мелкие идемпотентные шаги, а не пишут одним большим скриптом?A)Ради рестартуемости, наблюдаемости, параллелизма и гранулярных ретраев: упал шаг — перезапускают его, а не всё с нуляB)Мелкие шаги работают медленнее одного большого скрипта, зато их проще читать глазамиC)Дробление нужно чтобы обойти ограничение на максимальную длину SQL-запроса в СУБДD)Один большой шаг предпочтительнее, потому что его не получится случайно перезапустить не полностью
показать ответ и разбор
+A)Ради рестартуемости, наблюдаемости, параллелизма и гранулярных ретраев: упал шаг — перезапускают его, а не всё с нуля// разбор: Один монолитный шаг при сбое на середине приходится перезапускать целиком (дорого, и он мог оставить полусостояние). Разбиение на мелкие идемпотентные задачи даёт: рестартуемость (перезапустить только упавший шаг с его входа), наблюдаемость (видно, где именно встало, метрики по шагам), параллелизм (независимые шаги идут одновременно), точечные ретраи и понятные зависимости. Каждый шаг идемпотентен, поэтому повтор безопасен. Цена — больше оркестрации и накладных расходов на передачу данных между шагами, но на проде это окупается управляемостью.
- Зачем разделять extract, load и transform на отдельные шаги, а не делать всё в одном?A)Разделение нужно лишь для красоты архитектурной схемы, на практике всё можно делать в одном шагеB)Сохранённое сырьё → переиграть трансформацию без повторного extractC)Extract и transform обязаны идти вместе, иначе данные потеряют типы при сохранении в сыром видеD)Load должен идти строго перед extract, иначе трансформация не найдёт исходных данных
показать ответ и разбор
+B)Сохранённое сырьё → переиграть трансформацию без повторного extract// разбор: Слив всего в один шаг «выгрузил-преобразовал-записал» означает, что при любой правке логики или баге надо снова дёргать источник (нагрузка, лимиты, а иногда данные уже недоступны). Разделение на E, L, T сохраняет сырые данные (raw/bronze) отдельно, поэтому трансформацию (T) можно сколько угодно раз переигрывать на уже загруженном сырье, чинить и развивать логику без повторного extract. Это же даёт аудит (что реально пришло), разные витрины из одного сырья и рестарт с середины. Современный ELT именно поэтому грузит сырьё раньше трансформаций.
- Пайплайн должен стартовать, только когда апстрим доложил данные. Как выразить это надёжно?A)Просто поставить downstream-пайплайн по расписанию на пару часов после ожидаемого финиша апстримаB)Запускать downstream раньше апстрима и надеяться, что нужные данные подъедут по ходу выполненияC)Через готовность данных (sensor/датасет-триггер/зависимость), а не фиксированный сдвиг по времени от апстримаD)Объединить оба пайплайна в один гигантский шаг, чтобы вопрос готовности данных вообще не возникал
показать ответ и разбор
+C)Через готовность данных (sensor/датасет-триггер/зависимость), а не фиксированный сдвиг по времени от апстрима// разбор: Ставить свой пайплайн на «через 2 часа после апстрима» хрупко: апстрим задержался — считаешь по неполным/старым данным. Надёжно — гейт по факту готовности: сенсор ждёт появления файла/партиции/маркера, dataset-триггер запускает потребителя при обновлении датасета, либо явная межзадачная/меж-DAG зависимость. Так downstream стартует ровно когда данные реально готовы, а не когда «обычно» готовы. Дополняют freshness/SLA-проверкой, чтобы не ждать вечно и поднять алерт при опоздании апстрима.
- Зачем каждый прогон пайплайна логируют с run id, метриками и происхождением данных (lineage)?A)Логи прогонов нужны для юридического аудита и в ежедневной эксплуатации не применяютсяB)Достаточно логировать факт успеха или падения, метрики и происхождение данных избыточныC)Lineage данных можно восстановить вручную, поэтому фиксировать его при прогоне не нужноD)Для наблюдаемости и отладки: видно, что и когда обработано, откуда данные, где сбой — и можно точечно переиграть
показать ответ и разбор
+D)Для наблюдаемости и отладки: видно, что и когда обработано, откуда данные, где сбой — и можно точечно переиграть// разбор: Без следов прогон — чёрный ящик: непонятно, обработалось ли, сколько строк, откуда данные и почему цифры разъехались. Логирование run id, метрик (объём, длительность, DQ-результаты) и lineage (какие входы породили какой выход) даёт наблюдаемость: быстро найти сбойный прогон, понять радиус поражения бага (какие витрины затронуты), провести аудит и точечно переналить только пострадавшие партиции. Это фундамент эксплуатации: алерты, разбор инцидентов, backfill-трекинг. Особенно ценно, когда данные из пайплайна используют для решений и отчётности.
дальше
Теорию прочитали. Навык ставится повторением
В Сеньорчике эта подтема идёт в ежедневных сессиях: движок возвращает её, пока ответы не станут уверенными, и ведёт прогресс отдельно по каждой подтеме. Теория внутри тоже бесплатна, лимит только на количество вопросов в день.