GC Pressure: G1GC настройки, off-heap память и сигналы в логах

Почему Garbage Collector JVM тормозит Spark, как читать GC-логи, настраивать G1GC и использовать off-heap память для устранения Stop-the-World пауз

core internals optimization

Почему GC - проблема именно для Spark

PySpark - это Python-обёртка над JVM-процессом. Spark executor - это JVM, и все его данные: UnsafeRow буферы, shuffle structures, кэш DataFrame, промежуточные результаты агрегации - всё это живёт в Java heap (или рядом с ним, в off-heap). Когда heap переполняется - JVM останавливает все потоки для сборки мусора. В это время:

  • все Task'и executor'а заморожены
  • Driver считает executor'а висящим
  • если пауза длится дольше spark.network.timeout (120 сек) - executor считается потерянным
Типичный симптом в Spark UI (Executors tab):

Executor 3:
  Active Tasks:  4
  GC Time:       18 min из 22 min total  ← 82%!
  Task Time:     22 min

В Stage tab:
  Task 1234: Duration=45min, GC Time=38min  ← 84% времени в GC

GC pressure - это состояние, когда JVM тратит непропорционально много CPU на управление памятью вместо полезных вычислений. Порог: если GC Time > 10% Task Time - стоит разобраться. Если > 20% - серьёзная проблема.

Анатомия JVM Heap для Spark executor

Young Generation: быстрые Minor GC

Eden - место рождения объектов. Каждый new Object(), каждый промежуточный результат, каждая строка при сериализации - сначала попадает в Eden.

Minor GC (Young Collection): очищает Eden и Survivor пространства. Занимает 1–50 мс - приемлемо. Выжившие объекты переходят в Survivor, а при превышении возраста (tenuring threshold) - в Old Generation.

Проблема Spark: shuffle и агрегация порождают миллиарды короткоживущих объектов (промежуточные строки, буферы, итераторы). Eden переполняется быстро → Minor GC срабатывает часто → CPU занят GC, а не вычислениями.

Old Generation: дорогие Major/Full GC

Объекты, пережившие несколько Minor GC, попадают в Old Generation. Здесь живут:

  • Кэшированные DataFrame (.cache())
  • Broadcast переменные
  • Long-lived shuffle structures
  • Объекты из User Memory (словари UDF, аккумуляторы)

Major GC (G1 Mixed GC): очищает часть Old Generation. Занимает 100 мс – 2 сек при G1GC.

Full GC: очищает весь heap. Stop-the-World пауза 5–60 секунд. Это катастрофа для Spark: driver отмечает executor как dead, задачи перезапускаются.

Humongous Objects: особый класс проблем

G1GC делит heap на регионы одинакового размера (по умолчанию 1–32 МБ). Объект, занимающий > 50% размера региона, называется Humongous Object и выделяется напрямую в Old Generation, минуя Young.

Это критично для Spark потому что:

  • Shuffle буферы, строки Parquet, broadcast chunks - всё это может быть Humongous
  • Humongous объекты не компактируются → фрагментация Old Generation
  • Частые Humongous allocations вызывают преждевременные Full GC
GC лог: признак Humongous объекта:
[GC pause (G1 Humongous Allocation) 5120M->4890M(8192M), 0.8432 secs]
                                    ^^^^^^^^^
                                    Allocation failure из-за Humongous объекта

Эволюция GC в Spark

Почему G1GC - стандарт для Spark:

  • Предсказуемые паузы: цель MaxGCPauseMillis соблюдается в большинстве случаев
  • Инкрементальная очистка: не весь heap, а только выбранные регионы
  • Concurrent marking: большая часть работы идёт параллельно с приложением
  • Адаптивность: G1GC сам выбирает какие регионы очищать первыми (с наибольшим количеством мусора - отсюда "Garbage First")

Принцип работы G1GC

Фазы G1GC-цикла:

Фаза Тип Длительность Что делает
Young Collection STW 1–50 мс Очищает Eden, переводит в Survivor/Old
Concurrent Root Scan Concurrent мс Сканирует корни (stack, statics)
Concurrent Mark Concurrent 100–500 мс Помечает живые объекты по всему heap
Remark STW 1–10 мс Финализирует маркировку
Cleanup STW 1–5 мс Выбирает кандидатов на очистку
Mixed GC STW 50–200 мс Очищает Young + часть Old регионов
Full GC (аварийный) STW 5–60 сек Очищает всё, последний резерв

Как Spark создаёт GC pressure

Shuffle: фабрика временных объектов

При Sort Shuffle Write каждая запись проходит через:

1. Deserialize (Row object) → heap
2. Sort key extraction (new key object) → heap
3. Partition assignment (boxing int → Integer) → heap
4. Serialization to byte[] → heap
5. Write to buffer → heap

На 1 млн строк это 5+ млн промежуточных объектов → Eden переполняется сотни раз за Stage.

DataFrame vs RDD: кардинальная разница

# RDD: каждый элемент - Java-объект → огромное давление на GC
rdd_result = rdd.map(lambda r: (r[0], r[1] * 2)).filter(lambda r: r[1] > 100)
# Создаёт: 1 объект на каждую строку × 3 операции = 3 млрд объектов при 1 млрд строк

# DataFrame: UnsafeRow в Tungsten, Whole-Stage CodeGen
df_result = df.withColumn("v2", col("v") * 2).filter(col("v2") > 100)
# Создаёт: несколько батч-буферов на весь Stage
# Давление на GC - минимальное

Кэш в Old Generation

df.cache()  # сериализует DataFrame в Old Generation

# После cache():
# Old Generation занята: 5 ГБ из 6 ГБ
# G1GC начинает Mixed GC при 45% заполнения → срабатывает постоянно
# Minor GC стало чаще → Old Generation заполняется быстрее → порочный круг

Сигналы GC проблем в Spark UI

Executors Tab

Spark UI → Executors:

┌────┬─────────┬──────────┬──────────────┬───────────┬──────────┐
│ ID │ Status  │ RDD      │ Storage      │ Task Time │ GC Time  │
│    │         │ Blocks   │ Memory       │           │          │
├────┼─────────┼──────────┼──────────────┼───────────┼──────────┤
│  1 │ Active  │ 0        │ 2.1 GB/4.6GB │ 45 min    │ 2 min    │ ← 4%  OK
│  2 │ Active  │ 120      │ 4.4 GB/4.6GB │ 48 min    │ 14 min   │ ← 29% !!
│  3 │ Dead    │ 0        │ --           │ 32 min    │ 28 min   │ ← Потерян!
└────┴─────────┴──────────┴──────────────┴───────────┴──────────┘

Executor 2: GC Time 29% → нужна оптимизация
Executor 3: убит из-за GC → heartbeat timeout

Task Metrics

Stage → Tasks (Summary Metrics):

Metric        Min    25th   Median  75th   Max
─────────────────────────────────────────────────────
GC Time       0.1s   1.2s   3.8s    12s    145s  ← Max в 38x > Median
Duration      5s     8s     10s     18s    210s
Peak Exec Mem 1.2GB  1.8GB  2.1GB   2.3GB  2.3GB ← все упёрлись в потолок

Max GC Time >> Median → один executor сильно страдает
Peak Exec Mem = 2.3 ГБ у всех → пул исчерпан, GC срабатывает принудительно

Чтение GC логов

Включение GC логирования

# В spark-submit или SparkSession config:
spark = SparkSession.builder \
    .config("spark.executor.extraJavaOptions",
            "-XX:+UseG1GC "
            "-XX:+PrintGCDetails "
            "-XX:+PrintGCDateStamps "
            "-XX:+PrintGCTimeStamps "
            "-XX:+PrintAdaptiveSizePolicy "
            "-Xloggc:/tmp/gc-executor-%p.log") \
    .getOrCreate()

# Для Java 11+ (Spark 3.x):
spark = SparkSession.builder \
    .config("spark.executor.extraJavaOptions",
            "-XX:+UseG1GC "
            "-Xlog:gc*:file=/tmp/gc-executor-%p.log:time,uptime,level,tags") \
    .getOrCreate()

Анатомия GC лога

