Memory Fractions: spark.memory.fraction, storageFraction и tuning

Глубокий разбор On-Heap памяти Spark Executor: три региона (Reserved, User, Spark Memory), Unified Memory Manager и dynamic borrowing, Execution Memory для Shuffle/Sort/Join, Storage Memory для cache/broadcast, диагностика Spill через Spark UI, тюнинг fraction под ETL/ML/BI нагрузки.

optimization

1. Архитектурная карта On-Heap памяти Spark Executor

Предыдущий урок (Executor Sizing) научил правильно выбирать общий объём памяти Executor'а. Этот урок отвечает на следующий вопрос: что происходит с этой памятью внутри? Как Spark решает сколько отдать кешу, а сколько - вычислениям? И почему иногда кеш молча выбрасывается, а иногда запрос начинает неожиданно медленно писать данные на диск?

От простого к сложному: почему один параметр не решает проблему

Типичный сценарий: пайплайн падает с OOM или жутко тормозит. Инженер увеличивает spark.executor.memory с 8 GB до 32 GB. Проблема не исчезает. Почему?

Потому что executor.memory - это суммарный объём Heap. Spark делит этот Heap на несколько изолированных регионов с разными правилами управления, и именно соотношение между этими регионами определяет поведение в условиях нагрузки.

Три главных региона On-Heap памяти

Схема показывает полную картину On-Heap памяти для Executor с 20 GB. Три ключевых наблюдения:

  1. Reserved Memory (300 MB) абсолютно фиксирована - она не участвует ни в каких настройках
  2. User Memory определяется остатком после Spark Memory - чем больше memory.fraction, тем меньше User Memory
  3. Spark Memory делится на защищённую (Storage Floor) и динамическую зону - это и есть суть Unified Memory Manager

Математика расчёта регионов

def calculate_memory_regions(
    executor_memory_gb: float,
    memory_fraction: float = 0.6,
    storage_fraction: float = 0.5,
) -> dict:
    """
    Рассчитывает все регионы On-Heap памяти Spark Executor.

    Формулы:
    1. usable_heap = executor_memory - reserved
    2. spark_memory = usable_heap × memory_fraction
    3. user_memory  = usable_heap × (1 - memory_fraction)
    4. storage_floor = spark_memory × storage_fraction
    5. dynamic_zone  = spark_memory × (1 - storage_fraction)

    storage_floor - МИНИМУМ памяти для Storage (нельзя вытеснить Execution)
    dynamic_zone - может принадлежать Storage или Execution в любой момент
    """
    RESERVED_MB = 300
    executor_mb = executor_memory_gb * 1024

    usable_mb = executor_mb - RESERVED_MB
    spark_memory_mb = usable_mb * memory_fraction
    user_memory_mb = usable_mb * (1 - memory_fraction)
    storage_floor_mb = spark_memory_mb * storage_fraction
    dynamic_zone_mb = spark_memory_mb * (1 - storage_fraction)

    return {
        "executor_memory_mb": executor_mb,
        "reserved_mb": RESERVED_MB,
        "usable_mb": usable_mb,
        "spark_memory_mb": spark_memory_mb,
        "user_memory_mb": user_memory_mb,
        "storage_floor_mb": storage_floor_mb,
        "dynamic_zone_mb": dynamic_zone_mb,
        "max_execution_mb": spark_memory_mb,  # может занять всю Spark Memory
        "max_storage_mb": spark_memory_mb,    # тоже может занять всю Spark Memory
        "guaranteed_storage_mb": storage_floor_mb,  # минимально гарантировано
    }

# Дефолтные параметры для 20 GB Executor:
regions = calculate_memory_regions(20, 0.6, 0.5)
print(f"Executor Memory: {regions['executor_memory_mb']/1024:.1f} GB")
print(f"  Reserved:       {regions['reserved_mb']} MB (фикс.)")
print(f"  Usable Heap:    {regions['usable_mb']/1024:.2f} GB")
print(f"  Spark Memory:   {regions['spark_memory_mb']/1024:.2f} GB")
print(f"  User Memory:    {regions['user_memory_mb']/1024:.2f} GB")
print(f"  Storage Floor:  {regions['storage_floor_mb']/1024:.2f} GB (защищено)")
print(f"  Dynamic Zone:   {regions['dynamic_zone_mb']/1024:.2f} GB (заимствуется)")
Executor Memory: 20.0 GB
  Reserved:       300 MB (фикс.)
  Usable Heap:    19.71 GB
  Spark Memory:   11.82 GB
  User Memory:     7.88 GB
  Storage Floor:   5.91 GB (защищено)
  Dynamic Zone:    5.91 GB (заимствуется)

2. Опора первая: Execution Memory - память для вычислений

Execution Memory - это буферы, которые Spark использует для выполнения «тяжёлых» операций реляционной алгебры. Это самый критичный для производительности регион памяти.

