DAG Design: декомпозиция job, зависимости и параллелизм в Airflow

Принципы проектирования production Airflow DAG: атомарность, идемпотентность, декомпозиция по Medallion, Critical Path, Pools, Dynamic Task Mapping, cross-DAG зависимости, антипаттерны

streaming

Введение: DAG - это архитектурный документ

Начинающие data engineers воспринимают DAG как «скрипт с зависимостями». Опытные - как архитектурный документ, который описывает, как движутся данные через систему, какие гарантии надёжности обеспечивает пайплайн и как он восстанавливается после сбоев.

Плохо спроектированный DAG - это не просто «некрасивый код». Это источник:

  • Cascading failures: падение одной задачи блокирует весь пайплайн
  • Resource starvation: 20 Spark-сессий запускаются одновременно и убивают кластер
  • Невозможность selective rerun: чтобы перезапустить шаг 3, нужно перезапускать с шага 1
  • DAG spaghetti: зависимости превращаются в нечитаемую сеть, в которой никто не ориентируется
  • Непредсказуемое время завершения: пайплайн заканчивается в 7:00 вместо 3:00

Этот урок - о том, как проектировать DAG правильно: от первых принципов до production-кейса.


Часть 1. Философия проектирования DAG для Big Data

1.1. Три принципа надёжного DAG

Хороший DAG строится на трёх инженерных принципах. Они не специфичны для Airflow - они применимы к любой системе обработки данных.

Принцип 1: Атомарность (Atomicity)

Каждая задача в DAG должна делать одну логически завершённую операцию. «Скачать данные из источника А, очистить их и загрузить в Silver» - это три операции, а не одна.

Почему важно: если монолитная задача упала на этапе загрузки в Silver после 2 часов работы, при повторном запуске она скачает данные снова (лишние расходы), очистит снова (лишнее время) и только потом попробует загрузить. С атомарными задачами retry начинается ровно с того шага, где произошёл сбой.

Принцип 2: Идемпотентность (Idempotency)

Повторный запуск задачи с теми же параметрами должен давать тот же результат без побочных эффектов. Это единственная гарантия безопасного retry.

Нарушение идемпотентности: INSERT INTO table SELECT ... без DELETE перед этим - при повторе вставит дубли. Идемпотентный вариант: INSERT OVERWRITE PARTITION ... или MERGE INTO.

Принцип 3: Разделение вычислений (Separation of Compute)

Тяжёлые трансформации выполняются на Spark-кластере. Airflow Worker только инициирует задачи и отслеживает статус - не вычисляет ничего сам. Даже лёгкие операции (например, проверка количества строк) лучше выполнять в SQL через специализированный оператор, а не в Python-коде на воркере.

1.2. Эволюция от cron к event-driven оркестрации

Понимание эволюции помогает осознать, почему современный DAG устроен именно так:

Поколение Инструмент Проблема
1G: cron */5 * * * * script.sh Нет зависимостей, нет retry, нет мониторинга
2G: скриптовые цепочки script_a.sh && script_b.sh Нет параллелизма, один сбой = всё упало
3G: Airflow / Luigi DAG с зависимостями Scheduling-driven (по времени)
4G: Dataset-driven Airflow Datasets, Prefect Event-driven (по готовности данных)

Современный production-DAG сочетает элементы 3G и 4G: есть расписание (cron), но есть и механизмы ожидания внешних данных (Sensors, Datasets).


Часть 2. Декомпозиция Spark Job на задачи Airflow

2.1. Паттерн декомпозиции по Medallion Architecture

Medallion Architecture (Bronze → Silver → Gold) - естественная модель для декомпозиции DAG. Каждый слой - это граница идемпотентности: если задача Silver упала, Bronze уже записан и безопасен. Retry начинается с Silver.

Обратите внимание на несколько ключевых паттернов в этой схеме:

  1. Fan-out на Bronze: три источника извлекаются параллельно (нет зависимости друг от друга)
  2. Data Quality Gates: перед переходом к следующему слою - обязательная проверка
  3. Fail-Fast: если validate_bronze_a упал - Silver для source_a не запускается
  4. Fan-out на Gold: три Gold-задачи независимы, запускаются параллельно

2.2. Data Quality Gates: стратегия Fail-Fast

Встраивание проверок качества между слоями - один из важнейших паттернов. Без него испорченные данные могут незаметно добраться до Gold-слоя и отчётов.

"""
data_quality_gate.py
Пример реализации DQ-задачи через BranchPythonOperator.
"""
from airflow.operators.python import BranchPythonOperator, PythonOperator
from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator
from airflow.operators.empty import EmptyOperator
from airflow import DAG
from datetime import datetime

def check_bronze_quality(**context) -> str:
    """
    Проверяем качество Bronze-данных через Spark SQL запрос.

    Возвращает ID следующей задачи в зависимости от результата:
    - 'transform_silver' если данные прошли проверку
    - 'handle_dq_failure'  если качество недостаточное

    XCom используется только для передачи результата проверки (не данных!).
    """
    from airflow.providers.apache.spark.hooks.spark_sql import SparkSqlHook

    date = context["ds"]

    # SQL-запрос к Bronze таблице (выполняется Spark SQL)
    hook = SparkSqlHook(conn_id="spark_default")
    result = hook.run(f"""
        SELECT
            COUNT(*) AS total_rows,
            SUM(CASE WHEN event_id IS NULL THEN 1 ELSE 0 END) AS null_keys,
            COUNT(DISTINCT event_date) AS distinct_dates,
            MIN(event_date) AS min_date,
            MAX(event_date) AS max_date
        FROM catalog.bronze.events
        WHERE event_date = '{date}'
    """)

    total_rows = result[0]["total_rows"]
    null_keys = result[0]["null_keys"]

    # Правила качества
    if total_rows < 1_000_000:  # Ожидаем минимум 1M событий
        context["ti"].xcom_push(key="dq_failure_reason",
                                value=f"Too few rows: {total_rows} < 1,000,000")
        return "handle_dq_failure"

    if null_keys / total_rows > 0.01:  # Больше 1% NULL ключей
        context["ti"].xcom_push(key="dq_failure_reason",
                                value=f"Too many nulls: {null_keys / total_rows:.1%}")
        return "handle_dq_failure"

    # Данные прошли проверку
    context["ti"].xcom_push(key="dq_stats",
                            value={"total_rows": total_rows, "null_rate": null_keys / total_rows})
    return "transform_silver"


def send_dq_alert(**context):
    """Отправляет алерт в Slack при провале DQ-проверки."""
    reason = context["ti"].xcom_pull(key="dq_failure_reason", task_ids="check_bronze_quality")
    date = context["ds"]
    # В реальной системе: SlackOperator или отправка через webhook
    print(f"[DQ ALERT] Date={date}: {reason}")


with DAG("etl_with_dq_gate", start_date=datetime(2024, 1, 1), catchup=False) as dag:

    extract = SparkSubmitOperator(
        task_id="extract_bronze",
        application="s3://scripts/extract.py",
        conn_id="spark_default",
        application_args=["--date", "{{ ds }}"],
    )

    dq_check = BranchPythonOperator(
        task_id="check_bronze_quality",
        python_callable=check_bronze_quality,
        provide_context=True,
    )

    transform = SparkSubmitOperator(
        task_id="transform_silver",
        application="s3://scripts/transform.py",
        conn_id="spark_default",
        application_args=["--date", "{{ ds }}"],
    )

    dq_failure = PythonOperator(
        task_id="handle_dq_failure",
        python_callable=send_dq_alert,
        provide_context=True,
    )

    end = EmptyOperator(task_id="end", trigger_rule="none_failed_min_one_success")

    extract >> dq_check >> [transform, dq_failure]
    transform >> end
    dq_failure >> end

