5 мин на чтение

Airflow перед собеседованием: чеклист основных понятий

Тут давеча появилась возможность актуализировать знания по Apache Airflow. Заодно собрал чеклист терминов и определений, которые полезно повторить перед собеседованием.

Это именно опорный конспект: каждый пункт можно превратить в отдельную тему с примерами, но для первого прохода важнее увидеть общую картину и не путать DAG, Task и Task Instance.

Базовые сущности

1. Что такое Airflow

Apache Airflow — платформа для разработки, планирования и мониторинга пакетных workflow. Её часто используют для ETL/ELT-процессов, но одним ETL применение не ограничивается.

Workflow описывается кодом, а Airflow берёт на себя расписание, зависимости, повторные попытки, состояние запусков, логи и ручное управление через UI. При этом сами данные через Airflow обычно не «текут»: задачи читают и записывают их во внешние системы.

2. Что такое DAG

DAG — Directed Acyclic Graph, ориентированный ациклический граф. В нём определены задачи и зависимости между ними: что можно выполнять параллельно, а что должно дождаться предыдущего шага.

В Airflow определение DAG обычно пишется на Python. Важная оговорка: DAG описывает структуру workflow, но не является конкретным запуском этого workflow.

3. Что такое Task

Task — логическая единица работы внутри DAG. Например: выполнить Python-код, отправить SQL-запрос, дождаться файла или запустить другой DAG.

Task часто создают оператором вроде PythonOperator(...) или декоратором @task, поэтому её легко спутать с обычной функцией. Разница в том, что функция содержит бизнес-логику, а Task добавляет к ней семантику Airflow: зависимости, ретраи, таймаут, состояние и другие параметры выполнения.

4. Что такое Task Instance

Task Instance — конкретный экземпляр Task в рамках определённого DAG Run. Одна и та же Task запускается много раз по расписанию, и каждый такой запуск получает собственное состояние: scheduled, queued, running, success, failed, skipped и так далее.

Именно Task Instance мы видим в UI и перезапускаем после ошибки. Если одна задача упала, не обязательно повторять весь pipeline: можно устранить причину и очистить состояние нужного экземпляра с учётом его зависимостей.

Интерфейс Apache Airflow с графом DAG и задачами, завершившимися с разными статусами
Граф DAG в интерфейсе Airflow: зависимости и состояния отдельных задач. Источник: документация Apache Airflow.

Как Airflow выполняет задачи

5. Основные компоненты

  • UI / API Server — интерфейс для просмотра DAG, запусков, логов и ручных операций.
  • Scheduler — анализирует расписание и зависимости, создаёт DAG Run и переводит готовые Task Instance к выполнению.
  • Executor — определяет, каким способом и где запускать задачи.
  • Worker — непосредственно выполняет код задачи. Его устройство зависит от выбранного executor.
  • Metadata Database — хранит служебное состояние Airflow.

В распределённой установке также можно встретить DAG Processor, который разбирает файлы DAG, и Triggerer, который обслуживает отложенные задачи без занятого worker-слота.

6. Executors

Executor — настройка, которая определяет механизм исполнения Task Instance. Три названия, которые стоит знать:

  • LocalExecutor выполняет задачи локальными процессами;
  • CeleryExecutor распределяет их по Celery workers через брокер сообщений;
  • KubernetesExecutor запускает отдельный pod для каждого экземпляра задачи.

Выбор зависит от масштаба, инфраструктуры, требований к изоляции и стоимости эксплуатации. Сам executor не определяет порядок задач — этим занимается Scheduler на основании DAG и их состояний.

7. Жизненный цикл выполнения

Упрощённо цепочка выглядит так:

dag.py → DAG Processor → Scheduler → Task Instance → Executor → Worker
                                                      ↓
                                                Metadata DB

Metadata DB участвует не только в конце: компоненты постоянно читают и обновляют в ней служебное состояние.

8. Что хранится в Metadata DB

Это основная база Airflow — обычно PostgreSQL или MySQL. В ней хранятся сведения о DAG Run, Task Instance, подключениях, переменных, пулах и других объектах оркестратора.