Что использует Execution Memory

Hash Join (BroadcastHashJoin и SortMergeJoin partial): При построении hash-таблицы для JOIN, Spark хранит одну сторону JOIN в памяти как хэш-структуру. Для каждого ключа нужно ~50-200 байт в зависимости от типов данных. При 100 миллионах уникальных ключей это 5-20 GB только для хэш-таблицы.

Sort (OrderBy, SortMergeJoin): Буферы сортировки - Spark собирает в памяти данные для сортировки, затем пишет sorted runs. Размер In-Memory Sort buffer ограничен Execution Memory.

Hash Aggregation (groupBy с aggregation): При groupBy("key").count() Spark строит хэш-аггрегат в памяти: для каждого уникального ключа - аккумулятор. При высокой кардинальности ключа (groupBy("user_id") с миллионом уникальных пользователей) хэш-аггрегат может занять несколько гигабайт.

Window Functions: Оконные функции (RANK, LAG, LEAD) требуют держать в памяти весь PARTITION BY-буфер для вычисления функции. Для партиции с миллионом строк это несколько сотен MB.

Shuffle Write буферы: При записи Shuffle-данных каждый Task сначала буферизует данные в памяти (spark.shuffle.sort.bypassMergeThreshold управляет этим). Каждый открытый Shuffle Writer занимает spark.shuffle.file.buffer (дефолт 32 KB) × число Reduce-партиций.

Ключевое свойство: иммунитет к вытеснению

Execution Memory обладает принципиальным свойством: она не может быть вытеснена ради Storage Memory. Если Task выполняет JOIN и занял Execution Memory - Spark не может забрать эту память, чтобы отдать кешу.

Логика проста: вытеснение Execution Memory означало бы уничтожение промежуточного результата вычисления. Task пришлось бы пересчитывать всё с нуля - это хуже, чем Spill.

Что происходит при нехватке Execution Memory: Spill to Disk

Если Task требует больше Execution Memory чем доступно - Spill (сброс данных на диск):

Диаграмма показывает цену Spill: данные пишутся на диск и читаются обратно. Один Spill = минимум 2× лишних I/O. При нескольких Spill Run'ах overhead растёт линейно. Диск в 100-1000x медленнее RAM → Spill превращает секунды в минуты.


3. Опора вторая: Storage Memory - память для кеша и Broadcast

Storage Memory хранит данные, которые явно помещены в кеш командами .cache() или .persist(), а также данные Broadcast-переменных, разосланных всем Executor'ам.

Что хранится в Storage Memory

Кешированные DataFrame/RDD: Когда вы вызываете df.cache(), Spark материализует данные в Storage Memory в момент первого Action (count, show, write). Данные хранятся как блоки (Blocks) с указанными StorageLevel.

Broadcast переменные: При F.broadcast(lookup_table) или spark.sparkContext.broadcast(data), данные сначала сериализуются на Driver, затем рассылаются на все Executor'ы и хранятся в Storage Memory. Это позволяет избежать повторной передачи при каждом обращении.

Unroll буферы: Когда Spark десериализует блок из кеша (MEMORY_SER → MEMORY_ONLY), он временно использует Unroll буфер в Storage Memory для развёртывания данных.

StorageLevel: стратегии хранения и поведение при переполнении

from pyspark import StorageLevel

# MEMORY_ONLY: хранить в памяти, при нехватке - выбросить блок
# При следующем обращении пересчитать из источника
df.persist(StorageLevel.MEMORY_ONLY)

# MEMORY_AND_DISK: хранить в памяти, при нехватке - сбросить на диск
# Блок на диске медленнее, но не пересчитывается
df.persist(StorageLevel.MEMORY_AND_DISK)

# MEMORY_ONLY_SER: хранить в памяти в сериализованном виде (компактнее)
# CPU overhead при чтении (десериализация)
df.persist(StorageLevel.MEMORY_ONLY_SER)

# MEMORY_AND_DISK_SER: сериализованный + диск при переполнении
df.persist(StorageLevel.MEMORY_AND_DISK_SER)

# DISK_ONLY: только на диске (медленно, но не занимает RAM)
df.persist(StorageLevel.DISK_ONLY)

# Рекомендации по выбору:
# MEMORY_ONLY: данные легко пересчитываются (ранние Stage'ы, нет join'ов)
# MEMORY_AND_DISK: дорогие вычисления (много join'ов, сложные трансформации)
# MEMORY_ONLY_SER: много строк простых типов (числа, строки) → меньше RAM
# DISK_ONLY: очень большие данные, редко читаются

Что происходит при переполнении Storage Memory

Если Storage Memory заполнена и нужно кешировать новый блок:

  • MEMORY_ONLY: старый блок выбрасывается (не пишется на диск). При следующем обращении Spark пересчитает блок из источника. Если пересчёт дорогой (много трансформаций) - это очень медленно.
  • MEMORY_AND_DISK: старый блок записывается на диск, а в RAM занимает его место новый. При обращении к старому блоку - читается с диска. Медленнее RAM, но нет пересчёта.