2.3. Гранулярность задач: trade-off

Не существует универсального правила «как мелко дробить». Решение зависит от конкретного пайплайна:

Слишком крупные задачи Слишком мелкие задачи
При сбое - долгий retry Overhead Airflow Scheduler на планирование сотен задач
Плохая observability (где именно упало?) DAG spaghetti - нечитаемая сеть зависимостей
Нельзя параллелизировать независимые части Затруднён анализ end-to-end времени
Большой Impact при повторной работе Сложность поддержки

Практическое правило: задача должна соответствовать одной data boundary - естественной точке, в которой данные записываются на диск (в S3/HDFS). Например: скачать один источник → одна задача. Обработать один день данных → одна задача (если это атомарная операция). Одна DQ-проверка → одна задача.


Часть 3. Data Dependencies vs Execution Dependencies

3.1. Ключевое различие

Начинающие инженеры путают два типа зависимостей:

Execution dependency (task_a >> task_b): «задача B не должна запуститься до завершения задачи A». Это механическое ограничение Airflow.

Data dependency: «задача B читает данные, которые записала задача A». Это семантическая связь между данными.

Проблема: часто инженеры добавляют execution dependency там, где реальной data dependency нет. Это искусственно сериализует пайплайн и увеличивает время выполнения.

В правом варианте три загрузки запускаются параллельно (нет data dependency между ними), а build_silver ждёт все три. Общее время сократилось с 63 до 45 минут - на 28%.

3.2. Скрытые зависимости (Hidden Dependencies)

Скрытые зависимости - самая коварная проблема DAG design. Это зависимости, которых нет в коде DAG, но которые реально существуют:

Shared Tables: задача A и задача B читают из одной таблицы, в которую пишет задача C. Если запустить A и B без ожидания C - они прочитают устаревшие данные.

External Systems: задача A читает из API, который обновляется каждый час. Если запустить задачу раньше, чем API обновился - данные будут вчерашние.

Object Storage Mutations: задача A пишет файлы в S3, задача B читает оттуда. Если между ними нет execution dependency - B может прочитать неполные данные (запись ещё не завершилась).

# Выявление скрытых зависимостей: аудит data contracts
# Спросите себя о каждой паре задач:
# 1. Читает ли задача B данные, которые пишет задача A?
# 2. Читают ли обе задачи из общего источника, который мутирует?
# 3. Зависят ли задачи от одного внешнего API или сервиса?
# 4. Используют ли задачи общие Airflow Variables или Connections?

# Если ответ "да" хотя бы на один вопрос - нужна execution dependency

3.3. XCom: передача метаданных, не данных

XCom (Cross-Communication) - механизм Airflow для передачи небольших значений между задачами. Хранится в метабазе Airflow (PostgreSQL/MySQL).

XCom - только для метаданных. Никогда не для данных:

# АНТИПАТТЕРН: передача DataFrame через XCom
def bad_task(**context):
    df = spark.read.parquet("s3://data/users/")
    # Попытка сериализовать DataFrame в XCom - это катастрофа:
    # - Pickle DataFrame = сотни МБ в PostgreSQL
    # - Метабаза Airflow не предназначена для хранения данных
    # - Все другие задачи и планировщик будут работать медленнее
    context["ti"].xcom_push(key="users_df", value=df)  # НИКОГДА!


# ПРАВИЛЬНО: передаём только метаданные
def extract_task(**context):
    date = context["ds"]
    output_path = f"s3://silver/users/date={date}/"

    spark = create_spark_session()
    df = spark.read.parquet(f"s3://bronze/users/date={date}/")
    df.write.mode("overwrite").parquet(output_path)
    spark.stop()

    # Передаём только путь - строка в несколько байт
    context["ti"].xcom_push(key="output_path", value=output_path)
    context["ti"].xcom_push(key="row_count", value=df.count())


def transform_task(**context):
    # Получаем путь из XCom
    input_path = context["ti"].xcom_pull(
        key="output_path",
        task_ids="extract_task"
    )
    # Читаем данные из хранилища, а не из XCom
    df = spark.read.parquet(input_path)
    # ...

Что допустимо в XCom:

  • Пути к файлам (S3 URI, HDFS path)
  • Числовые метрики (количество строк, размер файла)
  • Временные метки (последний обработанный timestamp)
  • Имена партиций
  • Флаги статуса (bool)

Что нельзя в XCom:

  • DataFrames (любого размера)
  • Списки с тысячами элементов
  • JSON-объекты > 1 МБ
  • Бинарные данные

Часть 4. Параллелизм и оптимизация ресурсов

4.1. Fan-out и Fan-in паттерны

Fan-out: одна задача порождает несколько параллельных задач. Используется когда downstream-задачи независимы друг от друга.

Fan-in: несколько параллельных задач объединяются в одну. Используется для агрегации результатов нескольких источников.

4.2. Pools: контроль параллелизма

Pool в Airflow - это именованный семафор с ограниченным числом слотов. Задача, использующая pool, занимает слот при запуске и освобождает при завершении. Если все слоты заняты - задача ждёт в очереди.

Pools критично важны для предотвращения Resource Starvation - ситуации, когда слишком много Spark-задач одновременно конкурируют за ресурсы кластера.

# Создание pools через Airflow CLI
# airflow pools set spark_heavy_pool 3 "Тяжёлые Spark задачи: макс 3 одновременно"
# airflow pools set spark_light_pool 10 "Лёгкие Spark задачи: макс 10 одновременно"
# airflow pools set external_api_pool 5 "Запросы к внешнему API: макс 5 одновременно"

# Использование pool в задаче
from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator

# Тяжёлая задача - занимает 1 слот из spark_heavy_pool
# Одновременно может работать не более 3 таких задач
heavy_etl = SparkSubmitOperator(
    task_id="build_gold_aggregates",
    application="s3://scripts/gold_aggregates.py",
    conn_id="spark_default",
    num_executors=50,
    executor_memory="16g",
    pool="spark_heavy_pool",        # Ограничиваем параллелизм
    pool_slots=1,                   # Задача занимает 1 слот (по умолчанию)
    priority_weight=10,             # Приоритет в очереди (выше = важнее)
)

# Лёгкая задача (DQ-check) - не ограничивает тяжёлые задачи
dq_check = SparkSubmitOperator(
    task_id="validate_bronze",
    application="s3://scripts/validate.py",
    conn_id="spark_default",
    num_executors=5,
    executor_memory="4g",
    pool="spark_light_pool",        # Отдельный pool - не конкурирует с heavy
)

4.3. Конфигурация параллелизма на разных уровнях

Airflow предоставляет несколько уровней управления параллелизмом, и важно понимать иерархию:

# Уровень 1: Глобальный параллелизм Airflow
# airflow.cfg:
# [core]
# max_active_tasks_per_dag = 16  # Макс задач одного DAG одновременно
# parallelism = 32               # Макс задач всех DAG одновременно (глобально)

