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 нагрузки.
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. Три ключевых наблюдения:
- Reserved Memory (300 MB) абсолютно фиксирована - она не участвует ни в каких настройках
- User Memory определяется остатком после Spark Memory - чем больше
memory.fraction, тем меньше User Memory - 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:
-
Storage НИКОГДА не может вытеснить Execution. Даже если Storage Memory пустует, а Execution нужна память - Storage не «отнимает» у выполняющихся Task'ов.
-
Execution МОЖЕТ вытеснить Storage из Dynamic Zone. Если Task'е не хватает памяти, она вытеснит кешированные блоки из Dynamic Zone.
-
Storage Floor (storage.fraction × spark_memory) абсолютно защищена. Эту часть Execution не трогает никогда.
-
Максимально возможный размер обоих регионов = весь 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. Изменение без понимания проблемы может ухудшить производительность.
Алгоритм диагностики перед тюнингом:
- Открыть Spark UI → Stages - есть ли Spill (Disk) > 0?
- Если Spill есть → смотреть Peak Execution Memory на Task → рассчитать нужную Execution Memory
- Открыть Spark UI → Storage - насколько заполнен Storage, есть ли вытеснение кеша?
- Открыть Spark UI → Executors → GC Time > 10%? → проблема с Heap, а не с fractions
- Только после диагностики → применить
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 - это границы, внутри которых работает автоматика, а не жёсткие лимиты на каждую операцию.