# Minor GC (Young Collection) - нормально, быстро:
2024-01-15T10:23:45.123+0300: 125.456: [GC pause (G1 Evacuation Pause) (young)
 125.456: [G1Ergonomics (Heap Sizing) attempt to expand the heap] 
 (to space exhausted), 0.0234567 secs]
   [Eden: 512.0M(512.0M)->0.0B(512.0M)    ← Eden очищен
    Survivors: 64.0M->64.0M
    Heap: 2048.0M(8192.0M)->1600.0M(8192.0M)]  ← heap 2→1.6 ГБ, 23 мс
 [Times: user=0.15 sys=0.01, real=0.02 secs]
#          ^^^^               ^^^^
#          CPU time           Wall clock time
#          real << user → многопоточная сборка, OK

# Mixed GC (Old + Young) - настораживает если часто:
[GC pause (G1 Evacuation Pause) (mixed) 5200M->4800M(8192M), 0.3456 secs]
#                                ^^^^^ очищает и Old регионы
# 346 мс - это на границе допустимого

# Humongous Allocation - проблема:
[GC pause (G1 Humongous Allocation) 6800M->6600M(8192M), 0.8765 secs]
#                     ^^^^^^^^^^^^
# Объект > 50% размера G1-региона выделяется напрямую в Old

# Full GC - катастрофа:
[Full GC (Allocation Failure) 7800M->2100M(8192M), 12.3456 secs]
#         ^^^^^^^^^^^^^^^^^^ heap заполнился, concurrent mark не успел
# 12 секунд STW-паузы! Heartbeat timeout = 120 сек → при 10 Full GC executor умирает

Паттерны опасных логов

# Быстрый анализ GC лога:
grep "Full GC" gc.log | wc -l        # количество Full GC
grep "Full GC" gc.log | awk '{print $NF}'  # длительности Full GC

# Allocation rate (по Minor GC):
grep "young" gc.log | awk '{
    match($0, /([0-9.]+)M\(.*\)->[0-9.]+M/, arr)
    sum += arr[1]
} END {print sum/NR " MB/GC"}'

# Средняя пауза:
grep "real=" gc.log | awk '{
    match($0, /real=([0-9.]+)/, arr)
    sum += arr[1]; count++
} END {print sum/count " сек avg pause"}'

Настройка G1GC для Spark

Базовые параметры

# Рекомендуемые JVM-флаги для Spark executor (добавляются через spark.executor.extraJavaOptions):

-XX:+UseG1GC                          # явно выбрать G1GC (default в Java 9+)
-XX:MaxGCPauseMillis=200              # целевая макс пауза STW (не гарантия, цель)
-XX:InitiatingHeapOccupancyPercent=35 # начинать Concurrent Mark при 35% заполнения heap
                                       # (по умолч. 45% - для Spark лучше пониже)
-XX:G1HeapRegionSize=16m              # размер региона (1–32 МБ)
                                       # увеличить если много Humongous Allocation в логах
-XX:G1NewSizePercent=20               # минимальный размер Young Gen (% heap)
-XX:G1MaxNewSizePercent=40            # максимальный размер Young Gen
-XX:ParallelGCThreads=4               # потоки для STW-фаз (≈ cores/4, но не меньше 4)
-XX:ConcGCThreads=2                   # потоки для concurrent-фаз
-XX:+G1UseAdaptiveIHOP                # адаптивный IHOP (Spark 3.x+, Java 9+)
spark = SparkSession.builder \
    .config("spark.executor.extraJavaOptions",
        "-XX:+UseG1GC "
        "-XX:MaxGCPauseMillis=200 "
        "-XX:InitiatingHeapOccupancyPercent=35 "
        "-XX:G1HeapRegionSize=16m "
        "-XX:G1NewSizePercent=20 "
        "-XX:G1MaxNewSizePercent=40 "
        "-XX:ParallelGCThreads=4 "
        "-XX:ConcGCThreads=2 "
        "-XX:+G1UseAdaptiveIHOP") \
    .getOrCreate()

Настройка G1HeapRegionSize под Spark

G1HeapRegionSize определяет, какие объекты становятся Humongous. Для Spark с большими shuffle буферами нужно увеличить:

Правило: G1HeapRegionSize = heap / 2048
Минимум 1 МБ, максимум 32 МБ

heap = 8 ГБ → G1HeapRegionSize = 4 МБ (по умолчанию)
Humongous порог = 2 МБ

Shuffle буфер по умолчанию = 32 МБ → Humongous!

Решение: -XX:G1HeapRegionSize=32m
Новый Humongous порог = 16 МБ
Shuffle буфер 32 МБ → всё ещё Humongous, но Parquet row groups часто < 16 МБ → OK

Или: уменьшить spark.shuffle.file.buffer:
spark.conf.set("spark.shuffle.file.buffer", "2m")  # вместо 32 МБ → не Humongous

InitiatingHeapOccupancyPercent (IHOP)

ZGC для latency-sensitive приложений

# Java 15+ / Spark 3.3+: ZGC - паузы < 10 мс вне зависимости от размера heap
spark = SparkSession.builder \
    .config("spark.executor.extraJavaOptions",
        "-XX:+UseZGC "
        "-XX:ZAllocationSpikeTolerance=2 "
        "-XX:+ZGenerational") \   # Java 21+: generational ZGC
    .getOrCreate()

# Предупреждение: ZGC требует больше native memory (RSS)
# container limit = executor.memory + memoryOverhead
# при ZGC увеличить memoryOverhead на 10-20%
spark.conf.set("spark.executor.memoryOverhead", "2g")  # было 1 ГБ

Off-Heap: выводим данные из-под GC

Предыдущий урок о Tungsten уже объяснял механику off-heap. С точки зрения GC pressure off-heap критически важен: данные в off-heap полностью невидимы для GC. GC не тратит на них ни секунды.

# Настройка off-heap:
spark = SparkSession.builder \
    .config("spark.memory.offHeap.enabled", "true") \
    .config("spark.memory.offHeap.size", "8g") \   # дополнительно к executor.memory
    .config("spark.executor.memory", "4g") \
    .config("spark.executor.memoryOverhead", "2g") \
    # Итого на контейнер: 4 ГБ heap + 2 ГБ overhead + 8 ГБ off-heap = 14 ГБ
    .getOrCreate()

# Кэширование в off-heap:
from pyspark import StorageLevel
df.persist(StorageLevel.OFF_HEAP)

# Теперь кэш не занимает heap → GC не видит его → паузы короче

Когда off-heap помогает и когда нет

Сценарий Off-heap помогает? Почему
GC Time > 20% в Task Metrics Да Убирает shuffle буферы из GC-scan
Большой кэш DataFrame (cache()) Да (с OFF_HEAP) Кэш вне heap → Old Gen свободнее
OOM в executor из-за User Memory Нет User Memory - в heap, off-heap не помогает
Full GC при allocation failure Частично Меньше объектов в heap, но не полное решение
Python UDF создаёт много объектов Нет Python worker отдельный процесс
Metaspace overflow Нет Metaspace ≠ heap, другая проблема

Предостережение: off-heap память учитывается в container RSS. В K8s/YARN при превышении лимита контейнер убивается OOMKiller на уровне OS (не JVM OOM, а SIGKILL). Это harder to debug чем JVM OOM.

# Расчёт для K8s:
# container_limit = executor.memory + memoryOverhead + offHeap.size
# 4 ГБ + 2 ГБ + 8 ГБ = 14 ГБ → задать в resources.limits.memory = 14Gi

Serialization и GC pressure

Java Serialization (по умолчанию для RDD) создаёт много промежуточных объектов при сериализации/десериализации каждой строки. Kryo значительно компактнее:

# Включить Kryo для RDD (не влияет на DataFrame, там UnsafeRow)
spark = SparkSession.builder \
    .config("spark.serializer", "org.apache.spark.serializer.KryoSerializer") \
    .config("spark.kryo.registrationRequired", "false") \
    .config("spark.kryoserializer.buffer.max", "512m") \
    .getOrCreate()

# Зарегистрировать часто используемые классы (для максимальной эффективности)
# spark.kryo.classesToRegister = "com.example.MyClass,..."