# Уровень 2: DAG-level параллелизм
with DAG(
    dag_id="daily_etl",
    max_active_tasks=8,       # Максимум 8 задач этого DAG одновременно
    max_active_runs=2,        # Максимум 2 активных DagRun одновременно
    # max_active_runs=1 критично для Catchup: не запускать 100 дней параллельно
) as dag:
    pass

# Уровень 3: Task-level параллелизм (через Pool)
task = SparkSubmitOperator(
    task_id="heavy_task",
    pool="spark_heavy_pool",  # Ограничено size пула
    pool_slots=2,             # Эта задача занимает 2 слота (особо тяжёлая)
)

4.4. Dynamic Task Mapping: параллельная обработка множества таблиц

Dynamic Task Mapping (Airflow 2.3+) позволяет создавать задачи динамически в runtime, без необходимости заранее знать их количество. Это ключевой инструмент для параллельной обработки однотипных источников.

"""
Паттерн: параллельная ingestion 50 таблиц из PostgreSQL.
Вместо 50 захардкоженных задач - одна задача с .expand().
"""

from airflow import DAG
from airflow.decorators import task
from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator
from datetime import datetime


# ── Версия с @task decorator ──────────────────────────────────────────────────

@task
def get_tables_to_process(date: str) -> list[dict]:
    """
    Возвращает список таблиц для обработки.
    В реальной системе: читается из конфиг-таблицы в PostgreSQL или YAML.

    Каждый элемент списка - параметры для одной задачи ingestion.
    """
    # Пример: из конфиг-таблицы Airflow или внешней БД
    return [
        {"table": "users",    "schema": "crm",   "priority": "high"},
        {"table": "orders",   "schema": "sales",  "priority": "high"},
        {"table": "products", "schema": "catalog","priority": "medium"},
        # ... ещё 47 таблиц ...
    ]


@task
def ingest_table(table_config: dict, date: str) -> dict:
    """
    Выполняет ingestion одной таблицы.
    Вызывается параллельно для каждого элемента из get_tables_to_process().

    ВАЖНО: здесь не выполняем сам Spark! Запускаем Spark-задачу через hook
    или используем SparkSubmitOperator (см. альтернативную версию).
    """
    table = table_config["table"]
    schema = table_config["schema"]

    import subprocess
    # В реальной системе: SparkSubmitHook или KubernetesPodHook
    result = subprocess.run(
        [
            "spark-submit",
            f"s3://scripts/ingest_table.py",
            "--schema", schema,
            "--table", table,
            "--date", date,
        ],
        capture_output=True, text=True
    )

    if result.returncode != 0:
        raise RuntimeError(f"Failed to ingest {schema}.{table}: {result.stderr}")

    return {"table": f"{schema}.{table}", "status": "success", "date": date}


@task
def validate_all_ingested(results: list[dict]) -> None:
    """
    Fan-in: получает результаты от всех параллельных задач ingestion.
    Проверяет, что все таблицы обработаны успешно.
    """
    failed = [r for r in results if r.get("status") != "success"]
    if failed:
        raise ValueError(f"Failed ingestions: {[r['table'] for r in failed]}")
    print(f"All {len(results)} tables ingested successfully")


with DAG(
    dag_id="parallel_multi_table_ingestion",
    start_date=datetime(2024, 1, 1),
    schedule_interval="0 2 * * *",
    catchup=False,
    max_active_tasks=12,  # Не более 12 задач одновременно
) as dag:

    tables = get_tables_to_process(date="{{ ds }}")

    # .expand() создаёт одну задачу на каждый элемент списка tables
    # Параллелизм ограничен max_active_tasks на уровне DAG
    ingested = ingest_table.expand(
        table_config=tables,
        date="{{ ds }}",          # Константа - передаётся всем задачам
    )

    validate_all_ingested(ingested)
# ── Версия с SparkSubmitOperator через .expand_kwargs() ───────────────────────

from airflow import DAG
from airflow.decorators import task
from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator

@task
def generate_spark_configs(date: str) -> list[dict]:
    """
    Генерирует конфигурации для SparkSubmitOperator.
    Каждый dict - параметры для одного запуска Spark.
    """
    tables = ["users", "orders", "products", "sessions", "events"]
    return [
        {
            "application_args": [
                "--table", table,
                "--date", date,
                "--source", f"s3://bronze/{table}/",
                "--target", f"catalog.silver.{table}",
            ],
        }
        for table in tables
    ]


with DAG("spark_multi_table_etl", start_date=datetime(2024, 1, 1), catchup=False) as dag:

    configs = generate_spark_configs("{{ ds }}")

    # SparkSubmitOperator с expand_kwargs: создаёт N задач динамически
    spark_tasks = SparkSubmitOperator.partial(
        task_id="ingest_table",
        application="s3://scripts/ingest_single_table.py",
        conn_id="spark_default",
        num_executors=10,
        executor_memory="8g",
        pool="spark_light_pool",
    ).expand_kwargs(configs)

Часть 5. Critical Path Analysis

5.1. Что такое Critical Path

Critical Path - самая длинная цепочка зависимых задач в DAG. Именно она определяет минимально возможное время выполнения пайплайна. Уменьшить время выполнения DAG можно только ускорив задачи на Critical Path или распараллелив их.

Пример анализа:

В этом примере:

  • Critical Path: extract_eventsvalidate_eventstransform_silverbuild_gold_report = 45 минут
  • Ветки users и products завершаются раньше - их ускорение не уменьшит общее время
  • Чтобы ускорить DAG, нужно ускорить transform_silver (20 минут - самая долгая задача на Critical Path)

5.2. Практика: оптимизация Critical Path

from datetime import datetime, timedelta
from airflow import DAG
from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator

# Версия 1: Неоптимальный DAG (последовательный)
with DAG("unoptimized_etl", start_date=datetime(2024, 1, 1), catchup=False) as dag_v1:
    t1 = SparkSubmitOperator(task_id="extract_all",    ...)  # 15 мин
    t2 = SparkSubmitOperator(task_id="validate",       ...)  # 5 мин (можно параллелить)
    t3 = SparkSubmitOperator(task_id="transform",      ...)  # 20 мин
    t4 = SparkSubmitOperator(task_id="load_gold",      ...)  # 5 мин
    # Итого: 15+5+20+5 = 45 мин (всё последовательно)
    t1 >> t2 >> t3 >> t4

# Версия 2: Оптимизированный DAG (параллелизм где возможно)
with DAG("optimized_etl", start_date=datetime(2024, 1, 1), catchup=False) as dag_v2:
    # Extraction параллельно по источникам
    ex_events   = SparkSubmitOperator(task_id="extract_events",   ...)  # 15 мин
    ex_users    = SparkSubmitOperator(task_id="extract_users",    ...)  # 10 мин
    ex_products = SparkSubmitOperator(task_id="extract_products", ...)  # 8 мин

    # Validation параллельно (независимые DQ-проверки)
    val_events   = SparkSubmitOperator(task_id="validate_events",   ...)  # 5 мин
    val_users    = SparkSubmitOperator(task_id="validate_users",    ...)  # 3 мин
    val_products = SparkSubmitOperator(task_id="validate_products", ...)  # 2 мин

    # Transform - основная нагрузка (здесь Critical Path)
    transform = SparkSubmitOperator(task_id="transform_silver",    ...)  # 20 мин

    # Gold можно частично параллелить
    gold_kpi      = SparkSubmitOperator(task_id="build_kpi",      ...)  # 5 мин
    gold_segments = SparkSubmitOperator(task_id="build_segments", ...)  # 4 мин

    ex_events   >> val_events
    ex_users    >> val_users
    ex_products >> val_products

    [val_events, val_users, val_products] >> transform
    transform >> [gold_kpi, gold_segments]
    # Итого: max(15+5, 10+3, 8+2) + 20 + max(5,4) = 20+20+5 = 45 мин
    # Но извлечение идёт параллельно → только самая долгая ветка (20 мин) блокирует

