Serverless Spark на K8s: EMR Serverless, Dataproc Serverless, pod templates и auto-scaling

Serverless Spark на K8s: EMR Serverless, Dataproc Serverless, pod templates и auto-scaling

optimization

Эволюция инфраструктуры Spark: от Bare-Metal к Serverless

Чтобы понять, зачем нужен Serverless Spark, нужно вспомнить, как развивалась инфраструктура для аналитических вычислений - и почему каждое поколение пыталось решить проблемы предыдущего.

Первое поколение: постоянные YARN-кластеры

Классический Hadoop/Spark кластер - это набор физических или виртуальных машин, которые работают круглосуточно. ResourceManager (YARN) знает о всех нодах и распределяет задачи. Это надёжная, понятная модель. Но у неё есть критический изъян: кластер оплачивается постоянно, независимо от нагрузки.

Типичный паттерн нагрузки аналитического кластера:

00:00–06:00: 5% нагрузки (фоновые задачи)
06:00–08:00: 30% нагрузки (утренние ETL)
08:00–18:00: 60–80% нагрузки (рабочий день)
18:00–22:00: 40% нагрузки (вечерние батчи)
22:00–00:00: 10% нагрузки

Средняя утилизация: ~35%. Это значит, что 65% оплаченных ресурсов простаивают. В облаке, где платишь за час работы виртуальной машины, это прямые потери.

Дополнительно: постоянный кластер требует операционной поддержки. Кто-то должен обновлять ОС, мониторить ноды, обрабатывать падения, следить за version compatibility между YARN, Hadoop, Spark и Python. Это административная нагрузка, которая не создаёт бизнес-ценности.

Второе поколение: управляемые кластеры (EMR, Dataproc)

AWS EMR и GCP Dataproc решили проблему администрирования: облако берёт на себя установку Spark, управление версиями, патчинг. Вы просто указываете нужную версию Spark, выбираете размер нод - и кластер готов за 5–10 минут.

Но проблема экономики никуда не делась. Кластер по-прежнему нужно запустить, он по-прежнему работает, пока ты его не остановишь. Да, появился Auto-scaling - добавление/удаление нод в зависимости от нагрузки - но управлять этим нужно самостоятельно. И главное: сам процесс создания нового кластера занимает 5–15 минут - это неприемлемо для задач с жёстким SLA.

Третье поколение: Serverless Spark

Serverless Spark убирает понятие «кластер» из словаря пользователя. Вы описываете задачу (код, зависимости, ресурсы), отправляете её в облако - и облако само:

  1. Выделяет нужное количество вычислительных ресурсов (за секунды, не минуты)
  2. Запускает задачу
  3. Масштабирует ресурсы в ходе выполнения (вверх и вниз)
  4. Освобождает всё после завершения
  5. Выставляет счёт только за фактически потреблённые ресурсы (vCPU-секунды + GB-секунды памяти)

Это не просто удобство - это фундаментальное изменение экономики:

Классический кластер (10 нод × 8 vCPU × 32 GB, постоянный):
  Стоимость/час = $20/час
  Задача 2 часа/день: $20 × 24 = $480/день
  Из них полезная нагрузка: $20 × 2 = $40/день (8%)
  Оплата простоя: $440/день (92%)

Serverless (та же задача 2 часа/день):
  Стоимость только за выполнение = ~$40–60/день
  Экономия: ~90% относительно постоянного кластера

Kubernetes как фундамент современного Spark

Почему K8s стал стандартом

Kubernetes - это система оркестрации контейнеров. Spark поддерживает K8s как кластерный менеджер с версии 2.3 (2018), и с тех пор это стало де-факто стандартом для облачных Spark-деплоев.

Почему K8s вытесняет YARN для новых деплоев:

  • Изоляция через контейнеры: каждый Executor работает в отдельном Docker-контейнере. Никакой разделяемой библиотечной среды, никаких конфликтов версий Python.
  • Объявительная конфигурация: манифесты YAML описывают желаемое состояние; K8s сам достигает его. Никаких скриптов запуска и остановки.
  • Нативная эластичность: K8s умеет создавать и уничтожать поды за секунды. Это идеальная основа для динамического масштабирования Executor-ов.
  • Единая платформа: K8s используется для ML-обучения, микросервисов, баз данных - Spark органично вписывается в эту экосистему.

Spark на K8s: модель исполнения

В K8s-режиме каждый компонент Spark - это отдельный под (Pod):

  • Driver Pod: создаётся первым, содержит JVM с Spark Driver (DAGScheduler, SparkContext). Координирует выполнение задания.
  • Executor Pod: создаётся Driver-ом по мере необходимости через K8s API. Каждый Executor - изолированный контейнер с JVM. После завершения задачи под удаляется.

Взаимодействие с K8s API происходит через Service Account Executor-а. Driver формирует запрос к K8s API (POST /api/v1/namespaces/spark/pods), K8s планировщик находит подходящую ноду с достаточными ресурсами и запускает под. Это происходит быстро - обычно за 5–30 секунд для pre-pulled образов.

Полная архитектура Serverless Spark

На диаграмме показан ключевой принцип архитектуры: вычисление и хранение полностью разделены. Executor-поды эфемерны - они создаются и уничтожаются в ходе задачи. Данные хранятся в объектном хранилище (S3/GCS). Shuffle-данные вынесены в отдельный Remote Shuffle Service, чтобы уничтожение Executor-пода не приводило к потере промежуточных результатов.


Serverless Spark в публичных облаках: EMR vs Dataproc

AWS EMR Serverless

