Spill на диск: причины, диагностика в Spark UI и как избежать

Механика spill в Spark: почему данные сбрасываются на диск при нехватке Execution Memory, как диагностировать через Spark UI и какие техники устраняют spill в production

core internals optimization

Что такое 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.