Часть 6. Идемпотентность и retry-safe архитектура

6.1. Checkpoint Boundaries: атомарные точки записи

Каждая задача DAG должна заканчиваться атомарной записью в хранилище. Это и есть Checkpoint Boundary - точка, с которой можно безопасно перезапустить при сбое.

def run_silver_transform(date: str, source_table: str, target_table: str) -> None:
    """
    Идемпотентная трансформация Bronze → Silver.

    Атомарность достигается через replaceWhere:
    - Перезаписываем только партицию за указанную дату
    - Если задача запустится повторно - получим тот же результат
    - Другие партиции не затрагиваются

    Checkpoint Boundary: файлы в S3 после выполнения этой функции.
    """
    spark = create_spark_session()

    df_raw = (
        spark.table(source_table)
        .filter(f"event_date = '{date}'")
    )

    df_clean = (
        df_raw
        .dropDuplicates(["event_id"])
        .filter(F.col("event_id").isNotNull())
        .withColumn("_processed_at", F.current_timestamp())
        .withColumn("_pipeline_date", F.lit(date))
    )

    # Идемпотентная запись: overwrite конкретной партиции
    (
        df_clean.write
        .format("iceberg")
        .mode("overwrite")
        .option("partitionOverwriteMode", "dynamic")
        .option("replaceWhere", f"event_date = '{date}'")
        .saveAsTable(target_table)
    )

    spark.stop()


def non_idempotent_transform(date: str) -> None:
    """
    АНТИПАТТЕРН: НЕ идемпотентная трансформация.
    При повторном запуске вставит дубли!
    """
    spark = create_spark_session()
    df = spark.table("bronze.events").filter(f"event_date = '{date}'")

    # APPEND без предварительного удаления → дубли при retry
    df.write.mode("append").saveAsTable("silver.events")  # ОПАСНО!

    spark.stop()

6.2. Partial Failures: как обрабатывать частичные сбои

В fan-out паттерне несколько задач выполняются параллельно. Что делать, если одна упала, а остальные - нет?

from airflow.utils.trigger_rule import TriggerRule

# Стратегия 1: Fail-Fast - останавливаем всё при любом сбое
# (по умолчанию: trigger_rule="all_success")
transform_silver = SparkSubmitOperator(
    task_id="transform_silver",
    # Запускается только если ВСЕ upstream задачи успешны
    trigger_rule=TriggerRule.ALL_SUCCESS,
    # ...
)

# Стратегия 2: Best-Effort - продолжаем даже при частичном сбое
# Используется когда источники независимы, и частичный результат допустим
build_gold = SparkSubmitOperator(
    task_id="build_gold_partial",
    # Запускается если хотя бы одна upstream задача успешна
    trigger_rule=TriggerRule.ONE_SUCCESS,
    # ...
)

# Стратегия 3: Cleanup - всегда выполняется (для уборки ресурсов)
cleanup = PythonOperator(
    task_id="cleanup_temp_files",
    # Выполняется в любом случае (и при успехе, и при сбое)
    trigger_rule=TriggerRule.ALL_DONE,
    python_callable=lambda: print("Cleaning up temp files..."),
)

# Стратегия 4: Conditional - пропускаем если upstream был пропущен
# Используется с BranchPythonOperator
conditional_step = SparkSubmitOperator(
    task_id="optional_enrichment",
    trigger_rule=TriggerRule.NONE_FAILED_MIN_ONE_SUCCESS,
    # ...
)

6.3. Catchup и Backfill: обработка исторических данных

Catchup - автоматическое создание DagRun'ов для пропущенных интервалов при запуске нового DAG или после длительного downtime.

# Catchup включён: при запуске DAG Airflow создаст DagRun за каждый день
# с start_date до сегодня. Для DAG с start_date = 2024-01-01 и 180 днями
# это 180 DagRun'ов → 180 одновременных Spark-задач → КЛАСТЕР УПАДЁТ!
with DAG(
    "daily_etl",
    start_date=datetime(2024, 1, 1),
    catchup=True,    # ОПАСНО без max_active_runs
    max_active_runs=3,  # Максимум 3 DagRun одновременно → безопасный catchup
) as dag:
    pass

# Для production: catchup=False, backfill запускается вручную
with DAG(
    "daily_etl_prod",
    start_date=datetime(2024, 1, 1),
    catchup=False,   # Не догонять пропущенные запуски автоматически
    max_active_runs=1,
) as dag:
    pass

# Ручной backfill через CLI:
# airflow dags backfill daily_etl_prod \
#     --start-date 2024-01-01 \
#     --end-date 2024-03-31 \
#     --max-active-runs 5    # Ограничиваем параллелизм при backfill

Часть 7. Resource-Aware DAG Design для Spark

7.1. Многоуровневая стратегия ресурсов

Правильный resource-aware DAG учитывает ограничения на всех уровнях: от кластера до конкретной задачи.

7.2. Priority Weights: приоритизация задач

Когда пул занят и задачи ждут в очереди, priority_weight определяет порядок их запуска:

from airflow.utils.weight_rule import WeightRule

# Высокий приоритет: критичные Gold-задачи, от которых зависят отчёты
gold_revenue = SparkSubmitOperator(
    task_id="build_revenue_report",
    pool="spark_heavy_pool",
    priority_weight=100,   # Первыми в очереди
    weight_rule=WeightRule.ABSOLUTE,  # Используем абсолютный вес
)

# Средний приоритет: Silver трансформации
silver_transform = SparkSubmitOperator(
    task_id="transform_events",
    pool="spark_heavy_pool",
    priority_weight=50,
)

# Низкий приоритет: исторические пересчёты (backfill)
historical_backfill = SparkSubmitOperator(
    task_id="recalculate_history",
    pool="spark_heavy_pool",
    priority_weight=1,   # Последними в очереди
)

7.3. Dynamic Allocation vs Fixed Executors

Для задач с переменной нагрузкой (например, понедельник vs воскресенье) используйте Dynamic Allocation вместо фиксированного числа executor'ов:

# Фиксированное количество executor'ов - неэффективно
# В понедельник (пиковая нагрузка): 20 executor'ов мало
# В воскресенье (минимальная нагрузка): 20 executor'ов - расточительство
fixed_etl = SparkSubmitOperator(
    task_id="fixed_etl",
    num_executors=20,   # Всегда 20 - и в пик, и в спад
)