Сценарий тихой потери кеша (Cache Eviction без ошибки):

df1 = spark.read.parquet("hdfs://cluster/large_table_10GB/")
df1.persist(StorageLevel.MEMORY_ONLY)
df1.count()  # кешируем df1 в 10 GB Storage Memory

df2 = spark.read.parquet("hdfs://cluster/another_table_8GB/")
df2.persist(StorageLevel.MEMORY_ONLY)
df2.count()  # кешируем df2

# Если Storage Memory = 5.91 GB:
# df1 (10 GB) > Storage Memory → кешируется только часть!
# При кешировании df2 - часть блоков df1 ВЫБРАСЫВАЕТСЯ

# Никаких ошибок не выдаётся!
# При следующем обращении к df1 - часть пересчитывается из источника
result = df1.join(df2, "key")  # df1 частично пересчитывается → медленно

Для проверки что кеш действительно сохранился:

# Spark UI → Storage tab: показывает Fraction Cached
# Если Fraction Cached < 100% → часть блоков была вытеснена

# Программно:
catalog = spark.catalog
# В Spark 3.x: spark.catalog.isCached(table_name)
# Для DataFrame: нет прямого API, смотрим через UI

4. Управление пропорциями: Unified Memory Manager

До Spark 1.6 существовала статическая модель памяти: Execution и Storage имели жёсткие границы, заданные при запуске. Если Storage была заполнена, но Execution пустовала - нельзя было использовать свободную Storage память для вычислений. Это приводило к Spill даже при наличии формально «свободной» памяти.

Unified Memory Manager (Spark 1.6+) решил эту проблему через динамическое заимствование памяти.

Концепция Dynamic Borrowing

Критически важные правила Unified Memory Manager:

  1. Storage НИКОГДА не может вытеснить Execution. Даже если Storage Memory пустует, а Execution нужна память - Storage не «отнимает» у выполняющихся Task'ов.

  2. Execution МОЖЕТ вытеснить Storage из Dynamic Zone. Если Task'е не хватает памяти, она вытеснит кешированные блоки из Dynamic Zone.

  3. Storage Floor (storage.fraction × spark_memory) абсолютно защищена. Эту часть Execution не трогает никогда.

  4. Максимально возможный размер обоих регионов = весь Spark Memory. При пустом противнике можно занять всё.

Параметр spark.memory.fraction: граница Spark Memory

# spark.memory.fraction = доля от (executor.memory - reserved) под Spark Memory
# Дефолт: 0.6

# Что происходит с остальными 40%?
# → User Memory: ваши UDF объекты, RDD метаданные, кастомные структуры

# Когда увеличивать memory.fraction (например до 0.8):
# + Больше памяти для Execution и Storage
# - Меньше User Memory: риск OOM если UDF создаёт много объектов

# Когда уменьшать memory.fraction (например до 0.5):
# + Больше User Memory для тяжёлых UDF
# - Меньше памяти для Execution → больше Spill

spark = SparkSession.builder \
    .config("spark.memory.fraction", "0.75") \  # было 0.6, хотим больше для Execution
    .getOrCreate()

Параметр spark.memory.storageFraction: защита кеша

# spark.memory.storageFraction = доля Spark Memory под Storage Floor (защищённый кеш)
# Дефолт: 0.5

# Пример: spark_memory = 11.82 GB, storageFraction = 0.5
# Storage Floor = 5.91 GB (защищено от Execution)
# Dynamic Zone = 5.91 GB (заимствуется)

# Когда ПОДНЯТЬ storageFraction (например до 0.7):
# Использование: много .cache(), данные важно удержать в RAM
# Результат: 8.27 GB защищённого кеша, 3.55 GB Dynamic Zone
# Риск: Execution получает меньше "буфера" из Dynamic Zone → больше риск Spill

# Когда СНИЗИТЬ storageFraction (например до 0.2):
# Использование: нет .cache(), только тяжёлые Shuffle/Sort/Join
# Результат: 2.36 GB гарантированного кеша, 9.46 GB Dynamic Zone
# Результат: Execution имеет практически всю Spark Memory в своём распоряжении

spark = SparkSession.builder \
    .config("spark.memory.storageFraction", "0.3") \  # почти нет защищённого кеша
    .getOrCreate()

5. Стратегии тюнинга под разные типы нагрузок

Дефолтные настройки (0.6/0.5) - разумный компромисс для смешанных нагрузок. Но для специфических workload'ов тюнинг может дать значительное ускорение.

Полная математика тюнинга

