Airflow + Spark: SparkSubmitOperator vs KubernetesPodOperator

Оркестрация Spark-пайплайнов через Airflow: архитектурный выбор между SparkSubmitOperator и KubernetesPodOperator, управление зависимостями, мониторинг, retry-стратегии и production trade-offs

platform

Введение: почему оркестрация решает судьбу пайплайна

Представьте: вы написали отличный PySpark-скрипт, протестировали его локально, он работает. Но в production он должен:

  • Запускаться по расписанию (каждый час, каждый день в 3:00)
  • Повторно запускаться при сбое (retry при сетевой ошибке)
  • Не запускаться, если предыдущий шаг ещё не завершился (dependency management)
  • Уведомлять команду при падении (алерты)
  • Сохранять историю запусков для аудита (lineage, logs)
  • Останавливаться при достижении SLA (deadline awareness)

Ни одну из этих задач Spark не решает сам по себе. Это зона ответственности оркестратора - инструмента, который управляет жизненным циклом пайплайна: когда запустить, что передать в качестве параметров, куда отправить результат и что делать при ошибке.

Apache Airflow - де-факто стандарт оркестрации batch ETL в современных data platform. Он описывает пайплайны как DAG (Directed Acyclic Graph) на Python, позволяет настраивать зависимости между задачами, retry-стратегии, триггеры и мониторинг.

Ключевой вопрос этого урока: как правильно интегрировать Airflow и Spark? Выбор между SparkSubmitOperator и KubernetesPodOperator - это не вопрос синтаксиса, это архитектурное решение уровня всей платформы, влияющее на изоляцию, воспроизводимость, масштабирование и операционные расходы.


Часть 1. Архитектурный контекст: роль Airflow в Spark-пайплайнах

1.1. Separation of Concerns: оркестратор vs вычислительный движок

Фундаментальный принцип production data platform - разделение ответственности:

  • Airflow: решает когда и как запустить задачу, отслеживает статус, управляет зависимостями
  • Spark: решает что вычислить и как распределить вычисления по кластеру

Эти роли не должны смешиваться. Airflow не должен заниматься вычислениями, а Spark - управлением расписаниями.

1.2. Антипаттерн: Spark в памяти Airflow

Самая распространённая ошибка новичков - запустить PySpark-код прямо внутри PythonOperator:

# АНТИПАТТЕРН: НИКОГДА ТАК НЕ ДЕЛАЙТЕ
from airflow.operators.python import PythonOperator
from pyspark.sql import SparkSession

def run_etl():
    # Создаём SparkSession ПРЯМО НА ВОРКЕРЕ AIRFLOW
    spark = SparkSession.builder \
        .master("local[*]") \   # Использует все CPU воркера Airflow!
        .config("spark.driver.memory", "8g") \  # 8 ГБ из памяти воркера Airflow!
        .getOrCreate()

    df = spark.read.parquet("s3://bucket/data/")  # 50 ГБ данных
    df_result = df.groupBy("user_id").count()
    df_result.write.parquet("s3://bucket/output/")

    spark.stop()

task = PythonOperator(
    task_id="run_etl",
    python_callable=run_etl,
)

Почему это катастрофа:

  1. OOM на воркере: Spark-драйвер потребляет гигабайты RAM воркера Airflow. Когда несколько таких задач запускаются параллельно - воркер падает по Out-of-Memory, убивая все задачи на нём.
  2. Нет изоляции ресурсов: PySpark-задача потребляет CPU и память воркера, мешая другим задачам Airflow (в т.ч. простым SQL-запросам или операциям с API).
  3. Нет масштабирования: local[*] ограничен одной машиной - воркером Airflow. Никакой распределённости.
  4. Потеря задачи при рестарте воркера: если воркер Airflow упал, SparkSession безвозвратно потеряна. Нет возможности отследить состояние вычисления.
  5. Зависимости: воркер Airflow должен иметь установленные все Python-пакеты и Spark. Это кошмар поддержки.

Правильный подход: Airflow только инициирует запуск Spark на внешнем кластере и отслеживает статус. Вычисления происходят полностью за пределами Airflow.

1.3. Жизненный цикл задачи: от триггера до результата

Ключевой момент: воркер Airflow занят только polling - периодической проверкой статуса. Вся тяжёлая работа - на кластере. Это правильная архитектура.


Часть 2. Общая архитектура: компоненты системы

2.1. Компоненты production deployment

Полная production-система Airflow + Spark состоит из нескольких независимых слоёв:

2.2. Поддерживаемые Spark Cluster Managers

Airflow может оркестрировать Spark на любом кластерном менеджере:

Cluster Manager Запуск через Airflow Типичное применение
Apache YARN SparkSubmitOperator с --master yarn On-premise Hadoop кластеры
Spark Standalone SparkSubmitOperator с --master spark://host:7077 Небольшие on-premise установки
Kubernetes SparkSubmitOperator + KubernetesPodOperator + SparkKubernetesOperator Cloud-native платформы
Databricks DatabricksRunNowOperator / DatabricksSubmitRunOperator Managed Spark в облаке
AWS EMR EmrAddStepsOperator + EmrStepSensor AWS-ориентированные платформы

В этом уроке фокусируемся на наиболее универсальных и архитектурно интересных: SparkSubmitOperator и KubernetesPodOperator.


Часть 3. SparkSubmitOperator: классический подход

3.1. Механика работы

SparkSubmitOperator - это обёртка вокруг утилиты spark-submit. Когда Airflow Worker выполняет эту задачу, он буквально запускает процесс spark-submit на своём хосте, передавая ему все параметры: путь к скрипту, конфигурацию Spark, адрес кластера и аргументы приложения.

Spark-кластер при этом может быть YARN, Standalone или Kubernetes - воркер Airflow играет роль «клиента», который инициирует задачу. Важно понимать разницу между deployment modes:

Client mode (--deploy-mode client):

  • Spark Driver запускается прямо на воркере Airflow
  • Воркер Airflow занимает всю память, выделенную под Driver (часто 2–8 ГБ)
  • Если воркер упадёт - Driver умрёт, и вся Spark-задача потеряется
  • Не рекомендуется для production

Cluster mode (--deploy-mode cluster):

  • Spark Driver запускается на кластере (один из YARN NodeManagers или K8s pod)
  • Воркер Airflow только инициирует запуск и затем отслеживает статус
  • Если воркер Airflow упадёт - Driver продолжает работу на кластере (задача не потеряна)
  • Рекомендуется для production

3.2. Настройка Airflow Connection

Перед использованием SparkSubmitOperator нужно настроить Connection в Airflow, указав адрес кластера:

Через Airflow UI: Admin → Connections → New Connection:

  • Connection ID: spark_default
  • Connection Type: Spark
  • Host: yarn (для YARN) или spark://spark-master:7077 (для Standalone) или k8s://https://k8s-api:6443 (для K8s)

