Spark Config: тюнинг памяти, shuffle, AQE и PySpark-специфика
Иерархия конфигураций Spark, управление памятью Executor и Driver, shuffle.partitions, AQE, Arrow для PySpark, Parquet-настройки и чек-лист production-минимума.
Почему конфигурация определяет производительность¶
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")
Почему не ставить порог слишком высоко?
- Broadcast таблица хранится в памяти Driver и рассылается всем Executor
- При 100 Executor и 100 MB таблице: Driver отправит 100 × 100 MB = 10 GB по сети
- 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 в pandasspark.createDataFrame(pandas_df)- создание Spark DF из pandaspandas_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. Защита от случайного удаления данных |