Unified Memory Model: execution vs storage регионы и граница

Как Spark делит память executor между вычислениями и кэшем, механика динамического заимствования, spill to disk и практический tuning

core internals optimization

Эволюция управления памятью в Spark

До Spark 1.6 существовал Static Memory Manager: память executor жёстко делилась между тремя статическими регионами на этапе старта, и перераспределение между ними было невозможно во время работы.

Static Memory Manager (до Spark 1.6):
┌─────────────────────────────────────────────────────┐
│ Execution Memory  20%  │ Storage Memory  60%  │ 20% │
│ (shuffle, sort)        │ (cache, persist)     │     │
└─────────────────────────────────────────────────────┘
  Жёсткая граница: если Execution переполнено → OOM,
  даже если Storage пустой

Главная проблема: разные job'ы используют память принципиально по-разному. ETL-пайплайн с тяжёлыми shuffle почти не кэширует данные - ему нужно 80% Execution. Iterative ML-алгоритм (KMeans, ALS) активно переиспользует кэшированные данные - ему нужно 70% Storage. Со статическим менеджером нельзя настроить одну конфигурацию для обоих сценариев: либо OOM в Execution, либо бесполезно большой пустой Storage.

Unified Memory Manager (Spark 1.6+) решил это: Execution и Storage регионы - это части единого пула с мягкой (soft) границей, которая смещается динамически в зависимости от потребностей.

Карта памяти Executor

Общая схема памяти, доступной процессу executor:

Расчёт на конкретном примере

Возьмём типичный executor: spark.executor.memory = 8g.

heap = 8 ГБ = 8192 МБ

Reserved Memory = 300 МБ (фиксировано)
usable = 8192 - 300 = 7892 МБ

spark.memory.fraction = 0.6 (по умолчанию)
Spark Memory (Unified Pool) = 7892 × 0.6 = 4735 МБ ≈ 4.6 ГБ

User Memory = 7892 × 0.4 = 3157 МБ ≈ 3.1 ГБ

spark.memory.storageFraction = 0.5 (по умолчанию)
Storage soft limit = 4735 × 0.5 = 2368 МБ ≈ 2.3 ГБ
Execution soft limit = 4735 × 0.5 = 2368 МБ ≈ 2.3 ГБ
Регион Размер Назначение
Reserved Memory 300 МБ Внутренние нужды Spark (метаданные, системные структуры)
User Memory ~3.1 ГБ Пользовательские объекты: UDF-замыкания, Encoders, аккумуляторы, структуры данных в коде
Spark Memory ~4.6 ГБ Unified Pool: Execution + Storage (динамически)
JVM Overhead ~820 МБ Network buffers, JVM internals, Python worker memory

Важно: User Memory ≠ безопасная зона

User Memory - это память для пользовательских Java/Python объектов, которые Spark не контролирует. Если ваш UDF создаёт гигантский словарь в Python, он занимает JVM Overhead (Python worker) или User Memory через py4j. OOM здесь - ответственность разработчика.

# Это занимает User Memory - Spark не может evict это
lookup_table = {k: v for k, v in spark.table("big_table").collect()}  # опасно!

# Лучше: broadcast variable (управляется Spark как Storage)
lookup_bv = spark.sparkContext.broadcast(
    {row.k: row.v for row in spark.table("big_table").collect()}
)

Execution Memory: память для вычислений

Execution Memory - рабочая память для операций, требующих буферизации данных в процессе вычисления:

Ключевое свойство: Execution Memory выделяется per task. Если на executor работают 4 task'а параллельно, каждый получает ≈ 1/4 доступного Execution пула. Но это не жёсткое ограничение: если три task'а ждут (блокированы на I/O), четвёртый может использовать весь пул.

# Параллелизм задач на executor влияет на память каждой задачи
# cores = spark.executor.cores = 4 → 4 задачи одновременно
# Execution Memory / 4 = память на одну задачу
spark.conf.set("spark.executor.cores", "4")

# Увеличение cores → больше параллелизма, но меньше памяти на задачу
# Уменьшение cores → больше памяти на задачу, но меньше параллелизма

Storage Memory: кэш и broadcast