EMR Serverless - это managed сервис AWS, который предоставляет serverless Spark без необходимости управлять кластером. Под капотом AWS использует Firecracker (легковесные микровиртуальные машины) для изоляции задач.

Жизненный цикл задачи:

  1. Вы создаёте Application - логическую единицу Serverless-среды (тип runtime: Spark, версия, регион).
  2. Отправляете JobRun - конкретное задание с кодом, конфигурацией и IAM Role.
  3. EMR Serverless аллоцирует Driver и Executor-воркеры (занимает 10–60 секунд для cold start или < 5 секунд при Pre-initialized Capacity).
  4. Задача выполняется; EMR автоматически масштабирует воркеры.
  5. После завершения воркеры освобождаются. Платите только за time × vCPU × memory.

Ценообразование (2024):

  • Driver: $0.04730 за vCPU-час + $0.01847 за GB RAM-час
  • Executor: $0.04730 за vCPU-час + $0.01847 за GB RAM-час
  • Минимальная гранулярность: 1 минута
# Отправка задачи в EMR Serverless через boto3
import boto3

emr_client = boto3.client("emr-serverless", region_name="eu-west-1")

# Создание Application (один раз)
response = emr_client.create_application(
    name="spark-analytics",
    type="SPARK",
    releaseLabel="emr-6.15.0",
    # Pre-initialized capacity: держим "горячих" воркеров для быстрого старта
    initialCapacity={
        "DRIVER": {
            "workerCount": 1,
            "workerConfiguration": {
                "cpu": "4vCPU",
                "memory": "16GB",
                "disk": "100GB"
            }
        },
        "EXECUTOR": {
            "workerCount": 5,      # 5 воркеров всегда "горячие"
            "workerConfiguration": {
                "cpu": "4vCPU",
                "memory": "16GB",
                "disk": "200GB"
            }
        }
    },
    maximumCapacity={
        "cpu": "400vCPU",          # максимум автоскейлинга
        "memory": "3000GB"
    }
)
application_id = response["applicationId"]

# Отправка конкретного задания
job_response = emr_client.start_job_run(
    applicationId=application_id,
    executionRoleArn="arn:aws:iam::123456789:role/EMRServerlessRole",
    jobDriver={
        "sparkSubmit": {
            "entryPoint": "s3://my-bucket/scripts/gold_layer.py",
            "entryPointArguments": ["--date", "2024-01-15"],
            "sparkSubmitParameters": (
                "--conf spark.executor.cores=4 "
                "--conf spark.executor.memory=16g "
                "--conf spark.executor.instances=10 "
                "--conf spark.dynamicAllocation.enabled=true "
                "--conf spark.dynamicAllocation.minExecutors=2 "
                "--conf spark.dynamicAllocation.maxExecutors=50 "
                "--conf spark.sql.adaptive.enabled=true"
            )
        }
    },
    configurationOverrides={
        "monitoringConfiguration": {
            "s3MonitoringConfiguration": {
                "logUri": "s3://my-logs-bucket/emr-serverless/"
            }
        }
    }
)
print(f"Job submitted: {job_response['jobRunId']}")

Pre-initialized Capacity: решение проблемы Cold Start

Cold Start - это задержка между отправкой задачи и началом фактического выполнения. В Serverless-среде это время тратится на:

  • Запуск Firecracker micro-VM: ~2–5 секунд
  • Запуск JVM и инициализацию Spark: ~10–20 секунд
  • Загрузку Docker-образа (если не кэширован): ~30–120 секунд

Итого cold start может занимать 1–3 минуты. Для пайплайнов с жёстким SLA (например, отчёт должен быть готов в 8:00:30) это неприемлемо.

Pre-initialized Capacity решает это: EMR держит заданное число воркеров в «тёплом» состоянии - JVM уже запущена, образ загружен. Задача начинает выполнение за < 5 секунд. Но «тёплые» воркеры оплачиваются, даже если задача ещё не запущена - поэтому держать их имеет смысл только для критических pipeline с жёстким SLA.

GCP Dataproc Serverless

Dataproc Serverless - аналогичный сервис от Google Cloud. Ключевые отличия от EMR Serverless:

  • Глубокая интеграция с GCP-экосистемой: BigQuery, Cloud Spanner, Pub/Sub доступны через нативные коннекторы без дополнительных JAR.
  • Автоматическое определение ресурсов: Dataproc Serverless сам выбирает оптимальные размеры воркеров на основе характеристик задачи (не нужно вручную указывать --executor-memory 16g).
  • Управление через gcloud CLI: проще для одноразовых задач и скриптов.
# Отправка PySpark-задачи в Dataproc Serverless
gcloud dataproc batches submit pyspark gs://my-bucket/scripts/gold_layer.py \
  --region=europe-west4 \
  --project=my-gcp-project \
  --deps-bucket=gs://my-spark-staging/ \
  --jars=gs://spark-lib/bigtable/bigtable-dataproc-1.0.jar \
  --py-files=gs://my-bucket/libs/utils.zip \
  --properties=spark.executor.cores=4,spark.dynamicAllocation.enabled=true \
  --labels=env=production,team=analytics \
  -- --date=2024-01-15 --output-table=analytics.gold_sales
# Отправка через Python SDK (для Airflow DAG)
from google.cloud import dataproc_v1

