Spark Config: тюнинг памяти, shuffle, AQE и PySpark-специфика

Иерархия конфигураций Spark, управление памятью Executor и Driver, shuffle.partitions, AQE, Arrow для PySpark, Parquet-настройки и чек-лист production-минимума.

core

Почему конфигурация определяет производительность

Spark - это фреймворк с огромным числом настраиваемых параметров. Дефолтные значения рассчитаны на универсальную работу, а не на конкретную задачу. На маленьких данных дефолты выглядят нормально. На 100 GB или 1 TB они превращаются в серьёзные проблемы:

  • spark.sql.shuffle.partitions = 200 при агрегации 1 GB данных означает 200 задач с 5 MB каждая - накладные расходы на создание задачи превысят полезную работу. А при 10 TB - 200 партиций по 50 GB каждая - не хватит памяти.
  • Без настройки Arrow toPandas() в 7 раз медленнее оптимального.
  • Без правильного memoryOverhead контейнеры YARN/Kubernetes убивают Python-воркеры при работе с большими данными.

Умение читать explain() (прошлый урок) говорит что Spark делает. Умение настраивать конфигурацию говорит как быстро он это делает.


Иерархия конфигураций: от глобального к локальному

Spark применяет конфигурации в строгом порядке приоритета - более специфичные переопределяют более общие:

spark-defaults.conf - файл на кластере, настраивается администратором. Содержит базовые конфиги для всех Job, запускаемых на кластере. В managed-среде (Databricks, EMR) обычно недоступен для изменений.

spark-submit --conf - параметры при запуске Job. Удобно для CI/CD: разные профили для dev/staging/prod.

spark-submit \
    --conf spark.executor.memory=4g \
    --conf spark.sql.shuffle.partitions=100 \
    --conf spark.sql.adaptive.enabled=true \
    my_job.py

SparkSession.builder.config() - в коде приложения. Применяется при создании сессии. Не может переопределить некоторые кластерные настройки (например, количество и размер исполнителей в managed-окружениях).

spark = SparkSession.builder \
    .appName("MyPipeline") \
    .config("spark.sql.shuffle.partitions", "100") \
    .config("spark.sql.adaptive.enabled", "true") \
    .config("spark.sql.execution.arrow.pyspark.enabled", "true") \
    .getOrCreate()

spark.conf.set() - изменение во время выполнения. Работает только для динамически изменяемых конфигов (большинство spark.sql.*, но не ресурсные конфиги):

# Можно менять в runtime
spark.conf.set("spark.sql.shuffle.partitions", "50")
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "100mb")

# Нельзя менять в runtime (только при инициализации)
# spark.conf.set("spark.executor.memory", "8g")  → ошибка!

SQL SET - только для SQL-сессии:

SET spark.sql.shuffle.partitions = 50;
SET spark.sql.adaptive.enabled = true;

Как проверить применившиеся конфигурации

Spark UI → вкладка Environment → раздел Spark Properties - показывает все конфиги с их источником. Всегда проверяйте здесь, особенно когда конфиги задаются на нескольких уровнях.

# Программно: получить значение конфига
print(spark.conf.get("spark.sql.shuffle.partitions"))   # "200"
print(spark.conf.get("spark.executor.memory"))

# Все Spark-конфиги текущей сессии
spark.sparkContext.getConf().getAll()

# Только SQL-конфиги
spark.sql("SET").show(100, truncate=False)

Два мира в одном процессе: JVM и Python

Прежде чем разбирать конкретные параметры - важно понять архитектуру PySpark. Это принципиально отличает PySpark от нативного Scala/Java Spark.

JVM (Java Virtual Machine): здесь живёт настоящий Spark - Catalyst Optimizer, Shuffle, чтение/запись файлов, HashAggregate, Join операторы. Управляется кучей JVM через -Xmx.

Python Worker: отдельный Python-процесс, запускаемый на каждом Executor при необходимости выполнить Python-код (UDF, pandas_udf, mapInPandas). Управляется отдельными конфигами памяти.

Ключевой момент: большинство операций PySpark выполняется в JVM. Когда вы пишете df.filter(col("amount") > 100).groupBy("user_id").sum() - весь этот код выполняется в JVM, Python только описывает план через Py4J. Python-процессы запускаются только для UDF и явных Python-операций.

Это означает: конфиги памяти JVM (spark.executor.memory) и Python (spark.executor.pyspark.memory) - разные вещи и настраиваются независимо.


Управление памятью: анатомия Executor

Память одного Executor делится на несколько зон:

Total Container Memory (YARN/K8s)
├── spark.executor.memory = 4g
│   ├── Reserved Memory = 300MB (системный резерв, захардкожен)
│   └── Usable Memory = 3.7g
│       ├── spark.memory.fraction = 0.6 (60%)
│       │   ├── spark.memory.storageFraction = 0.5 (50% от unified)
│       │   │   └── Storage: cache, broadcast = 1.11g
│       │   └── Execution: shuffle, sort, join, agg = 1.11g (динамически)
│       └── User Memory (40%) = 1.48g
│           └── UDF-данные, Spark internal structures
└── spark.executor.memoryOverhead = 384MB (минимум 10%)
    └── Off-heap JVM: native memory, thread stacks, NIO buffers
    └── Python Worker: если pyspark.memory не задан

