Airflow + Spark: SparkSubmitOperator vs KubernetesPodOperator
Оркестрация Spark-пайплайнов через Airflow: архитектурный выбор между SparkSubmitOperator и KubernetesPodOperator, управление зависимостями, мониторинг, retry-стратегии и production trade-offs
Введение: почему оркестрация решает судьбу пайплайна¶
Представьте: вы написали отличный 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,
)
Почему это катастрофа:
- OOM на воркере: Spark-драйвер потребляет гигабайты RAM воркера Airflow. Когда несколько таких задач запускаются параллельно - воркер падает по Out-of-Memory, убивая все задачи на нём.
- Нет изоляции ресурсов: PySpark-задача потребляет CPU и память воркера, мешая другим задачам Airflow (в т.ч. простым SQL-запросам или операциям с API).
- Нет масштабирования:
local[*]ограничен одной машиной - воркером Airflow. Никакой распределённости. - Потеря задачи при рестарте воркера: если воркер Airflow упал, SparkSession безвозвратно потеряна. Нет возможности отследить состояние вычисления.
- Зависимости: воркер 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:
local[*]mode: Spark работает в рамках одного Pod (Driver = Executor). Подходит для небольших задач (< 50 ГБ данных).- Spark-on-Kubernetes: Pod является Spark Driver, который сам создаёт Executor Pods в K8s. Полноценное распределённое выполнение.
- 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, дают полную изоляцию и воспроизводимость. Сложнее в начальной настройке, но значительно проще в масштабировании и поддержке.
Ключевые правила:
- Никогда не запускайте PySpark в
PythonOperator- это OOM на воркере Airflow - Всегда используйте
--deploy-mode clusterвSparkSubmitOperator- Driver должен быть на кластере - Никогда не используйте
:latestтег вKubernetesPodOperatorв production - только фиксированные теги - Всегда проектируйте ETL-задачи идемпотентными - это единственная защита от проблем при retry
- Управляйте секретами через Kubernetes Secrets или внешние Secrets Manager - никакого хардкода в DAG или образах