def tune_memory_fractions(
    executor_memory_gb: float,
    workload_type: str,
    expected_cache_gb: float = 0.0,
    shuffle_data_per_task_gb: float = 0.0,
    n_tasks: int = 5,
) -> dict:
    """
    Рекомендует настройки memory.fraction и storageFraction
    для конкретного типа нагрузки.

    workload_type:
    - "shuffle_heavy": ETL с groupBy, JOIN, Window без кеша
    - "cache_heavy": ML, BI с многократным чтением одного датасета
    - "balanced": смешанная нагрузка (дефолт)
    - "broadcast_heavy": много JOIN с broadcast переменными
    """
    RESERVED_MB = 300
    executor_mb = executor_memory_gb * 1024
    usable_mb = executor_mb - RESERVED_MB

    # Базовые рекомендации по типу нагрузки
    configs = {
        "shuffle_heavy": {
            "memory_fraction": 0.80,
            "storage_fraction": 0.10,
            "description": "Максимум Execution Memory для Shuffle/Sort/Join. "
                          "Минимальный защищённый кеш (почти всё для Execution).",
        },
        "cache_heavy": {
            "memory_fraction": 0.70,
            "storage_fraction": 0.70,
            "description": "Большой защищённый Storage Floor для надёжного кеша. "
                          "ML training loops, BI с многократным чтением данных.",
        },
        "balanced": {
            "memory_fraction": 0.60,
            "storage_fraction": 0.50,
            "description": "Дефолт. Хороший баланс для смешанных нагрузок.",
        },
        "broadcast_heavy": {
            "memory_fraction": 0.65,
            "storage_fraction": 0.60,
            "description": "Умеренно повышенный Storage для Broadcast переменных.",
        },
    }

    config = configs.get(workload_type, configs["balanced"])
    mf = config["memory_fraction"]
    sf = config["storage_fraction"]

    spark_mem_mb = usable_mb * mf
    user_mem_mb = usable_mb * (1 - mf)
    storage_floor_mb = spark_mem_mb * sf
    dynamic_zone_mb = spark_mem_mb * (1 - sf)

    # Проверяем достаточно ли Execution Memory для данных задач
    max_execution_per_task_mb = spark_mem_mb / n_tasks
    needed_per_task_mb = shuffle_data_per_task_gb * 1024

    spill_risk = "ВЫСОКИЙ" if needed_per_task_mb > max_execution_per_task_mb \
                else "НИЗКИЙ" if needed_per_task_mb < max_execution_per_task_mb * 0.7 \
                else "УМЕРЕННЫЙ"

    cache_fits = storage_floor_mb >= expected_cache_gb * 1024
    cache_status = "ПОМЕСТИТСЯ ПОЛНОСТЬЮ" if cache_fits else "ВЫТЕСНЕНИЕ ВОЗМОЖНО"

    return {
        "workload_type": workload_type,
        "memory_fraction": mf,
        "storage_fraction": sf,
        "spark_memory_gb": spark_mem_mb / 1024,
        "user_memory_gb": user_mem_mb / 1024,
        "storage_floor_gb": storage_floor_mb / 1024,
        "dynamic_zone_gb": dynamic_zone_mb / 1024,
        "max_execution_per_task_gb": max_execution_per_task_mb / 1024,
        "spill_risk": spill_risk,
        "cache_status": cache_status,
        "description": config["description"],
    }


# Примеры
print("=== Тяжёлый ETL с GroupBy/JOIN (без кеша) ===")
result = tune_memory_fractions(
    executor_memory_gb=20,
    workload_type="shuffle_heavy",
    shuffle_data_per_task_gb=1.5,  # 1.5 GB данных на Task при сортировке
    n_tasks=5,
)
for k, v in result.items():
    print(f"  {k}: {v}")

print()
print("=== ML Training (многократное чтение одного датасета) ===")
result2 = tune_memory_fractions(
    executor_memory_gb=20,
    workload_type="cache_heavy",
    expected_cache_gb=8,  # датасет 8 GB нужно держать в RAM
    n_tasks=5,
)
for k, v in result2.items():
    print(f"  {k}: {v}")

Сценарий 1: ETL Shuffle-heavy (groupBy, JOIN, Window без кеша)

Характеристика нагрузки: пайплайн выполняет сложные аналитические запросы с несколькими groupBy, JOIN двух больших таблиц и WINDOW FUNCTION. Никакого кеширования - данные читаются один раз.

Проблема с дефолтом: storageFraction = 0.5 означает что половина Spark Memory «защищена» для кеша, который не используется. Это буквально 5+ GB, которые не могут стать буферами для Shuffle - они заблокированы Storage Floor.

Оптимальная настройка:

spark = SparkSession.builder \

    # Увеличиваем общую долю под Spark Memory
    # Больше памяти для вычислений в ущерб User Memory
    .config("spark.memory.fraction", "0.80") \

    # Снижаем защищённый кеш до минимума (кеш не используется)
    # Storage Floor = 10% × Spark Memory → почти всё для Execution
    .config("spark.memory.storageFraction", "0.10") \

    .getOrCreate()