Unified Memory (controlled by spark.memory.fraction) - главный ресурс для вычислений. Execution и Storage делят эту область динамически: если Storage не использует всю квоту - Execution может занять больше (и наоборот).

spark.memory.fraction (default: 0.6) - доля Usable Memory для Unified Memory. Увеличить если часто возникает spill на диск. Уменьшить если UDF или сторонние библиотеки требуют много User Memory.

spark.memory.storageFraction (default: 0.5) - доля Unified Memory, зарезервированная для Storage (кеш). Spark не вытеснит кешированные данные даже при нехватке Execution памяти, пока доля не превышена. Уменьшить если cache редко используется и нужно больше места для агрегаций.

# Агрессивная агрегация с большими join, мало кеша:
spark = SparkSession.builder \
    .config("spark.memory.fraction",        "0.7")   # больше unified
    .config("spark.memory.storageFraction", "0.3")   # меньше storage, больше execution
    .getOrCreate()

# Много кеша (iterative ML, повторные сканирования):
spark = SparkSession.builder \
    .config("spark.memory.fraction",        "0.6")
    .config("spark.memory.storageFraction", "0.7")   # больше storage
    .getOrCreate()

memoryOverhead: контейнерный overhead

spark.executor.memoryOverhead (default: max(executor_memory × 0.1, 384MB)) - дополнительная память сверх spark.executor.memory, которую запрашивает контейнер у YARN/Kubernetes. Включает:

  • Native off-heap память JVM (NIO буферы, thread stacks)
  • Память Python Worker (если spark.executor.pyspark.memory не задан явно)
  • Нативные библиотеки (PyArrow, NumPy, TensorFlow)

Типичная причина Container killed for exceeding memory limits: Python Worker потребляет больше памяти, чем отведено в overhead. Лечение:

# Вариант 1: Увеличить memoryOverhead
spark = SparkSession.builder \
    .config("spark.executor.memory",               "4g") \
    .config("spark.executor.memoryOverhead",       "1g")   # 25% от executor memory
    .getOrCreate()

# Вариант 2 (PySpark 3.x+): Явно задать память для Python-процесса
spark = SparkSession.builder \
    .config("spark.executor.memory",               "4g") \
    .config("spark.executor.pyspark.memory",       "512m") # для Python Worker
    .config("spark.executor.memoryOverhead",       "512m") # для JVM native
    .getOrCreate()

# Для Kubernetes:
spark = SparkSession.builder \
    .config("spark.kubernetes.memoryOverheadFactor", "0.25")  # 25%
    .getOrCreate()

Off-Heap память: когда она нужна

По умолчанию Spark хранит все данные в heap JVM. Off-heap (за пределами GC-управляемой кучи) полезен для:

  • Снижения давления на Garbage Collector при большом heap (> 8 GB)
  • Хранения данных без риска GC-пауз во время критических операций
spark = SparkSession.builder \
    .config("spark.memory.offHeap.enabled",  "true") \
    .config("spark.memory.offHeap.size",     "4g")   # дополнительные 4 GB off-heap
    .getOrCreate()

Off-heap не входит в spark.executor.memory - это дополнительная память. Общее потребление executor: executor.memory + offHeap.size + memoryOverhead.

Память Driver

Driver хранит: SparkContext, метаданные каталога, результаты collect(), планы всех DataFrame. Driver OOM обычно возникает при:

# ОПАСНО: collect() миллионов строк
huge_result = huge_df.collect()   # всё в память Driver

# ОПАСНО: broadcast очень большой таблицы
F.broadcast(big_df)   # broadcast данные хранятся в Driver

# ОПАСНО: toPandas() без limit
df.toPandas()   # всё собирается на Driver
spark = SparkSession.builder \
    .config("spark.driver.memory",         "4g")    # для сложных планов, метаданных
    .config("spark.driver.maxResultSize",  "1g")    # лимит на результаты collect()
    .getOrCreate()

spark.driver.maxResultSize = 1g - защита: если результат collect() превышает лимит, Spark бросит ошибку вместо OOM. Дефолт 1 GB.


spark.sql.shuffle.partitions: самый важный SQL-параметр

Число партиций после Shuffle - это spark.sql.shuffle.partitions (дефолт 200). Это значение влияет на каждый groupBy, join и window с Exchange.

Почему 200 - плохой дефолт для большинства задач:

Объём данных Идеальных партиций 200 партиций - проблема
100 MB 1–2 200 пустых/крошечных задач, overhead
10 GB ~80 200 × ~50 MB - приемлемо
100 GB ~800 200 × ~500 MB - риск OOM при агрегации
1 TB ~8000 200 × ~5 GB - точно OOM

Правило расчёта: целевой размер партиции после Shuffle - 100–200 MB.

