Spill на диск: причины, диагностика в Spark UI и как избежать
Механика spill в Spark: почему данные сбрасываются на диск при нехватке Execution Memory, как диагностировать через Spark UI и какие техники устраняют spill в production
Что такое Spill¶
Spill - это механизм защиты Spark от OOM: когда оператору (Sort, Aggregation, Join) не хватает Execution Memory для хранения промежуточных данных, он сбрасывает часть данных на диск, освобождает память, продолжает обработку, а в конце читает всё обратно и объединяет результаты.
Spill - не ошибка и не сбой. Spark продолжает работать корректно. Но это сигнал: памяти недостаточно, и за корректность приходится платить производительностью - дисковый I/O в 10–100 раз медленнее RAM.
Без Spill:
┌────────────────────────────────────────┐
│ Execution Memory (2 ГБ) │
│ [hash table: 1.8 ГБ] [free: 0.2 GB]│
│ │
│ Агрегация завершена полностью в памяти │
└────────────────────────────────────────┘
Время: 30 секунд
Со Spill:
┌────────────────────────────────────────┐
│ Execution Memory (2 ГБ) │
│ [batch 1: 1.9 ГБ] → сброс на диск │
│ [batch 2: 1.9 ГБ] → сброс на диск │
│ Merge spill files + in-memory data │
└────────────────────────────────────────┘
Время: 5–10 минут (merge с диска!)
Виды Spill в Spark¶
В Spark UI вы увидите две колонки: Spill (Memory) и Spill (Disk). Это разные величины одного события.
Spill (Memory) - сколько десериализованных данных было в памяти. Это «честный» объём: 500 МБ означает, что 500 МБ UnsafeRow данных не поместилось.
Spill (Disk) - сколько байт фактически записано на диск. Из-за LZ4-сжатия может быть в 3–5x меньше Spill (Memory). Именно эта метрика определяет нагрузку на диск и влияет на время I/O.
Три контекста возникновения Spill¶
Hash Aggregation vs Sort Aggregation¶
Это важнейшее разветвление при groupBy/agg с точки зрения памяти.
Hash Aggregation (предпочтительный путь):
groupBy("city").agg(sum("amount"))
In-memory hash table:
┌──────────────┬────────────┐
│ "Moscow" │ 1_234_567 │
│ "Berlin" │ 987_654 │
│ "Paris" │ 456_789 │
│ ... │ ... │
└──────────────┴────────────┘
Требует: O(distinct_keys × row_size) памяти
Если hash table не помещается в память - Spark автоматически переключается на Sort Aggregation:
Hash Agg переполнился → Sort-based fallback:
1. Сбросить hash table на диск (spill!)
2. Отсортировать данные по ключу
3. Итерировать: соседние строки с одним ключом агрегируются в потоке
Требует: O(partition_size) памяти (только буфер сортировки)
# В Spark UI (SQL tab): если видите переключение:
# ObjectHashAggregate → *(1) HashAggregate переключился на Sort Agg
# В explain() появится SortAggregateExec вместо HashAggregateExec
# Контроль через конфиг:
spark.conf.set("spark.sql.execution.useObjectHashAggregateExec", "true") # по умолчанию
# Лог executor'а при переключении:
# INFO HashAggregateExec: Falling back to sort-based aggregation due to memory pressure
Механика Spill при Sort¶
ExternalSorter - компонент, обрабатывающий сортировку и shuffle write. Именно он управляет spill:
Ключевой момент: каждый spill создаёт отдельный файл с частично отсортированными данными. Финальный merge читает все эти файлы одновременно через k-way merge с приоритетной очередью. Если было 50 spill-файлов - merge читает 50 файлов параллельно. Это дорого по I/O и CPU.
Причины Spill: от самых частых к редким¶
1. Data Skew - самая частая причина¶
Неравномерное распределение данных по партициям приводит к тому, что одна задача получает 95% данных, пока остальные простаивают.
Диагностика skew в Spark UI:
Stages → Stage 5 → Summary Metrics:
Task Metric Min 25th Median 75th Max
─────────────────────────────────────────────────────────
Duration 0.1s 0.3s 0.5s 1.2s 18m ← !
Spill(Memory) 0 0 0 0 9.5GB ← !
Spill(Disk) 0 0 0 0 2.1GB ← !
Records Read 1K 2K 3K 3.5K 950K ← !
Max >> Median → это skew, а не просто нехватка памяти
2. Слишком мало shuffle partitions¶
По умолчанию spark.sql.shuffle.partitions = 200. На 10 ТБ данных это даёт партиции по 50 ГБ каждая - гарантированный spill.
# Правило: 100–200 МБ на партицию
# Объём shuffle данных = 1 ТБ → нужно 1024 ГБ / 0.15 ГБ ≈ 6800 партиций
spark.conf.set("spark.sql.shuffle.partitions", "6800")
# AQE делает это автоматически (Spark 3.0+)
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
3. Exploding Join: декартово произведение¶
# Опасно: если orders имеет 1 млрд строк, а tags - 1000 строк
# результат = 1 трлн строк → гарантированный spill + OOM
result = orders.crossJoin(tags)
# Также опасно: join без правильного ключа
result = orders.join(events, orders.date == events.date)
# Если в один день миллиард событий → exploding partition
4. Слишком много кэша занимает Storage Memory¶
# Если Storage Memory взяла у Execution → Execution меньше места
df1.cache() # занял 3 ГБ Storage
df2.cache() # занял ещё 2 ГБ Storage
# Теперь Execution Memory урезана → агрегации идут в spill
# Диагностика: Spark UI → Storage tab → смотреть суммарный объём кэша
# vs доступный Execution Memory в Executors tab
5. Слишком мало executor cores (большие задачи)¶
# executor.memory = 16 ГБ, executor.cores = 8
# Execution Memory / 8 tasks = ~1 ГБ на задачу
# Если одна задача агрегирует 5 ГБ → spill
# Решение: уменьшить cores → больше памяти на задачу
# executor.memory = 16 ГБ, executor.cores = 2
# Execution Memory / 2 tasks = ~4 ГБ на задачу
Диагностика Spill в Spark UI¶
Stage-level диагностика¶
Spark UI → Stages → выбрать подозрительный Stage
В шапке таблицы Tasks:
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Input Output Shuffle Read Shuffle Write Spill(Mem) Spill(Disk)
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
120 GB -- 45 GB -- 8.5 GB 1.9 GB
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Spill(Mem) = 8.5 ГБ → серьёзная проблема
Spill(Disk) = 1.9 ГБ → сжатые данные на диске
Summary Metrics → смотреть на распределение между задачами:
Если Max >> Median → skew
Если все задачи имеют spill → нехватка памяти в целом
SQL Tab: находим оператор-виновник¶
Spark UI → SQL → выбрать Query → DAG
Искать операторы с аннотацией "spilled":
┌─────────────────────────────────────────────────────┐
│ Sort │
│ sort keys: [amount DESC] │
│ spill size: 4.5 GiB ← вот он! │
│ peak memory: 1.8 GiB │
│ ──────────────────────────────────────────────── │
│ Exchange (rangepartitioning) │
│ shuffle records written: 850M │
│ data size: 120 GiB │
└─────────────────────────────────────────────────────┘
Также смотреть:
- "number of output rows" после фильтра vs до join → exploding join?
- "data size" Exchange → если огромный → много shuffle data
Task-level диагностика¶
Stages → Stage → Tasks (список всех задач)
Сортировать по колонке "Spill (Memory)" - DESC.
Задачи с наибольшим spill - это либо горячие партиции (skew)
либо все задачи одинаково (нехватка памяти).
Также смотреть:
- Duration: skewed tasks будут в 10-100x дольше остальных
- GC Time: высокий GC + spill = double trouble
- Shuffle Read Size: большой → много данных на reducer
Логи executor'а¶
# В логах executor'а искать (grep):
grep "Spilling\|spilling\|spill" executor.log
# Типичные записи:
INFO ExternalSorter: Thread 47 spilling in-memory map of 1024.0 MB to disk (3 times so far)
INFO ExternalAppendOnlyMap: Spilling in-memory map of 512.0 MB to disk
INFO UnsafeExternalSorter: Thread 12 spilling sort data of 768.0 MB to disk (1 time so far)
Влияние Spill на производительность¶
Особенно болезненно на:
- Cloud object storage (S3, GCS) как
spark.local.dir: spill-файлы пишутся через сеть - в 1000x медленнее локального SSD - Медленные HDD: latency seek = 5–10 мс против 0.1 мс NVMe SSD
- Shared storage в контейнерах: spill конкурирует за I/O с другими процессами
# Убедиться что spark.local.dir указывает на локальный быстрый диск
spark.conf.get("spark.local.dir") # должен быть /local/disk1, не /mnt/s3/...
# Для K8s:
# spec.volumes → emptyDir с medium: "" (RAM) для маленьких spill
# или hostPath → локальный NVMe узла
Стратегии устранения Spill¶
1. Broadcast Join: убрать shuffle полностью¶
from pyspark.sql.functions import broadcast
# Было: Sort Merge Join → оба датасета shuffle → potential spill
result = large_orders.join(regions, "region_id")
# Стало: Broadcast Hash Join → regions рассылается, shuffle нет
result = large_orders.join(broadcast(regions), "region_id")
# Когда применимо:
# - малая сторона < autoBroadcastJoinThreshold (по умолчанию 10 МБ)
# - или явно через broadcast() если знаем что влезет в память
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "104857600") # 100 МБ
2. Salting: борьба со skewed groupBy¶
Если groupBy("city") даёт skew из-за одного доминирующего города:
from pyspark.sql.functions import col, concat, lit, floor, rand, explode, array
SALT_FACTOR = 20 # дробим горячие ключи на 20 частей
# Шаг 1: добавить соль к данным
df_salted = df.withColumn(
"salted_city",
concat(col("city"), lit("_"), (rand() * SALT_FACTOR).cast("int").cast("string"))
)
# Шаг 2: частичная агрегация по солёному ключу
partial = df_salted.groupBy("salted_city").agg({"amount": "sum"})
# Шаг 3: финальная агрегация - убираем соль, группируем по оригинальному ключу
from pyspark.sql.functions import regexp_replace
final = (partial
.withColumn("city", regexp_replace(col("salted_city"), "_\\d+$", ""))
.groupBy("city")
.agg({"sum(amount)": "sum"}))
3. Увеличение числа партиций¶
# До shuffle: убедиться что данные равномерно разбиты
df = df.repartition(1000, "region_id") # hash partition по ключу
# После shuffle: настроить число reduce-партиций
spark.conf.set("spark.sql.shuffle.partitions", "2000")
# Формула: target_partition_size = 100-200 МБ
# n_partitions = total_shuffle_data_bytes / 150_MB
import math
total_gb = 500 # оцениваем объём shuffle данных
n_partitions = math.ceil(total_gb * 1024 / 150) # ~3413
spark.conf.set("spark.sql.shuffle.partitions", str(n_partitions))
4. Ранняя фильтрация и проекция¶
# Плохо: join с полными таблицами, потом фильтр
result = (orders.join(customers, "customer_id")
.filter(col("order_date") > "2024-01-01")
.select("order_id", "customer_name", "amount"))
# Хорошо: фильтруем и проецируем ДО join (меньше данных → меньше памяти → меньше spill)
orders_filtered = (orders
.filter(col("order_date") > "2024-01-01")
.select("order_id", "customer_id", "amount"))
customers_slim = customers.select("customer_id", "customer_name")
result = orders_filtered.join(customers_slim, "customer_id")
5. AQE: автоматическая борьба со skew¶
# Включить AQE (Spark 3.0+, по умолчанию true в 3.2+)
spark.conf.set("spark.sql.adaptive.enabled", "true")
# AQE Skew Join Optimization: автоматически дробит skewed partitions
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")
# Порог: партиция считается skewed если > N * median И > min size
spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionFactor", "5")
spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes", "256MB")
# AQE Coalesce: уменьшает число пустых/маленьких партиций
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
spark.conf.set("spark.sql.adaptive.advisoryPartitionSizeInBytes", "128MB")
Как работает AQE Skew Join:
6. Tuning Execution Memory¶
# Освободить больше памяти для Execution:
# Способ 1: увеличить memory.fraction (меньше User Memory)
spark.conf.set("spark.memory.fraction", "0.75") # было 0.6
# Способ 2: уменьшить storageFraction (меньше Storage soft limit)
spark.conf.set("spark.memory.storageFraction", "0.3") # было 0.5
# Способ 3: меньше cores → больше памяти на задачу
# executor.memory=16g, cores=2 → ~4 ГБ Execution на задачу
# executor.memory=16g, cores=8 → ~1 ГБ Execution на задачу
# Способ 4: off-heap разгружает heap GC
spark.conf.set("spark.memory.offHeap.enabled", "true")
spark.conf.set("spark.memory.offHeap.size", "8g")
# Способ 5: unpersist неиспользуемый кэш
df_temp.unpersist() # освобождает Storage → Execution может взять больше
Практика: воспроизведение и устранение Spill¶
Шаг 1: Создаём skewed dataset¶
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, rand, when
spark = SparkSession.builder \
.appName("Spill-Demo") \
.config("spark.executor.memory", "2g") \ # намеренно мало
.config("spark.sql.shuffle.partitions", "10") \ # намеренно мало
.config("spark.sql.adaptive.enabled", "false") \ # отключить AQE для демо
.getOrCreate()
# Создаём skewed данные: 80% записей с city='Moscow'
df = (spark.range(5_000_000)
.withColumn("city",
when(rand() < 0.8, "Moscow")
.when(rand() < 0.5, "Berlin")
.otherwise("Paris"))
.withColumn("amount", (rand() * 10000).cast("double"))
.withColumn("payload", (rand() * 1000).cast("long").cast("string")))
df.write.mode("overwrite").parquet("/tmp/spill_demo")
Шаг 2: Запускаем проблемный запрос¶
df = spark.read.parquet("/tmp/spill_demo")
# Этот groupBy будет иметь сильный skew и spill
result = df.groupBy("city").agg({"amount": "sum", "payload": "count"})
result.show()
# Проверяем Spark UI → Stages → ищем Spill(Memory) > 0
# Ожидаем: Stage с groupBy покажет большой spill на задаче с 'Moscow'
Шаг 3: Исправляем через увеличение партиций + AQE¶
spark.conf.set("spark.sql.shuffle.partitions", "200")
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")
result = df.groupBy("city").agg({"amount": "sum", "payload": "count"})
result.show()
# Ожидаем: Spill(Memory) значительно меньше или 0
# Stage duration уменьшился
Шаг 4: Исправляем через Salting¶
from pyspark.sql.functions import concat, lit, regexp_replace
SALT = 50
# Добавляем соль
df_s = df.withColumn("city_s",
concat(col("city"), lit("_"), (rand() * SALT).cast("int").cast("string")))
# Частичная агрегация
partial = df_s.groupBy("city_s").agg({"amount": "sum", "payload": "count"})
# Убираем соль и финальная агрегация
final = (partial
.withColumn("city", regexp_replace("city_s", "_\\d+$", ""))
.groupBy("city")
.agg({"sum(amount)": "sum", "count(payload)": "sum"}))
final.show()
# Spill должен исчезнуть: каждая из 50 субпартиций "Moscow"
# теперь ≈ 80K записей вместо 4M
Шаг 5: Сравниваем метрики¶
# Смотреть в Spark UI:
# До оптимизации:
# - Stage Duration: 8 min (один straggler task)
# - Spill (Memory): 3.2 GB
# - Spill (Disk): 720 MB
# После AQE + partitions:
# - Stage Duration: 45 sec
# - Spill (Memory): 0
# - Spill (Disk): 0
Spill и shuffle: взаимосвязь¶
Shuffle сам по себе создаёт временные файлы на диске (shuffle files) - это не spill. Но если данные на reduce-стороне не помещаются при чтении и агрегации → это spill.
# Настройки для уменьшения давления на reduce-буфер:
spark.conf.set("spark.reducer.maxSizeInFlight", "96m") # было 48 МБ, увеличить при OOM
spark.conf.set("spark.shuffle.file.buffer", "1m") # буфер записи (было 32 КБ)
spark.conf.set("spark.shuffle.spill.batchSize", "10000") # элементов перед spill check
Best Practices¶
1. Мониторить Spill (Memory) как primary SLA-метрику:
# В production: настроить алерт если суммарный Spill (Disk) > 10 ГБ на job
# Prometheus + Spark metrics:
# metrics.conf → sink.prometheus.class=...
# spark_stage_spill_disk_bytes_total
2. Подбирать число партиций по формуле:
# Оценить объём shuffle данных (из предыдущего run в Spark UI)
shuffle_gb = 800 # из Stage metrics → Shuffle Write
target_mb_per_partition = 128
n = int(shuffle_gb * 1024 / target_mb_per_partition)
spark.conf.set("spark.sql.shuffle.partitions", str(n))
# Или доверить AQE, если enabled
3. Проверять skew до начала оптимизации:
# Быстрая диагностика распределения ключей
df.groupBy("join_key").count().orderBy(col("count").desc()).show(20)
# Если топ-ключ имеет > 10% всех записей → skew, нужен salting или AQE
4. Никогда не кэшировать без надобности:
# Каждый GB в Storage Memory = меньше Execution Memory = больше риск spill
df.cache() # только если используется 2+ раз И помещается в память
# Всегда unpersist по завершении:
df.unpersist()
5. Проверять spark.local.dir:
# Должен быть быстрый локальный диск
spark.conf.get("spark.local.dir") # /local/nvme0, /tmp (SSD)
# НЕ: /mnt/nfs, s3a://, hdfs:// - это катастрофа при spill
Итого¶
| Причина Spill | Диагностика | Решение |
|---|---|---|
| Data Skew | Max >> Median в Summary Metrics | Salting, AQE Skew Join |
| Мало партиций | Большие Shuffle Read Size | Увеличить shuffle.partitions |
| Мало памяти | Все задачи имеют spill | Увеличить executor.memory, уменьшить cores |
| Exploding Join | Огромное число output rows | Broadcast Join, фильтрация до join |
| Кэш забирает память | Storage > Execution в Executors tab | unpersist, уменьшить storageFraction |
| Медленный диск | Высокий Task Duration при малом Spill | Настроить spark.local.dir на SSD |
Spill (Memory) и Spill (Disk) в Spark UI - это ваш главный индикатор проблем с памятью. Если они ненулевые - разберитесь, это skew или нехватка памяти, и применяйте соответствующий рецепт.
Следующий урок рассматривает GC Pressure - что происходит когда heap переполнен Java-объектами, как читать GC логи и почему иногда нужно переходить на off-heap даже без видимого spill.