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 кластера.
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 или Spill0 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 ← простаивал!
Три проблемы видны сразу:
- Executor 2: GC Time = 49% от Task Time → критическое Memory Pressure
- Executor 3: 0 активных тасок при наличии 5 ядер → нагрузка не распределяется
- 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:
-
Нагрузка равномерна? — смотрите Task Time и Shuffle I/O по Executor'ам. Если один Executor в 5-10x нагружен больше остальных — Data Skew.
-
Памяти достаточно? — смотрите GC Time% (> 10% = проблема) и Storage Memory (> 90% = риск Eviction). Помните что Python Workers не видны в JVM метриках.
-
Инфраструктура стабильна? — смотрите 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)