def submit_dataproc_batch(
    project_id: str,
    region: str,
    batch_id: str,
    python_file: str,
    args: list,
) -> str:
    client = dataproc_v1.BatchControllerClient(
        client_options={"api_endpoint": f"{region}-dataproc.googleapis.com:443"}
    )

    batch = dataproc_v1.Batch()
    batch.pyspark_batch = dataproc_v1.PySparkBatch(
        main_python_file_uri=python_file,
        args=args,
        jar_file_uris=["gs://spark-lib/bigquery/spark-bigquery-latest_2.12.jar"],
    )
    batch.runtime_config = dataproc_v1.RuntimeConfig(
        properties={
            "spark.dynamicAllocation.enabled": "true",
            "spark.dynamicAllocation.executorAllocationRatio": "1.0",
            "spark.sql.adaptive.enabled": "true",
        }
    )

    operation = client.create_batch(
        request=dataproc_v1.CreateBatchRequest(
            parent=f"projects/{project_id}/locations/{region}",
            batch=batch,
            batch_id=batch_id,
        )
    )
    print(f"Submitted batch: {batch_id}")
    return operation.metadata.batch

# Использование
batch = submit_dataproc_batch(
    project_id="my-gcp-project",
    region="europe-west4",
    batch_id=f"gold-sales-{date.today().isoformat()}",
    python_file="gs://my-bucket/scripts/gold_layer.py",
    args=["--date", "2024-01-15"],
)

Сравнительная таблица: EMR Serverless vs Dataproc Serverless

Критерий EMR Serverless Dataproc Serverless
Облако AWS GCP
Изоляция Firecracker micro-VM GKE pods
Cold start 30–120s (cold) / < 5s (warm) 30–90s
Pre-warmed capacity Да (оплачивается) Нет (auto-managed)
Настройка ресурсов Явная (vCPU + GB) Автоматическая
Интеграция с хранилищем S3 + Glue Catalog GCS + BigQuery
Кастом образы Да (ECR) Да (Artifact Registry)
Мониторинг CloudWatch + Spark UI Cloud Logging + Spark History
Pricing Per vCPU-s + per GB-s Per shuffle-GB + per vCPU-s

Кастомизация окружения: Pod Templates

Зачем нужны Pod Templates

Стандартные параметры spark-submit (--executor-memory, --executor-cores) описывают ресурсы Spark. Но иногда нужно настроить сам K8s-под на более низком уровне:

  • Примонтировать дополнительный диск (SSD) для хранения shuffle-данных
  • Задать Security Context (запуск от non-root user)
  • Настроить Node Affinity (закрепить Driver на On-Demand ноде, Executor - на Spot)
  • Добавить sidecar-контейнер (например, Prometheus exporter)
  • Внедрить секреты (Kubernetes Secrets или AWS Secrets Manager)

Pod Template - это YAML-файл с частичным описанием K8s Pod. Spark merge-ит его с собственным описанием пода, давая возможность тонкой кастомизации.

Driver Pod Template: надёжность и наблюдаемость

Driver - критичный компонент: если он падает, задача падает целиком. Для Driver-а важны: стабильная нода (On-Demand, не Spot), достаточный диск для event log, доступ к Spark History Server.

# driver-pod-template.yaml
apiVersion: v1
kind: Pod
metadata:
  labels:
    app: spark-driver
    version: "3.5"
spec:
  # Tolerations: Driver допускается на ноды с taint "spark-driver"
  tolerations:
  - key: "spark-driver"
    operator: "Exists"
    effect: "NoSchedule"

  # Node Affinity: Driver ТОЛЬКО на On-Demand ноды (никогда не Spot!)
  affinity:
    nodeAffinity:
      requiredDuringSchedulingIgnoredDuringExecution:
        nodeSelectorTerms:
        - matchExpressions:
          - key: "node.kubernetes.io/lifecycle"
            operator: In
            values: ["on-demand"]
          - key: "node.kubernetes.io/instance-type"
            operator: In
            values: ["r6i.4xlarge", "r6i.8xlarge"]

  # Security Context: запуск как непривилегированный пользователь
  securityContext:
    runAsNonRoot: true
    runAsUser: 1000
    fsGroup: 1000

  containers:
  - name: spark-driver
    resources:
      requests:
        cpu: "4"
        memory: "16Gi"
      limits:
        cpu: "4"
        memory: "18Gi"    # немного больше request - для JVM overhead
    # Переменные окружения
    env:
    - name: SPARK_HISTORY_SERVER_URL
      value: "http://spark-history:18080"
    - name: AWS_REGION
      value: "eu-west-1"
    # Монтируем секреты через K8s Secret
    - name: DATABASE_PASSWORD
      valueFrom:
        secretKeyRef:
          name: spark-secrets
          key: db-password

    volumeMounts:
    # Локальный диск для Spark event log (Driver пишет его весь job)
    - name: event-log-volume
      mountPath: /tmp/spark-events

  volumes:
  # EmptyDir с medium: Memory - RAM-disk для быстрых временных данных
  - name: event-log-volume
    emptyDir:
      sizeLimit: 5Gi

Executor Pod Template: скорость и стоимость

Executor-поды - расходный материал: их много, они создаются и удаляются динамически. Главные приоритеты: скорость запуска, эффективное использование диска для shuffle, дешёвые Spot-инстансы.

# executor-pod-template.yaml
apiVersion: v1
kind: Pod
metadata:
  labels:
    app: spark-executor
