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

Airflow и DAG-и

Airflow: DAG и его семантика

Airflow планирует пайплайны как граф задач и ведёт состояние каждой: ретраи, SLA, алерты. Собес проверяет две вещи, на которых спотыкаются: смысл logical_date и требование атомарной идемпотентной задачи.

Стержень: запуск бежит за интервал данных, а logical_date это НАЧАЛО интервала, не время запуска.

// Формулировки: «что такое execution_date/logical_date?», «какой должна быть задача?», «что такое сенсор?».

DAG и logical_date

DAG (directed acyclic graph) - граф задач с зависимостями без циклов; Airflow планирует запуски по расписанию и хранит состояние каждой задачи, обеспечивая ретраи, SLA (service level agreement) и алерты.

Ключевая и коварная семантика: запуск привязан к интервалу данных, и logical_date (бывший execution_date) это начало интервала. Джоба за понедельник физически бежит во вторник ночью, когда данные понедельника уже собраны. Трактовать logical_date как «сейчас» - значит сдвинуть все данные на интервал.

// Отсюда правило: границы обрабатываемого периода берут из logical_date запуска, а не из now(). Тогда перезапуск за прошлую дату обработает именно её данные.

logical date
интервал данных, за который бежит запуск (не «сейчас»)
DAG
граф задач без циклов с расписанием и состоянием

Задача, оператор, сенсор

Задача - атомарная и идемпотентная единица ретрая: упала - перезапустилась без последствий. Две операции в одной задаче (truncate и insert) ломают это: truncate прошёл, insert упал, ретрай, а данных уже нет.