# Результат для 20 GB Executor:
# Spark Memory = 0.80 × 19.7 GB = 15.76 GB (было 11.82 GB)
# Storage Floor = 0.10 × 15.76 = 1.58 GB (было 5.91 GB)
# Dynamic Zone = 0.90 × 15.76 = 14.18 GB (было 5.91 GB)
# Max Execution (при пустом Storage) = 15.76 GB (было 11.82 GB)
# Прирост Execution: +3.94 GB → значительно меньше Spill на heavy JOIN

Сценарий 2: Iterative ML / BI (Cache-heavy)

Характеристика нагрузки: обучение ML-модели читает один и тот же датасет в каждой итерации (эпохе обучения). Или BI-дашборд многократно запрашивает один и тот же агрегированный датасет.

Проблема с дефолтом: если Execution Memory занята JOIN'ами и вытесняет кешированные данные из Dynamic Zone - каждая следующая итерация пересчитывает данные с нуля.

Оптимальная настройка:

spark = SparkSession.builder \
    # Умеренно увеличиваем Spark Memory
    .config("spark.memory.fraction", "0.70") \

    # Поднимаем Storage Floor: 70% Spark Memory защищено от Execution
    # Кеш гарантированно удержится в RAM при тяжёлых вычислениях
    .config("spark.memory.storageFraction", "0.70") \

    .getOrCreate()

# Результат для 20 GB Executor:
# Spark Memory = 0.70 × 19.7 = 13.79 GB
# Storage Floor = 0.70 × 13.79 = 9.65 GB (ЗАЩИЩЕНО от Execution)
# Dynamic Zone = 0.30 × 13.79 = 4.14 GB
# Датасет 8 GB → полностью умещается в Storage Floor → НЕ вытесняется
# Каждая итерация читает из RAM → в 100x быстрее чем без кеша

Сценарий 3: Mixed Lakehouse (Delta Lake / Iceberg)

Характеристика: Spark читает Delta Lake таблицы с Predicate Pushdown, выполняет JOIN с broadcast lookup'ами, затем пишет результат обратно.

spark = SparkSession.builder \
    # Близко к дефолту, чуть больше Execution для MERGE операций Delta
    .config("spark.memory.fraction", "0.65") \

    # Умеренно большой Storage Floor для Broadcast переменных (lookup таблицы)
    .config("spark.memory.storageFraction", "0.45") \

    # AQE для динамической адаптации
    .config("spark.sql.adaptive.enabled", "true") \

    # Off-Heap для векторизованного Parquet чтения (Tungsten)
    .config("spark.memory.offHeap.enabled", "true") \
    .config("spark.memory.offHeap.size", "2g") \

    .getOrCreate()

6. Диагностика Memory Spill и GC через Spark UI

Теоретические знания о регионах памяти бесполезны без умения диагностировать проблемы в реальных запросах. Spark UI предоставляет все необходимые инструменты.

Вкладка Stages: обнаружение Spill

В Spark UI → Stages → конкретный Stage видны агрегированные метрики по всем Task'ам:

Stage Details:
  Shuffle Read Size / Records: 2.4 GB / 1,234,567
  Shuffle Write Size / Records: 1.8 GB / 1,234,567

  Spill (Memory): 4.2 GB   ← данные сериализованы и ожидают записи
  Spill (Disk):   4.1 GB   ← данные записаны на диск

Spill (Memory) vs Spill (Disk):

  • Spill (Memory) - объём данных в несериализованном виде в памяти перед записью на диск. Обычно больше чем Spill (Disk) потому что сериализация сжимает данные.
  • Spill (Disk) - объём данных после сериализации и записи на диск. Это реальная нагрузка на дисковую подсистему.

Соотношение Spill(Memory) / Spill(Disk) показывает эффективность сжатия. Обычно 1.2-3x. Если соотношение близко к 1 - данные плохо сжимаются (случайные числа, уже сжатые данные).

Вкладка Stages → Tasks: детальный анализ

Для конкретного Stage кликаем «Show additional metrics» и смотрим:

  • Peak Execution Memory: максимальный объём Execution Memory для одной Task. Это ключевая метрика для расчёта нужного executor.memory.
  • Shuffle Read Fetch Wait Time: сколько Task ждала данных от других Executor'ов. Если > 30% Duration Stage - сетевой bottleneck.
  • GC Time: время сборки мусора. Если > 10% от Task Duration - проблема с Heap управлением.