spec:
  # Spot/Preemptible ноды - в 3x дешевле On-Demand
  # Потеря Executor - не катастрофа: задача перезапустит стейдж
  affinity:
    nodeAffinity:
      preferredDuringSchedulingIgnoredDuringExecution:
      - weight: 80
        preference:
          matchExpressions:
          - key: "node.kubernetes.io/lifecycle"
            operator: In
            values: ["spot", "preemptible"]
      - weight: 20
        preference:
          matchExpressions:
          - key: "node.kubernetes.io/lifecycle"
            operator: In
            values: ["on-demand"]

  # Anti-affinity: spread Executor-ы по разным нодам
  # Если одна нода упадёт - не теряем все Executor-ы сразу
    podAntiAffinity:
      preferredDuringSchedulingIgnoredDuringExecution:
      - weight: 50
        podAffinityTerm:
          labelSelector:
            matchLabels:
              app: spark-executor
          topologyKey: kubernetes.io/hostname

  tolerations:
  - key: "spot-instance"
    operator: "Exists"
    effect: "NoSchedule"

  securityContext:
    runAsNonRoot: true
    runAsUser: 1000

  containers:
  - name: spark-executor
    resources:
      requests:
        cpu: "4"
        memory: "16Gi"
      limits:
        cpu: "4"
        memory: "20Gi"    # Spark использует memoryOverhead сверх executor.memory

    volumeMounts:
    # Локальный NVMe SSD для shuffle-данных
    # Намного быстрее, чем запись shuffle на сетевой диск
    - name: shuffle-volume
      mountPath: /mnt/spark-local
    # tmpfs для временных буферов (RAM-based, очень быстрый)
    - name: tmp-volume
      mountPath: /tmp

  volumes:
  # hostPath: монтируем NVMe диск ноды напрямую в контейнер
  # (нода должна иметь /mnt/nvme предмонтированный)
  - name: shuffle-volume
    hostPath:
      path: /mnt/nvme/spark-shuffle
      type: DirectoryOrCreate
  # tmpfs: in-memory filesystem для temp файлов
  - name: tmp-volume
    emptyDir:
      medium: Memory
      sizeLimit: 4Gi

Подключение Pod Templates к задаче

from pyspark.sql import SparkSession

spark = (
    SparkSession.builder
    .appName("gold-layer-job")
    .master("k8s://https://k8s-api.company.com:6443")

    # Указываем Pod Template файлы
    .config("spark.kubernetes.driver.podTemplateFile",
            "s3://my-bucket/templates/driver-pod-template.yaml")
    .config("spark.kubernetes.executor.podTemplateFile",
            "s3://my-bucket/templates/executor-pod-template.yaml")

    # Образ с нашими зависимостями
    .config("spark.kubernetes.container.image",
            "123456789.dkr.ecr.eu-west-1.amazonaws.com/spark:3.5-analytics-v2.1")

    # Namespace и Service Account
    .config("spark.kubernetes.namespace", "spark-analytics")
    .config("spark.kubernetes.authenticate.driver.serviceAccountName",
            "spark-sa")

    # Пути к shuffle на локальном NVMe (монтируется через template)
    .config("spark.local.dir", "/mnt/spark-local/tmp")

    .getOrCreate()
)

Сборка кастомного Docker-образа

Для Serverless Spark необходим Docker-образ, содержащий всё необходимое: JVM, Spark, Python-зависимости, дополнительные JAR.

# Dockerfile для production Spark образа
# За основу берём официальный образ 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/*

# Копируем дополнительные JAR (iceberg, delta, s3a коннекторы)
COPY jars/iceberg-spark-runtime-3.5_2.12-1.4.2.jar $SPARK_HOME/jars/
COPY jars/delta-spark_2.12-3.1.0.jar $SPARK_HOME/jars/
COPY jars/aws-java-sdk-bundle-1.12.262.jar $SPARK_HOME/jars/
COPY jars/hadoop-aws-3.3.4.jar $SPARK_HOME/jars/

# Python зависимости: устанавливаем через requirements.txt
COPY requirements.txt /tmp/requirements.txt
RUN pip install --no-cache-dir -r /tmp/requirements.txt

# requirements.txt:
# pyarrow==14.0.2
# pandas==2.1.4
# boto3==1.34.0
# iceberg-python==0.5.1

# Копируем код приложения
COPY src/ /opt/spark/work-dir/src/

# Возвращаемся к непривилегированному пользователю
USER 1000

# Проверка: все JAR на месте
RUN ls $SPARK_HOME/jars/ | grep -E "(iceberg|delta|hadoop-aws)"
# CI/CD скрипт сборки и публикации образа
#!/bin/bash
set -euo pipefail

VERSION="3.5-analytics-v$(date +%Y%m%d-%H%M)"
ECR_REPO="123456789.dkr.ecr.eu-west-1.amazonaws.com/spark"

# Авторизация в ECR
aws ecr get-login-password --region eu-west-1 | \
  docker login --username AWS --password-stdin ${ECR_REPO}

# Сборка (многоплатформенная для ARM и x86)
docker buildx build \
  --platform linux/amd64,linux/arm64 \
  --tag ${ECR_REPO}:${VERSION} \
  --tag ${ECR_REPO}:latest \
  --push \
  .

echo "Published: ${ECR_REPO}:${VERSION}"

Экстремальный Auto-scaling и Dynamic Allocation

Как работает Dynamic Allocation на K8s

Spark Dynamic Allocation - это механизм автоматического управления количеством Executor-ов в ходе выполнения задачи. Принцип: если задачи ждут в очереди - добавить Executor-ов; если Executor-ы простаивают - удалить их.

На K8s Dynamic Allocation реализуется иначе, чем на YARN. В YARN Executor-ы возвращались в пул ресурсов. На K8s Executor-под удаляется (DELETE pod) - и создаётся заново при необходимости.