Storage Memory хранит:

  • Кэшированные DataFrame/RDD (.cache(), .persist(StorageLevel.MEMORY_AND_DISK))
  • Broadcast variables (широковещательные таблицы для Broadcast Hash Join)
  • Unroll буфер - временная область для разворачивания сжатых блоков при чтении из кэша
from pyspark import StorageLevel

# Разные уровни хранения → разный баланс памяти/диска/CPU
df.persist(StorageLevel.MEMORY_ONLY)          # только в Storage Memory (как UnsafeRow)
df.persist(StorageLevel.MEMORY_AND_DISK)      # в памяти, при нехватке - на диск
df.persist(StorageLevel.MEMORY_ONLY_SER)      # сериализованные байты (меньше памяти, CPU на десериализацию)
df.persist(StorageLevel.MEMORY_AND_DISK_SER)  # сериализованные, с spill на диск
df.persist(StorageLevel.DISK_ONLY)            # только на диск (обходит Storage Memory)
df.persist(StorageLevel.OFF_HEAP)             # в Tungsten off-heap (вне GC)

Broadcast variables:

# Broadcast переменная занимает Storage Memory на каждом executor
# Spark размещает её там при доставке через BroadcastExchange
small_dict = spark.sparkContext.broadcast({"US": "United States", "DE": "Germany"})

# Spark UI → Storage → видите broadcast_X_piece_Y в Storage Memory

RDD vs DataFrame: разница в объёме кэша

Одни и те же данные при кэшировании через RDD и DataFrame занимают принципиально разный объём памяти:

Одинаковый датасет (условно 10 ГБ CSV):
  rdd.cache()                           → ~60 ГБ в Storage Memory
  df.cache()                            → ~15 ГБ в Storage Memory
                                                 ↑ в 4 раза меньше!

Причина - способ хранения:

RDD cache DataFrame cache
Формат в памяти Java-объекты (heap) Tungsten UnsafeRow (бинарный off-heap или on-heap)
Накладные расходы Заголовок объекта (16 байт) + ссылки Плотный бинарный формат без объектных заголовков
GC влияние Каждая строка - отдельный Java-объект в Old Gen GC не видит данные (off-heap) или видит как один блок
Типичный коэффициент 4–6× больше baseline

Это одна из причин, почему DataFrame API предпочтительнее RDD: не только оптимизации Catalyst, но и в 4 раза меньшее потребление Storage Memory при кэшировании.

# Плохо: RDD-кэш дорог по памяти и не использует Tungsten
rdd = df.rdd
rdd.cache()
rdd.count()  # ~60 ГБ занято

# Хорошо: DataFrame-кэш компактен (UnsafeRow формат)
df.cache()
df.count()  # ~15 ГБ занято

Динамическая граница: механика заимствования

Это ключевая инновация Unified Memory Model. Граница между Execution и Storage - мягкая: регионы могут заимствовать друг у друга неиспользованную память.

Правила заимствования

Правило 1: Execution может взять память у Storage

Если Execution не хватает памяти и в Storage есть свободное место (сверх гарантированного минимума storageFraction), Execution берёт эту память.

Если Storage заполнен кэшем и Execution хочет больше - Spark принудительно вытесняет (evict) кэшированные блоки из Storage: сначала из памяти на диск (если MEMORY_AND_DISK), потом просто удаляет (если MEMORY_ONLY).

Правило 2: Storage может взять свободную память у Execution

Если Storage нужна память для кэширования и в Execution есть свободное место, Storage заимствует его.

Правило 3: Storage НЕ может выселить Execution

Если Execution занял всю память (в том числе взятую у Storage), Storage не может его выселить. Вычисления в приоритете над кэшем. Если Storage при этом нужна память - попытка кэширования просто не произойдёт (блок не закэшируется, придётся пересчитать).

Три сценария «перетягивания каната»

Сценарий 1: ETL с тяжёлыми join (много Execution, мало Storage)

Начало: Storage = 2.3 ГБ (пусто), Execution = 2.3 ГБ (частично занято)
Большой groupBy: Execution нужно 3.5 ГБ под hash table
→ Execution заимствует 1.2 ГБ из Storage (тот пустой)
→ Execution использует все 3.5 ГБ
Итог: Storage эффективно = 1.1 ГБ, Execution = 3.5 ГБ