Но главный совет: использовать DataFrame API вместо RDD. DataFrame с Tungsten/CodeGen создаёт в сотни раз меньше объектов, что радикально снижает GC pressure.

GC и размер executor: антипаттерн "жирного executor"

Интуитивно кажется: дать одному executor 64 ГБ heap - меньше executor'ов, меньше сети, быстрее. На практике - GC катастрофа.

Правило: executor heap > 32 ГБ - антипаттерн для Spark. Оптимальный диапазон: 8–16 ГБ на executor. Больше executor'ов по 8 ГБ лучше, чем один executor на 64 ГБ.

Есть ещё одна причина: heap > ~32 ГБ выходит за порог Compressed OOPs (Ordinary Object Pointers). Выше этого предела JVM использует 8-байтные вместо 4-байтных указателей → +50% к размеру объектов в heap.

executor.memory < 31 ГБ → Compressed OOPs active → указатели 4 байта
executor.memory > 31 ГБ → Compressed OOPs disabled → указатели 8 байт
                          → всё занимает больше памяти → больше GC pressure

Metaspace: скрытая причина GC

Metaspace хранит bytecode загруженных классов. В Spark есть особый источник классов - Janino (Whole-Stage CodeGen): для каждого уникального Physical Plan генерируется и компилируется новый Java-класс.

Spark UI → Executor → stderr лог:
java.lang.OutOfMemoryError: Metaspace
    at java.lang.ClassLoader.defineClass1(Native Method)
    at org.apache.spark.sql.execution.WholeStageCodegenExec$$anon$1...
# Если видите Metaspace OOM:

# Способ 1: увеличить Metaspace
spark.conf.set("spark.executor.extraJavaOptions",
    "-XX:MaxMetaspaceSize=512m "      # по умолчанию unlimited, но OS ограничивает
    "-XX:MetaspaceSize=128m")          # начальный размер (избегаем resize overhead)

# Способ 2: ограничить сложность выражений (меньше уникальных планов)
spark.conf.set("spark.sql.codegen.wholeStage", "false")  # отключить WSCG для debugging

# Способ 3: уменьшить число уникальных Physical Plans
# → не создавать DataFrame динамически в цикле без кэширования

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

Шаг 1: Быстрая диагностика из Spark UI

# Запустить job и открыть:
# 1. Executors tab → смотреть GC Time column
# 2. Stages tab → любой Stage → Task Metrics → GC Time (25th, Median, Max)
# 3. Если GC Time > 10% Duration → переходим к GC логам

# Простой тест allocation rate:
from pyspark.sql.functions import rand
df = spark.range(100_000_000).withColumn("v", rand())
df.groupBy((df.id % 100).alias("bucket")).agg({"v": "sum"}).collect()

# Потом Spark UI → Stages → Stage с этим groupBy → GC Time

Шаг 2: Анализ GC логов

# Включить логирование (локально):
export SPARK_EXECUTOR_OPTS="-XX:+UseG1GC \
  -Xlog:gc*:file=/tmp/gc-%p.log:time,level,tags"

# Запустить job, потом анализировать:
grep "GC pause" /tmp/gc-*.log | grep -v "young" | head -20  # Non-minor GC

# Считаем allocation rate (МБ/сек):
python3 - << 'EOF'
import re
with open("/tmp/gc-1234.log") as f:
    content = f.read()
pauses = re.findall(r"GC pause.*?(\d+\.\d+) secs", content)
print(f"GC pauses: {len(pauses)}")
print(f"Total GC time: {sum(float(p) for p in pauses):.1f} secs")
print(f"Max pause: {max(float(p) for p in pauses):.3f} secs")
EOF

Шаг 3: Применение tuning и сравнение

# Baseline конфигурация (без tuning):
spark_base = SparkSession.builder \
    .config("spark.executor.memory", "8g") \
    .config("spark.executor.cores", "8") \
    .getOrCreate()