Через переменную окружения:

# В docker-compose.yml или .env файле
AIRFLOW_CONN_SPARK_DEFAULT='spark://spark-master:7077'

# Для YARN:
AIRFLOW_CONN_SPARK_DEFAULT='{"conn_type": "spark", "host": "yarn", "extra": {"deploy-mode": "cluster"}}'

3.3. Базовый DAG с SparkSubmitOperator

"""
dag_spark_submit_etl.py
Базовый пример SparkSubmitOperator для production ETL.
"""

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

# Дефолтные аргументы применяются ко всем задачам в DAG
default_args = {
    "owner":            "data-engineering-team",
    "depends_on_past":  False,           # Не ждать успешного завершения предыдущего запуска
    "start_date":       datetime(2024, 1, 1),
    "email":            ["alerts@company.com"],
    "email_on_failure": True,            # Email при падении
    "email_on_retry":   False,           # Не спамить при retry
    "retries":          3,               # До 3 повторов
    "retry_delay":      timedelta(minutes=5),  # Пауза между повторами
    "retry_exponential_backoff": True,   # Экспоненциальный рост паузы
}

with DAG(
    dag_id="clickstream_silver_etl",
    default_args=default_args,
    description="Incremental clickstream ingestion: Bronze → Silver",
    schedule_interval="0 * * * *",  # Каждый час
    catchup=False,                   # Не запускать пропущенные запуски
    max_active_runs=1,               # Только один активный запуск одновременно
    tags=["etl", "clickstream", "silver"],
) as dag:

    # Шаг 1: Инкрементальная обработка кликстрима
    ingest_clickstream = SparkSubmitOperator(
        task_id="ingest_clickstream",

        # Путь к PySpark-скрипту (на HDFS или в Docker-образе)
        application="/opt/spark/scripts/clickstream_ingestion.py",

        # Spark Connection из Airflow
        conn_id="spark_default",

        # Режим деплоя: cluster = Driver на кластере, не на воркере Airflow
        deploy_mode="cluster",

        # Ресурсы кластера
        num_executors=20,
        executor_cores=4,
        executor_memory="8g",
        driver_memory="4g",

        # Spark конфигурация
        conf={
            "spark.sql.adaptive.enabled":                    "true",
            "spark.sql.adaptive.coalescePartitions.enabled": "true",
            "spark.sql.shuffle.partitions":                  "400",
            # Iceberg extensions
            "spark.sql.extensions": "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions",
            "spark.sql.catalog.catalog": "org.apache.iceberg.spark.SparkCatalog",
        },

        # JAR-зависимости (Iceberg, S3 коннектор)
        jars=(
            "s3://artifacts/iceberg-spark-runtime-3.4.jar,"
            "s3://artifacts/hadoop-aws-3.3.4.jar,"
            "s3://artifacts/aws-java-sdk-bundle-1.12.262.jar"
        ),

        # Аргументы передаются в скрипт как sys.argv[1], sys.argv[2], ...
        # {{ ds }} - Jinja-шаблон Airflow: дата запуска в формате YYYY-MM-DD
        # {{ prev_ds }} - дата предыдущего успешного запуска
        application_args=[
            "--date",        "{{ ds }}",           # Дата инкремента
            "--prev-date",   "{{ prev_ds }}",      # Начало окна
            "--source",      "s3://bronze/clickstream/",
            "--target",      "catalog.silver.clickstream",
            "--batch-id",    "{{ run_id }}",       # Уникальный ID запуска
        ],

        # Передача кастомного python-окружения (если нужно)
        # py_files="s3://artifacts/utils.zip",  # Вспомогательные Python-модули

        # Логи: передавать в Airflow (работает только в client mode!)
        # В cluster mode логи нужно смотреть в YARN UI или History Server
        verbose=True,

        # Таймаут: если задача висит дольше 2 часов - убить
        execution_timeout=timedelta(hours=2),
    )

    # Шаг 2: Проверка качества данных
    validate_data = SparkSubmitOperator(
        task_id="validate_data_quality",
        application="/opt/spark/scripts/data_quality_check.py",
        conn_id="spark_default",
        deploy_mode="cluster",
        num_executors=5,
        executor_memory="4g",
        driver_memory="2g",
        application_args=[
            "--date",   "{{ ds }}",
            "--table",  "catalog.silver.clickstream",
            "--checks", "not_null,row_count,freshness",
        ],
    )

    # Зависимость: сначала ingest, потом validate
    ingest_clickstream >> validate_data

3.4. PySpark-скрипт: получение аргументов из Airflow

"""
clickstream_ingestion.py
PySpark-скрипт, запускаемый через SparkSubmitOperator.
Получает параметры через аргументы командной строки.
"""

import argparse
import sys
import logging
from pyspark.sql import SparkSession, functions as F

logger = logging.getLogger(__name__)


def parse_args():
    """Парсим аргументы командной строки, переданные из Airflow."""
    parser = argparse.ArgumentParser(description="Clickstream Bronze → Silver ETL")
    parser.add_argument("--date",      required=True, help="Дата инкремента YYYY-MM-DD")
    parser.add_argument("--prev-date", required=True, help="Начало инкрементального окна")
    parser.add_argument("--source",    required=True, help="Путь к Bronze данным")
    parser.add_argument("--target",    required=True, help="Target Iceberg таблица")
    parser.add_argument("--batch-id",  required=True, help="Airflow run_id для трейсабилити")
    return parser.parse_args()


def main():
    args = parse_args()
    logger.info(f"Starting ETL for date={args.date}, batch_id={args.batch_id}")

    spark = SparkSession.builder \
        .appName(f"clickstream_silver_{args.date}") \
        .getOrCreate()

    # Читаем инкремент за нужную дату
    df = spark.read.parquet(f"{args.source}date={args.date}/")

    # Трансформации...
    df_clean = (
        df
        .filter(F.col("event_id").isNotNull())
        .withColumn("_batch_id", F.lit(args.batch_id))
        .withColumn("_ingested_at", F.current_timestamp())
    )

    # Записываем в Silver
    df_clean.writeTo(args.target).append()

    logger.info(f"ETL complete for date={args.date}")
    spark.stop()


if __name__ == "__main__":
    main()

3.5. Плюсы и минусы SparkSubmitOperator

Преимущества:

  • Зрелость: используется годами, хорошо изучен, много документации
  • Простота: не нужен Kubernetes, достаточно YARN или Standalone кластера
  • Нативная интеграция: spark-submit напрямую взаимодействует с Spark cluster manager, получая оптимальные настройки
  • Прозрачная модель: легко объяснить и дебажить (это буквально spark-submit + polling)
  • YARN UI: история выполнения задач видна в YARN Resource Manager

