Spark UI: вкладка Executors — Memory Breakdown и Task сводка

Полный разбор вкладки Executors в Spark UI: архитектура Executor и анатомия памяти, Storage vs Execution Memory и динамическое заимствование, Off-Heap и Python Workers как источник скрытых OOM, диагностика CPU утилизации и GC Time, Dead Executors и причины их гибели, дисковые метрики Shuffle, расчёт ROI кластера.

optimization

1. Архитектура Executor и назначение вкладки

Вкладка Executors — это финальный уровень диагностики в цепочке Jobs → Stages → SQL → Executors. Если предыдущие вкладки отвечали на вопросы «что», «как быстро» и «почему такой план», то вкладка Executors отвечает на вопрос «на каком железе и в какой памяти это выполнялось».

Executor — это изолированный JVM-процесс, запущенный на рабочем узле кластера. Каждый Executor имеет свой пул потоков (Threads), свою оперативную память и свой доступ к локальным дискам. Именно Executor'ы являются единицей вычислений: каждый Task выполняется в одном потоке одного Executor'а.

Анатомия Executor'а: что внутри JVM

Схема показывает полную картину памяти Executor'а. Критически важный момент: то что видит Spark UI в разделе Memory — это только JVM Heap. Python Workers и Off-Heap существуют вне JVM и могут вызвать OOM даже при «свободной» JVM памяти.

Что показывает вкладка Executors

Открыв вкладку, вы видите таблицу:

ID | Address   | State  | RDD Blocks | Storage Memory | Cores | Active Tasks | Failed Tasks
---|-----------|--------|------------|----------------|-------|--------------|-------------
0  | driver    | Active | 0          | 0 / 2.3 GB     | 1     | 0            | 0
1  | node1:45678| Active | 45         | 1.2 / 5.5 GB   | 5     | 4            | 0
2  | node2:45679| Active | 43         | 1.1 / 5.5 GB   | 5     | 5            | 0
3  | node3:45680| Active | 0          | 0 / 5.5 GB     | 5     | 0            | 0  ← простаивает!

Первое что вы должны заметить: Executor 3 не выполняет ни одной задачи и не кеширует данные. При 5 ядрах на остальных Executor'ах это означает что нагрузка могла распределяться более равномерно.


2. Memory Breakdown: Storage vs Execution

Наиболее информативная часть вкладки Executors — разбивка памяти. Понимание как память делится между хранением данных и вычислениями критично для диагностики Spill и Cache Eviction.

Unified Memory Manager: динамическое заимствование

Схема показывает два сценария. В сценарии A кеша нет, Execution свободно занимает всё место. В сценарии B кеш уже занят, и когда JOIN требует память — часть кеша принудительно вытесняется. Это Cache Eviction — скрытая причина того, что df.cache() не всегда работает ожидаемо.

Как читать колонку Storage Memory в Executors

Executor 1: Storage Memory: 1.2 GB / 5.5 GB
                              ↑           ↑
                              Используется   Доступно (unified pool)

Дробь X / Y означает: X GB использовано под кеш и broadcast, Y GB доступно суммарно для Storage и Execution.

Когда цифры говорят о проблеме:

  • 5.4 GB / 5.5 GB — почти всё занято под кеш. При следующем JOIN будет Eviction или Spill
  • 0 GB / 5.5 GB — кеша нет. Если вы вызывали .cache() — оно не сработало
  • Разные значения / Y на разных Executor'ах — неравномерная конфигурация или Dynamic Allocation

Программный расчёт параметров памяти

from pyspark.sql import SparkSession