spark = (
    SparkSession.builder

    # Включаем Dynamic Allocation
    .config("spark.dynamicAllocation.enabled", "true")

    # Начальное количество Executor-ов при старте задачи
    .config("spark.dynamicAllocation.initialExecutors", "2")

    # Минимум (не гасим всё до нуля - нужно время на перезапуск)
    .config("spark.dynamicAllocation.minExecutors", "1")

    # Максимум - ограничение масштабирования
    .config("spark.dynamicAllocation.maxExecutors", "50")

    # Через сколько секунд ожидания задач запрашиваем новый Executor
    # Маленькое значение = быстрый scale-up, но много ненужных Executor-ов
    .config("spark.dynamicAllocation.schedulerBacklogTimeout", "30s")

    # Через сколько секунд простоя Executor удаляется
    # Маленькое = агрессивная экономия, но overhead на пересоздание
    .config("spark.dynamicAllocation.executorIdleTimeout", "60s")

    # Shuffle tracking - критически важно для K8s (объяснение ниже)
    .config("spark.dynamicAllocation.shuffleTracking.enabled", "true")
    .config("spark.dynamicAllocation.shuffleTracking.timeout", "600s")

    .getOrCreate()
)

Проблема Shuffle Data Loss: самая опасная ловушка

Shuffle - это запись промежуточных данных на локальный диск Executor-а и их чтение другими Executor-ами. На YARN Executor-ы работали долго, и их shuffle-данные были доступны весь жизненный цикл задачи.

На K8s с Dynamic Allocation возникает опасная ситуация:

Stage 1: Executor A записал shuffle-файлы на свой локальный диск
Stage 2: Executor A простаивает → Dynamic Allocation гасит его (DELETE pod)
Stage 2 (другой Executor): пытается прочитать shuffle-файлы Executor A
          → Pod удалён → данных нет → FAILED!
          → Spark перезапускает Stage 1 (полностью!)

Это может приводить к бесконечным retry и значительному замедлению задачи.

Решение 1: Shuffle Tracking

spark.dynamicAllocation.shuffleTracking.enabled=true - Spark отслеживает, какие Executor-ы хранят shuffle-данные, нужные для текущих и будущих стейджей. Такие Executor-ы не удаляются до тех пор, пока их shuffle-данные не будут прочитаны. Это безопасно, но означает, что некоторые Executor-ы не смогут быть масштабированы вниз немедленно.

Решение 2: Remote Shuffle Service

Remote Shuffle Service (RSS) - это отдельный сервис, который принимает shuffle-данные от Executor-ов и хранит их вне Executor-контейнеров. Когда Executor удаляется, его shuffle-данные остаются в RSS и доступны другим Executor-ам.

Популярные реализации RSS:

  • Apache Celeborn (ex-Uber RSS): open-source, поддерживает Spark 3.x
  • Magnet (LinkedIn): интегрирован в некоторые корпоративные дистрибутивы
  • Crail: исследовательский проект CMU
# Конфигурация Apache Celeborn как Remote Shuffle Service
spark = (
    SparkSession.builder

    # Указываем Spark использовать RSS как shuffle manager
    .config("spark.shuffle.manager",
            "org.apache.spark.shuffle.celeborn.SparkShuffleManager")

    # Адрес Celeborn master
    .config("spark.celeborn.master.endpoints",
            "celeborn-master-1:9097,celeborn-master-2:9097")

    # Fallback на стандартный shuffle если RSS недоступен
    .config("spark.celeborn.shuffle.fallback.policy", "ALWAYS")

    # Теперь Dynamic Allocation может безопасно удалять Executor-ы
    # даже если у них есть непрочитанные shuffle-данные - они в RSS!
    .config("spark.dynamicAllocation.enabled", "true")
    .config("spark.dynamicAllocation.shuffleTracking.enabled", "false")

    .getOrCreate()
)

Kubernetes Cluster Autoscaler: масштабирование нод

Dynamic Allocation управляет количеством Executor-подов. Но если в кластере нет свободных нод для размещения новых подов - поды остаются в состоянии Pending.

Kubernetes Cluster Autoscaler (CA) решает это: он наблюдает за Pending-подами и автоматически запрашивает новые ноды у облачного провайдера (AWS EC2, GCP GCE). Новая нода появляется за 2–5 минут - именно это время является "реальным" cold start для масштабирования вверх.

# cluster-autoscaler-config.yaml (упрощённо)
apiVersion: apps/v1
kind: Deployment
metadata:
  name: cluster-autoscaler
spec:
  template:
    spec:
      containers:
      - name: cluster-autoscaler
        args:
        - --cloud-provider=aws
        - --node-group-auto-discovery=asg:tag=k8s.io/cluster-autoscaler/enabled,k8s.io/cluster-autoscaler/my-cluster
        # Ждём 2 минуты перед добавлением ноды (избегаем флаппинга)
        - --scale-up-from-zero=true
        - --new-pod-scale-up-delay=0
        # Ждём 10 минут перед удалением незагруженной ноды
        - --scale-down-delay-after-add=10m
        - --scale-down-unneeded-time=10m

Взаимодействие Dynamic Allocation и Cluster Autoscaler:

1. Spark нужно больше Executor-ов → запрашивает K8s API
2. K8s пробует разместить новые поды → не хватает нод → поды в Pending
3. CA видит Pending поды → запрашивает новые EC2/GCE ноды
4. Новые ноды появляются через 2–5 минут
5. Поды переходят в Running → Executor-ы присоединяются к задаче
6. После завершения задачи Executor-поды удаляются
7. Через 10 минут пустые ноды удаляются CA → расходы обнуляются

Тюнинг скорости масштабирования

Два параметра определяют «поведение» Dynamic Allocation:

schedulerBacklogTimeout - как долго задачи должны ждать в очереди перед запросом нового Executor-а. Маленькое значение (5–15 секунд) = агрессивное масштабирование вверх, больше мгновенных Executor-ов, но потенциально лишние если задача имеет переменную нагрузку. Большое значение (60+ секунд) = экономия, но задачи долго ждут.

executorIdleTimeout - через сколько секунд простоя удалять Executor. Маленькое значение (30–60 секунд) = быстрая экономия, но overhead на пересоздание если Executor снова понадобился. Большое значение (300+ секунд) = меньше overhead, но дольше оплачиваем простой.

Оптимальные значения зависят от характера задачи:

# Для batch ETL с длинными стейджами и мало shuffle:
.config("spark.dynamicAllocation.schedulerBacklogTimeout", "15s")
.config("spark.dynamicAllocation.executorIdleTimeout", "120s")

# Для задач с частым shuffle (shuffle-intensive): сохраняем Executor-ы дольше
.config("spark.dynamicAllocation.schedulerBacklogTimeout", "30s")
.config("spark.dynamicAllocation.executorIdleTimeout", "300s")
.config("spark.dynamicAllocation.shuffleTracking.timeout", "1800s")

# Для streaming micro-batch (маленькие батчи, частые запуски):
# Dynamic Allocation лучше отключить - overhead на создание/удаление > экономии
.config("spark.dynamicAllocation.enabled", "false")
.config("spark.executor.instances", "10")   # фиксированное количество

Лабораторная практика: Gold-слой в Serverless Spark

Сценарий: ежедневный расчёт Gold-витрины

Пайплайн: ежедневно в 07:00 UTC запускается расчёт агрегированных продаж за предыдущий день. Данные поступают в Silver-слой (Parquet на S3/GCS) из 15 источников. Задача должна:

  • Стартовать за < 2 минут от триггера Airflow
  • Масштабироваться до 50 Executor-ов в пике
  • Полностью гасить ресурсы по завершении
  • Стоить не более $15 за запуск

Шаг 1: PySpark-скрипт с оптимальной конфигурацией

# gold_layer.py - основной скрипт для Serverless-деплоя
import argparse
from datetime import date, timedelta
from pyspark.sql import SparkSession, functions as F
from pyspark.sql.window import Window

def create_spark_session() -> SparkSession:
    return (
        SparkSession.builder
        .appName("gold-sales-daily")

        # S3/GCS: используем instance credentials (IRSA/Workload Identity)
        # Явный ключ не нужен - IAM Role пода имеет доступ
        .config("spark.hadoop.fs.s3a.aws.credentials.provider",
                "com.amazonaws.auth.WebIdentityTokenCredentialsProvider")

        # Iceberg runtime
        .config("spark.sql.extensions",
                "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions")
        .config("spark.sql.catalog.glue", "org.apache.iceberg.spark.SparkCatalog")
        .config("spark.sql.catalog.glue.type", "glue")
        .config("spark.sql.catalog.glue.warehouse", "s3://my-lakehouse/warehouse/")

        # Dynamic Allocation с защитой shuffle
        .config("spark.dynamicAllocation.enabled", "true")
        .config("spark.dynamicAllocation.minExecutors", "2")
        .config("spark.dynamicAllocation.maxExecutors", "50")
        .config("spark.dynamicAllocation.schedulerBacklogTimeout", "20s")
        .config("spark.dynamicAllocation.executorIdleTimeout", "120s")
        .config("spark.dynamicAllocation.shuffleTracking.enabled", "true")

        # AQE: адаптивное планирование (особенно важно при переменном числе Executor)
        .config("spark.sql.adaptive.enabled", "true")
        .config("spark.sql.adaptive.coalescePartitions.enabled", "true")
        .config("spark.sql.adaptive.skewJoin.enabled", "true")

        # Оптимальный размер shuffle-партиций
        .config("spark.sql.adaptive.targetPostShuffleInputSize", "134217728")  # 128 MB

        .getOrCreate()
    )


def run_gold_layer(spark: SparkSession, process_date: str) -> None:
    """Расчёт Gold-витрины продаж за указанную дату."""

    prev_date = (date.fromisoformat(process_date) - timedelta(days=1)).isoformat()

    # Читаем Silver-данные (инкрементально за вчера)
    df_orders = (
        spark.read
        .format("iceberg")
        .load("glue.silver.orders")
        .filter(F.col("order_date") == process_date)
    )

    df_products = (
        spark.read.format("iceberg").load("glue.silver.products")
    )

    df_customers = (
        spark.read.format("iceberg").load("glue.silver.customers")
    )

    # Тяжёлая агрегация: именно здесь Dynamic Allocation будет масштабироваться
    w_customer = Window.partitionBy("customer_id").orderBy("order_date")

    gold = (
        df_orders
        .join(df_products.select("product_id", "category_l1", "brand"),
              "product_id")
        .join(df_customers.select("customer_id", "region", "segment"),
              "customer_id")
        .withColumn("running_revenue",
                    F.sum("amount").over(w_customer))
        .groupBy("region", "segment", "category_l1", "brand")
        .agg(
            F.sum("amount").alias("daily_revenue"),
            F.count("order_id").alias("order_count"),
            F.countDistinct("customer_id").alias("unique_customers"),
            F.avg("amount").alias("avg_order_value"),
            F.sum(F.when(F.col("is_return"), F.col("amount")).otherwise(0))
             .alias("return_amount"),
        )
        .withColumn("return_rate",
                    F.col("return_amount") / F.col("daily_revenue"))
        .withColumn("process_date", F.lit(process_date))
    )

    # Запись в Gold-таблицу через MERGE (upsert для идемпотентности)
    gold.createOrReplaceTempView("gold_updates")

    spark.sql(f"""
        MERGE INTO glue.gold.sales_daily AS target
        USING gold_updates AS source
        ON target.region = source.region
           AND target.segment = source.segment
           AND target.category_l1 = source.category_l1
           AND target.brand = source.brand
           AND target.process_date = source.process_date
        WHEN MATCHED THEN UPDATE SET *
        WHEN NOT MATCHED THEN INSERT *
    """)

    print(f"Gold layer updated for {process_date}: "
          f"{gold.count()} rows written")