# Dynamic Allocation - масштабируется автоматически
dynamic_etl = SparkSubmitOperator(
    task_id="dynamic_etl",
    conf={
        "spark.dynamicAllocation.enabled":            "true",
        "spark.dynamicAllocation.shuffleTracking.enabled": "true",
        "spark.dynamicAllocation.minExecutors":       "5",    # Минимум
        "spark.dynamicAllocation.maxExecutors":       "50",   # Максимум
        "spark.dynamicAllocation.initialExecutors":   "10",   # Стартовое значение
        # Убивать простаивающий executor через 60 секунд
        "spark.dynamicAllocation.executorIdleTimeout": "60s",
        # Добавлять executor если задача ждёт более 1 секунды
        "spark.dynamicAllocation.schedulerBacklogTimeout": "1s",
    },
)

Часть 8. Cross-DAG оркестрация

8.1. Проблема межпайплайновых зависимостей

В enterprise lakehouse редко существует один DAG. Обычно это десятки взаимосвязанных пайплайнов: Bronze DAG для каждого источника, общий Silver DAG, несколько Gold DAG для разных команд. Когда Gold DAG ждёт данные от трёх Bronze DAG из разных источников - нужен механизм coordination.

8.2. ExternalTaskSensor: ожидание завершения другого DAG

"""
Silver DAG, ожидающий завершения нескольких Bronze DAG.
"""

from datetime import datetime, timedelta
from airflow import DAG
from airflow.sensors.external_task import ExternalTaskSensor
from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator

with DAG(
    dag_id="silver_unified_pipeline",
    start_date=datetime(2024, 1, 1),
    schedule_interval="30 4 * * *",  # Запускается в 4:30, ждёт данных из Bronze
    catchup=False,
    max_active_runs=1,
) as dag:

    # Ждём завершения Bronze CRM за тот же день
    wait_crm = ExternalTaskSensor(
        task_id="wait_bronze_crm",
        external_dag_id="bronze_crm_pipeline",
        external_task_id=None,   # None = ждём завершения всего DagRun (не одной задачи)
        execution_delta=timedelta(hours=1),  # CRM DAG запускается в 3:30, мы - в 4:30
        mode="reschedule",   # Reschedule mode: не блокирует воркер, пока ждёт
        timeout=3600,        # Таймаут ожидания: 1 час
        poke_interval=60,    # Проверяем каждые 60 секунд
        allowed_states=["success"],    # Считаем готовым только при SUCCESS
        failed_states=["failed", "upstream_failed"],  # Провалить если upstream упал
    )

    # Ждём Bronze ERP (запускается ежедневно в 3:00, мы - в 4:30)
    wait_erp = ExternalTaskSensor(
        task_id="wait_bronze_erp",
        external_dag_id="bronze_erp_pipeline",
        external_task_id=None,
        execution_delta=timedelta(hours=1, minutes=30),  # ERP в 3:00, мы в 4:30
        mode="reschedule",
        timeout=7200,   # 2 часа - ERP может быть медленным
        poke_interval=120,
        allowed_states=["success"],
        failed_states=["failed"],
    )

    # Ждём Bronze Events (последний запуск до нашего старта)
    wait_events = ExternalTaskSensor(
        task_id="wait_bronze_events",
        external_dag_id="bronze_events_pipeline",
        external_task_id="load_to_bronze",  # Ждём конкретную задачу
        # execution_date_fn: кастомная функция для нестандартного маппинга дат
        execution_date_fn=lambda dt: dt.replace(minute=15),  # Последний запуск в XX:15
        mode="reschedule",
        timeout=1800,
        poke_interval=60,
        allowed_states=["success"],
    )

    # Трансформация запускается только когда все источники готовы
    transform = SparkSubmitOperator(
        task_id="transform_unified_silver",
        application="s3://scripts/silver_unified.py",
        conn_id="spark_default",
        num_executors=30,
        executor_memory="12g",
        application_args=["--date", "{{ ds }}"],
        pool="spark_heavy_pool",
    )

    [wait_crm, wait_erp, wait_events] >> transform

8.3. Dataset-Driven Scheduling (Airflow 2.4+)

Airflow 2.4 ввёл Datasets - event-driven альтернативу ExternalTaskSensor. DAG запускается не по расписанию, а когда другой DAG обновил определённый dataset.

"""
Dataset-driven orchestration: Silver DAG запускается автоматически
при обновлении Bronze datasets.
"""

from airflow import DAG, Dataset
from airflow.decorators import task
from datetime import datetime

# Определяем datasets (URI может быть любым идентификатором)
BRONZE_EVENTS_DS = Dataset("s3://bronze/events/")
BRONZE_USERS_DS  = Dataset("s3://bronze/users/")
SILVER_EVENTS_DS = Dataset("s3://silver/events/")

# Bronze DAG: продьюсер datasets
with DAG(
    dag_id="bronze_events_producer",
    start_date=datetime(2024, 1, 1),
    schedule_interval="*/15 * * * *",   # Каждые 15 минут
) as bronze_dag:

    @task(outlets=[BRONZE_EVENTS_DS])   # Объявляем что задача обновляет dataset
    def load_bronze_events(date: str):
        # Загрузка событий в Bronze
        print(f"Loading events for {date}")
        # ...после записи в S3 Airflow автоматически помечает dataset как updated

    load_bronze_events("{{ ds }}")

# Silver DAG: консьюмер datasets
# Запускается автоматически когда ОБА dataset обновятся
with DAG(
    dag_id="silver_unified_consumer",
    start_date=datetime(2024, 1, 1),
    schedule=[BRONZE_EVENTS_DS, BRONZE_USERS_DS],  # Ждём оба dataset!
    catchup=False,
) as silver_dag:

    @task(outlets=[SILVER_EVENTS_DS])
    def transform_silver():
        print("Both Bronze datasets updated, starting Silver transform...")
        # ...трансформация

    transform_silver()

Преимущества Dataset-driven подхода:

  • Нет polling overhead (нет постоянных запросов к Airflow DB как у ExternalTaskSensor)
  • Семантически точнее: «запустись когда данные готовы», а не «запустись в 4:30 и надейся»
  • Автоматическая обработка задержек источников без изменения расписания

Часть 9. Антипаттерны и production-проблемы

9.1. Top-level Airflow код

Это один из самых разрушительных антипаттернов. Код вне операторов выполняется при каждом парсинге DAG - а Airflow Scheduler парсит DAG файлы каждые 30–60 секунд.

# АНТИПАТТЕРН: SparkSession или DB-соединение на уровне модуля
from pyspark.sql import SparkSession

# Это выполняется при КАЖДОМ парсинге DAG-файла планировщиком!
# Scheduler парсит файлы каждые 30 секунд → 2 SparkSession в минуту
# → Airflow Scheduler зависает, потребляет всю память
spark = SparkSession.builder.getOrCreate()  # НИКОГДА!
tables = spark.catalog.listTables("silver")  # НИКОГДА!
df = spark.read.parquet("s3://data/")       # НИКОГДА!

# АНТИПАТТЕРН: HTTP-запросы на уровне модуля
import requests
config = requests.get("http://config-service/pipeline-config").json()  # НИКОГДА!

# ПРАВИЛЬНО: всё внутри операторов или callable-функций
from airflow.operators.python import PythonOperator

def fetch_config(**context):
    import requests  # Импорт внутри функции
    config = requests.get("http://config-service/pipeline-config").json()
    context["ti"].xcom_push(key="config", value=config)