def calculate_executor_memory_layout(
    executor_memory_gb: float,
    memory_fraction: float = 0.6,
    storage_fraction: float = 0.5,
    reserved_mb: int = 300,
) -> dict:
    """
    Рассчитывает реальное распределение памяти Executor'а.
    Это то, что вы должны видеть на вкладке Executors в Spark UI.

    Формулы Unified Memory Manager:
    usable = executor_memory - reserved_mb
    spark_memory = usable × memory_fraction
    user_memory  = usable × (1 - memory_fraction)
    storage_floor = spark_memory × storage_fraction  ← минимум, защищён
    execution_max = spark_memory                     ← максимум (если Storage пуст)
    storage_max   = spark_memory                     ← максимум (если Execution пуст)
    """
    usable_mb = executor_memory_gb * 1024 - reserved_mb
    spark_memory_mb = usable_mb * memory_fraction
    user_memory_mb = usable_mb * (1 - memory_fraction)
    storage_floor_mb = spark_memory_mb * storage_fraction

    return {
        "executor_memory_gb": executor_memory_gb,
        "reserved_mb": reserved_mb,
        "usable_gb": usable_mb / 1024,
        "spark_memory_gb": spark_memory_mb / 1024,
        "user_memory_gb": user_memory_mb / 1024,
        "storage_floor_gb": storage_floor_mb / 1024,
        "execution_max_gb": spark_memory_mb / 1024,
        "storage_max_gb": spark_memory_mb / 1024,
        "what_spark_ui_shows": f"{spark_memory_mb/1024:.1f} GB (X / this)",
    }


# Дефолтная конфигурация для executor.memory=8g:
layout = calculate_executor_memory_layout(8.0, 0.6, 0.5)
for k, v in layout.items():
    print(f"  {k}: {v}")

# Вывод:
# executor_memory_gb: 8.0
# reserved_mb: 300
# usable_gb: 7.71
# spark_memory_gb: 4.62     ← вот что показывает "/ X" в UI
# user_memory_gb: 3.08
# storage_floor_gb: 2.31    ← минимум для кеша, нельзя вытеснить
# execution_max_gb: 4.62    ← максимум для Execution (если Storage пуст)
# what_spark_ui_shows: 4.6 GB

3. Off-Heap память и PySpark Workers: скрытые источники OOM

Наиболее коварные OOM-ошибки в PySpark возникают не из-за переполнения JVM Heap, а из-за памяти вне JVM. Вкладка Executors показывает только JVM-часть — это принципиально важно понимать.

Полная формула памяти контейнера

Total Container Memory = spark.executor.memory
                       + spark.executor.memoryOverhead
                       + spark.memory.offHeap.size
                       + spark.executor.pyspark.memory (если задан)

YARN/K8s контролирует TOTAL. Если Total превышено → OOM Killer убивает контейнер!
Spark UI показывает только spark.executor.memory часть!

Что живёт в memoryOverhead

Главная ловушка PySpark: Python Worker'ы используют memoryOverhead (или pyspark.memory если задан). Если ваш UDF загружает большую ML-модель в каждый Worker — это не видно в JVM Heap на вкладке Executors, но контейнер будет убит OOM Killer'ом!

Диагностика: видим OOM в логах, но в UI память «свободна»

# Типичная ошибка Container killed by OOM Killer:
# В логах YARN/K8s:
Container killed on request. Exit code is 137.  ← 137 = SIGKILL = OOM Killer!
Container [pid=12345,containerID=xxx] is running beyond virtual memory limits.
Current usage: 16.0 GB of 15 GB physical memory used.

# НО! В Spark UI → Executors:
Storage Memory: 3.1 GB / 4.2 GB  ← "только" 74% JVM Heap использовано
# Кажется памяти достаточно, но контейнер убит?!

# Причина: Python Workers потребили 12+ GB вне JVM!
# JVM видит свои 4.2 GB, Python видит свои 12 GB
# Total = 4.2 + 12 = 16.2 GB > 15 GB limit → OOM!
# Решение: явно настроить память для Python Workers

spark = SparkSession.builder \

    # JVM часть — то что видит Spark UI
    .config("spark.executor.memory", "8g") \

    # Overhead для Python Workers и Netty:
    # Для PySpark с Pandas UDF: увеличиваем overhead
    # Правило: overhead = 20-50% executor.memory для тяжёлых PySpark UDF
    .config("spark.executor.memoryOverhead", "4g") \

    # Явный лимит на Python Workers (Spark 2.4+):
    # Если задан — Python Workers могут использовать ДО этого лимита из overhead
    # Если не задан — Python Workers могут занять весь overhead и больше!
    .config("spark.executor.pyspark.memory", "3g") \

    # Суммарно: 8 + 4 = 12 GB на контейнер
    # YARN/K8s должен выделять именно это количество!

    .getOrCreate()