# Формула:
# N = total_shuffle_data_bytes / target_partition_size_bytes
# Для 100 GB данных с целевым размером 128 MB:
# N = 100 * 1024 / 128 ≈ 800 партиций

spark.conf.set("spark.sql.shuffle.partitions", "800")

На практике total_shuffle_data заранее неизвестен. Стратегии:

# Стратегия 1: Зафиксировать под конкретный pipeline после профилирования
spark.conf.set("spark.sql.shuffle.partitions", "400")

# Стратегия 2: Использовать AQE (следующий раздел) для автоматики
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
spark.conf.set("spark.sql.adaptive.advisoryPartitionSizeInBytes", "128mb")

# Стратегия 3: Для dev/test - всегда уменьшать чтобы не гонять 196 пустых партиций
spark.conf.set("spark.sql.shuffle.partitions", "4")  # в тестах

Adaptive Query Execution (AQE): умный runtime-оптимизатор

AQE (включён по умолчанию в Spark 3.2+) позволяет Spark пересматривать физический план во время выполнения на основе реальной статистики. Три ключевых механизма:

Coalescing мелких партиций

После Shuffle многие партиции могут оказаться маленькими (особенно после агрессивной фильтрации). AQE автоматически объединяет их в более крупные:

spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled",     "true")
spark.conf.set("spark.sql.adaptive.advisoryPartitionSizeInBytes",   "128mb")  # целевой размер
spark.conf.set("spark.sql.adaptive.coalescePartitions.minPartitionNum", "1")  # минимум партиций

Как это работает: запрос начинается с 200 партициями, AQE видит что фактически данных всего 500 MB, объединяет в 4 партиции. Без AQE - 196 задач по 2 MB каждая.

Динамическое переключение join-стратегии

AQE может заменить SortMergeJoin на BroadcastHashJoin если после выполнения первой стадии выясняется, что одна из таблиц достаточно мала:

spark.conf.set("spark.sql.adaptive.autoBroadcastJoinThreshold", "30mb")
# AQE может broadcast таблицы до 30 MB, даже если статический threshold ниже

Обработка Data Skew

AQE автоматически разбивает перекошенные партиции на несколько при join:

spark.conf.set("spark.sql.adaptive.skewJoin.enabled",              "true")
spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionFactor", "5")   # партиция считается skewed если в N раз больше медианы
spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes", "256mb")

Полная конфигурация AQE для production:

spark = SparkSession.builder \
    .config("spark.sql.adaptive.enabled",                              "true") \
    .config("spark.sql.adaptive.coalescePartitions.enabled",           "true") \
    .config("spark.sql.adaptive.coalescePartitions.minPartitionNum",   "1") \
    .config("spark.sql.adaptive.advisoryPartitionSizeInBytes",         "134217728") \  # 128 MB
    .config("spark.sql.adaptive.skewJoin.enabled",                     "true") \
    .config("spark.sql.adaptive.skewJoin.skewedPartitionFactor",       "5") \
    .config("spark.sql.adaptive.autoBroadcastJoinThreshold",           "31457280") \  # 30 MB
    .getOrCreate()

Broadcast Join: настройка порога

spark.sql.autoBroadcastJoinThreshold (default: 10 MB) - если одна из сторон join меньше этого значения, Spark автоматически использует BroadcastHashJoin вместо SortMergeJoin. Нет Shuffle для большой таблицы - огромный выигрыш.

# Увеличить порог broadcast: если есть справочники до 100 MB
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "104857600")  # 100 MB

# Отключить автоматический broadcast (принудительно SMJ)
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1")

Почему не ставить порог слишком высоко?

  1. Broadcast таблица хранится в памяти Driver и рассылается всем Executor
  2. При 100 Executor и 100 MB таблице: Driver отправит 100 × 100 MB = 10 GB по сети
  3. Driver может упасть с OOM если broadcast таблица очень велика

spark.driver.maxResultSize защищает от OOM при collect, но не от broadcast.

Альтернатива: явный hint, когда автоматика не срабатывает:

# Явный hint без изменения глобального порога
orders.join(F.broadcast(small_reference), on="product_id", how="left")

Arrow: ускорение для PySpark

Apache Arrow - это ключевая оптимизация для операций, которые пересекают границу JVM-Python:

  • df.toPandas() - передача данных из Spark в pandas
  • spark.createDataFrame(pandas_df) - создание Spark DF из pandas
  • pandas_udf / mapInPandas - выполнение Python-функций над батчами
# Включить Arrow (по умолчанию выключен до PySpark 3.3)
spark.conf.set("spark.sql.execution.arrow.pyspark.enabled", "true")

# Безопасный fallback: если Arrow не может сериализовать тип - упасть в обычный режим
spark.conf.set("spark.sql.execution.arrow.pyspark.fallback.enabled", "true")

# Размер батча Arrow: сколько строк за раз передаётся в Python
spark.conf.set("spark.sql.execution.arrow.maxRecordsPerBatch", "10000")  # default

Насколько Arrow ускоряет? Бенчмарк для toPandas() на 10M строк:

Режим Время Механизм
Без Arrow ~85s pickle-сериализация по строкам через Py4J
С Arrow ~12s columnar binary transfer, zero-copy

Разница ~7x. Для pandas_udf Arrow даёт ~5x ускорение по сравнению с обычными UDF (данные не сериализуются построчно).

Ограничения Arrow: некоторые типы данных не поддерживаются Arrow-сериализацией (например, MapType с не-строковыми ключами). В таких случаях fallback.enabled=true позволяет автоматически вернуться к обычному режиму вместо ошибки.


Сериализация: Kryo вместо Java

Spark использует сериализацию при передаче объектов между Driver и Executor (broadcast переменные, closures), а также при Spill на диск.

По умолчанию - стандартная Java-сериализация: медленная и многословная.

Kryo - значительно быстрее (до 10x) и компактнее (до 2–3x меньший размер объекта):

spark = SparkSession.builder \
    .config("spark.serializer", "org.apache.spark.serializer.KryoSerializer") \
    .config("spark.kryo.unsafe",                "false")     # true - ещё быстрее, менее безопасно
    .config("spark.kryoserializer.buffer.max",  "512m")      # максимальный буфер для крупных объектов
    .getOrCreate()

Для кастомных классов нужна регистрация (иначе Kryo использует медленный fallback):

spark = SparkSession.builder \
    .config("spark.serializer",                        "org.apache.spark.serializer.KryoSerializer") \
    .config("spark.kryo.registrationRequired",         "false")    # true - ошибка если класс не зарегистрирован
    .config("spark.kryo.classesToRegister",
            "my.package.MyClass,my.package.AnotherClass") \
    .getOrCreate()

Kryo особенно важен при использовании RDD API и broadcast переменных с кастомными классами. Для чистого DataFrame/SQL API разница менее заметна - большинство операций выполняется в JVM без Python-сериализации.


Dynamic Allocation: автомасштабирование

Dynamic Allocation позволяет Spark автоматически запрашивать и освобождать Executor в зависимости от текущей нагрузки. Полезно для shared cluster где несколько Job работают одновременно.

spark = SparkSession.builder \
    .config("spark.dynamicAllocation.enabled",                    "true") \
    .config("spark.dynamicAllocation.minExecutors",               "2") \
    .config("spark.dynamicAllocation.maxExecutors",               "50") \
    .config("spark.dynamicAllocation.initialExecutors",           "5") \
    .config("spark.dynamicAllocation.executorIdleTimeout",        "60s") \   # убрать idle executor через 60s
    .config("spark.dynamicAllocation.cachedExecutorIdleTimeout",  "300s") \ # executor с кешем - через 5 min
    .config("spark.dynamicAllocation.schedulerBacklogTimeout",    "1s") \   # запросить новый executor если очередь > 1s
    .getOrCreate()

Важно для Dynamic Allocation: нужен External Shuffle Service - иначе при удалении Executor теряются Shuffle-файлы и Stage перезапускаются. На Kubernetes - Spark shuffle service или shuffle-aware node affinity.

Когда Dynamic Allocation лучше отключить:

  • Streaming Jobs - постоянная нагрузка, масштабирование не нужно
  • Jobs с тяжёлым кешем - частое добавление/удаление Executor инвалидирует кеш
  • Очень короткие Jobs - overhead на масштабирование больше полезного эффекта

Parquet и Data Lake: специфичные настройки

Размер партиций при чтении

spark.sql.files.maxPartitionBytes (default: 128 MB) - максимальный размер данных в одной Task при чтении файлов. Определяет число задач при чтении Parquet:

# Для больших файлов - увеличить чтобы уменьшить число задач
spark.conf.set("spark.sql.files.maxPartitionBytes", str(256 * 1024 * 1024))  # 256 MB

# Для маленьких файлов с высоким параллелизмом - уменьшить
spark.conf.set("spark.sql.files.maxPartitionBytes", str(64 * 1024 * 1024))   # 64 MB

spark.sql.files.openCostInBytes (default: 4 MB) - Spark суммирует файлы в одну Task пока их общий размер < maxPartitionBytes. Если файлы очень мелкие, увеличьте openCostInBytes для лучшей упаковки:

spark.conf.set("spark.sql.files.openCostInBytes", str(8 * 1024 * 1024))  # 8 MB

Partition Overwrite Mode: не затирать лишние партиции

Критически важная настройка для инкрементальной загрузки в партиционированные таблицы:

# STATIC (default): перезаписывает ВСЕ партиции - даже те, которых нет в новых данных!
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "static")
df.write.mode("overwrite").partitionBy("date").parquet("s3://silver/events/")
# ОПАСНО: если в df только date=2026-05-16, все другие даты стираются!

# DYNAMIC: перезаписывает ТОЛЬКО партиции, которые есть в новых данных
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")
df.write.mode("overwrite").partitionBy("date").parquet("s3://silver/events/")
# БЕЗОПАСНО: date=2026-05-15 и другие даты остаются нетронутыми

Это одна из самых распространённых причин потери данных в Spark ETL! static режим по умолчанию + невнимательность = стёртая история.