fetch_config_task = PythonOperator(
    task_id="fetch_config",
    python_callable=fetch_config,
)

9.2. Excessive Sensors (Sensoritis)

# АНТИПАТТЕРН: слишком много сенсоров с коротким poke_interval
# Каждый сенсор в poke_mode занимает слот воркера Airflow
# 50 сенсоров с poke_interval=30 = постоянный шторм запросов к БД

bad_sensor = ExternalTaskSensor(
    task_id="wait_upstream",
    external_dag_id="upstream_dag",
    mode="poke",            # ОПАСНО для долгих ожиданий!
    poke_interval=30,       # Каждые 30 секунд занимает воркер
    timeout=86400,          # 24 часа × 120 пробуждений/час = 2880 запросов к БД!
)

# ПРАВИЛЬНО: reschedule mode + разумный poke_interval
good_sensor = ExternalTaskSensor(
    task_id="wait_upstream",
    external_dag_id="upstream_dag",
    mode="reschedule",      # Освобождает воркер между проверками
    poke_interval=300,      # Проверяем каждые 5 минут
    timeout=7200,           # Таймаут 2 часа
)

# ЕЩЁ ЛУЧШЕ: Dataset-driven (Airflow 2.4+) - нет polling вообще

9.3. DAG Spaghetti: нечитаемые зависимости

# АНТИПАТТЕРН: 50+ задач в одном DAG без структуры
# Никто не понимает, что от чего зависит

with DAG("megadag") as dag:
    t1 = SparkSubmitOperator(task_id="extract_users", ...)
    t2 = SparkSubmitOperator(task_id="extract_events", ...)
    # ... ещё 48 задач ...
    t50 = SparkSubmitOperator(task_id="send_report", ...)

    t1 >> t5 >> t12 >> t23 >> t42 >> t50
    t2 >> t7 >> t15 >> t23
    t3 >> t7
    # ... ещё 50 строк зависимостей ...

# ПРАВИЛЬНО: Task Groups для логической группировки
from airflow.utils.task_group import TaskGroup

with DAG("structured_dag") as dag:

    with TaskGroup("bronze_extraction", tooltip="Извлечение из источников") as bronze_group:
        t_users   = SparkSubmitOperator(task_id="extract_users", ...)
        t_events  = SparkSubmitOperator(task_id="extract_events", ...)
        t_orders  = SparkSubmitOperator(task_id="extract_orders", ...)

    with TaskGroup("data_quality", tooltip="Проверки качества") as dq_group:
        dq_users  = SparkSubmitOperator(task_id="validate_users", ...)
        dq_events = SparkSubmitOperator(task_id="validate_events", ...)
        dq_orders = SparkSubmitOperator(task_id="validate_orders", ...)

    with TaskGroup("silver_transform", tooltip="Silver трансформации") as silver_group:
        s_events = SparkSubmitOperator(task_id="transform_events", ...)
        s_users  = SparkSubmitOperator(task_id="transform_users", ...)

    with TaskGroup("gold_marts", tooltip="Gold витрины") as gold_group:
        g_kpi     = SparkSubmitOperator(task_id="build_kpi", ...)
        g_funnel  = SparkSubmitOperator(task_id="build_funnel", ...)

    # Высокоуровневые зависимости - читаются как предложения
    bronze_group >> dq_group >> silver_group >> gold_group

9.4. Hardcoded Configuration

# АНТИПАТТЕРН: конфиги в коде DAG
task = SparkSubmitOperator(
    task_id="etl",
    application="hdfs://namenode:8020/scripts/etl.py",  # Жёстко! Изменить → деплой DAG
    num_executors=20,          # Не адаптируется к нагрузке
    executor_memory="8g",
    conf={"spark.master": "yarn"},  # Другой кластер → другой код
)

# ПРАВИЛЬНО: конфиги через Airflow Variables
from airflow.models import Variable

ETL_CONFIG = Variable.get("etl_spark_config", deserialize_json=True, default_var={
    "num_executors": 20,
    "executor_memory": "8g",
    "deploy_mode": "cluster",
    "application_path": "s3://scripts/etl.py",
})

task = SparkSubmitOperator(
    task_id="etl",
    application=ETL_CONFIG["application_path"],
    num_executors=ETL_CONFIG["num_executors"],
    executor_memory=ETL_CONFIG["executor_memory"],
    deploy_mode=ETL_CONFIG.get("deploy_mode", "cluster"),
)
# Изменить конфиг → обновить Variable в Airflow UI → без деплоя кода

9.5. Чек-лист production-ready DAG

"""
Чек-лист: перед деплоем DAG в production проверьте все пункты.
"""

CHECKLIST = {
    "atomicity": {
        "description": "Каждая задача делает одну логическую операцию",
        "verify": "Можно ли перезапустить задачу без влияния на другие?",
    },
    "idempotency": {
        "description": "Повторный запуск даёт тот же результат",
        "verify": "Используется OVERWRITE/MERGE вместо APPEND без деdup?",
    },
    "no_top_level_code": {
        "description": "Нет SparkSession/DB-соединений вне операторов",
        "verify": "Открыть файл и убедиться, что всё внутри callable/operators",
    },
    "max_active_runs": {
        "description": "max_active_runs=1 для production DAG",
        "verify": "Установлен max_active_runs?",
    },
    "pools_configured": {
        "description": "Тяжёлые Spark-задачи используют pools",
        "verify": "Все SparkSubmitOperator с num_executors>10 имеют pool?",
    },
    "retries_configured": {
        "description": "Настроены retries с exponential backoff",
        "verify": "retries >= 2, retry_delay >= 5 мин, retry_exponential_backoff=True?",
    },
    "execution_timeout": {
        "description": "Задачи имеют execution_timeout",
        "verify": "Каждая задача имеет execution_timeout для защиты от зависания?",
    },
    "no_xcom_data": {
        "description": "XCom используется только для метаданных",
        "verify": "Нет DataFrame/больших объектов в xcom_push()?",
    },
    "sensors_reschedule": {
        "description": "Сенсоры используют mode='reschedule'",
        "verify": "Нет ExternalTaskSensor/S3KeySensor с mode='poke' для длинных ожиданий?",
    },
    "catchup_disabled": {
        "description": "catchup=False для production DAG",
        "verify": "Установлен catchup=False или max_active_runs для безопасного catchup?",
    },
    "no_hardcoded_config": {
        "description": "Конфигурации из Airflow Variables",
        "verify": "Нет захардкоженных путей, адресов, паролей в коде DAG?",
    },
    "task_groups": {
        "description": "TaskGroup для DAG с более чем 10 задачами",
        "verify": "Логически связанные задачи сгруппированы в TaskGroup?",
    },
}

Часть 10. End-to-End Production Кейс: Daily Lakehouse Pipeline

10.1. Постановка задачи

Разработаем production DAG для ежедневного ETL-пайплайна аналитической платформы e-commerce:

  • 3 источника: CRM (PostgreSQL), Events (Kafka → S3), ERP (REST API)
  • Medallion Architecture: Bronze → Silver → Gold
  • SLA: Gold-данные должны быть готовы к 6:00 утра
  • Объём: ~500M событий/день, ~10M пользователей, ~2M заказов
  • Кластер: YARN, 100 executor'ов × 32 ГБ

10.2. Декомпозиция задач и критический путь