# Проверка: рассчитаем нужный overhead
def estimate_pyspark_overhead(
    pandas_udf_peak_gb: float,    # пиковое потребление Pandas в одном Worker
    n_cores: int = 5,              # ядер на Executor = Worker'ов одновременно
    model_size_gb: float = 0,      # если ML-модель загружается в каждый Worker
    arrow_buffer_multiplier: float = 2.5,  # Arrow buffers ≈ 2.5× данных
) -> float:
    """
    Оценивает необходимый spark.executor.memoryOverhead для PySpark.
    """
    # Память одного Python Worker
    per_worker_gb = (pandas_udf_peak_gb * arrow_buffer_multiplier
                     + model_size_gb
                     + 0.2)  # base Python interpreter

    # Одновременно работают n_cores Worker'ов
    python_total_gb = per_worker_gb * n_cores

    # Добавляем JVM overhead (Netty, JIT code)
    jvm_overhead_gb = 1.0

    return python_total_gb + jvm_overhead_gb


# Пример: Pandas UDF обрабатывает батчи по 100K строк × 200 байт = 20 MB
# Arrow буферы: 20 MB × 2.5 = 50 MB = 0.05 GB на Worker
# При 5 ядрах: 5 × 0.05 = 0.25 GB Python + 1 GB JVM = 1.25 GB минимум
recommended = estimate_pyspark_overhead(0.05, 5, model_size_gb=0.5)
print(f"Рекомендуемый memoryOverhead: {recommended:.1f} GB")
# 0.05 × 2.5 × 5 + 0.5 × 5 + 1.0 = 4.125 GB

4. CPU утилизация и GC Time: диагностика неэффективного использования

Вкладка Executors показывает не только память но и эффективность использования CPU. Анализ этих метрик позволяет определить является ли ваш кластер compute-bound или memory-bound.

Ключевые метрики CPU на вкладке

Executor | Cores | Active Tasks | Task Time | GC Time | Failed Tasks
---------|-------|--------------|-----------|---------|-------------
1        | 5     | 4/5          | 1h 23min  | 8min    | 0     ← GC 9.6% 
2        | 5     | 5/5          | 1h 31min  | 45min   | 0     ← GC 49%!
3        | 5     | 0/5          | 45min     | 2min    | 0     ← простаивал!

Три проблемы видны сразу:

  1. Executor 2: GC Time = 49% от Task Time → критическое Memory Pressure
  2. Executor 3: 0 активных тасок при наличии 5 ядер → нагрузка не распределяется
  3. Executor 1: GC = 9.6% → близко к опасному порогу (10%)

Расчёт Core Utilization Efficiency

from typing import NamedTuple


class ExecutorMetrics(NamedTuple):
    executor_id: str
    cores: int
    task_time_ms: int      # суммарное время всех Task
    gc_time_ms: int        # суммарное время GC
    shuffle_read_bytes: int
    shuffle_write_bytes: int
    failed_tasks: int
    active_tasks: int