Недостатки:

  • Spark на воркерах: воркеры Airflow должны иметь установленный Spark и Java - это тяжёлая зависимость
  • Версионный ад: если разные DAG используют разные версии Spark, конфигурировать воркер становится сложно
  • Python dependency drift: пакеты Python должны быть установлены на всех узлах YARN-кластера - сложно синхронизировать
  • Client mode риски: при случайном использовании client mode Driver поглощает ресурсы воркера Airflow
  • Логи в cluster mode: в cluster mode логи Driver находятся на кластере, а не в Airflow UI - нужна дополнительная настройка Log Forwarding

Часть 4. KubernetesPodOperator: Cloud-Native подход

4.1. Механика работы

KubernetesPodOperator работает принципиально иначе: Airflow не вызывает spark-submit, а создаёт Kubernetes Pod через Kubernetes API. Внутри этого Pod находится Docker-контейнер с вашим PySpark-скриптом и всеми его зависимостями.

Варианты запуска Spark внутри Pod:

  1. local[*] mode: Spark работает в рамках одного Pod (Driver = Executor). Подходит для небольших задач (< 50 ГБ данных).
  2. Spark-on-Kubernetes: Pod является Spark Driver, который сам создаёт Executor Pods в K8s. Полноценное распределённое выполнение.
  3. Spark в K8s через spark-submit: внутри Pod запускается spark-submit с --master k8s://.... Гибрид подходов.

4.2. Docker-образ для Spark ETL

Ключевое преимущество KubernetesPodOperator - все зависимости упакованы в Docker-образ. Нет больше «dependency hell», нет «у меня работает, на кластере не работает».

# Dockerfile - образ для Spark ETL задач
# Базируется на официальном образе Apache Spark

FROM apache/spark:3.5.1-python3

# Переключаемся на root для установки зависимостей
USER root