def analyze_stage_for_spill(spark, stage_id: int, history_server_url: str) -> dict:
    """
    Анализирует конкретный Stage на наличие Spill через History Server API.

    Возвращает рекомендации по настройке памяти.
    """
    import requests

    app_id = spark.sparkContext.applicationId
    url = f"{history_server_url}/api/v1/applications/{app_id}/stages/{stage_id}/0/taskList"

    tasks = requests.get(url, timeout=30).json()

    # Собираем метрики по всем Task'ам
    total_spill_memory = sum(
        t.get("taskMetrics", {}).get("memoryBytesSpilled", 0) for t in tasks
    )
    total_spill_disk = sum(
        t.get("taskMetrics", {}).get("diskBytesSpilled", 0) for t in tasks
    )
    peak_execution_memory = max(
        t.get("taskMetrics", {}).get("peakExecutionMemory", 0) for t in tasks
    )
    total_gc_time = sum(
        t.get("taskMetrics", {}).get("jvmGcTime", 0) for t in tasks
    )
    total_task_time = sum(t.get("duration", 0) for t in tasks)
    gc_pct = (total_gc_time / max(total_task_time, 1)) * 100

    return {
        "stage_id": stage_id,
        "total_tasks": len(tasks),
        "spill_memory_gb": total_spill_memory / 1024**3,
        "spill_disk_gb": total_spill_disk / 1024**3,
        "peak_execution_memory_gb": peak_execution_memory / 1024**3,
        "gc_time_pct": gc_pct,
        "diagnosis": {
            "has_spill": total_spill_disk > 0,
            "gc_pressure": gc_pct > 10,
            "recommendations": _generate_memory_recommendations(
                total_spill_disk, peak_execution_memory, gc_pct
            ),
        }
    }


def _generate_memory_recommendations(spill_disk: int, peak_exec: int, gc_pct: float) -> list:
    recommendations = []

    if spill_disk > 0:
        spill_gb = spill_disk / 1024**3
        recommendations.append(
            f"SPILL: {spill_gb:.1f} GB данных записано на диск. "
            f"Рассмотрите: увеличение executor.memory, снижение storageFraction, "
            f"или увеличение числа партиций (меньше данных на Task)."
        )

    if gc_pct > 10:
        recommendations.append(
            f"GC PRESSURE: {gc_pct:.0f}% времени Task тратит на GC. "
            f"Рассмотрите: уменьшение executor.memory (меньше Heap → быстрее GC), "
            f"включение Off-Heap (spark.memory.offHeap.enabled=true), "
            f"или переключение с G1GC на ZGC."
        )

    if peak_exec > 0:
        peak_gb = peak_exec / 1024**3
        recommendations.append(
            f"PEAK EXECUTION: {peak_gb:.1f} GB пиковая потребность. "
            f"Рекомендуемый executor.memory для нулевого Spill: "
            f"ceil({peak_gb:.1f} / memory_fraction / (1-storage_fraction)) GB."
        )

    return recommendations

Вкладка Executors: аудит использования памяти

Executors Tab:
  RDD Blocks:     234                    ← кешированных блоков в RAM
  Storage Memory: 4.2 GB / 5.9 GB       ← используется / доступно
  Peak JVM Memory: 18.4 GB              ← пиковое потребление JVM
  GC Time:        2.1 min (3.2%)        ← доля GC от общего времени

Storage Memory 4.2 / 5.9 GB:
  - 4.2 GB занято кешем
  - 1.7 GB свободно в Storage Floor
  - Dynamic Zone (5.9 GB) частично занята Execution

Если Storage Memory used близко к Storage Floor - кеш «заперт» в защищённой зоне. Это хорошо (кеш стабилен). Если Storage Memory used равно Spark Memory total - Storage занимает и Dynamic Zone, возможно вытесняя Execution при следующем тяжёлом запросе.


7. Практика: ликвидация дискового Spill

Полный пример диагностики и оптимизации

from pyspark.sql import SparkSession, functions as F
import time

# ── ШАГ 1: Воспроизводим проблему с дефолтными настройками ──────────────

spark_default = SparkSession.builder \
    .master("local[8]") \
    .appName("memory-fractions-before") \
    # Дефолтные memory fractions
    .config("spark.memory.fraction", "0.6") \
    .config("spark.memory.storageFraction", "0.5") \
    # Небольшая память чтобы увидеть Spill
    .config("spark.executor.memory", "4g") \
    .config("spark.sql.shuffle.partitions", "50") \
    .getOrCreate()

spark_default.sparkContext.setLogLevel("WARN")


def create_heavy_etl_data(spark: SparkSession):
    """Создаёт датасет имитирующий тяжёлый аналитический пайплайн."""
    n = 5_000_000

    fact = spark.range(n).select(
        F.col("id").alias("order_id"),
        (F.rand() * 50_000).cast("long").alias("customer_id"),
        (F.rand() * 5_000).cast("long").alias("product_id"),
        (F.rand() * 10_000).alias("revenue"),
        F.date_add(F.lit("2024-01-01"), (F.rand() * 365).cast("int")).alias("order_date"),
    )

    dim = spark.range(50_000).select(
        F.col("id").alias("customer_id"),
        F.concat(F.lit("Cust_"), F.col("id")).alias("customer_name"),
        F.array(F.lit("EU"), F.lit("US"), F.lit("APAC")).getItem(
            (F.col("id") % 3).cast("int")
        ).alias("region")
    )

    return fact, dim