def analyze_cluster_efficiency(executors: list[ExecutorMetrics]) -> dict:
    """
    Рассчитывает ROI кластера: насколько эффективно используются ресурсы.

    Core Utilization = чистое время вычислений / общее время × ядра

    Ideal: 80%+ это хорошо
    60-80%: допустимо, есть резервы
    <60%: кластер недоиспользуется или проблемы с памятью
    """
    total_task_time = sum(e.task_time_ms for e in executors)
    total_gc_time = sum(e.gc_time_ms for e in executors)
    total_failed = sum(e.failed_tasks for e in executors)
    total_cores = sum(e.cores for e in executors)
    total_shuffle_read = sum(e.shuffle_read_bytes for e in executors)
    total_shuffle_write = sum(e.shuffle_write_bytes for e in executors)

    # Чистое вычислительное время (без GC)
    compute_time = total_task_time - total_gc_time

    gc_pct = total_gc_time / max(total_task_time, 1) * 100

    # Дисбаланс нагрузки: std / mean задач по Executor'ам
    task_times = [e.task_time_ms for e in executors]
    mean_time = sum(task_times) / len(task_times)
    variance = sum((t - mean_time) ** 2 for t in task_times) / len(task_times)
    load_imbalance = (variance ** 0.5) / max(mean_time, 1)

    # Проблемные Executor'ы
    overloaded = [e for e in executors
                  if e.gc_time_ms / max(e.task_time_ms, 1) > 0.15]
    idle = [e for e in executors
            if e.task_time_ms < mean_time * 0.2]

    recommendations = []
    if gc_pct > 20:
        recommendations.append(
            f"ВЫСОКИЙ GC ({gc_pct:.0f}%): увеличьте executor.memory или "
            "переключитесь на MEMORY_ONLY_SER кеш и ZGC/G1GC с тюнингом."
        )
    if load_imbalance > 0.5:
        recommendations.append(
            f"ДИСБАЛАНС НАГРУЗКИ (imbalance={load_imbalance:.1f}): "
            "проверьте Data Skew или проблемы с Data Locality."
        )
    if idle:
        recommendations.append(
            f"{len(idle)} Executor'ов простаивают: "
            "проверьте конфигурацию Dynamic Allocation или данные не распределены."
        )
    if total_failed > 0:
        recommendations.append(
            f"{total_failed} Task упало: проверьте Dead Executors tab и логи."
        )

    return {
        "total_cores": total_cores,
        "gc_overhead_pct": gc_pct,
        "compute_efficiency_pct": compute_time / max(total_task_time, 1) * 100,
        "load_imbalance_score": load_imbalance,
        "overloaded_executors": [e.executor_id for e in overloaded],
        "idle_executors": [e.executor_id for e in idle],
        "total_failed_tasks": total_failed,
        "shuffle_read_gb": total_shuffle_read / 1024**3,
        "shuffle_write_gb": total_shuffle_write / 1024**3,
        "recommendations": recommendations,
    }

5. Dead Executors: диагностика потери узлов

Секция Dead Executors в нижней части вкладки хранит историю Executor'ов, которые прекратили работу. Это бесценный источник для постфактум-диагностики.

Два типа «смерти» Executor'а

Как определить причину смерти Executor'а

Dead Executor Log (что показывает Spark UI):
Executor 5 | node3:56789 | State: Dead
Finish Time: 14:23:45
Logs: [stderr] [stdout]

→ Нажимаем Logs → stderr:
"Container killed by YARN for exceeding memory limits.
 12.4 GB of 12.0 GB physical memory used."
→ Диагноз: OOM, нужно увеличить memoryOverhead

"ExecutorLost: Remote RPC client disassociated.
 Address of executor: SparkException: Lost executor 5 on node3."
→ Диагноз: Heartbeat timeout, скорее всего Full GC или сетевой сбой

"ExecutorLost: Executor heartbeat timed out after 120000 ms"
→ Диагноз: spark.network.timeout слишком маленький или Full GC
# Настройки которые снижают ложные срабатывания Heartbeat timeout

spark = SparkSession.builder \
    # Timeout сетевых операций (heartbeat включен)
    # Дефолт: 120 секунд. Для тяжёлых GC нужно больше!
    .config("spark.network.timeout", "800s") \

    # Специально heartbeat interval (должен быть << network.timeout)
    .config("spark.executor.heartbeatInterval", "60s") \

    # Ожидание перед повторным запуском Task на другом Executor
    .config("spark.task.maxFailures", "4") \

    # Dynamic Allocation таймауты:
    # Executor с кешем ждём дольше чем пустой
    .config("spark.dynamicAllocation.executorIdleTimeout", "180s") \
    .config("spark.dynamicAllocation.cachedExecutorIdleTimeout", "3600s") \

    .getOrCreate()

6. Дисковые метрики Shuffle: диагностика нагрузки на хранилище

Каждый Executor записывает промежуточные Shuffle данные на локальный диск (spark.local.dir). Вкладка Executors показывает агрегированные дисковые метрики.