Оператор исполняет действие, сенсор ждёт условия (файла, партиции, чужого DAG'а). Сенсоры в режиме reschedule не жгут слот воркера, освобождая его на время ожидания.

// Данные между задачами не гоняют через XCom это канал метаданных, килобайты (пути, параметры). Большие данные передают через хранилище; DataFrame в XCom раздувает мета-БД и падает на сериализации.

sensor
задача-ожидание внешнего условия (reschedule не жжёт слот)
XCom
обмен метаданными между задачами - не транспорт данных

Парсинг, catchup, ретраи

Код DAG-файла планировщик парсит постоянно, каждые несколько десятков секунд. Тяжёлые импорты или запрос к БД на верхнем уровне файла кладут scheduler - вся реальная работа должна быть внутри задач, а не в теле модуля.

catchup и backfill догоняют пропущенные интервалы. Но catchup=True (по умолчанию в старых версиях), случайно оставленный на новом DAG'е, зальёт сотни запусков за прошлые годы.

// Ретраи с паузой (retries + retry_delay/backoff) - дефолт для всего сетевого, а on_failure_callback и алерты обязательны: падение должно быть замечено сразу, а не найдено через неделю по кривым цифрам.

catchup / backfill
автодогон пропущенных интервалов / ручной прогон истории
on_failure_callback
хук оповещения о падении задачи

Как отвечать: «Что такое logical_date и почему это ловушка?»

logical_date, раньше execution_date, это не время, когда задача запустилась, а начало интервала данных, за который она бежит. Airflow работает по интервалам: чтобы обработать данные за понедельник, запуск стартует, когда понедельник закончился, то есть во вторник, но его logical_date указывает на понедельник. Ловушка в том, что интуитивно кажется, будто это «сейчас», и люди пишут внутри задачи фильтр по now() или today(). Тогда всё ломается двумя способами: во-первых, обрабатывается не тот интервал - сегодняшние данные вместо вчерашних; во-вторых, теряется идемпотентность - перезапуск вчерашнего запуска сегодня возьмёт уже другой кусок данных. Правильно - брать границы периода строго из logical_date запуска, тогда задача за любую дату всегда обрабатывает именно её данные, и бэкфил, и ретрай дают тот же результат.

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

На чём валят

  • Трактовать logical_date как «время запуска» - данные сдвинуты на интервал.
  • Неатомарная задача: truncate прошёл, insert упал, ретрай - данных нет.
  • DataFrame через XCom - раздутая мета-БД и падения сериализации.
  • Запрос к БД на верхнем уровне DAG-файла - scheduler парсит его каждые 30 секунд.
  • Включённый catchup на новом DAG'е - сотни неожиданных запусков за прошлые годы.

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

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

  1. #airflow_dags1 / 5
    Чем задача (оператор) отличается от DAG?
    A)Задача и DAG — это полные синонимы, разницы между ними нет
    B)DAG — граф целиком; задача (оператор) — один узел-шаг внутри него
    C)Задача больше DAG и содержит его внутри себя как часть
    D)DAG исполняется прямо на воркере, а задача планирует запуск других DAG-графов
    показать ответ и разбор
    +B)DAG — граф целиком; задача (оператор) — один узел-шаг внутри него

    // разбор: DAG — это весь граф; оператор описывает один шаг (SQL-запрос, Python-функция, сенсор ожидания). Задачи — узлы графа, связанные зависимостями. Такое разделение позволяет переиспользовать операторы и рассуждать о пайплайне как о наборе идемпотентных шагов, а не одном большом скрипте.

  2. #airflow_dags2 / 5
    Почему задача должна опираться на logical_date интервала, а не на «сейчас»?
    A)Использовать текущее время wall-clock в момент фактического запуска задачи на сервере планировщика оркестратора Airflow
    B)Дата запуска вообще не важна для корректности пайплайна данных
    C)logical_date — период, за который считаем; так ретрай и бэкфилл повторяемы
    D)Дата нужна лишь для красивого именования логов и артефактов запуска
    показать ответ и разбор
    +C)logical_date — период, за который считаем; так ретрай и бэкфилл повторяемы

    // разбор: Задача считает за свой интервал (logical/execution date), а не «за момент запуска». Тогда повторный ретрай вчерашнего запуска и бэкфилл за прошлый месяц дают ровно тот же результат — обработка детерминирована и идемпотентна. Привязка к wall-clock ломает повторяемость: результат зависит от того, когда именно запустили.

  3. #airflow_dags3 / 5
    Почему тяжёлую работу нельзя выносить в тело DAG-файла (top-level code)?
    A)Тяжёлый код в теле DAG-файла ускоряет выполнение задач на воркерах
    B)В теле DAG-файла можно спокойно делать запросы к боевой базе данных
    C)Код верхнего уровня в DAG-файле исполняется ровно один раз за всё время при самом первом развёртывании DAG в планировщик задач
    D)Файл парсится планировщиком постоянно — тяжёлый top-level код тормозит всё
    показать ответ и разбор
    +D)Файл парсится планировщиком постоянно — тяжёлый top-level код тормозит всё

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

  4. #airflow_dags4 / 5
    Что такое пайплайн данных (data pipeline)?
    A)Один SQL-запрос, выполняемый в базе вручную по нажатию кнопки оператором
    B)Последовательность шагов, переносящих и преобразующих данные из источника в приёмник
    C)График, визуализирующий бизнес-метрики на дашборде для руководства
    D)Формат сжатого хранения таблиц на диске в колоночном виде
    показать ответ и разбор
    +B)Последовательность шагов, переносящих и преобразующих данные из источника в приёмник

    // разбор: Пайплайн данных — цепочка шагов, которая извлекает данные из источников, преобразует и грузит в приёмник (extract → transform → load), обычно по расписанию, с зависимостями, ретраями и проверками. Оркестратор (Airflow) управляет порядком и повторами, чего голый набор cron-скриптов надёжно не даёт.

  5. #airflow_dags5 / 5
    Что задаёт cron-расписание для задачи?
    A)Периодичность автозапуска задачи по времени (например, каждый день в 3:00)
    B)Точный объём выделяемой оперативной памяти, резервируемой под задачу при запуске
    C)Порядок колонок в результирующей таблице после трансформации
    D)Число повторных попыток задачи после её падения с ошибкой
    показать ответ и разбор
    +A)Периодичность автозапуска задачи по времени (например, каждый день в 3:00)

    // разбор: Cron-выражение задаёт, когда запускать задачу (минуты, часы, дни), но не что она делает и не сколько ресурсов ей дать. Оркестраторы поверх cron добавляют зависимости между задачами, catchup (досчёт пропущенных интервалов), ретраи и сенсоры — то, чего голый cron не умеет.

дальше

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

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