Executor Sizing: 4-компонентная формула памяти и правило 5 cores
Инженерная математика Executor Sizing: иерархия Node→Container→Task, правило 5 cores и почему нельзя Fat Executors, 4-компонентная формула памяти (On-Heap + Overhead + Off-Heap + PySpark), архитектурный шлюз PySpark и ловушка Python workers, диагностика через Spark UI и пошаговый расчёт кластера под SLA.
1. Архитектурная анатомия воркера: Node → Container → Task¶
Прежде чем говорить о числах, необходимо понять физическую иерархию, на которой работает Spark. Большинство ошибок конфигурации происходит именно потому, что инженер не различает уровни этой иерархии и путает «сервер», «Executor» и «задачу».
Три уровня иерархии: физика и абстракции¶
Схема показывает три уровня: физический сервер содержит несколько Executor'ов (JVM-процессов), каждый Executor запускает несколько Task'ов (Java-потоков) параллельно.
Node (физический/виртуальный сервер) - это железо или VM. Он имеет конкретные CPU, RAM, диск, сеть. YARN NodeManager или K8s Kubelet занимает часть ресурсов для управления. Оставшееся распределяется между Executor'ами.
Executor (JVM-процесс) - это самостоятельный Java-процесс, который Spark запускает на каждом Node. У каждого Executor'а есть выделенная ему память и набор CPU-ядер. Executor'ы живут на протяжении всего приложения (если не включён Dynamic Allocation).
Task (Java-поток внутри Executor) - это минимальная единица работы. Каждый Task обрабатывает одну Partition данных. Task'и запускаются параллельно: число одновременных Task'ов = число ядер Executor'а.
Три антипаттерна конфигурации¶
Tiny Executors (один Executor на сервер, 1 core):
Запускаем 24 Executor'а на сервере с 24 ядрами, у каждого 1 core и 4 GB RAM. Что получаем:
- Overhead на 24 JVM-процесса: мониторинг, heartbeat, GC в каждом - потребляет до 20% ресурсов только на инфраструктуру
- Нет broadcast sharing: каждый Executor загружает свою копию broadcast-переменных (справочников)
- Нет Off-Heap преимуществ: 4 GB слишком мало для Tungsten
- Результат: 30-50% ресурсов уходит на накладные расходы вместо полезной работы
Fat Executors (один Executor на весь сервер, 24 cores):
Один Executor с 24 ядрами и 90 GB RAM. Выглядит логично - «всё в одной JVM, меньше overhead». Реальность:
- JVM GC обрабатывает 90 GB Heap. Full GC (Stop-The-World) на 90 GB может занимать 30-90 секунд. Весь Executor замирает - все 24 Task'а ждут.
- HDFS-клиент с 24 параллельными потоками пытается открыть 24 одновременных соединения с DataNode. Это вызывает троттлинг и конкуренцию за сетевые буферы.
- Если один Task упал → перезапускается весь Stage на том же огромном Executor'е.
- Результат: GC pauses убивают производительность при любой нагрузке с большими данными.
Правильный паттерн (5 cores, ~20 GB): разбираем в следующем разделе.
2. Золотое правило «5 Cores»: баланс многопоточности¶
Правило 5 ядер - эмпирический стандарт, выработанный на основе production-опыта тысяч Spark-кластеров. Понять его природу важнее, чем просто следовать цифре.
Почему именно 5, а не 3 и не 10¶
Параллелизм HDFS-клиента. Spark S3A/HDFS клиент внутри Executor'а имеет пул соединений для чтения блоков. При 5 потоках пять одновременных TCP-соединений к DataNode дают высокую пропускную способность без конкуренции. При 10+ потоках начинается contention: потоки конкурируют за одни и те же DataNode-соединения, сетевые буферы переполняются, появляются переотправки пакетов.
GC давление. 5 одновременных Task'ов генерируют умеренный «мусор» в Java Heap (промежуточные объекты при обработке данных). G1 GC или ZGC справляются с Minor GC за 50-200 ms без Stop-The-World. При 15+ Task'ах объём «мусора» растёт быстрее чем GC успевает убирать → накопление → Full GC.
CPU Cache Locality. 5 потоков на современном серверном CPU (обычно 2 сокета × 12 ядер = 24 ядра) умещаются в пределах одного NUMA-узла. Доступ к данным в L3-кеше одного NUMA-узла в 3-5 раз быстрее чем cross-NUMA. При 10+ потоках начинаются cross-NUMA операции.
Пошаговый алгоритм нарезки физического сервера¶
Рассмотрим конкретный сервер: 24 vCPU, 96 GB RAM на YARN-кластере.
# Параметры физического сервера
total_cpu = 24
total_ram_gb = 96
# Шаг 1: Резервируем ресурсы под ОС и YARN NodeManager
# Правило: 1 CPU + 1 GB для ОС, 1 CPU + 1 GB для NodeManager
os_cpu = 1
nm_cpu = 1
os_ram_gb = 1
nm_ram_gb = 1
available_cpu = total_cpu - os_cpu - nm_cpu # = 22
available_ram_gb = total_ram_gb - os_ram_gb - nm_ram_gb # = 94 GB
# Шаг 2: Вычисляем количество Executor'ов по правилу 5 cores
executor_cores = 5
# Число Executor'ов = floor(available_cpu / executor_cores)
num_executors_per_node = available_cpu // executor_cores # = 22 // 5 = 4
# Шаг 3: Вычисляем RAM на Executor
# available_ram / num_executors (с небольшим запасом)
ram_per_executor_raw = available_ram_gb / num_executors_per_node # = 94 / 4 = 23.5 GB
# Округляем вниз для безопасности
ram_per_executor_gb = 23 # GB
print(f"Сервер {total_cpu} vCPU / {total_ram_gb} GB:")
print(f" Executor'ов на ноду: {num_executors_per_node}")
print(f" Ядер на Executor: {executor_cores}")
print(f" RAM на Executor: {ram_per_executor_gb} GB")
# Проверка: суммарное потребление
total_used_cpu = num_executors_per_node * executor_cores + os_cpu + nm_cpu # 20 + 2 = 22 ✓
total_used_ram = num_executors_per_node * ram_per_executor_gb + os_ram_gb + nm_ram_gb # 92 + 2 = 94 ✓
print(f"Утилизация: CPU={total_used_cpu}/{total_cpu}, RAM={total_used_ram}/{total_ram_gb} GB")
Для кластера из 5 нод:
Узлы: 5
Executor'ов всего: 5 × 4 = 20 (один оставляем Driver'у)
Итого Executor'ов: 19 (рабочие) + 1 Driver
Общий параллелизм: 19 × 5 = 95 одновременных Task'ов
3. Четырёхкомпонентная формула памяти Spark Executor¶
Главное заблуждение начинающих инженеров: spark.executor.memory = 20GB означает, что YARN/K8s выделит контейнеру ровно 20 GB. Это не так. Параметр executor.memory - лишь один из четырёх компонентов общей памяти контейнера.
Полная формула¶
Total Container Memory =
spark.executor.memory (On-Heap JVM Heap)
+ spark.executor.memoryOverhead (Off-Heap нативная память)
+ spark.memory.offHeap.size (Tungsten Off-Heap, если включён)
+ spark.executor.pyspark.memory (Python Workers, только PySpark)
Схема показывает полную анатомию памяти контейнера. Контейнер = JVM Process + Overhead + PySpark. JVM Process = On-Heap (executor.memory) + Off-Heap Tungsten. On-Heap делится на Unified Memory (для Execution и Storage) и User Memory.
4. Глубокий разбор компонентов формулы¶
Компонент 1: spark.executor.memory (On-Heap JVM Heap)¶
Это основной пул памяти JVM-процесса Executor'а. Именно это значение указывается в параметре --executor-memory при spark-submit.
Внутри executor.memory Unified Memory Manager делит память на три части:
Reserved Memory (всегда 300 MB): зарезервирована жёстко для внутренних структур Spark (списки блоков, метаданные). Не настраивается.
User Memory = (1 - spark.memory.fraction) × (executor.memory - 300 MB): пространство для Java/Scala-объектов создаваемых пользовательским кодом: объекты в UDF, данные в рамках одной Task'и до Shuffle. Дефолт spark.memory.fraction = 0.6, то есть User Memory = 40% × (heap - 300 MB).
Unified Memory = spark.memory.fraction × (executor.memory - 300 MB): делится динамически между Execution и Storage Memory.
# Пример расчёта для executor.memory = 20 GB
executor_memory_mb = 20 * 1024 # = 20480 MB
reserved_mb = 300
usable_mb = executor_memory_mb - reserved_mb # = 20180 MB
memory_fraction = 0.6
unified_mb = usable_mb * memory_fraction # = 12108 MB ≈ 11.8 GB
user_mb = usable_mb * (1 - memory_fraction) # = 8072 MB ≈ 7.9 GB
print(f"executor.memory = 20 GB:")
print(f" Reserved: {reserved_mb} MB (фиксировано)")
print(f" Unified Memory: {unified_mb:.0f} MB ({unified_mb/1024:.1f} GB)")
print(f" User Memory: {user_mb:.0f} MB ({user_mb/1024:.1f} GB)")
print()
print(f"Unified Memory делится между:")
print(f" Execution Memory (Shuffle, Sort, Join) ↔ динамически ↔ Storage Memory (cache)")
Execution Memory используется для:
- Hash Join: хранение hash-таблицы для одной стороны JOIN
- Sort: буфер сортировки при SortMergeJoin или OrderBy
- Aggregation: хэш-таблица при Hash Aggregation
- Shuffle Write/Read буферы
Storage Memory используется для:
- Кешированные DataFrame через
.cache()или.persist() - Broadcast-переменные: small lookup-таблицы разосланные всем Executor'ам
- Unroll буферы при материализации кешированных партиций
Динамическое перераспределение: если Execution Memory нужно больше места и Storage Memory имеет свободное пространство - Execution может его «одолжить» (и наоборот). Но Execution Memory не может вытеснить уже занятый кеш (только если тот ещё не закреплён через StorageLevel.MEMORY_AND_DISK_SER).
Компонент 2: spark.executor.memoryOverhead (нативная JVM-память)¶
Это память вне Java Heap, необходимая для функционирования JVM и сетевых компонентов Spark.
# Формула расчёта по умолчанию:
import math
executor_memory_mb = 20 * 1024 # 20 GB = 20480 MB
overhead_fraction = 0.10 # 10%
overhead_mb = max(384, math.ceil(executor_memory_mb * overhead_fraction))
# = max(384, ceil(20480 * 0.10)) = max(384, 2048) = 2048 MB ≈ 2 GB
print(f"Memory Overhead = max(384MB, 10% × executor.memory)")
print(f"= max(384, {executor_memory_mb * 0.10:.0f}) = {overhead_mb} MB = {overhead_mb/1024:.1f} GB")
Что хранится в Memory Overhead:
- Netty NIO Direct Buffers: Spark использует Netty для сетевой передачи Shuffle-данных. NIO Direct Buffers - это буферы за пределами Heap, управляемые ОС напрямую. При большом Shuffle (например, джойн двух 100 GB таблиц) эти буферы могут вырасти до нескольких гигабайт.
- JVM thread stacks: каждый Java-поток (Task) потребляет ~1 MB для стека вызовов. При 5 Task'ах: 5 MB.
- JVM code cache: JIT-скомпилированный байткод методов Spark. При сложных планах выполнения может занимать 200-500 MB.
- Native library buffers: если используются нативные C++ библиотеки (ISA-L для Erasure Coding, zstd, lz4 компрессия) - их буферы живут в Overhead.
- Arrow buffers: при использовании Pandas UDF (
@udf(returnType=..., functionType=PandasUDFType....)) Apache Arrow сериализует данные в бинарные буферы вне Heap.
Когда увеличивать Overhead:
# Правило: base overhead × multiplier
def calculate_overhead(executor_memory_gb: float, workload_type: str) -> float:
"""
Рекомендуемый overhead для разных типов нагрузки.
"""
base_overhead_gb = max(0.375, executor_memory_gb * 0.10)
multipliers = {
"standard_scala": 1.0, # дефолт, стандартная Scala/Java ETL
"pyspark_basic": 1.5, # PySpark без тяжёлых UDF
"pyspark_pandas": 2.5, # PySpark с Pandas UDF / Apache Arrow
"heavy_shuffle": 1.5, # большой Shuffle, интенсивная сеть
"compression_ops": 2.0, # работа с zip/bzip2/lz4 нативно
"ml_tensorflow": 3.0, # ML фреймворки с нативными библиотеками
}
multiplier = multipliers.get(workload_type, 1.0)
return base_overhead_gb * multiplier
# Примеры:
for workload in ["standard_scala", "pyspark_pandas", "ml_tensorflow"]:
oh = calculate_overhead(20, workload)
print(f" {workload}: overhead = {oh:.1f} GB")
Симптом неправильного Overhead: Container killed by YARN for exceeding memory limits. X GB of Y GB physical memory used. Если X примерно равно executor.memory + overhead, но немного больше - нужно увеличить overhead.
Компонент 3: spark.memory.offHeap.size (Off-Heap Tungsten)¶
Project Tungsten - оптимизация Spark, которая размещает данные напрямую в нативной памяти ОС, полностью минуя JVM Garbage Collector.
spark = SparkSession.builder \
# Включить Off-Heap Tungsten память
.config("spark.memory.offHeap.enabled", "true") \
# Размер Off-Heap пула на Executor
# Рекомендация: 10-30% от executor.memory для начала
# Увеличьте если видите тяжёлый GC при сортировках/агрегациях
.config("spark.memory.offHeap.size", "4g") \
.getOrCreate()
Когда использовать Off-Heap:
- Тяжёлые сортировки (OrderBy над петабайтами): Tungsten хранит ключи сортировки в Off-Heap, что устраняет GC при работе с гигантскими Sort Runs
- Много долгоживущих объектов: если Spark кеширует в MEMORY_ONLY_2 - объекты долго живут в Heap и нагружают GC. Off-Heap кеш не создаёт GC-давления
- Векторизованные операции: Spark Vectorized Reader читает батчи данных в ColumnarBatch, которые хранятся в Off-Heap для работы с SIMD
Важно: offHeap.size не добавляется автоматически к memoryOverhead. YARN/K8s выделяет контейнеру сумму всех четырёх компонентов.
Компонент 4: spark.executor.pyspark.memory (Python Workers)¶
Это самый часто игнорируемый компонент - и именно он вызывает загадочные OOM в production PySpark пайплайнах.
Когда он работает: только в PySpark при наличии Python UDF, Pandas UDF, MapInPandas и других операций, требующих Python-интерпретатора на Executor'е.
Дефолтное значение: 0 - то есть Python Workers НЕ ограничены! Они могут занять любое количество памяти, и YARN/K8s убьёт контейнер по превышению лимита. Именно поэтому нужно явно устанавливать этот параметр.
Разбираем это в следующем разделе.
5. Ловушка PySpark: разрыв контекста JVM и Python Workers¶
Это одна из наиболее коварных областей PySpark-разработки. Архитектура PySpark предполагает два независимых процесса на каждом Executor'е, и незнание этого приводит к неожиданным OOM.
Анатомия архитектурного шлюза PySpark¶
Схема раскрывает ключевой момент: Python Workers - это отдельные процессы, не связанные с JVM Heap. Когда Python UDF читает данные через Arrow или Pandas, эти данные размещаются в Python Heap. Если Python Heap + JVM Heap превысят общий лимит контейнера - YARN/K8s убьёт контейнер.
Анатомия Python Worker-процесса¶
При каждом вызове Python UDF или Pandas UDF Spark запускает Python Worker из пула:
- Каждый Task'а, использующий Python код, получает свой Python Worker (или берёт из пула)
- Python Worker получает данные от JVM через сокет или Apache Arrow (нулевое копирование)
- Python код обрабатывает данные и возвращает результат обратно в JVM
Что потребляет память в Python Worker:
- Импорты библиотек:
import pandas,import numpy,import sklearn- каждый занимает 50-500 MB при импорте - Batch данных: если Pandas UDF получает батч из 10 000 строк × 50 колонок - это может быть 100-500 MB в pandas DataFrame
- Промежуточные вычисления: numpy операции создают временные массивы
- ML-модели: если UDF применяет sklearn модель - модель живёт в Python памяти
Конфигурация pyspark.memory и расчёт¶
spark = SparkSession.builder \
# spark.executor.pyspark.memory: сколько памяти разрешено КАЖДОМУ Python Worker'у
# Это НЕ суммарная память всех Python Workers, а лимит на один Worker-процесс!
# При 5 cores → потенциально 5 одновременных Python Workers
# Суммарная Python память = pyspark.memory × cores ≈ 2 GB × 5 = 10 GB
# Правило расчёта:
# 1. Оцените максимальный размер pandas batch: batch_size × avg_row_size
# 2. Добавьте overhead библиотек: ~500 MB для стандартного стека
# 3. Умножьте на 2 (промежуточные вычисления удваивают потребление)
.config("spark.executor.pyspark.memory", "2g") \
# Суммарная память контейнера теперь:
# executor.memory (20g) + memoryOverhead (3g) + pyspark.memory (2g × 5 cores)
# = 20 + 3 + 10 = 33 GB на контейнер
.getOrCreate()
Практический расчёт для PySpark с Pandas UDF:
def calculate_pyspark_memory(
batch_size_rows: int, # число строк в одном Pandas batch
avg_row_size_bytes: int, # средний размер одной строки в байтах
n_cores: int = 5, # cores на Executor
library_overhead_mb: int = 512, # pandas + numpy + прочее
) -> dict:
"""
Рассчитывает рекомендуемый spark.executor.pyspark.memory.
batch_size_rows: задаётся через spark.sql.execution.arrow.maxRecordsPerBatch
(дефолт: 10000 строк).
"""
# Размер одного batch в MB
batch_mb = (batch_size_rows * avg_row_size_bytes) / (1024 * 1024)
# Промежуточные данные при обработке (удваиваем)
processing_mb = batch_mb * 2
# Итого на один Worker
per_worker_mb = batch_mb + processing_mb + library_overhead_mb
# Рекомендуем с запасом 20%
recommended_mb = int(per_worker_mb * 1.2)
# Суммарная Python память на весь Executor
total_python_mb = recommended_mb * n_cores
return {
"batch_size_mb": batch_mb,
"per_worker_recommended_mb": recommended_mb,
"total_python_memory_mb": total_python_mb,
"suggested_config": f"{recommended_mb // 1024}g" if recommended_mb >= 1024
else f"{recommended_mb}m"
}
# Пример: Pandas UDF обрабатывает транзакции
result = calculate_pyspark_memory(
batch_size_rows=10000, # дефолт Arrow batch
avg_row_size_bytes=200, # 200 байт/строка (10 колонок × 20 байт)
n_cores=5,
library_overhead_mb=512 # pandas + numpy
)
print("Рекомендуемая pyspark.memory:")
print(f" Per Worker: {result['per_worker_recommended_mb']} MB")
print(f" Total (5 cores): {result['total_python_memory_mb']} MB")
print(f" Config: spark.executor.pyspark.memory = {result['suggested_config']}")
6. Аудит утилизации памяти через Spark UI¶
Правильно настроить Executor Sizing - это половина дела. Нужно уметь проверить что настройки работают и диагностировать проблемы в production.
Вкладка Executors: главный экран диагностики¶
В Spark UI (порт 4040) вкладка Executors показывает состояние каждого Executor'а в реальном времени. Ключевые колонки:
Storage Memory: показывает used / total. total = Unified Memory (то что доступно для кеша и Execution). Если used стабильно близко к total и при этом в Stages видны Spill → памяти недостаточно.
Task Time: суммарное время CPU всех Task'ов на этом Executor'е. Нормальная утилизация ≈ 70-90% от Task Time / Wall Clock Time.
GC Time: критически важная метрика. Правило: GC Time > 10% от Task Time - сигнал тревоги. При 20%+ GC - Executor страдает от Memory Pressure. Причины: слишком много данных в Heap, слишком большой executor.memory для G1GC, тяжёлые объекты в User Memory.
Peak Execution Memory: сколько памяти пиково использовалось для вычислений (Shuffle, Sort). Если близко к Unified Memory total → при следующем подобном запросе возможен Spill.
Вкладка Stages: мониторинг Spill¶
В вкладке Stages для каждого Stage видны агрегированные метрики Task'ов. Ключевые для диагностики:
Shuffle Spill (Memory): объём данных сброшенных из памяти на диск перед записью в Shuffle файл. Это бесплатный Spill - данные ещё в памяти Executor'а.
Shuffle Spill (Disk): объём данных уже записанных на диск из-за нехватки памяти. Это дорогостоящий Spill - диск в 100-1000x медленнее RAM.
Если видите Spill (Disk) > 0 - памяти Executor'а не хватает для данного Stage. Решения:
# Вариант 1: Увеличить executor.memory
spark.conf.set("spark.executor.memory", "32g")
# Вариант 2: Увеличить число партиций (меньше данных на Task)
spark.conf.set("spark.sql.shuffle.partitions", "400")
# Вариант 3: Освободить Storage Memory (uncache неиспользуемое)
spark.catalog.uncacheTable("old_cached_table")
# Вариант 4: AQE coalesce (поднять advisoryPartitionSizeInBytes)
spark.conf.set("spark.sql.adaptive.advisoryPartitionSizeInBytes", "128MB")
Программный сбор метрик: SparkContext StatusTracker¶
from pyspark.sql import SparkSession
import json
spark = SparkSession.builder.getOrCreate()
def get_executor_stats(spark: SparkSession) -> list[dict]:
"""
Собирает метрики Executor'ов через SparkContext StatusTracker.
Полезно для programmatic мониторинга в production.
"""
sc = spark.sparkContext
status = sc.statusTracker()
executor_infos = status.getExecutorInfos()
stats = []
for executor in executor_infos:
# Получаем метрики через REST API History Server (если доступен)
# Здесь упрощённая версия через доступные API
stats.append({
"executor_id": executor.executorId(),
"host": executor.host(),
"num_running_tasks": executor.numRunningTasks(),
"num_failed_tasks": executor.numFailedTasks(),
"num_completed_tasks": executor.numCompletedTasks(),
})
return stats
def check_gc_pressure(spark: SparkSession, gc_threshold_pct: float = 10.0) -> list[str]:
"""
Проверяет GC pressure на всех Executor'ах.
Возвращает список Executor'ов с проблемным GC.
"""
# Для доступа к GC Time используем History Server REST API:
# GET /api/v1/applications/{appId}/executors
# Поле: "totalGCTime", "totalDuration"
# Здесь показываем логику:
problem_executors = []
# Пример данных из History Server API response:
executor_metrics = [
{"executorId": "1", "totalGCTime": 5000, "totalDuration": 100000},
{"executorId": "2", "totalGCTime": 25000, "totalDuration": 100000}, # 25% GC!
]
for ex in executor_metrics:
if ex["totalDuration"] > 0:
gc_pct = ex["totalGCTime"] / ex["totalDuration"] * 100
if gc_pct > gc_threshold_pct:
problem_executors.append(
f"Executor {ex['executorId']}: GC={gc_pct:.1f}% (порог={gc_threshold_pct}%)"
)
return problem_executors
7. Практика: математический расчёт кластера под SLA¶
Условие задачи¶
Характеристики кластера: 5 серверов × (24 vCPU, 96 GB RAM) Нагрузка: PySpark ETL-пайплайн с Pandas UDF (трансформации через numpy/pandas) Требование: стабильная работа без OOM, минимальный Spill
Шаг 1: Инвентаризация ресурсов¶
class ClusterCalculator:
"""
Калькулятор конфигурации Spark Executor'ов.
Реализует стандартную методологию Executor Sizing.
"""
def __init__(
self,
nodes: int,
cpu_per_node: int,
ram_per_node_gb: int,
workload_type: str = "pyspark_pandas",
):
self.nodes = nodes
self.cpu_per_node = cpu_per_node
self.ram_per_node_gb = ram_per_node_gb
self.workload_type = workload_type
def calculate(self) -> dict:
# ── Шаг 1: Резервирование под ОС и менеджер узлов ───────────────
# ОС: 1 CPU + 1 GB RAM
# NodeManager/Kubelet: 1 CPU + 1 GB RAM
# Итого резерв: 2 CPU + 2 GB RAM
os_cpu = 1
nm_cpu = 1
os_ram_gb = 1
nm_ram_gb = 1
available_cpu = self.cpu_per_node - os_cpu - nm_cpu # 24 - 2 = 22
available_ram_gb = self.ram_per_node_gb - os_ram_gb - nm_ram_gb # 96 - 2 = 94 GB
# ── Шаг 2: Правило 5 cores ────────────────────────────────────────
executor_cores = 5
executors_per_node = available_cpu // executor_cores # 22 // 5 = 4
unused_cpu = available_cpu % executor_cores # 22 % 5 = 2 CPU (потеря)
# ── Шаг 3: Базовая RAM на Executor ───────────────────────────────
raw_ram_per_ex_gb = available_ram_gb / executors_per_node # 94 / 4 = 23.5 GB
# Округляем вниз: 23 GB для executor.memory + overhead
total_ram_per_ex_gb = int(raw_ram_per_ex_gb) # 23 GB на все 4 компонента
# ── Шаг 4: Рассчитываем Memory Overhead ──────────────────────────
# PySpark + Pandas → overhead × 2.5
overhead_multipliers = {
"standard_scala": 1.0,
"pyspark_basic": 1.5,
"pyspark_pandas": 2.5,
"heavy_shuffle": 1.5,
"ml_tensorflow": 3.0,
}
overhead_base_gb = max(0.375, total_ram_per_ex_gb * 0.10) # max(375MB, 2.3GB) = 2.3 GB
overhead_multiplier = overhead_multipliers.get(self.workload_type, 1.0)
overhead_gb = round(overhead_base_gb * overhead_multiplier, 1) # 2.3 × 2.5 = 5.75 ≈ 5.8 GB
# ── Шаг 5: pyspark.memory (Python Workers) ───────────────────────
# PySpark pandas: 2 GB × cores (консервативно)
pyspark_per_worker_gb = 2.0
pyspark_total_gb = pyspark_per_worker_gb * executor_cores # 2 × 5 = 10 GB
# НО pyspark.memory - это лимит на один Worker, не суммарный!
pyspark_config_gb = pyspark_per_worker_gb # 2 GB per worker
# ── Шаг 6: executor.memory = total - overhead - pyspark × cores ──
# Мы хотим: executor.memory + overhead + pyspark_total ≤ total_ram_per_ex_gb
executor_memory_gb = total_ram_per_ex_gb - overhead_gb - pyspark_total_gb
# = 23 - 5.8 - 10 = 7.2 GB
# Маловато! Нужно пересмотреть.
# Корректировка: уменьшаем pyspark.memory
# Для batch_size=10K строк × 200 байт = ~2 MB batch
# 2 MB × 3 (overhead) + 512 MB (libs) ≈ 518 MB → округляем до 1 GB
pyspark_per_worker_gb = 1.0
pyspark_total_gb = pyspark_per_worker_gb * executor_cores # 1 × 5 = 5 GB
executor_memory_gb = total_ram_per_ex_gb - overhead_gb - pyspark_total_gb
# = 23 - 5.8 - 5 = 12.2 → берём 12 GB
executor_memory_gb = 12
# ── Шаг 7: Проверка ───────────────────────────────────────────────
actual_total = executor_memory_gb + overhead_gb + pyspark_total_gb
# = 12 + 5.8 + 5 = 22.8 GB ≤ 23 GB ✓
# ── Шаг 8: Итоговые параметры кластера ───────────────────────────
# Один Executor оставляем под Driver на головной ноде
total_executors = self.nodes * executors_per_node
working_executors = total_executors - 1 # минус Driver
total_parallel_tasks = working_executors * executor_cores
return {
"nodes": self.nodes,
"executors_per_node": executors_per_node,
"total_executors": total_executors,
"working_executors": working_executors,
"executor_cores": executor_cores,
"executor_memory_gb": executor_memory_gb,
"memory_overhead_gb": overhead_gb,
"pyspark_memory_per_worker_gb": pyspark_per_worker_gb,
"total_container_memory_gb": actual_total,
"total_parallel_tasks": total_parallel_tasks,
"unused_cpu_per_node": unused_cpu,
"spark_submit_params": {
"executor-cores": executor_cores,
"executor-memory": f"{executor_memory_gb}g",
"num-executors": working_executors,
"conf_memory_overhead": f"{int(overhead_gb * 1024)}m",
"conf_pyspark_memory": f"{pyspark_per_worker_gb}g",
}
}
def print_report(self) -> None:
result = self.calculate()
print(f"\n{'='*60}")
print(f"Executor Sizing Report")
print(f"Кластер: {result['nodes']} нод × {self.cpu_per_node} CPU × {self.ram_per_node_gb} GB RAM")
print(f"Нагрузка: {self.workload_type}")
print(f"{'='*60}")
print(f"\nКонфигурация Executor:")
print(f" executor-cores: {result['executor_cores']}")
print(f" executor-memory: {result['executor_memory_gb']} GB")
print(f" memoryOverhead: {result['memory_overhead_gb']} GB")
print(f" pyspark.memory: {result['pyspark_memory_per_worker_gb']} GB / worker")
print(f" Total Container: {result['total_container_memory_gb']:.1f} GB")
print(f"\nРазмещение на кластере:")
print(f" Executors на ноду: {result['executors_per_node']}")
print(f" Всего Executors: {result['total_executors']}")
print(f" Рабочих Executors: {result['working_executors']} (1 для Driver)")
print(f" Параллельных Tasks: {result['total_parallel_tasks']}")
print(f"\nspark-submit параметры:")
params = result['spark_submit_params']
print(f" --executor-cores {params['executor-cores']} \\")
print(f" --executor-memory {params['executor-memory']} \\")
print(f" --num-executors {params['num-executors']} \\")
print(f" --conf spark.executor.memoryOverhead={params['conf_memory_overhead']} \\")
print(f" --conf spark.executor.pyspark.memory={params['conf_pyspark_memory']}")
# Запускаем расчёт для задания
calc = ClusterCalculator(
nodes=5,
cpu_per_node=24,
ram_per_node_gb=96,
workload_type="pyspark_pandas"
)
calc.print_report()
Шаг 2: Применение в SparkSession¶
from pyspark.sql import SparkSession
# Параметры из калькулятора
spark = SparkSession.builder \
.master("yarn") \
.appName("etl-pyspark-production") \
# ── Executor ресурсы ──────────────────────────────────────────────
.config("spark.executor.instances", "19") \
.config("spark.executor.cores", "5") \
.config("spark.executor.memory", "12g") \
# ── Memory Overhead: PySpark + Pandas требует 2.5× дефолта ────────
.config("spark.executor.memoryOverhead", "5836m") \
# ── Python Workers ────────────────────────────────────────────────
.config("spark.executor.pyspark.memory", "1g") \
# ── Driver (на головной ноде) ─────────────────────────────────────
.config("spark.driver.memory", "12g") \
.config("spark.driver.cores", "4") \
.config("spark.driver.memoryOverhead", "2g") \
# ── AQE: для динамической оптимизации ────────────────────────────
.config("spark.sql.adaptive.enabled", "true") \
.config("spark.sql.adaptive.coalescePartitions.enabled", "true") \
.config("spark.sql.adaptive.advisoryPartitionSizeInBytes", "128mb") \
# ── Arrow для Pandas UDF: оптимизация передачи данных ────────────
.config("spark.sql.execution.arrow.pyspark.enabled", "true") \
.config("spark.sql.execution.arrow.maxRecordsPerBatch", "10000") \
# ── GC: G1GC для умеренных heap (< 32 GB) ────────────────────────
.config("spark.executor.extraJavaOptions",
"-XX:+UseG1GC -XX:G1HeapRegionSize=16m "
"-XX:+PrintGCDetails -XX:+PrintGCDateStamps") \
.enableHiveSupport() \
.getOrCreate()
Шаг 3: Сравнение дефолтной и оптимизированной конфигурации¶
# Сравнительная таблица: дефолты vs оптимизированная конфигурация
configurations = {
"Дефолтная (Fat Executor)": {
"executor_cores": 24,
"executor_memory_gb": 90,
"overhead_gb": 9, # 10%
"pyspark_memory_gb": 0, # не задан!
"num_executors": 5,
"parallel_tasks": 120,
"container_gb": 99,
"risks": [
"Full GC 30-90 сек при 90GB Heap",
"Python workers неограничены → OOM",
"HDFS 24-поточный троттлинг",
]
},
"Правильная (5-core)": {
"executor_cores": 5,
"executor_memory_gb": 12,
"overhead_gb": 5.8,
"pyspark_memory_gb": 1, # per worker
"num_executors": 19,
"parallel_tasks": 95,
"container_gb": 22.8,
"risks": [
"Стабильная работа G1GC",
"Python Workers ограничены 1 GB/worker",
"HDFS оптимальный параллелизм",
]
}
}
print(f"{'Параметр':35s} {'Fat Executor':20s} {'5-core Optimal':20s}")
print("-" * 75)
print(f"{'executor.cores':35s} {'24':20s} {'5':20s}")
print(f"{'executor.memory':35s} {'90 GB':20s} {'12 GB':20s}")
print(f"{'memoryOverhead':35s} {'9 GB':20s} {'5.8 GB':20s}")
print(f"{'pyspark.memory':35s} {'0 (опасно!)':20s} {'1 GB / worker':20s}")
print(f"{'num.executors':35s} {'5':20s} {'19':20s}")
print(f"{'Параллельных Tasks':35s} {'120':20s} {'95':20s}")
print(f"{'Container Memory':35s} {'99 GB':20s} {'22.8 GB':20s}")
print(f"{'GC Stop-The-World':35s} {'30-90 сек':20s} {'< 1 сек':20s}")
print(f"{'OOM риск':35s} {'ВЫСОКИЙ':20s} {'НИЗКИЙ':20s}")
print(f"{'Fault Tolerance':35s} {'Плохая':20s} {'Хорошая':20s}")
Шаг 4: Диагностика после запуска¶
# Проверить что контейнеры запустились с нужными параметрами
yarn application -status <APPLICATION_ID>
# Смотрим логи конкретного Executor
yarn logs -applicationId <APPLICATION_ID> -containerId <CONTAINER_ID> 2>&1 | grep -E "GC|OOM|killed"
# Мониторим в реальном времени
watch -n 5 'yarn application -status <APP_ID> | grep -E "State|Progress|Container"'
# Программная проверка через SparkContext
def verify_executor_config(spark: SparkSession) -> None:
"""Проверяет что применились нужные конфигурации."""
sc = spark.sparkContext
configs_to_check = [
"spark.executor.cores",
"spark.executor.memory",
"spark.executor.memoryOverhead",
"spark.executor.pyspark.memory",
"spark.executor.instances",
]
print("Активные конфигурации Executor:")
for key in configs_to_check:
try:
value = spark.conf.get(key)
print(f" {key} = {value}")
except Exception:
print(f" {key} = НЕ ЗАДАН (дефолт)")
verify_executor_config(spark)
Итоги: золотые правила Executor Sizing¶
Правило 1: 5 cores per Executor - золотой стандарт. Отступайте от него осознанно: для коротких задач с быстрым Data Engineering можно 4-6 cores, для ML-инференса - 8-10 cores если нет тяжёлого HDFS I/O.
Правило 2: 4 компонента памяти, не 1. Всегда считайте: executor.memory + overhead + pyspark.memory × cores ≤ available_ram_per_executor. Забытый overhead или pyspark.memory - главная причина Container killed by YARN.
Правило 3: Резервируйте для ОС. Никогда не отдавайте Spark 100% ресурсов сервера. 2 CPU + 2-4 GB под ОС и NodeManager - обязательный минимум.
Правило 4: GC Time > 10% = сигнал тревоги. Смотрите в Spark UI → Executors → GC Time. Если GC > 10% от Task Time - либо уменьшите heap (больше Executor'ов с меньшей памятью), либо включите Off-Heap, либо переключитесь с G1GC на ZGC.
Правило 5: PySpark требует явного pyspark.memory. Дефолт 0 означает «неограничено». Для любого production PySpark пайплайна с Pandas UDF - устанавливайте явно.
Правило 6: Пересчитывайте при смене нагрузки. Правильный Executor для batch ETL - не то же самое что для ML-тренинга или Structured Streaming. Конфигурация должна соответствовать характеру задачи.