Сжатие Parquet

# Кодек сжатия для Parquet (default: snappy)
spark.conf.set("spark.sql.parquet.compression.codec", "snappy")   # баланс скорость/размер
spark.conf.set("spark.sql.parquet.compression.codec", "zstd")     # лучшее сжатие, CPU дороже
spark.conf.set("spark.sql.parquet.compression.codec", "lz4")      # быстрее snappy, хуже сжатие
spark.conf.set("spark.sql.parquet.compression.codec", "gzip")     # лучшее сжатие, медленно
spark.conf.set("spark.sql.parquet.compression.codec", "uncompressed")  # без сжатия

# Для CPU-ограниченных Job (много данных, мало CPU):
spark.conf.set("spark.sql.parquet.compression.codec", "snappy")

# Для I/O-ограниченных Job (медленные диски/S3, много CPU):
spark.conf.set("spark.sql.parquet.compression.codec", "zstd")
spark.conf.set("spark.io.compression.zstd.level",     "3")   # уровень сжатия 1-22

Практика выбора: для горячих данных (часто читаются) - snappy или lz4. Для архивных данных (редко читаются, важен размер) - zstd или gzip.

Schema merge для Parquet (осторожно!)

# Автоматически объединять схемы файлов с разными схемами (очень медленно!)
spark.conf.set("spark.sql.parquet.mergeSchema", "false")  # default и рекомендуется

# При чтении конкретного датасета с эволюцией схемы:
df = spark.read \
    .option("mergeSchema", "true") \  # только для конкретного чтения
    .parquet("s3://bucket/data_with_evolving_schema/")

mergeSchema=true вызывает дополнительный pass по всем файлам для чтения схем - очень дорого на больших датасетах. Включайте только явно и только когда нужно.


Timezone и время: часто забываемый параметр

# Timezone для интерпретации timestamp-значений (default: JVM timezone)
spark.conf.set("spark.sql.session.timeZone", "UTC")

# КРИТИЧЕСКИ ВАЖНО для production: всегда UTC!
# Иначе данные разных сессий с разными timezone несовместимы
spark = SparkSession.builder \
    .config("spark.sql.session.timeZone", "UTC") \
    .getOrCreate()

Без явного UTC разные пользователи (разные часовые пояса) или разные серверы (кластер в EU, клиент в US) могут получать разные результаты при работе с timestamp.


Speculative Execution: защита от медленных задач

Speculative execution - Spark запускает дублирующую копию медленной задачи на другом Executor. Первая завершившаяся копия принимается, вторая отменяется. Защита от straggler-задач (медленных нод, GC-пауз, сетевых проблем):

spark.conf.set("spark.speculation",              "true")
spark.conf.set("spark.speculation.interval",     "100ms")  # как часто проверять
spark.conf.set("spark.speculation.multiplier",   "1.5")    # запустить копию если задача в 1.5x медленнее медианы
spark.conf.set("spark.speculation.quantile",     "0.75")   # когда 75% задач Stage завершились

Осторожно: speculative execution не подходит для задач с side effects (запись в БД, Kafka). Если задача записывает данные дважды - это проблема. Для идемпотентных операций (перезапись Parquet по partition) - безопасно.


Executor sizing: стратегия выбора конфигурации

Нет универсально "правильного" размера Executor. Компромисс между параллелизмом, памятью и накладными расходами:

Правило 5 ядер на Executor: эмпирически найдено, что 5 CPU на Executor даёт хороший баланс. При большем числе ядер GC останавливает все задачи одновременно (Stop-the-World GC на 16-ядерном Executor - дорогая пауза). При меньшем - слишком много JVM instances.

# Рекомендуемая конфигурация для 40-ядерной ноды:
# 8 Executor × 5 cores (1 ядро на OS/Daemon)
# memory: (node_memory - OS_reserve) / 8 Executor

spark-submit \
    --num-executors 8 \
    --executor-cores 5 \
    --executor-memory 19g \   # (160GB - 8GB OS) / 8 = ~19 GB
    --driver-memory 4g \
    my_job.py

Диагностика через Spark UI

Spark UI - главный инструмент для проверки эффективности конфигурации:

Вкладка Environment → Spark Properties: все активные конфиги. Проверьте, что ваши настройки применились.

Вкладка Stages → задачи:

  • Ровные времена выполнения задач → хороший баланс, нет skew
  • Одна задача в 10x медленнее остальных → Data Skew
  • GC Time > 20% от Task Time → GC pressure, нужно больше памяти или off-heap

Вкладка SQL → Jobs:

  • Exchange с огромным Shuffle Write → много партиций или неэффективная агрегация
  • InMemoryTableScan → кеш работает
  • BroadcastHashJoin → broadcast сработал

Вкладка Executors:

  • Spill (Memory) > 0 → нехватка Execution Memory
  • Spill (Disk) > 0 → нехватка Execution Memory, данные льются на диск
  • GC Time % → если высокий, нужно больше heap или off-heap

Чек-лист: production-минимум конфигурации