def run_analytics_pipeline(spark: SparkSession, fact, dim) -> tuple[float, dict]:
    """Запускает тяжёлый аналитический запрос с JOIN и Window."""
    t0 = time.time()

    # JOIN + Window + GroupBy - всё что нагружает Execution Memory
    result = fact.join(dim, "customer_id") \
        .withColumn(
            "revenue_rank",
            F.rank().over(
                F.Window.partitionBy("region")
                .orderBy(F.desc("revenue"))
            )
        ) \
        .filter(F.col("revenue_rank") <= 1000) \
        .groupBy("region", F.month("order_date").alias("month")) \
        .agg(
            F.sum("revenue").alias("total_revenue"),
            F.count("*").alias("order_count"),
            F.countDistinct("customer_id").alias("unique_customers"),
        ) \
        .orderBy("region", "month")

    row_count = result.count()
    elapsed = time.time() - t0

    return elapsed, {"rows": row_count}


print("=== ДО оптимизации (дефолтные fractions) ===")
fact, dim = create_heavy_etl_data(spark_default)
fact.cache()
dim.cache()
fact.count()
dim.count()

t_before, stats_before = run_analytics_pipeline(spark_default, fact, dim)
print(f"Время выполнения: {t_before:.1f}с")
print("Откройте Spark UI (http://localhost:4040) → Stages")
print("Найдите Spill (Memory) и Spill (Disk) метрики")
print()

spark_default.stop()

# ── ШАГ 2: Диагностируем проблему ───────────────────────────────────────
# Из Spark UI видим (пример):
# Stage 3 (Window): Spill (Memory) = 2.1 GB, Spill (Disk) = 1.4 GB
# Stage 5 (Sort):   Spill (Memory) = 1.8 GB, Spill (Disk) = 1.2 GB
# Peak Execution Memory: ~1.5 GB per Task

# Вычисляем нужную Execution Memory:
# Peak per Task = 1.5 GB
# 8 Tasks параллельно = 12 GB Execution Memory нужно
# executor.memory = 4 GB → spark_memory = 4 × 0.6 × (1 - 0.5) = 1.2 GB доступно Execution
# 1.2 GB << 12 GB → огромный Spill!

# ── ШАГ 3: Пересчитываем fractions ──────────────────────────────────────
# Цель: дать Execution Memory >= Peak per Task × N concurrent tasks
# executor.memory = 4 GB
# usable = 4 GB - 0.3 GB = 3.7 GB
# Нужно для Execution: ~1.5 GB per Task (но при 8 Tasks одновременно
# и 4 GB executor: 4/8 cores = 0.5 GB per Task - мало)
# Решение: поднять memory.fraction до 0.85, storageFraction до 0.05

# ── ШАГ 4: Применяем оптимизацию ────────────────────────────────────────
spark_optimized = SparkSession.builder \
    .master("local[8]") \
    .appName("memory-fractions-after") \

    # Оптимизированные fractions для Shuffle-heavy нагрузки без кеша
    .config("spark.memory.fraction", "0.85") \      # Максимум для Spark Memory
    .config("spark.memory.storageFraction", "0.05") \  # Минимум для защищённого кеша

    .config("spark.executor.memory", "4g") \
    .config("spark.sql.shuffle.partitions", "50") \
    .getOrCreate()

spark_optimized.sparkContext.setLogLevel("WARN")

print("=== ПОСЛЕ оптимизации (memory.fraction=0.85, storageFraction=0.05) ===")

# Расчёт нового распределения:
optimized_regions = calculate_memory_regions(4, 0.85, 0.05)
print(f"Spark Memory: {optimized_regions['spark_memory_mb']/1024:.2f} GB "
      f"(было {calculate_memory_regions(4, 0.6, 0.5)['spark_memory_mb']/1024:.2f} GB)")
print(f"Storage Floor: {optimized_regions['storage_floor_mb']/1024:.2f} GB "
      f"(было {calculate_memory_regions(4, 0.6, 0.5)['storage_floor_mb']/1024:.2f} GB)")
print(f"Dynamic Zone: {optimized_regions['dynamic_zone_mb']/1024:.2f} GB "
      f"(было {calculate_memory_regions(4, 0.6, 0.5)['dynamic_zone_mb']/1024:.2f} GB)")

fact2, dim2 = create_heavy_etl_data(spark_optimized)
t_after, stats_after = run_analytics_pipeline(spark_optimized, fact2, dim2)
print(f"Время выполнения: {t_after:.1f}с")
print(f"Ускорение: {t_before/t_after:.1f}x")
print()
print("Проверьте Spark UI:")
print("  Spill (Memory) и Spill (Disk) должны быть 0 или значительно снижены!")

spark_optimized.stop()

Калькулятор оптимальных Memory Fractions