Что означают дисковые метрики

Executor | Shuffle Read  | Shuffle Write | Disk Used
---------|---------------|---------------|----------
1        | 2.3 GB        | 2.1 GB        | 45 MB
2        | 2.1 GB        | 2.0 GB        | 38 MB
3        | 8.9 GB        | 8.7 GB        | 1.2 GB   ← в 4x больше!
4        | 2.2 GB        | 2.2 GB        | 41 MB

Executor 3 имеет в 4 раза больше дискового трафика — классический признак Data Skew. Горячий ключ join'а попал именно в партиции этого Executor'а.

Как дисковые метрики помогают обосновать переход на NVMe

def analyze_shuffle_io(executors_data: list[dict],
                       disk_throughput_mb_per_sec: float = 100) -> dict:
    """
    Анализирует дисковую нагрузку Shuffle и рассчитывает overhead от медленных дисков.

    disk_throughput_mb_per_sec:
    - HDD 7200 RPM: ~100-150 MB/s
    - SATA SSD: ~500 MB/s
    - NVMe SSD: ~3000-7000 MB/s

    Если Shuffle I/O Bound, переход HDD→NVMe даёт 10-50x ускорение!
    """
    total_shuffle_gb = sum(
        e.get("shuffle_write_bytes", 0) + e.get("shuffle_read_bytes", 0)
        for e in executors_data
    ) / 1024**3

    total_shuffle_time_sec = total_shuffle_gb * 1024 / disk_throughput_mb_per_sec

    # Дисбаланс Shuffle по Executor'ам
    shuffle_amounts = [
        (e.get("shuffle_write_bytes", 0) + e.get("shuffle_read_bytes", 0)) / 1024**3
        for e in executors_data
    ]
    mean_shuffle = sum(shuffle_amounts) / len(shuffle_amounts)
    max_shuffle = max(shuffle_amounts)
    skew_ratio = max_shuffle / max(mean_shuffle, 0.001)

    return {
        "total_shuffle_gb": total_shuffle_gb,
        "estimated_io_time_sec_hdd": total_shuffle_gb * 1024 / 100,
        "estimated_io_time_sec_nvme": total_shuffle_gb * 1024 / 3500,
        "nvme_speedup_factor": 3500 / 100,
        "shuffle_skew_ratio": skew_ratio,
        "recommendation": (
            f"Shuffle I/O Time: {total_shuffle_time_sec:.0f}s (HDD) vs "
            f"{total_shuffle_gb * 1024 / 3500:.0f}s (NVMe). "
            f"{'Переход на NVMe даст ~35x ускорение Shuffle!' if total_shuffle_gb > 10 else 'Shuffle небольшой, NVMe необязателен.'}"
        )
    }

7. Сводная таблица и расчёт ROI кластера

Вкладка Executors позволяет рассчитать бизнес-эффективность использования инфраструктуры. Именно здесь данные переводятся в деньги.

Формула Core Utilization Efficiency

Core Utilization Efficiency = (Task Compute Time) / (Total Executor Time × Cores)

где:
Task Compute Time = Total Task Time - GC Time - Shuffle Wait Time
Total Executor Time = время жизни кластера × число Executor'ов × ядер на Executor

Идеал: 70%+
Проблема: < 50% → кластер работает вхолостую
Критично: < 30% → деньги тратятся без эффекта

Полная программная диагностика через History Server API

import requests
import json
from dataclasses import dataclass


@dataclass
class ExecutorSummary:
    executor_id: str
    host: str
    cores: int
    total_duration_ms: int
    total_gc_time_ms: int
    total_tasks: int
    failed_tasks: int
    storage_memory_used: int
    storage_memory_total: int
    max_memory: int
    shuffle_read: int
    shuffle_write: int
    is_active: bool