# Системные зависимости
RUN apt-get update && apt-get install -y --no-install-recommends \
    curl \
    && rm -rf /var/lib/apt/lists/*

# Python зависимости
# Используем requirements.txt для воспроизводимости
COPY requirements.txt /opt/app/requirements.txt
RUN pip install --no-cache-dir -r /opt/app/requirements.txt

# Spark JARs: Iceberg, S3, Delta
COPY jars/ /opt/spark/jars/

# Наши ETL-скрипты
COPY src/ /opt/app/src/
COPY config/ /opt/app/config/

# Возвращаемся на spark пользователя (безопасность)
USER spark

# Рабочая директория
WORKDIR /opt/app

# По умолчанию - help (переопределяется в KubernetesPodOperator)
CMD ["python", "-m", "src.main", "--help"]
# requirements.txt
pyspark==3.5.1
delta-spark==3.1.0
pyiceberg==0.7.1
boto3==1.34.0
great-expectations==0.18.0
python-dateutil==2.9.0
pydantic==2.6.0
# Сборка и публикация образа (в CI/CD pipeline)
docker build -t company/spark-etl:v1.2.3 .
docker push company-registry.example.com/spark-etl:v1.2.3

# Тег latest для разработки (не использовать в prod!)
docker tag company/spark-etl:v1.2.3 company/spark-etl:latest

4.3. DAG с KubernetesPodOperator

"""
dag_kubernetes_pod_etl.py
KubernetesPodOperator для production ETL с полной изоляцией зависимостей.
"""

from datetime import datetime, timedelta
from airflow import DAG
from airflow.providers.cncf.kubernetes.operators.kubernetes_pod import KubernetesPodOperator
from kubernetes.client import models as k8s

default_args = {
    "owner":            "data-engineering-team",
    "depends_on_past":  False,
    "start_date":       datetime(2024, 1, 1),
    "email":            ["alerts@company.com"],
    "email_on_failure": True,
    "retries":          3,
    "retry_delay":      timedelta(minutes=5),
}

with DAG(
    dag_id="clickstream_silver_etl_k8s",
    default_args=default_args,
    description="Clickstream ETL via KubernetesPodOperator",
    schedule_interval="0 * * * *",
    catchup=False,
    max_active_runs=1,
    tags=["etl", "clickstream", "kubernetes"],
) as dag:

    ingest_clickstream = KubernetesPodOperator(
        task_id="ingest_clickstream",

        # Имя Pod (должно быть уникальным в K8s namespace)
        name="clickstream-ingestion",

        # Namespace в K8s для запуска Pod
        namespace="data-platform",

        # Docker образ: СТРОГО фиксированный тег (не :latest!)
        # :latest в production - грубая ошибка (нет воспроизводимости)
        image="company-registry.example.com/spark-etl:v1.2.3",

        # Команда внутри контейнера (переопределяет CMD в Dockerfile)
        cmds=["python", "-m", "src.clickstream_ingestion"],

        # Аргументы, передаваемые в скрипт
        # Jinja-шаблоны работают так же, как в SparkSubmitOperator
        arguments=[
            "--date",     "{{ ds }}",
            "--prev-date","{{ prev_ds }}",
            "--source",   "s3://bronze/clickstream/",
            "--target",   "catalog.silver.clickstream",
            "--batch-id", "{{ run_id }}",
        ],

        # Ресурсные запросы/лимиты Pod
        # requests: гарантированные ресурсы (K8s зарезервирует)
        # limits: максимум (Pod будет убит при превышении)
        container_resources=k8s.V1ResourceRequirements(
            requests={"memory": "8Gi", "cpu": "2"},
            limits={"memory": "12Gi", "cpu": "4"},
        ),

        # Переменные окружения (незашифрованные конфиги)
        env_vars={
            "SPARK_DRIVER_MEMORY": "6g",
            "SPARK_EXECUTOR_INSTANCES": "20",
            "APP_ENV": "production",
        },

        # Секреты из Kubernetes Secrets (пароли, токены, ключи S3)
        secrets=[
            k8s.V1EnvFromSource(
                secret_ref=k8s.V1SecretEnvSource(name="aws-credentials")
            ),
            k8s.V1EnvFromSource(
                secret_ref=k8s.V1SecretEnvSource(name="iceberg-catalog-credentials")
            ),
        ],

        # Service Account с нужными правами в K8s
        # (чтение секретов, создание дочерних Pod для Executor'ов)
        service_account_name="spark-etl-sa",

        # Политика pull image: Always гарантирует свежий образ
        # (но замедляет старт на крупных кластерах без cache)
        image_pull_policy="Always",

        # Pull secret для приватного registry
        image_pull_secrets=[k8s.V1LocalObjectReference(name="registry-credentials")],

        # Получаем логи из Pod в Airflow (по умолчанию True)
        # Это работает и в "cluster mode" (в отличие от SparkSubmitOperator)
        get_logs=True,

        # Удалять ли Pod после завершения
        # True: чисто, но теряем возможность пост-мортем
        # False: Pod остаётся для отладки, но нужна уборка
        is_delete_operator_pod=True,

        # Если Pod уже существует - что делать?
        termination_grace_period=30,

        # Tolerations: запускать Pod на специальных нодах (spark-nodes)
        tolerations=[
            k8s.V1Toleration(
                key="spark-workload",
                operator="Equal",
                value="true",
                effect="NoSchedule",
            )
        ],

        # Node affinity: предпочитать ноды с SSD и достаточным объёмом RAM
        affinity=k8s.V1Affinity(
            node_affinity=k8s.V1NodeAffinity(
                preferred_during_scheduling_ignored_during_execution=[
                    k8s.V1PreferredSchedulingTerm(
                        weight=100,
                        preference=k8s.V1NodeSelectorTerm(
                            match_expressions=[
                                k8s.V1NodeSelectorRequirement(
                                    key="node.kubernetes.io/instance-type",
                                    operator="In",
                                    values=["m5.8xlarge", "m5.16xlarge"],
                                )
                            ]
                        ),
                    )
                ]
            )
        ),

        # Таймаут (после чего Pod убивается и задача считается неудачной)
        execution_timeout=timedelta(hours=2),
    )

    validate_data = KubernetesPodOperator(
        task_id="validate_data_quality",
        name="clickstream-dq-check",
        namespace="data-platform",
        image="company-registry.example.com/spark-etl:v1.2.3",
        cmds=["python", "-m", "src.data_quality"],
        arguments=[
            "--date",   "{{ ds }}",
            "--table",  "catalog.silver.clickstream",
        ],
        container_resources=k8s.V1ResourceRequirements(
            requests={"memory": "4Gi", "cpu": "1"},
            limits={"memory": "6Gi", "cpu": "2"},
        ),
        service_account_name="spark-etl-sa",
        get_logs=True,
        is_delete_operator_pod=True,
    )

    ingest_clickstream >> validate_data

4.4. Управление секретами в KubernetesPodOperator

Секреты (пароли БД, AWS-ключи, токены) никогда не должны быть в коде DAG или в переменных окружения Docker-образа. Правильный подход - Kubernetes Secrets:

# kubernetes/secrets/aws-credentials.yaml
# НИКОГДА не коммитить в git с реальными значениями!
# Используйте Helm Secrets, Sealed Secrets, External Secrets Operator
apiVersion: v1
kind: Secret
metadata:
  name: aws-credentials
  namespace: data-platform
type: Opaque
stringData:
  AWS_ACCESS_KEY_ID: "AKIA..."        # В prod: из Vault или AWS Secrets Manager
  AWS_SECRET_ACCESS_KEY: "xxxxx..."   # В prod: из Vault или AWS Secrets Manager
  AWS_DEFAULT_REGION: "eu-west-1"
# Применение через kubectl (в CI/CD pipeline, не в git)
kubectl apply -f aws-credentials.yaml -n data-platform

# Лучший подход: External Secrets Operator
# (автоматически синхронизирует из AWS Secrets Manager / HashiCorp Vault)
# ExternalSecret: автоматически создаёт K8s Secret из AWS Secrets Manager
apiVersion: external-secrets.io/v1beta1
kind: ExternalSecret
metadata:
  name: aws-credentials
  namespace: data-platform
spec:
  refreshInterval: 1h
  secretStoreRef:
    kind: ClusterSecretStore
    name: aws-secrets-manager
  target:
    name: aws-credentials
  data:
    - secretKey: AWS_ACCESS_KEY_ID
      remoteRef:
        key: prod/data-platform/aws-credentials
        property: access_key_id
    - secretKey: AWS_SECRET_ACCESS_KEY
      remoteRef:
        key: prod/data-platform/aws-credentials
        property: secret_access_key

4.5. Spark-on-Kubernetes внутри Pod

Если задача слишком велика для одного Pod (local mode), Driver Pod создаёт дополнительные Executor Pods:

"""
src/clickstream_ingestion.py - PySpark с Spark-on-K8s internal cluster
"""
import os
import argparse
from pyspark.sql import SparkSession, functions as F

def main():
    args = parse_args()

    # Определяем: запускаемся в K8s-кластере или локально
    k8s_namespace = os.environ.get("SPARK_K8S_NAMESPACE", "data-platform")
    k8s_service_account = os.environ.get("SPARK_K8S_SERVICE_ACCOUNT", "spark-etl-sa")

    spark = (
        SparkSession.builder
        .appName(f"clickstream_silver_{args.date}")

        # Executor Pods используют тот же Docker образ, что и Driver
        .config("spark.kubernetes.container.image",
                os.environ.get("SPARK_IMAGE", "company-registry.example.com/spark-etl:v1.2.3"))

        # K8s namespace для создания Executor Pods
        .config("spark.kubernetes.namespace", k8s_namespace)
        .config("spark.kubernetes.authenticate.driver.serviceAccountName", k8s_service_account)

        # Executor ресурсы
        .config("spark.executor.instances", "20")
        .config("spark.executor.memory", "8g")
        .config("spark.executor.cores", "4")

        # Динамическое масштабирование (если нужно)
        .config("spark.dynamicAllocation.enabled", "true")
        .config("spark.dynamicAllocation.shuffleTracking.enabled", "true")
        .config("spark.dynamicAllocation.minExecutors", "5")
        .config("spark.dynamicAllocation.maxExecutors", "50")

        # S3 настройки
        .config("spark.hadoop.fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem")
        .config("spark.hadoop.fs.s3a.aws.credentials.provider",
                "com.amazonaws.auth.EnvironmentVariableCredentialsProvider")

        # Iceberg
        .config("spark.sql.extensions",
                "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions")

        .getOrCreate()
    )

    # ETL логика (идентична локальному коду - никаких изменений)
    df_raw = spark.read.parquet(f"{args.source}date={args.date}/")

    df_clean = (
        df_raw
        .filter(F.col("event_id").isNotNull())
        .withColumn("_batch_id", F.lit(args.batch_id))
        .withColumn("_ingested_at", F.current_timestamp())
    )

    df_clean.writeTo(args.target).append()

    spark.stop()

Часть 5. Управление зависимостями: главная боль

5.1. Dependency Hell в SparkSubmitOperator

При использовании SparkSubmitOperator зависимости нужно распространять на все узлы YARN-кластера:

5.2. Решения для SparkSubmitOperator

Когда нельзя перейти на K8s, управление Python-зависимостями в SparkSubmitOperator решается несколькими способами:

# Вариант 1: --py-files - передаём zip-архив с зависимостями
# Работает для небольших зависимостей (нет C-расширений)
ingest = SparkSubmitOperator(
    task_id="ingest",
    application="s3://scripts/clickstream_ingestion.py",
    py_files="s3://artifacts/utils.zip",  # Вспомогательные Python-модули
    packages="org.apache.iceberg:iceberg-spark-runtime-3.4_2.12:1.4.0",
    conn_id="spark_default",
)

# Вариант 2: conda/virtualenv - упакованное Python-окружение
# spark-submit --conf spark.yarn.dist.archives=environment.tar.gz#environment
ingest_conda = SparkSubmitOperator(
    task_id="ingest_conda",
    application="s3://scripts/clickstream_ingestion.py",
    conf={
        "spark.yarn.dist.archives": "s3://artifacts/environment.tar.gz#environment",
        "spark.pyspark.python": "environment/bin/python",
    },
    conn_id="spark_default",
)

# Создание архива окружения (в CI/CD pipeline):
# conda create -n myenv python=3.11 pyspark pandas iceberg
# conda pack -n myenv -o environment.tar.gz
# aws s3 cp environment.tar.gz s3://artifacts/

# Вариант 3: --packages - Maven-координаты для Java JAR зависимостей
ingest_packages = SparkSubmitOperator(
    task_id="ingest_packages",
    application="s3://scripts/clickstream_ingestion.py",
    packages=(
        "org.apache.iceberg:iceberg-spark-runtime-3.4_2.12:1.4.0,"
        "io.delta:delta-core_2.12:3.0.0,"
        "org.apache.hadoop:hadoop-aws:3.3.4"
    ),
    # Кастомный Maven репозиторий (если нет доступа к Maven Central)
    repositories="https://artifacts.internal.company.com/maven2",
    conn_id="spark_default",
)

5.3. Управление образами в KubernetesPodOperator: CI/CD pipeline

# .github/workflows/build-spark-image.yml
# CI/CD: автосборка и публикация Docker образа при изменении кода

name: Build and Push Spark ETL Image

on:
  push:
    branches: [main]
    paths:
      - 'src/**'
      - 'requirements.txt'
      - 'Dockerfile'
      - 'jars/**'

jobs:
  build:
    runs-on: ubuntu-latest
    steps:
      - uses: actions/checkout@v4

      - name: Set up Docker Buildx
        uses: docker/setup-buildx-action@v3

      - name: Login to Registry
        uses: docker/login-action@v3
        with:
          registry: company-registry.example.com
          username: ${{ secrets.REGISTRY_USER }}
          password: ${{ secrets.REGISTRY_PASSWORD }}

      - name: Build and push
        uses: docker/build-push-action@v5
        with:
          context: .
          push: true
          # Тег = git SHA для воспроизводимости
          tags: |
            company-registry.example.com/spark-etl:${{ github.sha }}
            company-registry.example.com/spark-etl:latest
          # Кэширование слоёв для ускорения сборки
          cache-from: type=registry,ref=company-registry.example.com/spark-etl:buildcache
          cache-to: type=registry,ref=company-registry.example.com/spark-etl:buildcache,mode=max

      - name: Update DAG with new image tag
        # Обновляем DAG, чтобы он использовал новый тег
        run: |
          sed -i "s|spark-etl:.*|spark-etl:${{ github.sha }}|g" dags/clickstream_etl.py
          git config user.email "ci@company.com"
          git config user.name "CI Bot"
          git commit -am "Update spark-etl image to ${{ github.sha }}"
          git push

Часть 6. Мониторинг, логирование и observability

6.1. Логирование: главное отличие операторов

Логирование - одно из ключевых различий операторов в production:

Аспект SparkSubmitOperator (client mode) SparkSubmitOperator (cluster mode) KubernetesPodOperator
Driver логи в Airflow UI ✅ Да ❌ Нет (только на кластере) ✅ Да (get_logs=True)
Executor логи в Airflow UI ❌ Нет ❌ Нет ❌ Нет (только через Spark UI)
Spark UI доступен ✅ (пока Driver жив) ✅ (пока Driver жив) ✅ (через pod port-forward)
History Server ✅ После завершения ✅ После завершения ✅ После завершения
Потеря логов при crash ❌ Частичная ❌ Да ✅ Нет (K8s хранит)

6.2. Настройка centralized logging

"""
Добавляем ссылку на Spark History Server в логи Airflow задачи.
Это позволяет быстро переходить от Airflow к Spark UI при отладке.
"""

from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator
from airflow.models import Variable

# Адрес History Server хранится в Airflow Variables
HISTORY_SERVER_URL = Variable.get("spark_history_server_url",
                                  default_var="http://spark-history:18080")

def get_spark_ui_link(application_id: str) -> str:
    """Возвращает ссылку на конкретную задачу в Spark History Server."""
    return f"{HISTORY_SERVER_URL}/history/{application_id}/jobs/"


# В callback-функции для on_success_callback / on_failure_callback
def on_task_failure(context):
    """Вызывается при падении задачи. Логирует ссылку на Spark UI."""
    ti = context["task_instance"]
    log = ti.log

    # В реальной системе application_id берётся из метаданных задачи
    # (например, через XCom или из логов воркера)
    log.error(
        f"Task {ti.task_id} FAILED on {ti.execution_date}. "
        f"Check Spark History Server: {HISTORY_SERVER_URL}"
    )
    # Здесь можно отправить алерт в Slack/PagerDuty

# Использование в SparkSubmitOperator:
task = SparkSubmitOperator(
    task_id="ingest",
    application="s3://scripts/etl.py",
    conn_id="spark_default",
    on_failure_callback=on_task_failure,
)

6.3. Мониторинг через Prometheus + Grafana

Для системного мониторинга Spark-задач используется связка Prometheus + Grafana:

# В SparkSubmitOperator: включаем Prometheus метрики
ingest = SparkSubmitOperator(
    task_id="ingest",
    application="s3://scripts/etl.py",
    conn_id="spark_default",
    conf={
        # Prometheus Pushgateway: Spark отправляет метрики туда
        "spark.metrics.conf.*.sink.prometheuspushgateway.class":
            "org.apache.spark.metrics.sink.PrometheusSink",
        "spark.metrics.conf.*.sink.prometheuspushgateway.pushgateway-address":
            "prometheus-pushgateway:9091",
        "spark.metrics.conf.*.sink.prometheuspushgateway.job-name":
            "spark_clickstream_etl",
        # Отправляем метрики каждые 30 секунд
        "spark.metrics.conf.*.sink.prometheuspushgateway.period": "30",
    },
)

Ключевые метрики для Dashboard в Grafana:

  • spark_executor_shuffleWriteBytes: объём shuffle - индикатор тяжёлых операций
  • spark_executor_jvmHeapMemory: использование heap - индикатор потенциального OOM
  • spark_executor_cpuTime: CPU utilization на executor'ах
  • spark_driver_dagScheduler_failedStages: количество упавших stage - индикатор проблем
  • airflow_dagrun_duration_success: время выполнения DAG run - SLA мониторинг

Часть 7. Retry-стратегии и "зомби-задачи"

7.1. Что происходит при сбое воркера Airflow

Один из наиболее опасных сценариев в production - воркер Airflow упал в середине выполнения Spark-задачи. Поведение зависит от оператора и deploy mode:

7.2. Стратегия retry в Airflow

Грамотная retry-стратегия учитывает характер возможных ошибок:

from datetime import timedelta
from airflow.utils.trigger_rule import TriggerRule

# Для ETL-задач с идемпотентной записью (OVERWRITE или MERGE):
# - Безопасно повторять неограниченное число раз
# - Используем экспоненциальный backoff для сетевых ошибок

etl_task_defaults = {
    "retries":                        3,
    "retry_delay":                    timedelta(minutes=5),
    "retry_exponential_backoff":      True,   # 5 мин → 10 мин → 20 мин
    "max_retry_delay":                timedelta(hours=1),  # Не больше часа паузы
    "execution_timeout":              timedelta(hours=3),  # Убиваем зависшую задачу
}

# Для задач, которые НЕ идемпотентны (APPEND без деdup):
# - Нужно тщательно продумать retry
# - Лучше использовать идемпотентную запись (MERGE/OVERWRITE)
append_only_task_defaults = {
    "retries": 1,  # Минимум retry для append-задач
}

# Trigger rule: запускать DQ-check даже если ETL упал
# (чтобы зафиксировать факт падения в DQ отчёте)
validate_task = KubernetesPodOperator(
    task_id="validate",
    # ...
    trigger_rule=TriggerRule.ALL_DONE,  # Запускаем независимо от статуса upstream
)

7.3. Идемпотентность как защита от retry-проблем

Retry-стратегия безопасна только если задача идемпотентна - повторный запуск с теми же параметрами даёт тот же результат без дублирования данных:

def write_idempotent(spark, df, target_table: str, partition_date: str):
    """
    Идемпотентная запись: можно запускать повторно без дублей.

    Использует OVERWRITE для конкретной партиции.
    Если партиция уже записана - перезаписывает (не дублирует).
    """
    (
        df.write
        .format("iceberg")
        .mode("overwrite")
        # Перезаписываем только конкретную партицию, не всю таблицу
        .option("partitionOverwriteMode", "dynamic")
        .option("replaceWhere", f"event_date = '{partition_date}'")
        .saveAsTable(target_table)
    )


# Пример использования в скрипте, вызываемом из Airflow:
def main():
    args = parse_args()
    spark = create_spark_session()

    df = process_data(spark, args.date)

    # Безопасно при retry: если упал и перезапустился —
    # партиция перезапишется корректно
    write_idempotent(spark, df, "catalog.silver.clickstream", args.date)

Часть 8. Архитектурный выбор: SparkSubmitOperator vs KubernetesPodOperator

8.1. SparkKubernetesOperator: третий вариант

Помимо двух основных операторов, существует SparkKubernetesOperator из проекта spark-on-k8s-operator. Он использует Custom Resource Definition (CRD) Kubernetes - ресурс SparkApplication, который описывает всю Spark-задачу в K8s-нативном формате.

# spark-application.yaml - CRD для SparkKubernetesOperator
apiVersion: sparkoperator.k8s.io/v1beta2
kind: SparkApplication
metadata:
  name: clickstream-etl
  namespace: data-platform
spec:
  type: Python
  pythonVersion: "3"
  mode: cluster
  image: "company-registry.example.com/spark-etl:v1.2.3"
  imagePullPolicy: Always
  mainApplicationFile: "local:///opt/app/src/clickstream_ingestion.py"
  arguments:
    - "--date=2024-01-15"
    - "--target=catalog.silver.clickstream"
  sparkVersion: "3.5.1"
  driver:
    cores: 2
    memory: "4g"
    serviceAccount: spark-etl-sa
  executor:
    cores: 4
    memory: "8g"
    instances: 20
  restartPolicy:
    type: OnFailure
    onFailureRetries: 2
    onFailureRetryInterval: 10
# Использование SparkKubernetesOperator в DAG
from airflow.providers.cncf.kubernetes.operators.spark_kubernetes import SparkKubernetesOperator

spark_task = SparkKubernetesOperator(
    task_id="run_spark_etl",
    namespace="data-platform",
    application_file="spark-applications/clickstream-etl.yaml",
    kubernetes_conn_id="k8s_default",
    # Airflow polling пока SparkApplication не COMPLETED
    attach_log=True,
)

8.2. Граф принятия решений

8.3. Сводная матрица сравнения

Критерий SparkSubmitOperator KubernetesPodOperator SparkKubernetesOperator
Инфраструктура YARN / K8s / Standalone Kubernetes Kubernetes + spark-operator
Dependency management Сложно (conda-pack, --py-files) Просто (Docker image) Просто (Docker image)
Изоляция Python ❌ Низкая ✅ Полная ✅ Полная
Логи в Airflow UI Только client mode ✅ Всегда ✅ Всегда
Устойчивость к сбою воркера ❌ client mode / ✅ cluster mode ✅ Всегда ✅ Всегда
Сложность настройки Низкая Средняя Высокая
CI/CD интеграция Умеренная ✅ Отличная ✅ Отличная
Масштабирование YARN autoscaling K8s autoscaling K8s + Spark dynamic alloc
Зрелость ✅ Очень высокая Высокая Средняя
Подходит для On-premise Hadoop Cloud-native K8s-first платформы

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

9.1. Монолитный DAG

# АНТИПАТТЕРН: огромный монолитный DAG
# Все ETL-задачи - в одном файле, все зависимости перепутаны
with DAG("giant_etl_dag", ...) as dag:
    task_1 = SparkSubmitOperator(task_id="bronze_users", ...)
    task_2 = SparkSubmitOperator(task_id="bronze_events", ...)
    task_3 = SparkSubmitOperator(task_id="bronze_orders", ...)
    task_4 = SparkSubmitOperator(task_id="silver_users", ...)
    # ... ещё 47 задач ...
    task_52 = SparkSubmitOperator(task_id="gold_report", ...)

    # Зависимости - спагетти
    task_1 >> task_4 >> task_15 >> task_32 >> task_52
    task_2 >> task_5 >> task_16 >> task_33 >> task_52

# ПРАВИЛЬНО: разбить на независимые DAG
# dag_bronze_ingestion.py  - только ingestion
# dag_silver_transformation.py  - только трансформации
# dag_gold_aggregation.py  - только агрегации
# Связь между DAG через sensors или TriggerDagRunOperator

9.2. Hardcoded конфигурация Spark

# АНТИПАТТЕРН: жёстко заданные конфиги прямо в DAG
task = SparkSubmitOperator(
    task_id="etl",
    application="s3://scripts/etl.py",
    num_executors=20,          # Захардкожено!
    executor_memory="8g",      # Захардкожено!
    conf={"spark.sql.shuffle.partitions": "400"},  # Захардкожено!
)

# ПРАВИЛЬНО: конфиги в Airflow Variables или config-файлах
from airflow.models import Variable

spark_configs = Variable.get("spark_etl_config", deserialize_json=True)

task = SparkSubmitOperator(
    task_id="etl",
    application="s3://scripts/etl.py",
    num_executors=spark_configs["num_executors"],
    executor_memory=spark_configs["executor_memory"],
    conf=spark_configs.get("spark_conf", {}),
)

9.3. Mutable Docker image tags

# АНТИПАТТЕРН: использование :latest в production
task = KubernetesPodOperator(
    task_id="etl",
    image="company/spark-etl:latest",  # Что именно запускается?!
    image_pull_policy="Always",
    # ...
)
# Проблемы:
# - Нет воспроизводимости: разные запуски могут использовать разный код
# - При откате (rollback) непонятно, на что откатываться
# - Нарушает принцип immutable artifacts

# ПРАВИЛЬНО: конкретный тег = конкретный commit
task = KubernetesPodOperator(
    task_id="etl",
    image="company/spark-etl:sha-a1b2c3d4",  # Git SHA - точно известно, что запускается
    image_pull_policy="IfNotPresent",          # Не перекачивать образ при каждом запуске
    # ...
)

9.4. Zombie Spark Jobs

Zombie-задача - это Spark Job, который продолжает работать на кластере, хотя Airflow уже считает задачу завершённой (или повторно её запустил после timeout).

# Защита от zombie jobs:

# 1. Cleanup hook при retry
from airflow.providers.apache.spark.hooks.spark_submit import SparkSubmitHook

def kill_zombie_spark_jobs(context):
    """
    Вызывается при retry задачи.
    Убивает предыдущую Spark-задачу, если она ещё висит.
    """
    ti = context["task_instance"]
    # В реальной системе: получаем application_id из XCom
    # и убиваем его через YARN API или Spark REST API
    pass

task = SparkSubmitOperator(
    task_id="etl",
    application="s3://scripts/etl.py",
    conn_id="spark_default",
    on_retry_callback=kill_zombie_spark_jobs,
    execution_timeout=timedelta(hours=3),  # Жёсткий таймаут = защита от зависания
)

# 2. Для KubernetesPodOperator - Pod автоматически удаляется
# Kubernetes сам управляет lifecycle Pod → нет zombie по определению
task_k8s = KubernetesPodOperator(
    task_id="etl",
    image="company/spark-etl:v1.2.3",
    is_delete_operator_pod=True,   # Удалять Pod при завершении
    termination_grace_period=60,   # 60 секунд на graceful shutdown
)

9.5. Resource Starvation в multi-tenant кластерах

# АНТИПАТТЕРН: запустить 10 тяжёлых Spark задач одновременно
with DAG("parallel_etl") as dag:
    tasks = [
        SparkSubmitOperator(
            task_id=f"etl_{table}",
            num_executors=50,      # 50 × 10 задач = 500 executor'ов одновременно!
            executor_memory="8g",  # 500 × 8 ГБ = 4 ТБ RAM!
        )
        for table in ["users", "events", "orders", "products", ...]  # 10 таблиц
    ]

# ПРАВИЛЬНО: ограничиваем параллелизм через pool
from airflow.models import Pool

# Создаём pool с ограничением на параллельные Spark задачи
# airflow pools set spark_etl_pool 3 "Max 3 concurrent Spark jobs"

tasks_with_pool = [
    SparkSubmitOperator(
        task_id=f"etl_{table}",
        num_executors=20,      # Разумное количество
        executor_memory="8g",
        pool="spark_etl_pool", # Максимум 3 из этих задач выполняется одновременно
    )
    for table in ["users", "events", "orders", "products", ...]
]

Часть 10. End-to-End Production Кейс: один пайплайн, два оператора

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

Реализуем один и тот же ETL-пайплайн двумя способами, чтобы наглядно сравнить операционный опыт:

Задача: инкрементальная обработка кликстрима - Bronze → Silver → Gold (Daily Report).

  • Bronze: сырые события, Parquet, партиционированные по часам
  • Silver: дедуплицированные и обогащённые события, Iceberg, партиционированные по дням
  • Gold: агрегированные KPI по дням и сегментам пользователей, Iceberg

10.2. Версия A: SparkSubmitOperator (YARN)

"""
dag_clickstream_yarn.py
ETL через SparkSubmitOperator на YARN-кластере.
Подходит для on-premise инфраструктуры.
"""

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

# Конфиги вынесены в Airflow Variables (управляются через UI/CLI)
YARN_CLUSTER_CONFIGS = Variable.get("yarn_spark_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},
})

# Адрес History Server для логов (хранится в Variable)
HISTORY_SERVER = Variable.get("spark_history_server_url",
                              default_var="http://spark-history:18080")

SPARK_BASE_CONF = {
    "spark.sql.adaptive.enabled":                    "true",
    "spark.sql.adaptive.coalescePartitions.enabled": "true",
    "spark.sql.extensions":
        "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions",
    "spark.sql.catalog.catalog":
        "org.apache.iceberg.spark.SparkCatalog",
    "spark.sql.catalog.catalog.type": "hive",
    "spark.sql.catalog.catalog.uri":
        "thrift://hive-metastore:9083",
}

def make_spark_submit(
    task_id: str,
    script: str,
    size: str,
    app_args: list,
    **kwargs
) -> SparkSubmitOperator:
    """Фабричная функция для создания SparkSubmitOperator с общими настройками."""
    cfg = YARN_CLUSTER_CONFIGS[size]
    return SparkSubmitOperator(
        task_id=task_id,
        application=f"s3://scripts/{script}",
        conn_id="spark_yarn_prod",
        deploy_mode="cluster",
        conf={**SPARK_BASE_CONF},
        jars="s3://artifacts/iceberg-spark-runtime-3.4.jar",
        application_args=app_args,
        **cfg,
        **kwargs,
    )


with DAG(
    dag_id="clickstream_etl_yarn",
    default_args={
        "owner": "data-team",
        "retries": 3,
        "retry_delay": timedelta(minutes=10),
        "retry_exponential_backoff": True,
        "execution_timeout": timedelta(hours=4),
        "start_date": datetime(2024, 1, 1),
    },
    schedule_interval="30 2 * * *",   # Каждый день в 2:30 UTC
    catchup=False,
    max_active_runs=1,
    tags=["etl", "clickstream", "yarn"],
) as dag:

    bronze_to_silver = make_spark_submit(
        task_id="bronze_to_silver",
        script="clickstream_silver.py",
        size="large",
        app_args=[
            "--date",     "{{ ds }}",
            "--source",   "s3://bronze/clickstream/",
            "--target",   "catalog.silver.clickstream",
            "--batch-id", "{{ run_id }}",
        ],
    )

    silver_to_gold = make_spark_submit(
        task_id="silver_to_gold",
        script="clickstream_gold.py",
        size="medium",
        app_args=[
            "--date",       "{{ ds }}",
            "--source",     "catalog.silver.clickstream",
            "--target",     "catalog.gold.clickstream_daily_kpi",
        ],
    )

    data_quality = make_spark_submit(
        task_id="data_quality_check",
        script="dq_check.py",
        size="small",
        app_args=[
            "--date",   "{{ ds }}",
            "--tables", "catalog.silver.clickstream,catalog.gold.clickstream_daily_kpi",
        ],
    )

    bronze_to_silver >> silver_to_gold >> data_quality

10.3. Версия B: KubernetesPodOperator (K8s)

"""
dag_clickstream_k8s.py
Тот же ETL через KubernetesPodOperator.
Подходит для cloud-native инфраструктуры.
"""

from datetime import datetime, timedelta
from airflow import DAG
from airflow.providers.cncf.kubernetes.operators.kubernetes_pod import KubernetesPodOperator
from kubernetes.client import models as k8s
from airflow.models import Variable

# Версия образа - единственная точка изменения при деплое новой версии кода
IMAGE_TAG = Variable.get("spark_etl_image_tag", default_var="v1.2.3")
IMAGE = f"company-registry.example.com/spark-etl:{IMAGE_TAG}"
NAMESPACE = "data-platform"
SERVICE_ACCOUNT = "spark-etl-sa"

# Общие секреты (AWS, Iceberg catalog credentials)
COMMON_SECRETS = [
    k8s.V1EnvFromSource(secret_ref=k8s.V1SecretEnvSource(name="aws-credentials")),
    k8s.V1EnvFromSource(secret_ref=k8s.V1SecretEnvSource(name="iceberg-catalog-creds")),
]

def make_kpo(
    task_id: str,
    module: str,
    resources: dict,
    app_args: list,
) -> KubernetesPodOperator:
    """Фабричная функция для создания KubernetesPodOperator с общими настройками."""
    return KubernetesPodOperator(
        task_id=task_id,
        name=task_id.replace("_", "-"),
        namespace=NAMESPACE,
        image=IMAGE,
        image_pull_policy="IfNotPresent",
        cmds=["python", "-m", f"src.{module}"],
        arguments=app_args,
        container_resources=k8s.V1ResourceRequirements(**resources),
        secrets=COMMON_SECRETS,
        service_account_name=SERVICE_ACCOUNT,
        get_logs=True,
        is_delete_operator_pod=True,
        tolerations=[
            k8s.V1Toleration(key="spark-workload", operator="Exists", effect="NoSchedule")
        ],
    )


with DAG(
    dag_id="clickstream_etl_k8s",
    default_args={
        "owner": "data-team",
        "retries": 3,
        "retry_delay": timedelta(minutes=10),
        "retry_exponential_backoff": True,
        "execution_timeout": timedelta(hours=4),
        "start_date": datetime(2024, 1, 1),
    },
    schedule_interval="30 2 * * *",
    catchup=False,
    max_active_runs=1,
    tags=["etl", "clickstream", "kubernetes"],
) as dag:

    bronze_to_silver = make_kpo(
        task_id="bronze_to_silver",
        module="clickstream_silver",
        resources={
            "requests": {"memory": "8Gi",  "cpu": "4"},
            "limits":   {"memory": "16Gi", "cpu": "8"},
        },
        app_args=[
            "--date",     "{{ ds }}",
            "--source",   "s3://bronze/clickstream/",
            "--target",   "catalog.silver.clickstream",
            "--batch-id", "{{ run_id }}",
        ],
    )

    silver_to_gold = make_kpo(
        task_id="silver_to_gold",
        module="clickstream_gold",
        resources={
            "requests": {"memory": "4Gi", "cpu": "2"},
            "limits":   {"memory": "8Gi", "cpu": "4"},
        },
        app_args=[
            "--date",   "{{ ds }}",
            "--source", "catalog.silver.clickstream",
            "--target", "catalog.gold.clickstream_daily_kpi",
        ],
    )

    data_quality = make_kpo(
        task_id="data_quality_check",
        module="dq_check",
        resources={
            "requests": {"memory": "2Gi", "cpu": "1"},
            "limits":   {"memory": "4Gi", "cpu": "2"},
        },
        app_args=[
            "--date",   "{{ ds }}",
            "--tables", "catalog.silver.clickstream,catalog.gold.clickstream_daily_kpi",
        ],
    )

    bronze_to_silver >> silver_to_gold >> data_quality

10.4. Сравнение операционного опыта

Первый деплой:

  • YARN версия: настройка spark-submit на воркерах (30–60 мин), создание Airflow Connection (5 мин)
  • K8s версия: сборка Docker-образа (15–20 мин), настройка K8s RBAC и secrets (30 мин), создание Airflow Connection (5 мин)

Деплой новой версии кода:

  • YARN версия: обновить скрипт на S3, обновить JAR/zip зависимости (возможно - на всех нодах)
  • K8s версия: docker build && docker push, обновить IMAGE_TAG в Airflow Variable - готово

Отладка упавшей задачи:

  • YARN версия: открыть Airflow UI → перейти в YARN Resource Manager → найти application_id → открыть логи
  • K8s версия: открыть Airflow UI → Task Logs (логи прямо здесь, get_logs=True)

Масштабирование под нагрузку:

  • YARN версия: добавить ноды в YARN кластер (операция уровня инфраструктуры)
  • K8s версия: Kubernetes Cluster Autoscaler автоматически добавит ноды при нехватке ресурсов (Scale to Zero при простое)

Воспроизводимость:

  • YARN версия: зависит от состояния нод кластера (Python-версии, установленных пакетов)
  • K8s версия: 100% воспроизводимо (всё в Docker-образе, образ immutable)

Итоги урока

Выбор между SparkSubmitOperator и KubernetesPodOperator - это не технический, а архитектурный выбор:

SparkSubmitOperator - для зрелых on-premise Hadoop-платформ, где уже есть YARN-кластер, команда знает Spark-экосистему, и переход на K8s не планируется. Простота в настройке, но сложность в управлении зависимостями.

KubernetesPodOperator - для cloud-native платформ, где инфраструктура строится на Kubernetes. Docker-контейнеры решают проблему dependency hell, дают полную изоляцию и воспроизводимость. Сложнее в начальной настройке, но значительно проще в масштабировании и поддержке.

Ключевые правила:

  1. Никогда не запускайте PySpark в PythonOperator - это OOM на воркере Airflow
  2. Всегда используйте --deploy-mode cluster в SparkSubmitOperator - Driver должен быть на кластере
  3. Никогда не используйте :latest тег в KubernetesPodOperator в production - только фиксированные теги
  4. Всегда проектируйте ETL-задачи идемпотентными - это единственная защита от проблем при retry
  5. Управляйте секретами через Kubernetes Secrets или внешние Secrets Manager - никакого хардкода в DAG или образах