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 и логи.

optimization platform

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.

Три ключевых правила:

  1. Никогда не оставляйте spark.local.dir на /tmp (системном диске) в production
  2. Для PySpark всегда явно задавайте memoryOverhead — дефолтные 10% катастрофически мало
  3. Apache Arrow (spark.sql.execution.arrow.pyspark.enabled=true) — обязательная настройка для любого PySpark с Pandas UDF