GC Pressure: G1GC настройки, off-heap память и сигналы в логах
Почему Garbage Collector JVM тормозит Spark, как читать GC-логи, настраивать G1GC и использовать off-heap память для устранения Stop-the-World пауз
Почему 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 оптимизирует этот обмен.