def optimal_memory_fractions(
    executor_memory_gb: float,
    peak_execution_per_task_gb: float,
    n_concurrent_tasks: int,
    cache_size_gb: float = 0.0,
    user_memory_min_gb: float = 1.0,
) -> dict:
    """
    Вычисляет оптимальные memory.fraction и storageFraction
    на основе реальных метрик из Spark UI.

    Входные данные берутся из Spark UI:
    - peak_execution_per_task_gb: из Stages → Peak Execution Memory
    - n_concurrent_tasks: из Executor cores
    - cache_size_gb: из Storage tab
    """
    RESERVED_MB = 300
    usable_gb = executor_memory_gb - RESERVED_MB / 1024

    # Минимум Execution Memory для нулевого Spill
    min_execution_gb = peak_execution_per_task_gb * n_concurrent_tasks

    # Минимум Storage Floor для надёжного кеша
    min_storage_gb = cache_size_gb

    # Минимум User Memory
    min_user_gb = user_memory_min_gb

    # Нужная Spark Memory
    min_spark_gb = min_execution_gb + min_storage_gb
    if min_spark_gb + min_user_gb > usable_gb:
        # Не хватает памяти - нужно увеличить executor.memory
        return {
            "error": f"Недостаточно памяти! Нужно минимум "
                     f"{min_spark_gb + min_user_gb:.1f} GB + 300 MB reserved, "
                     f"но executor.memory = {executor_memory_gb} GB. "
                     f"Увеличьте до {min_spark_gb + min_user_gb + RESERVED_MB/1024:.0f} GB."
        }

    # Вычисляем оптимальные fractions
    # memory.fraction: обеспечивает min_spark_gb
    opt_memory_fraction = min(0.9, max(0.5, min_spark_gb / usable_gb))
    # Округляем до одного знака
    opt_memory_fraction = round(opt_memory_fraction, 1)

    actual_spark_gb = usable_gb * opt_memory_fraction

    # storage.fraction: обеспечивает min_storage_gb внутри Spark Memory
    if actual_spark_gb > 0 and min_storage_gb > 0:
        opt_storage_fraction = min(0.8, max(0.05, min_storage_gb / actual_spark_gb))
        opt_storage_fraction = round(opt_storage_fraction, 1)
    else:
        opt_storage_fraction = 0.1 if cache_size_gb == 0 else 0.5

    return {
        "recommended_memory_fraction": opt_memory_fraction,
        "recommended_storage_fraction": opt_storage_fraction,
        "resulting_spark_memory_gb": actual_spark_gb,
        "resulting_storage_floor_gb": actual_spark_gb * opt_storage_fraction,
        "resulting_user_memory_gb": usable_gb * (1 - opt_memory_fraction),
        "expected_spill": peak_execution_per_task_gb * n_concurrent_tasks > actual_spark_gb,
        "config": f"spark.memory.fraction={opt_memory_fraction}, "
                  f"spark.memory.storageFraction={opt_storage_fraction}",
    }


# Пример: используем реальные метрики из Spark UI
result = optimal_memory_fractions(
    executor_memory_gb=20,
    peak_execution_per_task_gb=1.8,   # из Stages → Peak Execution Memory
    n_concurrent_tasks=5,              # executor.cores
    cache_size_gb=6.0,                 # планируем кешировать 6 GB датасет
)
print("Рекомендации:")
for k, v in result.items():
    print(f"  {k}: {v}")

Итоги: когда трогать Memory Fractions, а когда нет

Дефолтные значения (0.6 / 0.5) правильны для большинства случаев. Не меняйте их без диагностики в Spark UI. Изменение без понимания проблемы может ухудшить производительность.

Алгоритм диагностики перед тюнингом:

  1. Открыть Spark UI → Stages - есть ли Spill (Disk) > 0?
  2. Если Spill есть → смотреть Peak Execution Memory на Task → рассчитать нужную Execution Memory
  3. Открыть Spark UI → Storage - насколько заполнен Storage, есть ли вытеснение кеша?
  4. Открыть Spark UI → Executors → GC Time > 10%? → проблема с Heap, а не с fractions
  5. Только после диагностики → применить tune_memory_fractions() калькулятор

Правила тюнинга в одной строке:

  • Много Spill, нет кеша → увеличить memory.fraction, снизить storageFraction
  • Кеш вытесняется → увеличить storageFraction (или executor.memory)
  • GC > 10% → уменьшить executor.memory (меньше Heap → быстрее GC) или Off-Heap
  • OOM → проверить User Memory (может UDF создаёт огромные объекты?)

Ключевое понимание: Unified Memory Manager - это система автоматического заимствования, а не статический разделитель. В большинстве реальных нагрузок она сама правильно балансирует между Execution и Storage. Параметры memory.fraction и storageFraction - это границы, внутри которых работает автоматика, а не жёсткие лимиты на каждую операцию.