Critical Path: extract_events_s3 (30) → validate_events (8) → transform_events_dedup (35) → validate_silver_all (10) → build_daily_kpi (15) = 98 минут - DAG завершится в 3:38.

10.3. Полная реализация DAG

"""
dag_daily_lakehouse_pipeline.py
Production DAG: daily e-commerce lakehouse ETL.
Реализует Best Practices: pools, task groups, idempotency, DQ gates.
"""

from datetime import datetime, timedelta
from airflow import DAG
from airflow.decorators import task
from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator
from airflow.operators.python import BranchPythonOperator, PythonOperator
from airflow.operators.empty import EmptyOperator
from airflow.utils.task_group import TaskGroup
from airflow.utils.trigger_rule import TriggerRule
from airflow.models import Variable

# ── Конфигурация ──────────────────────────────────────────────────────────────

# Конфиги хранятся в Airflow Variables, не в коде
SPARK_CONFIGS = Variable.get("spark_lakehouse_configs", deserialize_json=True, default_var={
    "small":  {"num_executors": 5,  "executor_memory": "4g",  "executor_cores": 2},
    "medium": {"num_executors": 20, "executor_memory": "8g",  "executor_cores": 4},
    "large":  {"num_executors": 50, "executor_memory": "16g", "executor_cores": 8},
})

BASE_SPARK_CONF = {
    "spark.sql.adaptive.enabled":                    "true",
    "spark.sql.adaptive.coalescePartitions.enabled": "true",
    "spark.sql.extensions":
        "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions",
    "spark.sql.catalog.lakehouse": "org.apache.iceberg.spark.SparkCatalog",
}

SCRIPTS_BASE = "s3://company-scripts/etl/"
DQ_THRESHOLDS = Variable.get("dq_thresholds", deserialize_json=True, default_var={
    "min_events_per_day": 100_000_000,   # 100M событий минимум
    "max_null_rate":      0.01,           # Максимум 1% NULL ключей
    "min_orders_per_day": 500_000,        # 500K заказов минимум
})

# ── Вспомогательные функции ───────────────────────────────────────────────────

def make_spark_task(
    task_id: str,
    script: str,
    size: str,
    app_args: list,
    pool: str = "spark_medium_pool",
    **kwargs,
) -> SparkSubmitOperator:
    """Фабричная функция: создаёт SparkSubmitOperator с общими настройками."""
    cfg = SPARK_CONFIGS.get(size, SPARK_CONFIGS["medium"])
    return SparkSubmitOperator(
        task_id=task_id,
        application=f"{SCRIPTS_BASE}{script}",
        conn_id="spark_yarn_prod",
        deploy_mode="cluster",
        conf={**BASE_SPARK_CONF},
        pool=pool,
        execution_timeout=timedelta(hours=3),
        application_args=app_args,
        **cfg,
        **kwargs,
    )


def make_dq_branch(task_id: str, check_fn, pool: str = "spark_light_pool"):
    """Фабричная функция: создаёт BranchPythonOperator для DQ-проверки."""
    return BranchPythonOperator(
        task_id=task_id,
        python_callable=check_fn,
        provide_context=True,
        pool=pool,
        retries=0,  # DQ-проверки не ретраим - сразу падаем
    )


# ── DQ Functions ──────────────────────────────────────────────────────────────

def check_bronze_events(**context) -> str:
    """DQ-проверка Bronze Events: row count + null rate."""
    from airflow.providers.apache.spark.hooks.spark_sql import SparkSqlHook
    date = context["ds"]
    hook = SparkSqlHook(conn_id="spark_yarn_prod")
    row = hook.run(
        f"SELECT COUNT(*) AS n, SUM(CASE WHEN event_id IS NULL THEN 1 ELSE 0 END) AS nulls "
        f"FROM catalog.bronze.events WHERE event_date = '{date}'"
    )[0]
    total, nulls = row["n"], row["nulls"]
    if total < DQ_THRESHOLDS["min_events_per_day"]:
        context["ti"].xcom_push(key="dq_fail", value=f"events: {total} < min")
        return "dq_gate.handle_dq_failure"
    if total > 0 and nulls / total > DQ_THRESHOLDS["max_null_rate"]:
        context["ti"].xcom_push(key="dq_fail", value=f"events null_rate: {nulls/total:.2%}")
        return "dq_gate.handle_dq_failure"
    return "silver_transform.transform_events_dedup"


def check_silver_integrity(**context) -> str:
    """DQ-проверка Silver: ссылочная целостность events → users."""
    from airflow.providers.apache.spark.hooks.spark_sql import SparkSqlHook
    date = context["ds"]
    hook = SparkSqlHook(conn_id="spark_yarn_prod")
    # Проверяем: все user_id в events присутствуют в users
    orphan_row = hook.run(
        f"SELECT COUNT(*) AS orphans FROM catalog.silver.events e "
        f"LEFT ANTI JOIN catalog.silver.users u ON e.user_id = u.user_id "
        f"WHERE e.event_date = '{date}'"
    )[0]
    if orphan_row["orphans"] > 1000:
        context["ti"].xcom_push(key="dq_fail", value=f"orphan events: {orphan_row['orphans']}")
        return "gold_gate.handle_silver_dq_failure"
    return "gold_marts.build_daily_kpi"


def send_failure_alert(**context) -> None:
    """Отправляет алерт при провале DQ."""
    ti = context["task_instance"]
    reason = ti.xcom_pull(key="dq_fail") or "Unknown DQ failure"
    date = context["ds"]
    print(f"[ALERT] DQ FAILURE on {date}: {reason}")
    # В реальности: SlackWebhookOperator или PagerDuty API


# ── Основной DAG ──────────────────────────────────────────────────────────────

default_args = {
    "owner":                    "data-platform-team",
    "depends_on_past":          False,
    "start_date":               datetime(2024, 1, 1),
    "email":                    ["data-oncall@company.com"],
    "email_on_failure":         True,
    "email_on_retry":           False,
    "retries":                  3,
    "retry_delay":              timedelta(minutes=10),
    "retry_exponential_backoff": True,
    "max_retry_delay":          timedelta(hours=1),
}