Набор конфигураций, которые должны быть в любом enterprise PySpark-пайплайне:

spark = SparkSession.builder \
    .appName("MyProductionPipeline") \

    # ─── Ресурсы ───
    .config("spark.executor.memory",              "8g") \
    .config("spark.executor.cores",               "5") \
    .config("spark.executor.memoryOverhead",      "2g") \
    .config("spark.driver.memory",               "4g") \
    .config("spark.driver.maxResultSize",        "2g") \

    # ─── SQL оптимизация ───
    .config("spark.sql.shuffle.partitions",       "200") \  # настроить под объём!
    .config("spark.sql.adaptive.enabled",         "true") \
    .config("spark.sql.adaptive.coalescePartitions.enabled",  "true") \
    .config("spark.sql.adaptive.skewJoin.enabled","true") \
    .config("spark.sql.adaptive.advisoryPartitionSizeInBytes", "134217728") \  # 128 MB

    # ─── Broadcast ───
    .config("spark.sql.autoBroadcastJoinThreshold", "31457280") \  # 30 MB

    # ─── PySpark / Arrow ───
    .config("spark.sql.execution.arrow.pyspark.enabled",          "true") \
    .config("spark.sql.execution.arrow.pyspark.fallback.enabled", "true") \

    # ─── Parquet ───
    .config("spark.sql.parquet.compression.codec",          "snappy") \
    .config("spark.sql.sources.partitionOverwriteMode",     "dynamic") \
    .config("spark.sql.parquet.mergeSchema",                "false") \

    # ─── Timezone (критично!) ───
    .config("spark.sql.session.timeZone", "UTC") \

    # ─── Сериализация ───
    .config("spark.serializer", "org.apache.spark.serializer.KryoSerializer") \

    .getOrCreate()

Антипаттерны конфигурации

1. Слепое копирование конфигов из Stack Overflow. Каждый пайплайн уникален. Конфиги под 100 GB join не подходят для 1 GB lookup. Всегда профилируйте через Spark UI.

2. Гигантские Executor (64 ядра, 256 GB). GC Stop-the-World на 256 GB heap - это секунды паузы. Все задачи Executor ждут. Лучше много средних Executor.

3. shuffle.partitions = 200 для всех задач. Для маленьких данных - 196 пустых задач. Для больших - задачи по 5 GB с OOM. Используйте AQE или настраивайте явно под объём.

4. partitionOverwriteMode = static для инкрементальной загрузки. Потеря исторических данных. Всегда dynamic для партиционированных таблиц при инкрементальной загрузке.

5. Не задать UTC timezone. Расчёты с timestamp дают разные результаты в разных средах. Всегда spark.sql.session.timeZone = UTC.

6. broadcast порог слишком высок без контроля Driver. Aggressive broadcast при autoBroadcastJoinThreshold = 1g может убить Driver OOM при 100 Executor.


Справочник конфигурационных параметров

Полные таблицы всех параметров для быстрого поиска при отладке и настройке продакшн-пайплайнов.


PySpark / Spark Core

Параметры управляют памятью JVM и Python-процессов, сериализацией, сетью, отказоустойчивостью и динамическим выделением ресурсов.

Память Driver и Executor

Параметр Дефолт Описание
spark.driver.memory 1g Память Driver process. Нужна для .collect() и построения графа задач
spark.driver.cores 1 Количество CPU для Driver
spark.driver.maxResultSize 1g Максимальный объём данных, возвращаемых в Driver. 0 = без ограничений
spark.executor.instances - Фиксированное количество Executor (при отключённом dynamic allocation)
spark.executor.memory - Память JVM-кучи на один Executor. Здесь живут кэш и внутренние структуры Spark
spark.executor.cores 1 (YARN) CPU cores на один Executor
spark.executor.memoryOverhead max(executor_memory × 0.1, 384MB) Дополнительная память поверх heap (JVM overhead, NIO буферы). Защищает контейнер от OOM по лимитам ОС
spark.executor.pyspark.memory 0 (нет лимита) Лимит физической памяти для Python-воркера на Executor. При превышении процесс убивается
spark.executor.extraJavaOptions - Дополнительные JVM-флаги для Executor. Сюда прописывают GC-настройки: -XX:+UseG1GC -XX:InitiatingHeapOccupancyPercent=35
spark.executor.heartbeatInterval 10s Интервал heartbeat Executor → Driver. Должен быть существенно меньше spark.network.timeout

Unified Memory Manager

Параметр Дефолт Описание
spark.memory.fraction 0.6 Доля heap под Execution + Storage memory. Остальные 40% - пользовательские объекты и внутренние нужды Spark
spark.memory.storageFraction 0.5 Доля внутри memory.fraction, защищённая от вытеснения кэша операциями shuffle/join
spark.memory.offHeap.enabled false Включение off-heap памяти вне JVM GC. Снижает нагрузку на Garbage Collector при тяжёлых sort/join
spark.memory.offHeap.size 0 Размер off-heap памяти в байтах. Требует offHeap.enabled=true
spark.cleaner.periodicGC.interval 30min Интервал принудительного System.gc() на Driver. Предотвращает утечки при долгоживущих сессиях