if __name__ == "__main__":
    parser = argparse.ArgumentParser()
    parser.add_argument("--date", required=True, help="Processing date YYYY-MM-DD")
    args = parser.parse_args()

    spark = create_spark_session()
    run_gold_layer(spark, args.date)
    spark.stop()

Шаг 2: Airflow DAG для оркестрации

# dags/gold_layer_daily.py
from datetime import datetime, timedelta
from airflow import DAG
from airflow.providers.amazon.aws.operators.emr import (
    EmrServerlessStartJobOperator,
    EmrServerlessStopApplicationOperator,
)
from airflow.providers.amazon.aws.sensors.emr import EmrServerlessJobSensor

EMR_APPLICATION_ID = "00fv43iq0jrsgf0l"
EMR_EXECUTION_ROLE = "arn:aws:iam::123456789:role/EMRServerlessExecutionRole"
S3_BUCKET = "my-lakehouse"

default_args = {
    "owner": "data-engineering",
    "retries": 2,
    "retry_delay": timedelta(minutes=5),
    "email_on_failure": True,
}

with DAG(
    dag_id="gold_layer_daily",
    default_args=default_args,
    schedule_interval="0 7 * * *",     # 07:00 UTC ежедневно
    start_date=datetime(2024, 1, 1),
    catchup=False,
    tags=["gold", "analytics"],
) as dag:

    # Запуск Spark-задачи в EMR Serverless
    submit_job = EmrServerlessStartJobOperator(
        task_id="submit_gold_spark_job",
        application_id=EMR_APPLICATION_ID,
        execution_role_arn=EMR_EXECUTION_ROLE,
        job_driver={
            "sparkSubmit": {
                "entryPoint": f"s3://{S3_BUCKET}/scripts/gold_layer.py",
                "entryPointArguments": [
                    "--date", "{{ ds }}"   # дата запуска DAG (YYYY-MM-DD)
                ],
                "sparkSubmitParameters": " ".join([
                    # Kubernetes pod templates
                    f"--conf spark.kubernetes.driver.podTemplateFile="
                    f"s3://{S3_BUCKET}/templates/driver-pod-template.yaml",
                    f"--conf spark.kubernetes.executor.podTemplateFile="
                    f"s3://{S3_BUCKET}/templates/executor-pod-template.yaml",
                    # Образ
                    "--conf spark.kubernetes.container.image="
                    "123456789.dkr.ecr.eu-west-1.amazonaws.com/spark:3.5-analytics-v2.1",
                    # Ресурсы Driver
                    "--conf spark.driver.cores=4",
                    "--conf spark.driver.memory=16g",
                    # Ресурсы Executor
                    "--conf spark.executor.cores=4",
                    "--conf spark.executor.memory=16g",
                    "--conf spark.executor.memoryOverhead=4g",
                ])
            }
        },
        configuration_overrides={
            "monitoringConfiguration": {
                "s3MonitoringConfiguration": {
                    "logUri": f"s3://{S3_BUCKET}/logs/emr-serverless/"
                }
            }
        },
        aws_conn_id="aws_default",
    )

    # Ожидание завершения задачи
    wait_for_job = EmrServerlessJobSensor(
        task_id="wait_gold_job_completion",
        application_id=EMR_APPLICATION_ID,
        job_run_id="{{ task_instance.xcom_pull('submit_gold_spark_job') }}",
        aws_conn_id="aws_default",
        # Таймаут: если задача не завершилась за 2 часа - считаем failed
        timeout=7200,
        poke_interval=60,     # проверяем каждую минуту
    )

    submit_job >> wait_for_job

Шаг 3: Мониторинг Auto-scaling

После запуска задачи отслеживаем поведение Dynamic Allocation в Spark UI:

Вывод логов Driver (динамика Executor-ов):
07:02:15 INFO DynamicAllocation: Requesting 2 executors (initial)
07:02:45 INFO ExecutorAllocationManager: Requesting 8 more executors (backlog: 45 tasks)
07:03:15 INFO ExecutorAllocationManager: Requesting 15 more executors (backlog: 62 tasks)
07:05:30 INFO ExecutorAllocationManager: Requesting 25 more executors (peak shuffle)
07:08:00 INFO ExecutorAllocationManager: Cancelling 10 executors (idle > 120s)
07:12:00 INFO ExecutorAllocationManager: Cancelling 15 executors (idle > 120s)
07:14:00 INFO ExecutorAllocationManager: Cancelling 5 executors (job finishing)
07:14:45 INFO SparkContext: Successfully stopped SparkContext
# Скрипт для получения метрик масштабирования после выполнения
import boto3
import json

