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.

optimization

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 из пула:

  1. Каждый Task'а, использующий Python код, получает свой Python Worker (или берёт из пула)
  2. Python Worker получает данные от JVM через сокет или Apache Arrow (нулевое копирование)
  3. 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. Конфигурация должна соответствовать характеру задачи.