Сценарий 2: ML с кэшем (много Storage, потом тяжёлый shuffle)

Начало: ML кэширует training data → Storage = 4 ГБ (взял у Execution)
Запускается обучение: shuffle нужно 2.5 ГБ в Execution
→ Execution пытается взять у Storage
→ Storage занят кэшем, Execution начинает evict'ить блоки
→ Кэш training data вылетает в DISK_ONLY или удаляется
→ Следующая итерация перечитывает данные с диска (медленно!)

Сценарий 3: Кэширование при занятом Execution

Execution занят тяжёлым join = 4 ГБ
df.cache() - пользователь хочет закэшировать 1 ГБ данных
→ Storage пытается взять память у Execution
→ Execution занят (не может быть выселен)
→ Блоки df не закэшируются в память
→ При следующем обращении df будет пересчитан заново

Spill to Disk: когда памяти не хватает

Когда Execution не может получить достаточно памяти (eviction Storage не помог, вся память занята), оператор сбрасывает (spill) промежуточные данные на диск.

Механика Spill при Shuffle

Что теряется при Spill:

Метрика Без Spill Со Spill
Disk I/O 0 Spill size × 2 (write + read)
CPU Только вычисления Сериализация + сортировка для spill
Время Baseline 10–100x медленнее
Тип данных In-memory UnsafeRow Сжатый (LZ4) бинарный файл

Spark сжимает spill-файлы через LZ4 по умолчанию (spark.shuffle.spill.compress=true), что уменьшает объём I/O, но добавляет CPU.

Типы Spill в Spark UI

Task Metrics в Spark UI → Stages → (выбрать Stage):

Spill (Memory): объём данных, которые были в памяти до сброса на диск
Spill (Disk):   фактический объём записанных на диск данных (после сжатия)
Shuffle Write:  объём shuffle-данных
GC Time:        время GC (косвенный индикатор memory pressure)
# Spill виден также в логах executor'а:
# INFO ExternalSorter: Thread 1 spilling in-memory map of 512.0 MB to disk (1 time so far)
# INFO ExternalAppendOnlyMap: Spilling in-memory map of 256.0 MB to disk

Наблюдение через Spark UI

Storage Tab

Spark UI → Storage

Таблица кэшированных датасетов:
┌──────────────────┬────────────┬─────────────┬──────────────┬──────────────┐
│ RDD Name         │ Storage    │ Cached      │ Size in      │ Size on      │
│                  │ Level      │ Partitions  │ Memory       │ Disk         │
├──────────────────┼────────────┼─────────────┼──────────────┼──────────────┤
│ DataFrame[...]   │ Disk Memory│ 200/200     │ 1.8 GB       │ 0.0 B        │  ← всё в памяти
│ broadcast_0      │ Memory     │ 1/1         │ 45.0 MB      │ 0.0 B        │  ← broadcast var
│ DataFrame[...]   │ Disk Memory│ 180/200     │ 1.2 GB       │ 600.0 MB     │  ← часть spill'd
└──────────────────┴────────────┴─────────────┴──────────────┴──────────────┘

Третья строка: 20 партиций были вытеснены на диск!

Executors Tab

Spark UI → Executors

┌────┬─────────┬──────────┬──────────────────────────────┬──────────┬──────────┐
│ ID │ Address │ Cores    │ Memory Usage                 │ GC Time  │ Tasks    │
├────┼─────────┼──────────┼──────────────────────────────┼──────────┼──────────┤
│  1 │ host1   │ 4 / 4    │ 3.1 GB / 4.6 GB used         │ 2.1 min  │ 145      │
│  2 │ host2   │ 4 / 4    │ 4.4 GB / 4.6 GB used ← 96%! │ 8.3 min  │ 143      │
└────┴─────────┴──────────┴──────────────────────────────┴──────────┴──────────┘

Executor 2: почти весь пул занят + высокое GC Time → кандидат на OOM

Stages Tab: Spill метрики

Spark UI → Stages → Stage 4 (Sort) → Summary Metrics:

Task Metric         Min      25th    Median   75th    Max
─────────────────────────────────────────────────────────────
Duration            12s      18s     22s      31s     180s     ← большой разброс
GC Time             0.1s     0.5s    1.2s     3.1s    45s      ← высокое GC
Spill (Memory)      0        0       512 MB   1.2 GB  3.8 GB   ← spill есть!
Spill (Disk)        0        0       89 MB    210 MB  650 MB   ← сжатый spill
Peak Exec Memory    800 MB   1.1 GB  1.8 GB   2.3 GB  2.3 GB

Когда Spill (Memory) ≠ 0 и Spill (Disk) ≠ 0 - это сигнал, что executor'у не хватает Execution Memory.

Настройки Unified Memory Model

Ключевые параметры

# Доля JVM heap, отводимая под Spark Memory (Unified Pool)
# Остаток (1 - fraction) = User Memory
spark.conf.set("spark.memory.fraction", "0.6")       # по умолчанию 0.6

# Доля Unified Pool, гарантированная Storage (soft limit)
# Остаток = soft limit для Execution
spark.conf.set("spark.memory.storageFraction", "0.5") # по умолчанию 0.5

# Off-heap память (вне GC, для Tungsten)
spark.conf.set("spark.memory.offHeap.enabled", "true")
spark.conf.set("spark.memory.offHeap.size", "4g")

# Overhead сверх executor.memory (native memory, Python worker)
spark.conf.set("spark.executor.memoryOverhead", "1g")  # или 10% автоматически

Когда и что менять

Типичные конфигурации под разные нагрузки

# ETL с тяжёлыми join и shuffle (ML-модели, агрегации)
# Много Execution, мало Storage (кэш почти не используется)
spark = SparkSession.builder \
    .config("spark.executor.memory", "16g") \
    .config("spark.executor.cores", "4") \
    .config("spark.memory.fraction", "0.75") \       # больше Spark Memory
    .config("spark.memory.storageFraction", "0.2") \ # меньше Storage soft limit
    .config("spark.memory.offHeap.enabled", "true") \
    .config("spark.memory.offHeap.size", "8g") \
    .getOrCreate()

# Iterative ML (KMeans, ALS, gradient boosting)
# Много Storage (переиспользуем training data), умеренный Execution
spark = SparkSession.builder \
    .config("spark.executor.memory", "16g") \
    .config("spark.executor.cores", "2") \           # меньше cores → больше памяти на task
    .config("spark.memory.fraction", "0.7") \
    .config("spark.memory.storageFraction", "0.6") \ # больше Storage soft limit
    .getOrCreate()

# Streaming (небольшие micro-batches, минимальный кэш)
spark = SparkSession.builder \
    .config("spark.executor.memory", "4g") \
    .config("spark.memory.fraction", "0.6") \
    .config("spark.memory.storageFraction", "0.3") \
    .config("spark.executor.memoryOverhead", "1g") \ # Python workers
    .getOrCreate()

Off-Heap и Unified Memory Model

Off-heap Memory существует параллельно с Unified Pool - это отдельный пул памяти вне Java heap, управляемый Tungsten напрямую через sun.misc.Unsafe:

Total Executor Memory:
┌─────────────────────────────────────────────────────────────────┐
│ JVM Heap (executor.memory = 8 ГБ)                               │
│ ┌─────────────┬────────────────────────────────────────────┐    │
│ │Reserved     │ Unified Pool (4.6 ГБ)          │User Mem   │    │
│ │300 МБ       │ ┌────────────┬───────────────┐  │ 3.1 ГБ    │    │
│ │             │ │Execution   │ Storage       │  │           │    │
│ │             │ │(dynamic)   │ (dynamic)     │  │           │    │
│ │             │ └────────────┴───────────────┘  │           │    │
│ └─────────────┴────────────────────────────────────────────┘    │
│                                                                   │
│ Off-Heap (spark.memory.offHeap.size = 4 ГБ) - вне GC             │
│ ┌──────────────────────────────────────────────────────────┐    │
│ │ Tungsten UnsafeRow buffers + OFF_HEAP persist blocks     │    │
│ └──────────────────────────────────────────────────────────┘    │
│                                                                   │
│ JVM Overhead (memoryOverhead = ~820 МБ)                           │
│ ┌──────────────────────────────────────────────────────────┐    │
│ │ Native memory: network buffers, JVM internals, Python    │    │
│ └──────────────────────────────────────────────────────────┘    │
└─────────────────────────────────────────────────────────────────┘