def fetch_executor_summary(app_id: str, history_server: str) -> list[ExecutorSummary]:
    """Получает данные Executors через History Server REST API."""
    url = f"{history_server}/api/v1/applications/{app_id}/allexecutors"
    response = requests.get(url, timeout=30)
    executors_json = response.json()

    return [
        ExecutorSummary(
            executor_id=e["id"],
            host=e.get("hostPort", "unknown"),
            cores=e.get("totalCores", 0),
            total_duration_ms=e.get("totalDuration", 0),
            total_gc_time_ms=e.get("totalGCTime", 0),
            total_tasks=e.get("completedTasks", 0),
            failed_tasks=e.get("failedTasks", 0),
            storage_memory_used=e.get("memoryUsed", 0),
            storage_memory_total=e.get("maxMemory", 0),
            max_memory=e.get("maxMemory", 0),
            shuffle_read=e.get("totalShuffleRead", 0),
            shuffle_write=e.get("totalShuffleWrite", 0),
            is_active=e.get("isActive", False),
        )
        for e in executors_json
        if e["id"] != "driver"  # исключаем Driver
    ]


def generate_cluster_health_report(executors: list[ExecutorSummary],
                                    hourly_cost_usd: float = 5.0) -> str:
    """
    Генерирует читаемый отчёт о здоровье кластера.
    Конвертирует технические метрики в бизнес-показатели.
    """
    active = [e for e in executors if e.is_active]
    dead = [e for e in executors if not e.is_active]

    total_task_ms = sum(e.total_duration_ms for e in executors)
    total_gc_ms = sum(e.total_gc_time_ms for e in executors)
    total_failed = sum(e.failed_tasks for e in executors)
    total_shuffle_gb = sum(e.shuffle_read + e.shuffle_write for e in executors) / 1024**3
    gc_pct = total_gc_ms / max(total_task_ms, 1) * 100

    # Расчёт эффективности
    compute_pct = (total_task_ms - total_gc_ms) / max(total_task_ms, 1) * 100

    # Memory utilization
    max_mem_used = max(
        e.storage_memory_used / max(e.storage_memory_total, 1)
        for e in active
    ) * 100 if active else 0

    # Estimated costs
    runtime_hours = total_task_ms / 1000 / 3600 / max(len(executors), 1)
    wasted_hours = runtime_hours * (1 - compute_pct / 100)
    wasted_cost = wasted_hours * hourly_cost_usd

    lines = [
        "=" * 60,
        "CLUSTER HEALTH REPORT (Spark UI → Executors)",
        "=" * 60,
        f"Active Executors:    {len(active)}",
        f"Dead Executors:      {len(dead)} {'← INVESTIGATE!' if dead else ''}",
        f"Total Failed Tasks:  {total_failed} {'← CHECK LOGS!' if total_failed > 5 else ''}",
        "",
        "MEMORY:",
        f"  Peak Storage Usage: {max_mem_used:.0f}%",
        f"  {'⚠️  HIGH: cache eviction likely' if max_mem_used > 85 else '✅ OK'}",
        "",
        "CPU EFFICIENCY:",
        f"  Compute Time:  {compute_pct:.0f}%",
        f"  GC Overhead:   {gc_pct:.0f}% {'❌ CRITICAL: tune memory!' if gc_pct > 20 else '⚠️  HIGH' if gc_pct > 10 else '✅ OK'}",
        f"  Status: {'🔴 OPTIMIZE URGENTLY' if compute_pct < 50 else '🟡 IMPROVE' if compute_pct < 70 else '🟢 HEALTHY'}",
        "",
        "SHUFFLE I/O:",
        f"  Total Shuffle: {total_shuffle_gb:.1f} GB",
        "",
        f"ESTIMATED WASTE: ${wasted_cost:.2f}/run (at ${hourly_cost_usd}/hour)",
        "=" * 60,
    ]

    if gc_pct > 20:
        lines.append("ACTION: Increase executor.memory or switch to G1GC/ZGC")
    if dead:
        lines.append(f"ACTION: Check logs for {len(dead)} dead executors (OOM or node failure?)")
    if max_mem_used > 85:
        lines.append("ACTION: Reduce cache size or increase executor.memory")

    return "\n".join(lines)


# Пример использования:
# executors = fetch_executor_summary("app-20240115-143022-0001",
#                                     "http://history-server:18080")
# print(generate_cluster_health_report(executors, hourly_cost_usd=8.50))