Сериализация и производительность

Параметр Дефолт Описание
spark.serializer JavaSerializer Механизм сериализации для shuffle и кэша. KryoSerializer работает в 2–10× быстрее
spark.kryoserializer.buffer.max 64m Максимальный буфер Kryo при сериализации. Увеличить при ошибке Buffer overflow
spark.default.parallelism all cores Базовое количество partitions для RDD-операций (не DataFrame). Влияет на sc.parallelize()

Сеть и отказоустойчивость

Параметр Дефолт Описание
spark.network.timeout 120s Таймаут сетевых взаимодействий между узлами. При тяжёлом shuffle поднимают до 300–600s
spark.task.maxFailures 4 Максимальное число retry задачи перед отменой всего Stage
spark.locality.wait 3s Сколько ждать data locality перед тем, как отправить задачу на другой узел
spark.speculation false Запуск дублирующих копий медленных задач (speculative execution)
spark.speculation.multiplier 1.5 Задача считается медленной, если её время > медиана × multiplier

Dynamic Allocation

Параметр Дефолт Описание
spark.dynamicAllocation.enabled false Динамическое выделение Executor - Spark запрашивает воркеры по мере необходимости
spark.dynamicAllocation.minExecutors 0 Минимальное количество Executor при dynamic allocation
spark.dynamicAllocation.maxExecutors Максимальное количество Executor при dynamic allocation
spark.dynamicAllocation.initialExecutors = minExecutors Стартовое количество Executor
spark.shuffle.service.enabled false External Shuffle Service. Обязателен при dynamicAllocation.enabled=true

Логирование и Catalog

Параметр Дефолт Описание
spark.app.name - Имя Spark application в Spark UI и логах
spark.master - Тип cluster manager: local[*], yarn, k8s://..., spark://...
spark.submit.deployMode client Где запускается Driver: client (локально) или cluster (на кластере)
spark.sql.warehouse.dir $PWD/spark-warehouse Путь для managed Spark SQL tables (Hive warehouse)
spark.eventLog.enabled false Включение записи event logs для Spark History Server
spark.eventLog.dir - Путь хранения event logs (HDFS или S3)
spark.history.fs.logDirectory - Директория, которую читает Spark History Server

S3 / Cloud Object Storage

Параметр Дефолт Описание
spark.hadoop.fs.s3a.fast.upload true Буферизация загрузки на S3 в памяти. Радикально ускоряет финальный .write.save()
spark.hadoop.fs.s3a.* - Группа настроек подключения к S3-compatible storage (endpoint, credentials, buffer size)
spark.hadoop.mapreduce.fileoutputcommitter.algorithm.version 1 Алгоритм записи output files. Версия 2 ускоряет commit на S3, но менее безопасна при сбоях

Spark SQL

Shuffle и параллелизм

Параметр Дефолт Описание
spark.sql.shuffle.partitions 200 Количество партиций после shuffle (join/groupBy/window). Дефолт 200 - главная причина тормозов: на малых данных пустые таски, на больших - OOM
spark.sql.files.maxPartitionBytes 128MB Максимальный размер данных в одну партицию при чтении файлов. Регулирует начальный параллелизм джобы
spark.sql.execution.sortBeforeRepartition true Сортировка данных перед repartition() для улучшения locality при записи

Adaptive Query Execution (AQE)

Параметр Дефолт Описание
spark.sql.adaptive.enabled true (Spark 3.2+) Включение AQE - Spark перестраивает физический план во время выполнения на основе реальной статистики
spark.sql.adaptive.coalescePartitions.enabled true Автообъединение мелких shuffle-партиций. Устраняет оверхед на тысячи пустых задач
spark.sql.adaptive.advisoryPartitionSizeInBytes 64MB Целевой размер партиции при AQE coalesce. На мощных кластерах поднимают до 128–256 MB
spark.sql.adaptive.maxShuffledPartitions 500 Верхний лимит партиций для AQE. При терабайтных объёмах нужно увеличить
spark.sql.adaptive.skewJoin.enabled true Автоматическое разбиение перекошенных (skewed) партиций в join
spark.sql.adaptive.skewJoin.skewedPartitionFactor 5 Партиция считается перекошенной, если её размер > медиана × factor
spark.sql.adaptive.localShuffleReader.enabled true Локальное чтение shuffle-партиций без сетевого запроса, когда данные уже на том же узле

Оптимизатор и джоины

Параметр Дефолт Описание
spark.sql.autoBroadcastJoinThreshold 10MB Максимальный размер таблицы для автоматического BroadcastHashJoin. -1 отключает авто-broadcast
spark.sql.broadcastTimeout 300s Таймаут ожидания broadcast. При медленной сети или большой таблице - увеличить
spark.sql.join.preferSortMergeJoin true Предпочитать SortMergeJoin вместо ShuffleHashJoin
spark.sql.cbo.enabled false Cost-Based Optimizer - использует статистику таблиц для выбора оптимального плана
spark.sql.statistics.histogram.enabled false Сбор гистограмм столбцов для CBO. Включать вместе с cbo.enabled
spark.sql.optimizer.dynamicPartitionPruning.enabled true Dynamic Partition Pruning - отсечение партиций на основе результатов sub-query