Off-heap особенно полезен при:

  • GC Time > 10% в Spark UI → переводим Tungsten buffers в off-heap
  • Большом кэше с StorageLevel.OFF_HEAP → данные вне GC
  • Контейнерной среде (K8s/YARN) с жёсткими лимитами на heap
# Суммарная память контейнера при off-heap:
# container_mem = executor.memory + executor.memoryOverhead + offHeap.size
# YARN: запрашиваем executor.memory + executor.memoryOverhead (overhead включает offHeap автоматически в YARN)
# K8s: явно задаём все три компонента

Практика: диагностика и tuning

Шаг 1: Симуляция memory pressure

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, rand

spark = SparkSession.builder \
    .appName("MemoryModel-Demo") \
    .config("spark.executor.memory", "2g") \     # намеренно мало для демонстрации
    .config("spark.memory.fraction", "0.6") \
    .config("spark.memory.storageFraction", "0.5") \
    .getOrCreate()

# Создаём датасет ~500 МБ
df = spark.range(5_000_000) \
    .withColumn("key", (col("id") % 1000).cast("int")) \
    .withColumn("value", rand() * 1000) \
    .withColumn("payload", col("id").cast("string"))

# Кэшируем - занимает Storage Memory
df.cache()
df.count()  # materialise cache

print("Проверьте Spark UI → Storage: сколько занято в памяти vs диске")

Шаг 2: Провоцируем конфликт Execution vs Storage

# Тяжёлый join - нужно много Execution Memory
from pyspark.sql.functions import broadcast

large_df = spark.range(10_000_000) \
    .withColumn("key", (col("id") % 1000).cast("int")) \
    .withColumn("value", rand() * 10000)

# Этот join потребует много Execution Memory
# и может вытеснить df из Storage
result = df.join(large_df, "key") \
           .groupBy("key") \
           .agg({"value": "sum"}) \
           .count()

print(result)
print("Проверьте Storage tab: остался ли df в памяти или был evict'ирован?")

Шаг 3: Чтение метрик spill

# Проверить spill через Spark UI → Stages → Stage with Sort/Shuffle
# или через SparkContext listener (advanced)

# Быстрая проверка: если задача очень долгая и диск активен → spill
# В логах executor ищем:
# "Spilling in-memory map" или "ExternalSorter: Thread spilling"

# Диагностика через конфиг:
spark.conf.set("spark.sql.execution.sortBeforeRepartition", "true")
spark.conf.set("spark.shuffle.spill.numElementsForceSpillThreshold", "1000000")

Шаг 4: Сравнение storageFraction

import time

# Много Storage (кэш приоритетен)
spark.conf.set("spark.memory.storageFraction", "0.7")
df.cache()
t0 = time.time()
df.join(large_df, "key").groupBy("key").count().show()
print(f"storageFraction=0.7: {time.time()-t0:.1f}s")
df.unpersist()

# Мало Storage (Execution приоритетен)
spark.conf.set("spark.memory.storageFraction", "0.2")
df.cache()
t0 = time.time()
df.join(large_df, "key").groupBy("key").count().show()
print(f"storageFraction=0.2: {time.time()-t0:.1f}s")
df.unpersist()

Чеклист диагностики OOM

# 1. Определить тип OOM из stacktrace:

# java.lang.OutOfMemoryError: Java heap space
#   → нехватка heap: увеличить executor.memory

# java.lang.OutOfMemoryError: GC overhead limit exceeded
#   → слишком много объектов в heap: включить off-heap, уменьшить кэш,
#     использовать MEMORY_AND_DISK_SER

# java.lang.OutOfMemoryError: Direct buffer memory
#   → нехватка native memory: увеличить executor.memoryOverhead