Итоговый алгоритм диагностики через все 4 вкладки Spark UI

Мы завершили изучение всех четырёх основных вкладок Spark UI. Объединим их в единую методику:


Практика: полное расследование проблемного Job

from pyspark.sql import SparkSession, functions as F
import time

spark = SparkSession.builder \
    .master("local[4]") \
    .appName("executors-demo") \
    .config("spark.executor.memory", "1g") \  # мало, для демонстрации Spill
    .config("spark.ui.enabled", "true") \
    .getOrCreate()

spark.sparkContext.setLogLevel("WARN")

print("Spark UI: http://localhost:4040")
print()

# ── ДЕМОНСТРАЦИЯ 1: Нормальная нагрузка ─────────────────────────────────────
print("=== ДЕМОНСТРАЦИЯ 1: Нормальная нагрузка ===")
df = spark.range(5_000_000).select(
    F.col("id"),
    (F.rand() * 100).cast("int").alias("key"),
    (F.rand() * 500).alias("value"),
)
df.groupBy("key").sum("value").count()
print("Проверьте Executors → Active Tasks, Task Time, GC Time")
print()

# ── ДЕМОНСТРАЦИЯ 2: Memory Pressure с GC ────────────────────────────────────
print("=== ДЕМОНСТРАЦИЯ 2: Cache pressure ===")
large_df = spark.range(2_000_000).select(
    F.col("id"),
    F.concat(F.lit("item_"), F.col("id").cast("string")).alias("str_key"),
    (F.rand() * 100).alias("value"),
)
large_df.cache()
large_df.count()  # материализуем кеш

# Теперь делаем join который потребует Execution Memory
dim = spark.range(1000).select(
    (F.col("id") % 100).alias("key"),
    F.concat(F.lit("cat_"), F.col("id").cast("string")).alias("category"),
)
result = large_df.withColumn("key", (F.col("id") % 100).cast("int")) \
    .join(dim, "key") \
    .groupBy("category") \
    .count()
result.collect()

print("Проверьте Executors → Storage Memory (заполнен кешем)")
print("При нехватке: смотрите Stages → Spill метрики")
print()

print("""
ЗАДАНИЕ: Откройте http://localhost:4040 → Executors

1. Колонка Storage Memory:
   - Какая доля занята кешем?
   - Если заполнено >90% — риск Cache Eviction

2. Колонка Task Time vs GC Time:
   - GC% = GC Time / Task Time × 100
   - > 10% → проблема с памятью

3. Колонка Active Tasks:
   - Все ядра заняты? Или простаивают?

4. Shuffle Write / Shuffle Read:
   - Распределены равномерно между Executors?
   - Или один Executor имеет в 5-10x больше других? → Data Skew

5. Прокрутите вниз → Dead Executors:
   - Есть ли мёртвые Executors?
   - Красный цвет → аварийное завершение (OOM, Node failure)
   - Серый → штатное завершение (Dynamic Allocation)
""")

spark.stop()

Итоги: вкладка Executors как финальный уровень диагностики

Вкладка Executors завершает цикл диагностики Spark UI. Она переводит абстрактные метрики производительности в конкретные физические узлы, память и диски.

Три главных вопроса которые решает вкладка Executors:

  1. Нагрузка равномерна? — смотрите Task Time и Shuffle I/O по Executor'ам. Если один Executor в 5-10x нагружен больше остальных — Data Skew.

  2. Памяти достаточно? — смотрите GC Time% (> 10% = проблема) и Storage Memory (> 90% = риск Eviction). Помните что Python Workers не видны в JVM метриках.

  3. Инфраструктура стабильна? — смотрите Dead Executors. Красные записи с Exit Code 137 = OOM Killer. «Lost» = Heartbeat timeout (длинный GC или сетевой сбой).

Единая методика диагностики Spark UI (полный цикл):

Jobs → найти медленный Job
Stages → найти проблемный Stage (Skew? Spill? Fetch Wait?)
SQL → проверить план (BHJ vs SMJ, WSCG, Partition Pruning)
Executors → подтвердить диагноз на уровне узлов (GC, OOM, Imbalance)