Чтение файлов: Parquet и ORC

Параметр Дефолт Описание
spark.sql.parquet.filterPushdown true Predicate Pushdown для Parquet - читаются только строки, удовлетворяющие фильтру
spark.sql.parquet.enableVectorizedReader true Векторизованное чтение Parquet блоками (batches). Ускоряет десериализацию в 3–5×
spark.sql.parquet.mergeSchema false Автоматический merge schema при чтении нескольких Parquet-файлов. Замедляет чтение - держать false
spark.sql.orc.filterPushdown true Predicate Pushdown для ORC-файлов
spark.sql.hive.metastorePartitionPruning true Отсечение партиций на уровне Hive Metastore/Glue - не сканируются директории, не подходящие под WHERE
spark.sql.files.ignoreCorruptFiles false Игнорировать повреждённые файлы при чтении
spark.sql.files.ignoreMissingFiles false Игнорировать отсутствующие файлы (например, удалённые между планированием и выполнением)

Запись и перезапись

Параметр Дефолт Описание
spark.sql.sources.partitionOverwriteMode static Режим перезаписи партиций. static удаляет все партиции, dynamic - только те, что в записываемом DataFrame. Критично для инкрементальной загрузки
spark.sql.parquet.compression.codec snappy Кодек сжатия Parquet: snappy (быстро), zstd (баланс), gzip (максимальное сжатие)

PySpark / Arrow / Pandas

Параметр Дефолт Описание
spark.sql.execution.arrow.pyspark.enabled false Arrow-оптимизация передачи данных между JVM и Python. Обязательно включить - даёт ~7× ускорение
spark.sql.execution.arrow.maxRecordsPerBatch 10000 Размер Arrow-батча в записях. Увеличить при широких схемах (много колонок)
spark.sql.execution.pandas.convertToArrowArraySafely false Безопасное преобразование типов при Arrow-конвертации. Бросает исключение вместо тихого truncation

SQL-семантика и Catalog

Параметр Дефолт Описание
spark.sql.caseSensitive false Case-sensitive режим SQL-парсера и column resolution
spark.sql.session.timeZone JVM default Timezone сессии для timestamp операций. Всегда устанавливать UTC
spark.sql.ansi.enabled false ANSI SQL: строгая типизация, запрет implicit cast, деление на ноль - исключение
spark.sql.catalogImplementation in-memory Тип Catalog: in-memory (временный) или hive (постоянный через Hive Metastore)
spark.sql.hive.metastore.version - Версия Hive Metastore для совместимости
spark.sql.legacy.timeParserPolicy EXCEPTION Политика парсинга legacy datetime форматов (LEGACY, CORRECTED, EXCEPTION)
spark.sql.debug.maxToStringFields 25 Количество полей схемы в explain и debug-выводе. Увеличить при широких схемах

Кэш в памяти

Параметр Дефолт Описание
spark.sql.inMemoryColumnarStorage.compressed true Сжатие колоночного кэша (.cache()) в памяти
spark.sql.inMemoryColumnarStorage.batchSize 10000 Размер батча колоночного хранилища кэша

Structured Streaming

Параметр Дефолт Описание
spark.sql.streaming.checkpointLocation - Checkpoint path для гарантий Exactly-Once в Structured Streaming
spark.sql.streaming.stateStore.providerClass HDFSBackedStateStoreProvider Backend для хранения состояний stateful операций. RocksDBStateStoreProvider выносит state на локальный диск, устраняя OOM
spark.sql.streaming.minBatchesToRetain 100 Количество хранимых метаданных прошлых батчей. Уменьшить при OOM Driver на долгоживущих стримах
spark.sql.streaming.forceDeleteTempCheckpointLocation false Автоочистка временных checkpoint-метаданных. Помогает избежать замусоривания дисков
spark.sql.streaming.schemaInference false Разрешить schema inference в streaming source (не рекомендуется в продакшне)

Расширения и Delta Lake

Параметр Дефолт Описание
spark.sql.extensions - Подключение Spark SQL extensions, например Delta Lake: io.delta.sql.DeltaSparkSessionExtension
spark.sql.catalog.spark_catalog - Замена default catalog: org.apache.spark.sql.delta.catalog.DeltaCatalog для Delta Lake
spark.databricks.delta.optimizeWrite.enabled false Auto-optimize write для Delta Lake (только Databricks Runtime)
spark.databricks.delta.autoCompact.enabled false Автоматический compaction Delta-файлов после записи (только Databricks Runtime)
spark.databricks.delta.schema.autoMerge.enabled false Auto schema evolution при записи в Delta Lake (аналог mergeSchema=true)
spark.databricks.delta.retentionDurationCheck.enabled true Проверка retention policy при VACUUM. Защита от случайного удаления данных