Hardware Tuning: spark.local.dir на NVMe и memoryOverhead для Python
Физика дисковых операций Spark и архитектура памяти PySpark: роль локального диска в Shuffle и Spill, влияние NVMe vs HDD на производительность, многодисковый round-robin, 4-компонентная память PySpark-контейнера, Exit Code 137 и OOMKilled, spark.executor.pyspark.memory, Python UDF vs Pandas UDF, Kubernetes emptyDir и настройка в YARN, диагностика через Spark UI и логи.
1. Архитектура дискового ввода-вывода в Spark: роль локального диска¶
Apache Spark позиционируется как движок для вычислений в оперативной памяти (In-Memory Processing). Но это только часть правды. На практике Spark активно использует локальный диск Executor'а для двух критических операций, и именно производительность этого диска часто определяет общую скорость работы.
Два ключевых сценария использования диска¶
Сценарий 1: Shuffle Write/Read. При каждой Shuffle-операции (groupBy, join, repartition) каждый Executor записывает свои выходные партиции на локальный диск как «Shuffle Files». Затем другие Executor'ы читают нужные им Shuffle Partitions с этого диска по сети.
Для крупного ETL с 10 TB данных объём Shuffle Write может достигать 2-5 TB суммарно на весь кластер. Если диски медленные — каждый Stage ждёт завершения I/O.
Сценарий 2: Disk Spill. Когда Execution Memory переполнена (при сортировке, хэш-агрегации, Join), Spark сбрасывает излишки данных на диск в виде Sorted Runs. Это называется Disk Spill. После обработки эти данные читаются обратно и мержатся.
Схема показывает где именно происходят дисковые операции и почему выбор spark.local.dir критически важен. По умолчанию /tmp часто размещается на том же диске что и ОС — медленном HDD или tmpfs с ограниченным объёмом.
I/O Wait: как медленный диск парализует CPU¶
Когда диск не успевает обрабатывать Shuffle Write, процессоры Executor'а уходят в состояние I/O Wait — ожидание завершения дисковой операции. В этом состоянии ядра CPU не делают полезной работы, но и не освобождают слот задачи.
# Наблюдение I/O Wait на Executor-узле (в реальном времени):
iostat -x 1
# Device r/s w/s rMB/s wMB/s await util%
# sda 12 8534 0.1 812.4 345 100% ← диск загружен на 100%!
# await=345ms означает: CPU ждёт диск по 345ms на каждую I/O операцию
# top/htop: смотреть колонку %wa (I/O Wait)
# Нормально: %wa < 5%
# Проблема: %wa > 20% — диск стал bottleneck
2. Конфигурация spark.local.dir и её влияние на производительность¶
spark.local.dir — параметр, который задаёт список директорий на локальных узлах Executor'ов, куда Spark записывает все временные файлы.
Дефолтное поведение и его проблемы¶
# Дефолт: /tmp (часто системный диск)
# Проблемы:
# 1. Системный диск = HDD → медленно
# 2. /tmp часто tmpfs (RAM-диск) → ограниченный объём
# 3. /tmp делится с ОС, логами и другими процессами → конкуренция
# Проверить текущее значение:
current = spark.conf.get("spark.local.dir")
print(f"Текущий spark.local.dir: {current}")
Правильная настройка: выделенные быстрые диски¶
# ── Один NVMe диск ────────────────────────────────────────────────────
spark = SparkSession.builder \
.config("spark.local.dir", "/mnt/nvme0/spark-local") \
.getOrCreate()
# ── Несколько NVMe дисков: Round-Robin striping ───────────────────────
# Spark чередует запись между дисками последовательно:
# файл 1 → /mnt/nvme0, файл 2 → /mnt/nvme1, файл 3 → /mnt/nvme2, ...
# Суммарная пропускная способность: ~3 × пропускная способность одного диска!
spark = SparkSession.builder \
.config("spark.local.dir",
"/mnt/nvme0/spark-local,"
"/mnt/nvme1/spark-local,"
"/mnt/nvme2/spark-local") \
.getOrCreate()
# ── Через spark-submit ────────────────────────────────────────────────
# spark-submit \
# --conf spark.local.dir=/mnt/nvme0/spark,/mnt/nvme1/spark \
# your_job.py
# ── Через spark-defaults.conf (для всех приложений на кластере) ───────
# В $SPARK_HOME/conf/spark-defaults.conf:
# spark.local.dir /mnt/nvme0/spark,/mnt/nvme1/spark,/mnt/nvme2/spark
# ── ВАЖНО: директории должны существовать! ─────────────────────────────
# Spark НЕ создаёт их автоматически
import subprocess
for disk in ["/mnt/nvme0/spark-local", "/mnt/nvme1/spark-local"]:
subprocess.run(["mkdir", "-p", disk])
subprocess.run(["chmod", "777", disk])
Диагностика: куда реально пишет Spark¶
# 1. Через логи Executor'а
grep "spark.local.dir\|shuffle\|local dir" /var/log/spark/executor.log | head -20
# 2. Через Spark UI (порт 4040)
# Stages → конкретный Stage → Tasks → "Shuffle Write Location"
# Показывает путь файла shuffle для каждой Task'и
# 3. Смотреть файловую систему во время выполнения
ls -la /mnt/nvme0/spark-local/
# blockmgr-xxx/ ← блоки shuffle
# spark-xxx/ ← временные файлы задач
# 4. Объём shuffle файлов
du -sh /mnt/nvme0/spark-local/
3. Эволюция дисков: перенос spark.local.dir на NVMe/SSD¶
Разрыв в производительности между типами дисков огромен. Для Shuffle-интенсивных операций это напрямую определяет время выполнения Stage.
Сравнение производительности дисков для Spark Shuffle¶
Облачные инстансы: Local NVMe vs сетевые диски¶
# ── Рекомендации для облачных провайдеров ─────────────────────────────
# AWS:
# ✅ Рекомендуются инстансы с Local NVMe: i3, i4i, d3, d3en, r6i.nvme
# ❌ НЕ использовать EBS (gp2/gp3) для spark.local.dir!
# EBS gp3 лимиты: 1000 MB/s, 16000 IOPS — исчерпываются при массовом shuffle
# Local NVMe i3.8xlarge: 4× NVMe, суммарно ~8+ GB/s
# GCP:
# ✅ Инстансы с Local SSD (375 GB NVMe per disk)
# Можно монтировать несколько Local SSD
# Yandex Cloud:
# ✅ Инстансы с NVMe-дисками типа network-ssd-nonreplicated (до 20k IOPS/TB)
# Лучше: выделенные ноды с Direct-attached NVMe
# Пример конфигурации для AWS i3.4xlarge (2× NVMe):
SPARK_CONF = {
"spark.local.dir": "/mnt/nvme0,/mnt/nvme1",
"spark.executor.cores": "8",
"spark.executor.memory": "30g",
"spark.executor.memoryOverhead": "5g",
}
Настройка диска в продакшн: автоматическое монтирование¶
#!/bin/bash
# setup_spark_nvme.sh
# Запускается при старте каждого Executor'а (через init script)
# Находим доступные NVMe диски
NVME_DISKS=$(ls /dev/nvme*n1 2>/dev/null)
i=0
SPARK_DIRS=""
for disk in $NVME_DISKS; do
MOUNT_POINT="/mnt/nvme${i}/spark-local"
# Форматируем и монтируем (если не смонтирован)
if ! mountpoint -q /mnt/nvme${i}; then
mkfs.xfs -f $disk
mkdir -p /mnt/nvme${i}
mount -o noatime,nodiratime $disk /mnt/nvme${i}
fi
mkdir -p $MOUNT_POINT
chmod 777 $MOUNT_POINT
if [ -n "$SPARK_DIRS" ]; then
SPARK_DIRS="${SPARK_DIRS},"
fi
SPARK_DIRS="${SPARK_DIRS}${MOUNT_POINT}"
i=$((i+1))
done
echo "spark.local.dir = $SPARK_DIRS"
# Передаём в spark-submit через --conf или в spark-defaults.conf
4. Анатомия памяти PySpark-контейнера: JVM против Python¶
При работе с PySpark архитектура памяти значительно сложнее, чем кажется на первый взгляд. На каждом Executor'е работают два принципиально разных процесса с независимым управлением памятью.
Четыре компонента памяти PySpark контейнера¶
Схема показывает полную картину памяти PySpark-контейнера. Критически важно: Python Workers живут в memoryOverhead, а не в executor.memory. Если не выделить достаточно Overhead — контейнер будет убит операционной системой несмотря на то что JVM Heap ещё свободен.
Что происходит при нехватке памяти¶
Сценарий: spark.executor.memory=8g, spark.executor.memoryOverhead=384m (дефолт)
Pod Limit: 8g + 384m + 300m = ~8.7 GB
Во время выполнения Pandas UDF:
JVM Heap: 6.2 GB (использовано)
Python Worker: 2.5 GB (pandas DataFrame в памяти)
JVM Native: 0.3 GB
ОС: 0.3 GB
Итого: 9.3 GB >> Pod Limit 8.7 GB
Результат: OOMKilled (Exit Code 137)
К8s: Pod убит. Spark: Task пересчитывается 3 раза, Stage падает.
5. Параметр memoryOverhead: защита от OOM Killer¶
spark.executor.memoryOverhead — это память вне JVM Heap, которую YARN/Kubernetes выделяет контейнеру Executor'а сверх spark.executor.memory.
Дефолтное значение и его проблема¶
# Дефолт: max(executorMemory × 0.1, 384 MB)
# Для spark.executor.memory = 8g:
# overhead = max(8192 × 0.1, 384) = max(819, 384) = 819 MB
# Проблема: 819 MB катастрофически мало для PySpark с Pandas/NumPy!
# pandas.read_parquet() на 500 MB файле = ~2 GB в RAM Pandas DataFrame
# + numpy операции добавляют временные массивы
# + импорт тяжёлых библиотек (scikit-learn, torch) = 1-2 GB
# Итого Python Worker легко занимает 3-5 GB
Как рассчитать правильный overhead¶
def calculate_pyspark_overhead(
executor_memory_gb: float,
num_cores: int,
use_pandas_udf: bool = False,
use_ml_libraries: bool = False,
typical_pandas_df_size_mb: float = 0,
) -> dict:
"""
Рассчитывает рекомендуемый spark.executor.memoryOverhead.
Учитывает:
- JVM Native Memory (Netty, JIT code cache, JVM metadata)
- Python Worker Memory (отдельный Python процесс на каждый core)
- Apache Arrow transfer buffers (при Pandas UDF)
- ML библиотеки (scikit-learn, tensorflow, torch)
"""
# JVM Native overhead (Netty, code cache, JVM metadata)
jvm_native_gb = max(0.5, executor_memory_gb * 0.05)
# Python Worker overhead
if not use_pandas_udf and not use_ml_libraries:
# Простые Python UDF: небольшой overhead на Python interpreter
python_per_worker_gb = 0.2
elif use_pandas_udf:
# Pandas UDF: pandas + numpy + arrow buffers
# typical_pandas_df_size_mb × 2 (input + output) × safety factor
pandas_ram = typical_pandas_df_size_mb / 1024 * 2 * 1.5
python_per_worker_gb = max(0.5, pandas_ram)
else:
# No Python UDF
python_per_worker_gb = 0.1
if use_ml_libraries:
# scikit-learn + модели: 0.5-2 GB дополнительно
ml_overhead_gb = 1.5
else:
ml_overhead_gb = 0
# Суммарный Python overhead (один Worker на core)
total_python_gb = python_per_worker_gb * num_cores + ml_overhead_gb
# Итого overhead с запасом 20%
recommended_gb = (jvm_native_gb + total_python_gb) * 1.2
# Минимальный безопасный overhead
recommended_gb = max(recommended_gb, 1.0)
total_container_gb = executor_memory_gb + recommended_gb + 0.3 # + ОС
return {
"executor_memory_gb": executor_memory_gb,
"jvm_native_gb": jvm_native_gb,
"python_total_gb": total_python_gb,
"recommended_overhead_gb": round(recommended_gb, 1),
"total_container_gb": round(total_container_gb, 1),
"config_string": f"spark.executor.memoryOverhead = {recommended_gb:.1f}g",
}
# Примеры:
print("=== Scala ETL без Python ===")
r = calculate_pyspark_overhead(8, 4, use_pandas_udf=False)
print(f" Overhead: {r['recommended_overhead_gb']} GB")
print(f" Container: {r['total_container_gb']} GB")
# Overhead: 1.0 GB, Container: 9.3 GB
print("\n=== PySpark с Pandas UDF (500 MB батчи) ===")
r2 = calculate_pyspark_overhead(8, 4,
use_pandas_udf=True,
typical_pandas_df_size_mb=500)
print(f" Overhead: {r2['recommended_overhead_gb']} GB")
print(f" Container: {r2['total_container_gb']} GB")
# Overhead: 4.2 GB, Container: 12.5 GB
print("\n=== PySpark с ML (scikit-learn) ===")
r3 = calculate_pyspark_overhead(16, 5,
use_pandas_udf=True,
use_ml_libraries=True,
typical_pandas_df_size_mb=1000)
print(f" Overhead: {r3['recommended_overhead_gb']} GB")
print(f" Container: {r3['total_container_gb']} GB")
# Overhead: 7.2 GB, Container: 23.5 GB
Настройка для разных типов нагрузки¶
# ── Чистый Scala/Java ETL ────────────────────────────────────────────
spark_scala = SparkSession.builder \
.config("spark.executor.memory", "8g") \
.config("spark.executor.memoryOverhead", "1g") \
.getOrCreate()
# Контейнер: 8 + 1 + 0.3 = 9.3 GB
# ── PySpark с простыми Python UDF ────────────────────────────────────
spark_simple_py = SparkSession.builder \
.config("spark.executor.memory", "8g") \
.config("spark.executor.memoryOverhead", "2g") \
.getOrCreate()
# Контейнер: 8 + 2 + 0.3 = 10.3 GB
# ── PySpark с Pandas UDF (Apache Arrow) ─────────────────────────────
spark_pandas = SparkSession.builder \
.config("spark.executor.memory", "8g") \
.config("spark.executor.memoryOverhead", "4g") \
# Arrow включён для Pandas UDF
.config("spark.sql.execution.arrow.pyspark.enabled", "true") \
.config("spark.sql.execution.arrow.maxRecordsPerBatch", "10000") \
.getOrCreate()
# Контейнер: 8 + 4 + 0.3 = 12.3 GB
# ── PySpark с ML (scikit-learn, PyTorch inference) ───────────────────
spark_ml = SparkSession.builder \
.config("spark.executor.memory", "16g") \
.config("spark.executor.memoryOverhead", "8g") \
.getOrCreate()
# Контейнер: 16 + 8 + 0.3 = 24.3 GB
# ── Универсальная формула (правило 80/20) ────────────────────────────
# 80% total memory → executor.memory
# 20% total memory → memoryOverhead
# Пример: 20 GB pod
# executor.memory = 16g, memoryOverhead = 4g
6. Специализированный параметр pyspark.memory¶
Начиная со Spark 2.4, появился параметр spark.executor.pyspark.memory для точного контроля максимального потребления каждого Python Worker.
Как это работает и чем отличается от memoryOverhead¶
spark = SparkSession.builder \
# Общий контейнер
.config("spark.executor.memory", "8g") \
# Полный overhead (включая Python)
.config("spark.executor.memoryOverhead", "4g") \
# ДОПОЛНИТЕЛЬНО: явный лимит на Python Worker
# Если Python процесс превысит этот лимит — MemoryError в Python
# (вместо OOMKilled Pod'а)
# Это защита: лучше упасть с MemoryError в одном Task'е
# чем убить весь Container и потерять все Task'и
.config("spark.executor.pyspark.memory", "2g") \
.getOrCreate()
# Что происходит при превышении pyspark.memory:
# 1. Python Worker получает SIGKILL с MemoryError
# 2. Конкретный Task падает (retry)
# 3. НО Pod/Container остаётся живым!
# 4. Другие Tasks на том же Executor продолжают работу
# Без pyspark.memory при нехватке overhead:
# 1. Container убивает YARN/K8s (OOMKilled)
# 2. ВСЕ Tasks на этом Executor теряются
# 3. Spark пересчитывает Stage заново
# ПРАВИЛО:
# pyspark.memory < memoryOverhead (pyspark.memory — это лимит одного Worker)
# Несколько Workers на Executor (cores > 1) каждый может занять pyspark.memory
# Поэтому: memoryOverhead > pyspark.memory × executor_cores
7. Опасности Python UDF и раздувание памяти¶
Классические Python UDF — главный триггер неконтролируемого роста потребления памяти. Понимание физики их работы позволяет избежать типичных ловушек.
Физика передачи данных между JVM и Python¶
Практическое сравнение: UDF vs Pandas UDF¶
from pyspark.sql import SparkSession, functions as F
from pyspark.sql.functions import udf, pandas_udf
from pyspark.sql.types import DoubleType
import pandas as pd
import numpy as np
import time
spark = SparkSession.builder \
.config("spark.sql.execution.arrow.pyspark.enabled", "true") \
.getOrCreate()
df = spark.range(1_000_000).select(
F.col("id").cast("double").alias("value")
)
df.cache()
df.count()
# ── Вариант 1: Классический Python UDF ───────────────────────────────
@udf(returnType=DoubleType())
def normalize_classic(x: float) -> float:
"""Каждая строка = один вызов Python функции."""
return (x - 500000.0) / 288675.0 # стандартизация
t0 = time.time()
result_udf = df.withColumn("normalized", normalize_classic(F.col("value")))
count_udf = result_udf.count()
t_udf = time.time() - t0
print(f"Classic UDF: {t_udf:.2f}s ({count_udf:,} rows)")
# ── Вариант 2: Pandas UDF (Apache Arrow) ─────────────────────────────
@pandas_udf(DoubleType())
def normalize_vectorized(x: pd.Series) -> pd.Series:
"""Весь батч = одна векторная операция NumPy."""
return (x - 500000.0) / 288675.0 # numpy SIMD операция!
t0 = time.time()
result_vudf = df.withColumn("normalized", normalize_vectorized(F.col("value")))
count_vudf = result_vudf.count()
t_vudf = time.time() - t0
print(f"Pandas UDF: {t_vudf:.2f}s ({count_vudf:,} rows)")
print(f"\nPandas UDF быстрее в {t_udf/t_vudf:.1f}x")
print(f"(при 1M строк: {t_udf:.1f}с vs {t_vudf:.1f}с)")
# ── Вариант 3: Без UDF (лучший вариант!) ─────────────────────────────
t0 = time.time()
result_sql = df.withColumn(
"normalized",
(F.col("value") - F.lit(500000.0)) / F.lit(288675.0)
)
count_sql = result_sql.count()
t_sql = time.time() - t0
print(f"\nSQL/DataFrame (без UDF): {t_sql:.2f}s ({count_sql:,} rows)")
print(f"SQL быстрее UDF в {t_udf/t_sql:.1f}x")
# SQL работает полностью в JVM Tungsten, Python вообще не задействован!
8. Настройка в Kubernetes: emptyDir и Resource Limits¶
Kubernetes требует специального подхода к настройке spark.local.dir и памяти.
Правильная конфигурация emptyDir для NVMe¶
# executor-pod-template.yaml
apiVersion: v1
kind: Pod
metadata:
name: spark-executor
spec:
containers:
- name: spark-executor
resources:
requests:
memory: "20Gi" # spark.executor.memory + overhead + буфер
cpu: "5"
limits:
memory: "20Gi" # Limits = Requests для Guaranteed QoS класса!
cpu: "5"
volumeMounts:
- mountPath: /mnt/nvme0/spark-local
name: spark-local-0
- mountPath: /mnt/nvme1/spark-local
name: spark-local-1
volumes:
- name: spark-local-0
emptyDir:
# medium: "" — использовать локальный диск ноды (NVMe если нода его имеет)
# medium: "Memory" — tmpfs (RAM-диск, быстро но ограничен по RAM)
medium: ""
sizeLimit: "400Gi" # Ограничиваем размер чтобы не заполнить весь диск
- name: spark-local-1
emptyDir:
medium: ""
sizeLimit: "400Gi"
# Гарантируем что Pod размещается на ноде с NVMe
nodeSelector:
cloud.google.com/gke-local-ssd: "true" # GKE
node.kubernetes.io/instance-type: "i3.4xlarge" # EKS
# Или через node affinity для более гибкого выбора
# Spark конфигурация для K8s
spark = SparkSession.builder \
.master("k8s://https://kubernetes.default.svc") \
.config("spark.executor.memory", "12g") \
.config("spark.executor.memoryOverhead", "6g") \
.config("spark.executor.pyspark.memory", "4g") \
.config("spark.kubernetes.executor.podTemplateFile",
"/opt/spark/conf/executor-pod-template.yaml") \
.config("spark.local.dir",
"/mnt/nvme0/spark-local,/mnt/nvme1/spark-local") \
.config("spark.kubernetes.executor.limit.cores", "5") \
.getOrCreate()
Расчёт Resource Limits для K8s Pod¶
def calculate_k8s_pod_resources(
executor_memory_gb: float,
executor_cores: int,
memory_overhead_gb: float,
use_pandas_udf: bool = False,
) -> dict:
"""
Рассчитывает Resource Requests/Limits для Kubernetes Pod.
КРИТИЧНО: в K8s использовать Guaranteed QoS class:
requests == limits для обеих (memory и cpu).
Это предотвращает убийство Pod'а OOM Killer'ом при burst нагрузке.
"""
# Общая память контейнера
os_overhead_gb = 0.3
total_memory_gb = executor_memory_gb + memory_overhead_gb + os_overhead_gb
# Rounded up to nearest 0.5 GB
pod_memory_limit_gb = round(total_memory_gb * 2) / 2 + 0.5 # +0.5 запас
# CPU: по executor cores + небольшой запас на системные процессы
cpu_request = executor_cores
cpu_limit = executor_cores # requests == limits для Guaranteed QoS
return {
"executor_memory": f"{executor_memory_gb}g",
"memory_overhead": f"{memory_overhead_gb}g",
"pod_memory_request": f"{pod_memory_limit_gb}Gi",
"pod_memory_limit": f"{pod_memory_limit_gb}Gi",
"cpu_request": str(cpu_request),
"cpu_limit": str(cpu_limit),
"qos_class": "Guaranteed (requests == limits)",
"yaml_snippet": f"""
resources:
requests:
memory: "{pod_memory_limit_gb}Gi"
cpu: "{cpu_request}"
limits:
memory: "{pod_memory_limit_gb}Gi"
cpu: "{cpu_limit}"
"""
}
# Пример для PySpark с Pandas UDF:
config = calculate_k8s_pod_resources(
executor_memory_gb=12,
executor_cores=4,
memory_overhead_gb=6,
use_pandas_udf=True,
)
print(config["yaml_snippet"])
# resources:
# requests:
# memory: "19.0Gi"
# cpu: "4"
# limits:
# memory: "19.0Gi"
# cpu: "4"
9. Отладка проблем через Spark UI и логи¶
Диагностика проблем с диском¶
# ── Spark UI анализ Shuffle метрик ────────────────────────────────────
# Вкладка: Stages → конкретный Stage → Summary Metrics
# Ключевые метрики:
# - Shuffle Write Time: сколько времени тратится на запись shuffle файлов
# - Shuffle Read Time: сколько времени тратится на чтение shuffle файлов
# - Spill (Memory): объём данных сброшенных с RAM (до сжатия)
# - Spill (Disk): объём данных записанных на диск
# Интерпретация:
# Shuffle Write Time / Stage Duration > 30% → диск узкое место!
# Spill (Disk) > 0 → увеличьте executor.memory или shuffle.partitions
# ── Программный мониторинг через SparkContext ──────────────────────────
def get_shuffle_io_stats(spark) -> dict:
"""
Получает метрики Shuffle I/O из активных Stage'ов.
"""
sc = spark.sparkContext
status = sc.statusTracker()
stats = {
"active_stages": len(status.getActiveStageIds()),
"active_executors": len(status.getExecutorInfos()),
}
return stats
# ── Диагностика через логи ────────────────────────────────────────────
# YARN: Container killed by YARN (insufficient memoryOverhead)
# grep "killed by YARN\|memory limit\|OOMKilled\|memoryOverhead" /var/log/yarn/
# Kubernetes OOMKilled:
# kubectl describe pod spark-executor-xxx | grep -A5 "OOMKilled\|Reason"
# Events:
# OOMKilled: 1
# Exit Code: 137 ← Exit Code 137 = 128 + 9 (SIGKILL)
# Python Worker OOM внутри Spark (более мягкий):
# grep "MemoryError\|Cannot allocate memory\|PyArrow\|pyspark.memory" executor.log
Признаки проблем и их решения¶
DIAGNOSTICS = {
"exit_code_137": {
"symptom": "Container killed (Exit Code 137 / OOMKilled)",
"cause": "spark.executor.memoryOverhead слишком мало",
"solution": [
"Увеличить spark.executor.memoryOverhead (до 30-40% от executor.memory)",
"Добавить spark.executor.pyspark.memory",
"Оптимизировать Python UDF: перейти на Pandas UDF",
"Уменьшить spark.sql.execution.arrow.maxRecordsPerBatch",
]
},
"high_shuffle_write_time": {
"symptom": "Shuffle Write Time > 30% Stage Duration",
"cause": "Медленный диск под spark.local.dir",
"solution": [
"Перенести spark.local.dir на NVMe SSD",
"Использовать несколько дисков для striping",
"Увеличить spark.sql.shuffle.partitions (меньше данных на Task)",
"Увеличить executor.memory (меньше Spill)",
]
},
"high_spill": {
"symptom": "Большой Spill (Disk) в Spark UI",
"cause": "Нехватка Execution Memory",
"solution": [
"Увеличить spark.executor.memory",
"Уменьшить spark.memory.storageFraction",
"Увеличить spark.sql.shuffle.partitions",
"Убедиться что spark.local.dir на быстром NVMe",
]
},
"gc_pressure": {
"symptom": "GC Time > 10% Task Duration",
"cause": "Слишком большой JVM Heap или много мелких объектов",
"solution": [
"Уменьшить executor.memory (меньше Heap → быстрее GC)",
"Включить spark.memory.offHeap.enabled=true",
"Включить KryoSerializer",
"Перейти с RDD API на DataFrame API",
]
},
}
def diagnose_problem(symptom: str) -> None:
"""Выводит диагностику и решения для конкретного симптома."""
if symptom in DIAGNOSTICS:
d = DIAGNOSTICS[symptom]
print(f"Симптом: {d['symptom']}")
print(f"Причина: {d['cause']}")
print("Решения:")
for s in d["solution"]:
print(f" • {s}")
10. Production чек-лист: настройка железа под PySpark-пайплайны¶
Полная конфигурация production PySpark¶
import os
from pyspark.sql import SparkSession
def create_production_spark_session(
app_name: str,
workload_type: str = "etl", # etl, ml, streaming
num_nodes: int = 10,
cores_per_node: int = 5,
ram_per_node_gb: int = 20,
has_nvme: bool = True,
nvme_paths: list[str] = None,
) -> SparkSession:
"""
Создаёт production-ready SparkSession с корректными
настройками железа и памяти для PySpark.
Включает:
- Оптимальный spark.local.dir на NVMe
- Корректный memoryOverhead под Python Workers
- Arrow для Pandas UDF
- AQE для динамической оптимизации
"""
if nvme_paths is None:
nvme_paths = ["/mnt/nvme0/spark-local", "/mnt/nvme1/spark-local"]
# ── Расчёт параметров памяти ──────────────────────────────────────
# Резервируем 2 CPU + 4 GB под ОС и YARN NodeManager
available_ram = ram_per_node_gb - 4
# Распределение памяти в зависимости от нагрузки
workload_config = {
"etl": {
"executor_memory_gb": int(available_ram * 0.75),
"overhead_multiplier": 0.20, # 20% на overhead
},
"ml": {
"executor_memory_gb": int(available_ram * 0.60),
"overhead_multiplier": 0.35, # 35% под Python ML libs
},
"streaming": {
"executor_memory_gb": int(available_ram * 0.70),
"overhead_multiplier": 0.25,
},
}
config = workload_config.get(workload_type, workload_config["etl"])
executor_memory_gb = config["executor_memory_gb"]
overhead_gb = max(1, int(available_ram * config["overhead_multiplier"]))
pyspark_memory_gb = max(1, overhead_gb - 1) # немного меньше overhead
# ── Настройка дисков ──────────────────────────────────────────────
if has_nvme:
local_dir = ",".join(nvme_paths)
else:
local_dir = "/tmp/spark-local"
builder = SparkSession.builder \
.appName(app_name) \
# ── Диски ─────────────────────────────────────────────────────
.config("spark.local.dir", local_dir) \
# ── Память JVM ────────────────────────────────────────────────
.config("spark.executor.memory", f"{executor_memory_gb}g") \
.config("spark.executor.cores", str(cores_per_node)) \
.config("spark.executor.instances", str(num_nodes - 1)) \
# ── Память Python ─────────────────────────────────────────────
.config("spark.executor.memoryOverhead", f"{overhead_gb}g") \
.config("spark.executor.pyspark.memory", f"{pyspark_memory_gb}g") \
# ── Memory Fractions ──────────────────────────────────────────
.config("spark.memory.fraction", "0.7") \
.config("spark.memory.storageFraction", "0.3") \
# ── Arrow для Pandas UDF ──────────────────────────────────────
.config("spark.sql.execution.arrow.pyspark.enabled", "true") \
.config("spark.sql.execution.arrow.maxRecordsPerBatch", "10000") \
# ── AQE ───────────────────────────────────────────────────────
.config("spark.sql.adaptive.enabled", "true") \
.config("spark.sql.adaptive.coalescePartitions.enabled", "true") \
.config("spark.sql.adaptive.skewJoin.enabled", "true") \
# ── Shuffle ───────────────────────────────────────────────────
.config("spark.sql.shuffle.partitions",
str(num_nodes * cores_per_node * 4)) \
spark = builder.getOrCreate()
print(f"SparkSession создан для '{workload_type}' нагрузки:")
print(f" spark.local.dir: {local_dir}")
print(f" spark.executor.memory: {executor_memory_gb}g")
print(f" spark.executor.memoryOverhead: {overhead_gb}g")
print(f" spark.executor.pyspark.memory: {pyspark_memory_gb}g")
print(f" Container суммарно: {executor_memory_gb + overhead_gb + 0.3:.1f}g")
return spark
# ── Использование ─────────────────────────────────────────────────────
# ETL на кластере с NVMe (10 нод × 24 CPU × 96 GB RAM)
spark_etl = create_production_spark_session(
app_name="daily-etl",
workload_type="etl",
num_nodes=10,
cores_per_node=5,
ram_per_node_gb=96,
has_nvme=True,
nvme_paths=["/mnt/nvme0/spark", "/mnt/nvme1/spark", "/mnt/nvme2/spark"]
)
# ML инференс с тяжёлыми Python библиотеками
spark_ml = create_production_spark_session(
app_name="ml-inference",
workload_type="ml",
num_nodes=5,
cores_per_node=4,
ram_per_node_gb=64,
has_nvme=True,
nvme_paths=["/mnt/nvme0/spark"]
)
Итоги: физическая оптимизация не менее важна чем логическая¶
spark.local.dir на NVMe — одна из самых недооценённых оптимизаций. При тяжёлых ETL с большим объёмом Shuffle перенос временных файлов с HDD на NVMe может сократить время Stage в 5-20 раз без единой строки изменения в коде. Всегда указывайте несколько NVMe-дисков для Round-Robin striping.
spark.executor.memoryOverhead — главный источник загадочных падений PySpark в production. Container killed с Exit Code 137 при свободном JVM Heap — это всегда сигнал о недостаточном Overhead. Для PySpark с Pandas UDF устанавливайте Overhead не менее 30-40% от executor.memory.
Три ключевых правила:
- Никогда не оставляйте
spark.local.dirна/tmp(системном диске) в production - Для PySpark всегда явно задавайте
memoryOverhead— дефолтные 10% катастрофически мало - Apache Arrow (
spark.sql.execution.arrow.pyspark.enabled=true) — обязательная настройка для любого PySpark с Pandas UDF