Сами обрабатываемые таблицы, файлы и датасеты туда складывать не следует. Metadata DB хранит состояние работы Airflow, а данные ETL остаются в хранилищах, объектных бакетах и внешних БД.

Зависимости и параметры запуска

9. Как задаются зависимости

Зависимости можно описать методами set_downstream() и set_upstream(), но чаще используется более читаемый синтаксис >> и <<:

source_ready >> extract_marker >> quality_marker >> merge >> validate
task1 >> [task2, task3]
[task1, task2] >> task3

Запись [a, b] >> c означает, что у c два upstream-родителя. Задача c всё равно будет запущена один раз в рамках текущего DAG Run. С правилом trigger_rule="all_success" по умолчанию Scheduler дождётся успешного завершения и a, и b.

10. Ретраи, задержка и таймаут

Эти параметры можно задавать задаче напрямую или передавать через default_args:

from datetime import timedelta

default_args = {
    "retries": 3,
    "retry_delay": timedelta(minutes=5),
    "execution_timeout": timedelta(minutes=30),
}
  • retries=3 означает три повторные попытки после первого запуска, то есть максимум четыре попытки;
  • retry_delay задаёт паузу между ними;
  • execution_timeout ограничивает длительность выполнения Task Instance: после превышения лимита задача завершается с ошибкой и может перейти к ретраю.

11. Расписание и catchup

Параметр schedule задаёт периодичность запуска:

with DAG(
    dag_id="daily_orders",
    schedule="@daily",
    catchup=False,
):
    ...

catchup=True заставляет Scheduler создать пропущенные интервалы между start_date и текущим временем. При catchup=False исторические интервалы не догоняются — планируются только актуальные новые запуски.

Важно помнить про интервальную модель Airflow: ежедневный запуск обычно планируется после завершения соответствующего суточного интервала, а не в момент его начала.

12. Основные пресеты расписания

Пресет Cron-эквивалент Значение
@once один запуск
@hourly 0 * * * * каждый час
@daily 0 0 * * * каждый день
@weekly 0 0 * * 0 каждую неделю
@monthly 0 0 1 * * каждый месяц
@quarterly 0 0 1 */3 * раз в квартал
@yearly 0 0 1 1 * раз в год

Фактическое время зависит от часового пояса DAG, поэтому на собеседовании полезно упомянуть timezone и границы data interval.

Обмен данными и типы задач

13. XCom

XCom — механизм обмена небольшими сериализуемыми значениями между задачами. В TaskFlow API возвращаемое значение автоматически сохраняется в XCom и передаётся следующей задаче:

from airflow.decorators import task

@task
def extract():
    return "orders.csv"

@task
def load(filename: str):
    print(filename)

load(extract())

Большие DataFrame или файлы через стандартный XCom передавать не стоит. Практичнее сохранить данные во внешнем хранилище, а через XCom передать небольшой идентификатор, URI или путь.

14. Sensor

Sensor — разновидность Task, которая ожидает выполнения условия: появления файла, готовности внешней задачи, записи в базе или ответа сервиса.

Обычный sensor в режиме poke удерживает worker-слот между проверками. Режим reschedule или deferrable-оператор освобождает ресурс на время ожидания — это важное отличие для загруженной системы.

15. Какие ещё бывают Task

Операторов и провайдеров много, но для базового разговора хватит нескольких примеров:

  • PythonOperator — выполнить Python-функцию;
  • EmptyOperator — создать пустой узел, например для группировки веток;
  • ShortCircuitOperator — пропустить downstream-задачи, если условие ложно;
  • TriggerDagRunOperator — запустить другой DAG;
  • SQLExecuteQueryOperator — выполнить SQL через настроенное подключение;
  • EmailOperator — отправить письмо.

Итог

Минимальная картина такая: DAG описывает граф, Task — узел графа, Task Instance — конкретное выполнение этого узла, Scheduler решает, когда задача готова, а Executor определяет способ её запуска. Состояние всего процесса хранится в Metadata DB, небольшие значения между задачами передаются через XCom.

Если эти связи не путаются, дальше уже проще обсуждать pools, branching, dynamic task mapping, datasets, deferrable operators и выбор executor под конкретную инфраструктуру.