def get_job_metrics(application_id: str, job_run_id: str) -> dict:
    client = boto3.client("emr-serverless")
    response = client.get_job_run(
        applicationId=application_id,
        jobRunId=job_run_id
    )
    job = response["jobRun"]

    print(f"Job State: {job['state']}")
    print(f"Duration: {job.get('totalExecutionDurationSeconds', 0) / 60:.1f} minutes")

    # Billing: vCPU-seconds и memory-GB-seconds
    total_resource_utilization = job.get("totalResourceUtilization", {})
    vcpu_hours = total_resource_utilization.get("vCPUHour", 0)
    memory_gb_hours = total_resource_utilization.get("memoryGBHour", 0)

    # Расчёт стоимости (цены 2024)
    vcpu_cost = vcpu_hours * 0.04730
    memory_cost = memory_gb_hours * 0.01847
    total_cost = vcpu_cost + memory_cost

    print(f"vCPU-hours consumed: {vcpu_hours:.2f}")
    print(f"Memory-GB-hours: {memory_gb_hours:.2f}")
    print(f"Estimated cost: ${total_cost:.2f}")

    return {"duration_min": job.get('totalExecutionDurationSeconds', 0) / 60,
            "cost_usd": total_cost}

# После завершения задачи:
# Job State: SUCCESS
# Duration: 12.4 minutes
# vCPU-hours consumed: 18.7
# Memory-GB-hours: 74.8
# Estimated cost: $2.27

Матрица выбора: Serverless vs постоянный кластер

Когда Serverless - правильный выбор

Нерегулярные batch-пайплайны: задачи, запускаемые 1–4 раза в день или по событию, когда кластер большую часть времени простаивает. Экономия: 70–90% стоимости инфраструктуры.

Ad-hoc аналитика: разовые запросы аналитиков, которые нельзя предсказать. Serverless даёт ресурсы мгновенно, не нужно ждать свободного слота в постоянном кластере.

Резкие скачки объёмов: например, обработка данных за квартал в конце квартала - задача в 10× крупнее обычной. Serverless масштабируется до нужного размера автоматически.

Разные требования к ресурсам: если одни задачи требуют 2 Executor-а, другие - 50, а третьи - GPU. Serverless выделяет именно то, что нужно каждой задаче.

Когда постоянный кластер выгоднее

Structured Streaming с постоянной нагрузкой: micro-batch запускается каждые 5–30 секунд. Cold-start overhead для каждого micro-batch недопустим. Executor-ы должны работать непрерывно. Serverless для streaming не подходит.

Круглосуточная высокая нагрузка: если кластер загружен на 80%+ круглосуточно, постоянный кластер дешевле - нет overhead на создание/удаление контейнеров, нет платы за Pre-initialized Capacity.

State-heavy задачи: если задача использует большой RDD.cache() или DataFrame.persist(), постоянный кластер сохраняет кэш между запусками. Serverless теряет всё в памяти после завершения.

Интерактивные запросы с низкой задержкой: Jupyter Notebook, BI-инструменты, которые ожидают ответа за < 5 секунд. Serverless cold-start сделает первый запрос медленным.

Сценарий Serverless Постоянный
Batch ETL 1–4 раза/день Да Нет (платим за простой)
Structured Streaming 24/7 Нет Да
Ad-hoc аналитика Да Нет (нужно резервировать)
Загрузка >80% круглосуточно Нет Да (дешевле)
ML feature engineering (пакетный) Да Нет
Кэширование данных между запросами Нет Да
Переменный объём данных (×10 сезонность) Да Нет (overprovisioning)

Домашнее задание

Условие:

Вам дан тяжёлый PySpark-скрипт обработки исторических данных о транзакциях за 3 года (500 GB Parquet в Silver-слое). Задача запускается ежедневно в 06:00, занимает 45 минут на постоянном EMR-кластере (6 × r5.4xlarge: 16 vCPU, 128 GB RAM каждый). Стоимость кластера: $0.904/час × 6 нод × 24 часа = $130.2/день (кластер работает круглосуточно).

Вы должны спроектировать манифест для перевода этой задачи на AWS EMR Serverless.

Задание состоит из трёх частей:

Часть 1: Конфигурационный манифест

Составьте полную конфигурацию spark-submit с параметрами для Serverless-деплоя:

  • spark.dynamicAllocation.* с обоснованием каждого числа (minExecutors, maxExecutors, schedulerBacklogTimeout, executorIdleTimeout)
  • Размеры Driver и Executor (vCPU, RAM, disk) с обоснованием от характеристик задачи
  • Нужна ли Pre-initialized Capacity? Если да - сколько воркеров и почему?
  • Включить Shuffle Tracking или Remote Shuffle Service? Аргументируйте.

Часть 2: Pod Template для Executor-ов

Напишите executor-pod-template.yaml с:

  • Spot/Preemptible нодами (с обоснованием: почему это безопасно для данной задачи)
  • Монтированием NVMe SSD для shuffle-данных (с расчётом необходимого объёма)
  • Anti-affinity правилом для распределения Executor-ов по разным нодам

Часть 3: FinOps-обоснование

Рассчитайте стоимость ежедневного запуска в EMR Serverless:

  • Сколько vCPU-часов потребляет задача (Driver + пиковое число Executor-ов × продолжительность)?
  • По ценам EMR Serverless 2024 ($0.04730 за vCPU-час + $0.01847 за GB RAM-час) - итоговая стоимость за один запуск?
  • Сравните с текущей стоимостью постоянного кластера ($130.2/день)?
  • При каком количестве запусков в день постоянный кластер становится выгоднее Serverless?

Ожидаемый результат: три YAML/текстовых артефакта с расчётами и обоснованием каждого решения. Итоговая ежедневная стоимость в Serverless должна быть < $20 при той же задаче в 45 минут.