Serverless Spark на K8s: EMR Serverless, Dataproc Serverless, pod templates и auto-scaling
Serverless Spark на K8s: EMR Serverless, Dataproc Serverless, pod templates и auto-scaling
Эволюция инфраструктуры 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 убирает понятие «кластер» из словаря пользователя. Вы описываете задачу (код, зависимости, ресурсы), отправляете её в облако - и облако само:
- Выделяет нужное количество вычислительных ресурсов (за секунды, не минуты)
- Запускает задачу
- Масштабирует ресурсы в ходе выполнения (вверх и вниз)
- Освобождает всё после завершения
- Выставляет счёт только за фактически потреблённые ресурсы (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 (легковесные микровиртуальные машины) для изоляции задач.
Жизненный цикл задачи:
- Вы создаёте Application - логическую единицу Serverless-среды (тип runtime: Spark, версия, регион).
- Отправляете JobRun - конкретное задание с кодом, конфигурацией и IAM Role.
- EMR Serverless аллоцирует Driver и Executor-воркеры (занимает 10–60 секунд для cold start или < 5 секунд при Pre-initialized Capacity).
- Задача выполняется; EMR автоматически масштабирует воркеры.
- После завершения воркеры освобождаются. Платите только за 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 минут.