DAG Design: декомпозиция job, зависимости и параллелизм в Airflow
Принципы проектирования production Airflow DAG: атомарность, идемпотентность, декомпозиция по Medallion, Critical Path, Pools, Dynamic Task Mapping, cross-DAG зависимости, антипаттерны
Введение: 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.
Обратите внимание на несколько ключевых паттернов в этой схеме:
- Fan-out на Bronze: три источника извлекаются параллельно (нет зависимости друг от друга)
- Data Quality Gates: перед переходом к следующему слою - обязательная проверка
- Fail-Fast: если
validate_bronze_aупал - Silver для source_a не запускается - 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_events→validate_events→transform_silver→build_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
Это возможно потому, что:
- Каждый шаг идемпотентен - повторный запуск не создаёт дубли
- Каждый шаг атомарно записывает в хранилище - предыдущие шаги не нужно повторять
- Checkpoint Boundaries чётко определены - каждая задача читает из стабильного источника (Bronze/Silver)
Итоги урока¶
Хороший DAG - это не просто Python-файл с операторами. Это архитектурная модель, отражающая реальное движение данных через систему.
Ключевые принципы:
- Атомарность: одна задача - одна логическая операция. Это минимизирует стоимость retry.
- Идемпотентность: OVERWRITE/MERGE вместо APPEND. Это делает retry безопасным.
- Реальные зависимости: dependency graph отражает движение данных, а не удобство написания кода. Параллелизация там, где нет реальной data dependency.
- Data Quality Gates: проверки между слоями с Fail-Fast стратегией. Испорченные данные не должны добираться до Gold.
- Resource Awareness: Pools, priority_weight и max_active_runs защищают кластер от перегрузки.
- Critical Path: понимание узкого места позволяет сфокусировать оптимизацию там, где это реально влияет на SLA.
- No top-level code: ни SparkSession, ни HTTP-запросов, ни тяжёлых вычислений вне операторов.
- Cross-DAG coordination: ExternalTaskSensor (или Datasets в Airflow 2.4+) для связи между пайплайнами разных команд.