# Tuned конфигурация:
spark_tuned = SparkSession.builder \
    .config("spark.executor.memory", "8g") \
    .config("spark.executor.cores", "4") \        # меньше cores → меньше конкуренции за GC
    .config("spark.memory.offHeap.enabled", "true") \
    .config("spark.memory.offHeap.size", "4g") \
    .config("spark.executor.extraJavaOptions",
        "-XX:+UseG1GC "
        "-XX:MaxGCPauseMillis=200 "
        "-XX:InitiatingHeapOccupancyPercent=35 "
        "-XX:G1HeapRegionSize=16m") \
    .getOrCreate()

# Сравниваем GC Time в Spark UI между двумя прогонами

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

# Вопрос 1: Какой процент времени в GC?
# Spark UI → Executors → GC Time / Task Time
# < 5%:  норма
# 5-10%: стоит следить
# > 10%: оптимизировать
# > 20%: срочно

# Вопрос 2: Что в GC логах?
# Full GC → нехватка Old Generation
# Humongous Allocation → увеличить G1HeapRegionSize или уменьшить буферы
# Allocation Failure → Young Gen слишком мал или allocation rate слишком высок
# Metaspace OOM → увеличить MaxMetaspaceSize

# Вопрос 3: RDD или DataFrame?
# RDD + много объектов → перейти на DataFrame
# DataFrame + высокий GC → shuffle heavy workload → off-heap

# Вопрос 4: Большой heap?
# executor.memory > 32 ГБ → разбить на несколько executor'ов

# Вопрос 5: Много кэша в heap?
# Storage tab → много данных in-memory → использовать MEMORY_AND_DISK_SER или OFF_HEAP

Best Practices

1. Стандартная конфигурация G1GC для ETL:

gc_opts = (
    "-XX:+UseG1GC "
    "-XX:MaxGCPauseMillis=200 "
    "-XX:InitiatingHeapOccupancyPercent=35 "
    "-XX:G1HeapRegionSize=16m "
    "-XX:G1NewSizePercent=20 "
    "-XX:G1MaxNewSizePercent=40 "
    "-XX:+G1UseAdaptiveIHOP "
    "-XX:ParallelGCThreads=4 "
    "-XX:ConcGCThreads=2 "
    "-XX:+PrintGCDetails "
    "-XX:+PrintGCDateStamps "
    "-Xloggc:/tmp/gc-%p.log"
)

spark.conf.set("spark.executor.extraJavaOptions", gc_opts)

2. Off-heap при GC Time > 15%:

spark.conf.set("spark.memory.offHeap.enabled", "true")
spark.conf.set("spark.memory.offHeap.size", "4g")  # ≈ половину от executor.memory

3. Не создавать executor'ы > 16 ГБ heap:

# Лучше: 4 executor × 8 ГБ
spark.conf.set("spark.executor.memory", "8g")
spark.conf.set("spark.executor.instances", "4")

# Хуже: 1 executor × 32 ГБ
# spark.conf.set("spark.executor.memory", "32g")

4. DataFrame > RDD для минимизации объектов:

# Избегать rdd.map(lambda ...) в пользу df.withColumn(...)
# Избегать rdd.collect() → обрабатывать через df.write.*

5. Сериализация кэша:

from pyspark import StorageLevel
# Вместо MEMORY_ONLY → меньше heap, но CPU на десериализацию:
df.persist(StorageLevel.MEMORY_AND_DISK_SER)
# Без давления на GC:
df.persist(StorageLevel.OFF_HEAP)

Итого

Симптом Причина Решение
GC Time 20–50% Много short-lived объектов Использовать DataFrame, Kryo, off-heap
Full GC каждые 30–60 сек Old Generation переполнен Снизить IHOP, уменьшить кэш, off-heap
Humongous Allocation Объекты > G1HeapRegionSize/2 Увеличить G1HeapRegionSize=16m–32m
Executor Dead (heartbeat timeout) Full GC > network.timeout Срочно tuning GC, уменьшить heap
Metaspace OOM Много generated классов (CodeGen) -XX:MaxMetaspaceSize=512m
GC time >> 50% одного executor Skew + GC cumulative Salting + GC tuning + off-heap

Следующий урок рассматривает Py4J и Python Workers - как именно Python-код взаимодействует с JVM Spark, где теряется производительность при пересечении границы Python/JVM и как Arrow оптимизирует этот обмен.