with DAG(
    dag_id="daily_lakehouse_pipeline",
    default_args=default_args,
    description="Daily e-commerce lakehouse: Bronze → Silver → Gold",
    schedule_interval="0 2 * * *",   # Каждый день в 02:00 UTC
    catchup=False,
    max_active_runs=1,               # Только один активный запуск
    tags=["lakehouse", "daily", "production"],
) as dag:

    # ── Bronze Extraction (параллельно, Fan-out) ──────────────────────────────

    with TaskGroup("bronze_extraction", tooltip="Параллельное извлечение из источников") as bronze_group:

        extract_crm = make_spark_task(
            task_id="extract_crm_users",
            script="bronze/extract_crm.py",
            size="medium",
            app_args=["--date", "{{ ds }}", "--source", "postgresql://crm-db/users"],
            pool="spark_medium_pool",
        )

        extract_events = make_spark_task(
            task_id="extract_events_s3",
            script="bronze/extract_events.py",
            size="large",    # Самый большой источник → large
            app_args=["--date", "{{ ds }}", "--source", "s3://raw-events/"],
            pool="spark_heavy_pool",   # Тяжёлая задача - занимает слот в heavy pool
        )

        extract_erp = make_spark_task(
            task_id="extract_erp_orders",
            script="bronze/extract_erp.py",
            size="medium",
            app_args=["--date", "{{ ds }}", "--source", "https://erp-api/orders"],
            pool="spark_medium_pool",
        )
        # Задачи внутри группы независимы → запускаются параллельно

    # ── DQ Gate: Bronze ───────────────────────────────────────────────────────

    with TaskGroup("dq_gate", tooltip="Data Quality проверки Bronze данных") as dq_gate_group:

        # BranchPythonOperator: проверяет качество Events (критичный источник)
        dq_events_check = make_dq_branch(
            task_id="check_events_quality",
            check_fn=check_bronze_events,
        )

        dq_failure_handler = PythonOperator(
            task_id="handle_dq_failure",
            python_callable=send_failure_alert,
            provide_context=True,
            trigger_rule=TriggerRule.ONE_FAILED,
        )

        # DQ завершился (успешно или через failure handler)
        dq_end = EmptyOperator(
            task_id="dq_end",
            trigger_rule=TriggerRule.NONE_FAILED_MIN_ONE_SUCCESS,
        )

        dq_events_check >> [dq_failure_handler, dq_end]
        dq_failure_handler >> dq_end

    # ── Silver Transformation (частично параллельно) ──────────────────────────

    with TaskGroup("silver_transform", tooltip="Silver трансформации") as silver_group:

        # Events - на Critical Path, запускается первым (прошёл DQ)
        transform_events = make_spark_task(
            task_id="transform_events_dedup",
            script="silver/transform_events.py",
            size="large",
            app_args=["--date", "{{ ds }}", "--target", "catalog.silver.events"],
            pool="spark_heavy_pool",
        )

        # Users SCD Type 2 - независим от events
        transform_users = make_spark_task(
            task_id="transform_users_scd2",
            script="silver/transform_users.py",
            size="medium",
            app_args=["--date", "{{ ds }}", "--target", "catalog.silver.users"],
            pool="spark_medium_pool",
        )

        # Orders MERGE - независим от events
        transform_orders = make_spark_task(
            task_id="transform_orders_merge",
            script="silver/transform_orders.py",
            size="medium",
            app_args=["--date", "{{ ds }}", "--target", "catalog.silver.orders"],
            pool="spark_medium_pool",
        )

    # ── DQ Gate: Silver ───────────────────────────────────────────────────────

    with TaskGroup("gold_gate", tooltip="Проверка Silver перед Gold") as gold_gate_group:

        silver_dq_check = make_dq_branch(
            task_id="check_silver_integrity",
            check_fn=check_silver_integrity,
        )

        silver_dq_failure = PythonOperator(
            task_id="handle_silver_dq_failure",
            python_callable=send_failure_alert,
            provide_context=True,
        )

        silver_dq_end = EmptyOperator(
            task_id="silver_dq_end",
            trigger_rule=TriggerRule.NONE_FAILED_MIN_ONE_SUCCESS,
        )

        silver_dq_check >> [silver_dq_failure, silver_dq_end]
        silver_dq_failure >> silver_dq_end

    # ── Gold Marts (параллельно, Fan-out) ─────────────────────────────────────

    with TaskGroup("gold_marts", tooltip="Построение Gold витрин") as gold_group:

        build_kpi = make_spark_task(
            task_id="build_daily_kpi",
            script="gold/build_kpi.py",
            size="medium",
            app_args=["--date", "{{ ds }}", "--target", "catalog.gold.daily_kpi"],
            pool="spark_medium_pool",
        )

        build_segments = make_spark_task(
            task_id="build_user_segments",
            script="gold/build_segments.py",
            size="medium",
            app_args=["--date", "{{ ds }}", "--target", "catalog.gold.user_segments"],
            pool="spark_medium_pool",
        )

        build_funnel = make_spark_task(
            task_id="build_funnel_report",
            script="gold/build_funnel.py",
            size="small",
            app_args=["--date", "{{ ds }}", "--target", "catalog.gold.funnel_report"],
            pool="spark_light_pool",
        )

    # ── Завершение ────────────────────────────────────────────────────────────

    @task
    def notify_downstream(date: str) -> None:
        """
        Уведомляет downstream системы о готовности данных.
        Может использоваться как Dataset trigger для других DAG.
        """
        print(f"[OK] Daily lakehouse pipeline completed for {date}")
        # В реальности: отправка в Slack, обновление status-таблицы,
        # публикация в EventBridge для event-driven downstream DAG

    notify = notify_downstream("{{ ds }}")

    # ── Граф зависимостей ─────────────────────────────────────────────────────

    # Bronze → DQ Gate
    bronze_group >> dq_gate_group

    # DQ Gate (успешный end) → Silver трансформации
    # Каждый источник начинает трансформацию когда его DQ прошёл
    extract_crm    >> transform_users    # CRM → Users
    extract_erp    >> transform_orders   # ERP → Orders
    dq_gate_group  >> transform_events  # Events (через DQ) → Events transform

    # Silver → DQ Gate → Gold
    silver_group >> gold_gate_group >> gold_group >> notify

10.4. Selective Rerun: перезапуск только нужного шага

Одно из главных преимуществ правильно декомпозированного DAG - возможность перезапустить только упавший шаг:

# Перезапустить только Silver трансформацию за конкретный день
# (без повторного извлечения Bronze - данные уже на месте)
airflow tasks clear daily_lakehouse_pipeline \
    --task-regex "silver_transform.*" \
    --start-date 2024-03-15 \
    --end-date 2024-03-15 \
    --yes

# Перезапустить только Gold витрины
airflow tasks clear daily_lakehouse_pipeline \
    --task-regex "gold_marts.*" \
    --start-date 2024-03-15 \
    --end-date 2024-03-15 \
    --yes

# Перезапустить конкретную задачу
airflow tasks run daily_lakehouse_pipeline gold_marts.build_daily_kpi 2024-03-15 \
    --local

Это возможно потому, что:

  1. Каждый шаг идемпотентен - повторный запуск не создаёт дубли
  2. Каждый шаг атомарно записывает в хранилище - предыдущие шаги не нужно повторять
  3. Checkpoint Boundaries чётко определены - каждая задача читает из стабильного источника (Bronze/Silver)

Итоги урока

Хороший DAG - это не просто Python-файл с операторами. Это архитектурная модель, отражающая реальное движение данных через систему.

Ключевые принципы:

  1. Атомарность: одна задача - одна логическая операция. Это минимизирует стоимость retry.
  2. Идемпотентность: OVERWRITE/MERGE вместо APPEND. Это делает retry безопасным.
  3. Реальные зависимости: dependency graph отражает движение данных, а не удобство написания кода. Параллелизация там, где нет реальной data dependency.
  4. Data Quality Gates: проверки между слоями с Fail-Fast стратегией. Испорченные данные не должны добираться до Gold.
  5. Resource Awareness: Pools, priority_weight и max_active_runs защищают кластер от перегрузки.
  6. Critical Path: понимание узкого места позволяет сфокусировать оптимизацию там, где это реально влияет на SLA.
  7. No top-level code: ни SparkSession, ни HTTP-запросов, ни тяжёлых вычислений вне операторов.
  8. Cross-DAG coordination: ExternalTaskSensor (или Datasets в Airflow 2.4+) для связи между пайплайнами разных команд.