# org.apache.spark.memory.SparkOutOfMemoryError: Unable to acquire X bytes
#   → нехватка Execution Memory: уменьшить cores/task,
#     увеличить memory.fraction, проверить на spill

# PySpark: Python worker killed by SIGKILL (OOM)
#   → Python worker превысил память: увеличить memoryOverhead,
#     переписать UDF на built-in функции или Pandas UDF

Кэш: когда использовать, когда избегать

Когда кэш ускоряет

# Многократное переиспользование одного датасета (ML итерации)
training_data = spark.table("features").filter("split = 'train'")
training_data.cache()  # первое обращение - чтение с диска + кэш
training_data.count()  # materialise

for iteration in range(100):
    model = train_iteration(training_data)  # каждая итерация из памяти → быстро

Когда кэш мешает

# Линейный пайплайн без переиспользования - кэш только занимает память
df1 = spark.table("raw").filter(...)    # используется один раз
df2 = df1.join(ref, "key")             # используется один раз
df3 = df2.groupBy("country").sum(...)  # используется один раз

# Плохо: кэшировать df1 и df2 - бессмысленно
df1.cache()  # займёт Storage, будет evict'ирован при join
df2.cache()  # занимает Storage при groupBy

# Хорошо: не кэшировать, Spark сам оптимизирует пайплайн
result = df3.write.parquet("/output")

Правило: когда кэшировать

  • Датасет используется 2+ раз в вычислительном графе
  • Стоимость пересчёта значительно выше, чем чтение из кэша
  • Датасет помещается в память (иначе MEMORY_AND_DISK или не кэшировать)
  • В ML-итерациях - всегда кэшировать training/validation sets

Best Practices

1. Подбирать executor.memory и executor.cores сбалансированно:

# Правило: память на задачу = executor.memory × fraction / cores
# Пример: 8 ГБ × 0.6 / 4 cores = 1.2 ГБ на задачу для Execution
# Если задаче нужно 3 ГБ → уменьшить cores до 2 (3 ГБ / task)

2. Следить за GC Time в Spark UI:

GC Time < 5% - норма
GC Time 5-15% - стоит включить off-heap
GC Time > 15% - серьёзная проблема, оптимизировать

3. Использовать правильный StorageLevel:

# Для повторно используемых данных
df.persist(StorageLevel.MEMORY_AND_DISK)  # безопаснее MEMORY_ONLY
# При высоком GC
df.persist(StorageLevel.MEMORY_AND_DISK_SER)  # меньше объектов в heap
# При нехватке heap, но достаточно off-heap
df.persist(StorageLevel.OFF_HEAP)

4. Явно unpersist когда кэш больше не нужен:

training_data.cache()
model = train(training_data)
training_data.unpersist()  # освобождаем Storage Memory для следующих шагов

5. Не кэшировать результаты shuffle (они уже на диске):

# После shuffle данные на диске в shuffle files - кэш не нужен
result = df.groupBy("key").sum("value")  # shuffle произошёл
result.cache()  # бессмысленно если result используется один раз

Итого

Регион Размер Назначение Конфиг
Reserved Memory 300 МБ (фиксировано) Внутренние нужды Spark -
User Memory (1 - fraction) × (heap - 300 МБ) UDF объекты, Encoders spark.memory.fraction
Unified Pool fraction × (heap - 300 МБ) Execution + Storage spark.memory.fraction
Storage soft limit storageFraction × unified Cache, Broadcast spark.memory.storageFraction
Execution soft limit (1 - storageFraction) × unified Shuffle, Sort, Join spark.memory.storageFraction
Off-Heap offHeap.size Tungsten buffers, OFF_HEAP cache spark.memory.offHeap.*
JVM Overhead max(384 МБ, 10%) Native memory, Python spark.executor.memoryOverhead

Правила заимствования:

  • Execution может взять свободное место у Storage, и принудительно evict'ировать кэш
  • Storage может взять только свободное место у Execution - вычисления выселить нельзя
  • Если нет места даже после всех перестановок - Execution делает Spill to Disk

Следующий урок рассматривает Spill to Disk подробнее: как работают External Sorter и External Hash Map, когда и почему spill катастрофически замедляет job, и как избежать spill через правильный tuning.