Spark UI: вкладка Stages — Task Skew, Shuffle Write/Read и GC Time
Полный разбор вкладки Stages в Spark UI: таблица квантилей и верхнеуровневый аудит, Task Skew и диагностика по гистограмме, механика Shuffle Write/Read и Fetch Wait Time, Spill (Memory) и Spill (Disk) как сигнал нехватки памяти, GC Time и стратегии снижения, декомпозиция Task Duration и анализ Task Table.
1. Метрики стадии и верхнеуровневый аудит Tasks¶
Если вкладка Jobs — это диспетчерский пульт всего приложения, то вкладка Stages — это операционный зал, где видны детали каждого отдельного этапа работы. Именно сюда попадают опытные инженеры после того, как на вкладке Jobs найдён медленный Job и нужно понять почему он медленный.
Каждая стадия в Spark — это совокупность параллельных Task'ов, каждый из которых обрабатывает ровно одну Partition данных на одном ядре CPU одного Executor'а. Вкладка Stages предоставляет агрегированные метрики по всем Task'ам стадии в виде таблицы квантилей.
Таблица квантилей: главный диагностический инструмент¶
Схема показывает три паттерна которые нужно уметь различать с первого взгляда. Первый — идеальный: все квантили близки, Max не сильно превышает Median. Второй — Data Skew: Max в сотни или тысячи раз больше Median. Третий — Memory Pressure: умеренный разброс Duration, но огромный Spill и высокий GC.
Как читать каждую колонку квантильной таблицы¶
Min — самая быстрая Task. Если Min аномально маленький (почти 0), возможно некоторые Task'и получили пустые партиции.
25th Percentile — нижняя квартиль. 25% Task выполнились быстрее этого значения.
Median (50th) — средняя Task. Это наиболее репрезентативная метрика. Именно с ней нужно сравнивать Max.
75th Percentile — верхняя квартиль. Если 75th близко к Median — хорошо. Если 75th в 2-3 раза больше Median — умеренный skew.
Max — самая медленная Task. Это ваш главный индикатор. Правило:
- Max < 2× Median → нормально
- Max 2–10× Median → исследуйте
- Max > 10× Median → активная проблема (Data Skew или аппаратный сбой)
Метрики Infrastructure Overhead¶
Две метрики, которые часто игнорируют начинающие инженеры, но которые могут объяснить непонятные замедления:
Task Deserialization Time — время на десериализацию Task'а перед выполнением. Если эта метрика > 1-2 секунды, проблема в:
- Слишком большое замыкание (closure) функции — Python UDF захватывает огромный объект
- Тяжёлые broadcast-переменные (>100 MB), которые нужно десериализовать на каждом Executor'е
- Слишком много мелких Task'ов (партиций), и overhead на инициализацию каждого доминирует
Scheduler Delay — время ожидания Task'а в очереди планировщика. Если Scheduler Delay > 100-200ms и это происходит массово:
- Кластер перегружен: слишком много параллельных Job'ов
- Слишком много Task'ов — планировщик Driver'а не успевает их распределять
- Проблема с сетью между Driver'ом и Executor'ами
# Диагностика: когда Task Deserialization Time высокий
# Типичная причина — большой объект в замыкании Python UDF
# ❌ Антипаттерн: ML-модель захватывается в closure
import pickle
with open("model.pkl", "rb") as f:
model = pickle.load(f) # 500 MB модель!
@udf(returnType=FloatType())
def predict(features):
# model захвачен в closure → 500 MB сериализуется в КАЖДЫЙ Task!
return float(model.predict([features])[0])
# ✅ Правильно: SparkContext.broadcast() для больших объектов
sc = spark.sparkContext
model_broadcast = sc.broadcast(model) # один раз, не в каждый Task
@udf(returnType=FloatType())
def predict_correct(features):
m = model_broadcast.value # читается из broadcast, не из closure
return float(m.predict([features])[0])
2. Феномен Task Skew и диагностика по гистограмме¶
Data Skew — неравномерное распределение данных между партициями — является одной из самых распространённых причин катастрофической деградации производительности в production Spark-кластерах.
Физика Task Skew: как 1 из 200 Task убивает весь Stage¶
Схема показывает реальную катастрофу: 150 Executor'ов простаивают 47 минут, потребляя ресурсы кластера и деньги, пока один Task обрабатывает весь «горячий» ключ (NULL значения в join-ключе). Это классический NULL Skew.
Три источника Data Skew¶
1. NULL или дефолтные значения в ключе:
# NULL в join-ключе: все NULL строки попадают в одну партицию!
orders.join(customers, orders.customer_id == customers.id, "left")
# Если 30% orders.customer_id = NULL → одна гигантская партиция
# Решение: фильтровать NULL до join или заменять на уникальный ключ
orders_clean = orders.filter(orders.customer_id.isNotNull())
orders_with_null = orders.filter(orders.customer_id.isNull())
result = orders_clean.join(customers, "customer_id").union(
orders_with_null.withColumn("customer_name", F.lit("Unknown"))
)
2. Hot Keys — несколько пользователей с миллиардами строк:
# Крупный корпоративный клиент с 50% всех заказов
# hash("BigCorp_ID") % 200 = 42 → всегда одна партиция
events.groupBy("company_id").count()
# Партиция 42: 500M строк, остальные: ~100K строк
3. Временные паттерны — пиковые часы/дни:
# Чёрная пятница: 70% событий за 6 часов
events.groupBy(F.hour("timestamp")).count()
# Партиция hour=18, hour=19: гигантские, остальные маленькие
Диагностика Skew по метрикам Spark UI¶
# Инструмент программной диагностики: анализируем распределение данных
# ДО запуска дорогого join/groupBy
def diagnose_skew(df, key_column: str, sample_fraction: float = 0.01) -> dict:
"""
Быстрая диагностика потенциального Data Skew перед дорогой операцией.
Использует sample для скорости — не требует полного scan.
Возвращает рекомендацию: нужен ли AQE/Salting.
"""
from pyspark.sql import functions as F
# Быстрый анализ через 1% выборку
sample = df.sample(fraction=sample_fraction, seed=42)
key_counts = sample.groupBy(key_column).count()
# Статистика распределения
stats = key_counts.agg(
F.count("count").alias("unique_keys"),
F.max("count").alias("max_count"),
F.min("count").alias("min_count"),
F.percentile_approx("count", 0.5).alias("median_count"),
F.percentile_approx("count", 0.99).alias("p99_count"),
).first()
skew_ratio = stats["max_count"] / max(stats["median_count"], 1)
# Топ-10 горячих ключей
top_keys = key_counts.orderBy(F.desc("count")).limit(10).collect()
result = {
"unique_keys": stats["unique_keys"],
"skew_ratio_max_vs_median": skew_ratio,
"p99_vs_median": stats["p99_count"] / max(stats["median_count"], 1),
"top_keys": [(row[key_column], row["count"]) for row in top_keys],
}
# Рекомендации
if skew_ratio < 5:
result["recommendation"] = "OK: Skew незначительный. AQE достаточно."
elif skew_ratio < 50:
result["recommendation"] = "УМЕРЕННЫЙ SKEW: включите AQE Skew Join."
elif skew_ratio < 500:
result["recommendation"] = "ТЯЖЁЛЫЙ SKEW: AQE + ручной Salting для hot keys."
else:
result["recommendation"] = "КРИТИЧЕСКИЙ SKEW: требуется ручная обработка NULL/hot keys."
return result
Что видеть в Spark UI при Data Skew¶
На вкладке Stages нажмите на нужный Stage. Сразу видите:
- Summary Metrics → Duration: Max >> Median (в 10-1000+ раз)
- Summary Metrics → Shuffle Read Size: у одного Task объём в сотни раз больше
- Кнопка "Event Timeline" → одна полоска Task намного длиннее остальных
Прокрутите вниз до Task Table:
- Сортируйте по
Duration(убывание) → первая строка — Straggler Task - Посмотрите колонку
Records Readдля этого Task — должно быть в разы больше медианы - Посмотрите
Input Sizeдля этого Task — если это Shuffle Read, проблема в распределении ключей
3. Shuffle Read и Shuffle Write: оценка сетевого шторма¶
Shuffle — это самая дорогостоящая операция в распределённых вычислениях. Когда данные перераспределяются между Executor'ами (при groupBy, join, distinct, repartition), они сначала записываются на локальный диск Map-стороны (Shuffle Write), потом читаются по сети Reduce-стороной (Shuffle Read).
Механика Shuffle: что за данными стоит¶
Схема показывает полный путь данных при Shuffle. Каждый Task Map-стороны сортирует свои данные по хэш-ключу и пишет в отдельные файлы на локальный диск. Reduce-сторона потом читает «своё» хэш-значение со всех Map-узлов по сети. Метрика Fetch Wait Time — это время ожидания пока блоки придут от удалённых узлов.
Интерпретация метрик Shuffle в Spark UI¶
Shuffle Write Size — объём данных, записанных на локальный диск Map-стороной до передачи. Если это значение большое (>100 GB), это сигнал к уменьшению числа партиций или применению partial aggregation.
Shuffle Read Size — объём данных, полученных Reduce-стороной по сети. В идеале Shuffle Read ≈ Shuffle Write (небольшая разница из-за компрессии).
Shuffle Fetch Wait Time — критически важная метрика. Показывает сколько времени Task провёл в ожидании доставки блоков от других Executor'ов.
# Практическая диагностика Shuffle через анализ в Spark UI
# Сценарий 1: Избыточный Shuffle — слишком мало партиций
# Признаки в UI: Shuffle Read Size = 500 MB / partition (гигантские партиции)
# → каждый Task читает 500 MB = очень медленно
# Лечение: увеличить число партиций
spark.conf.set("spark.sql.shuffle.partitions", "400") # было 200
# Или через AQE (автоматически):
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
# AQE объединит маленькие партиции и разобьёт большие
# Сценарий 2: Высокий Fetch Wait Time → сетевой bottleneck
# Признаки в UI: Fetch Wait Time > 50% от Task Duration
# Причины:
# - Map-узлы перегружены: диски медленные → отдают данные медленно
# - Сеть переполнена: много параллельных Shuffle
# - spark.local.dir на медленном диске (HDD вместо NVMe)
# Диагностика:
# 1. Смотрим вкладку Executors → Shuffle Write Time
# 2. Если Write Time высокий → проблема в дисках Map-стороны
# 3. Если Write Time нормальный → проблема в сети или перегрузке Reduce-стороны
# Оптимизация:
# Переместить spark.local.dir на NVMe диски:
spark.conf.set("spark.local.dir", "/mnt/nvme0/spark,/mnt/nvme1/spark")
# Сценарий 3: External Shuffle Service — уменьшает Fetch Wait
# При Dynamic Allocation без External Shuffle Service
# когда Executor убивается, его Shuffle файлы теряются!
# Включаем External Shuffle Service (на стороне YARN/K8s):
spark.conf.set("spark.shuffle.service.enabled", "true")
# Теперь Shuffle файлы хранятся независимо от жизни Executor'а
# Сценарий 4: Reduce-сторона перегружает память
# Признаки: высокий Shuffle Fetch Wait + Spill (Disk) > 0
# Task читает слишком много данных одновременно
# Ограничение буфера чтения на Reduce-стороне:
spark.conf.set("spark.reducer.maxSizeInFlight", "48m") # дефолт 48 MB
# Меньше = меньше давление на Execution Memory, но больше запросов к Map-стороне
Когда Shuffle Write/Read — это нормально¶
Не все большие значения Shuffle плохи. Важен контекст:
| Shuffle объём | Контекст | Оценка |
|---|---|---|
| 10 GB Shuffle Write, 10 GB Read | Обработка 100 GB данных | Норма (10% shuffle overhead) |
| 500 GB Shuffle Write, 500 GB Read | 1 TB агрегация | Может быть нормально, зависит от запроса |
| 10 MB Shuffle Write, 10 GB Read | — | Аномалия! Откуда взялся Read в 1000x больше Write? |
| 100 GB Write, 100 GB Read, 45 мин | groupBy с 5 уникальными значениями | Проблема! Мало партиций → нужен repartition |
4. Борьба с дисковыми проливами: Spill (Memory) и Spill (Disk)¶
Когда Executor не может удержать промежуточные данные в Execution Memory, Spark сбрасывает их на локальный диск — это называется Spill. Spill не вызывает ошибку (в отличие от OOM), но радикально замедляет выполнение Stage.
Что за Spill скрывается в цифрах UI¶
Два значения Spill на вкладке Stages означают разные вещи:
Spill (Memory) — объём данных в несжатом виде в памяти непосредственно перед сбросом на диск. Это «честный» размер данных с точки зрения оперативной памяти.
Spill (Disk) — объём данных после компрессии на диске. Обычно в 2-5 раз меньше Spill (Memory) благодаря сжатию (Snappy, LZ4, ZSTD).
Схема показывает цикличность Spill: буфер заполняется, Spark пишет на диск, освобождает память, снова наполняет буфер. В финале читает все Spill-файлы обратно. Это огромный overhead по дисковым операциям.
Как рассчитать нужную память для устранения Spill¶
# Данные из Spark UI:
spill_memory_gb = 45 # Spill (Memory) в GB
shuffle_read_gb = 8 # Shuffle Read Size в GB
executor_memory_gb = 8 # spark.executor.memory
memory_fraction = 0.6 # spark.memory.fraction
n_cores = 5 # spark.executor.cores
# Текущая Execution Memory на Task:
usable_gb = executor_memory_gb - 0.3 # вычитаем reserved
spark_memory_gb = usable_gb * memory_fraction # unified memory
execution_per_task_gb = spark_memory_gb / n_cores # при 5 cores каждая Task ~0.9 GB
print(f"Execution Memory per Task: {execution_per_task_gb:.1f} GB")
print(f"Shuffle Read per Task: {shuffle_read_gb:.1f} GB")
print(f"Нехватка: {shuffle_read_gb - execution_per_task_gb:.1f} GB → Spill неизбежен!")
# Решение 1: увеличить executor.memory
# Нужно: shuffle_read × 2.5 (2x для sort buffer + 0.5 overhead)
needed_per_task_gb = shuffle_read_gb * 2.5
needed_total_gb = needed_per_task_gb * n_cores / memory_fraction
print(f"\nРекомендуемый executor.memory: {needed_total_gb:.0f} GB")
# Решение 2: увеличить число партиций (уменьшить данные на Task)
# Если текущих партиций 200 → нужно:
current_partitions = 200
needed_partitions = int(current_partitions * shuffle_read_gb / execution_per_task_gb) + 1
print(f"Или увеличить partitions до: {needed_partitions}")
Что делать со Spill: практические решения¶
from pyspark.sql import SparkSession
spark = SparkSession.builder.getOrCreate()
# ── Решение 1: Увеличить число партиций Shuffle ────────────────────────
# Дефолт 200 → часто недостаточно для больших данных
# Больше партиций = меньше данных на Task = меньше Spill
spark.conf.set("spark.sql.shuffle.partitions", "800")
# Или через AQE (предпочтительнее):
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.advisoryPartitionSizeInBytes", "128MB")
# AQE само выберет нужное число партиций
# ── Решение 2: Увеличить executor.memory ──────────────────────────────
# Но осторожно с GC! Большой heap → более длинные GC паузы
# Правило: не больше 32 GB на Executor (G1GC деградирует при > 32 GB heap)
# ── Решение 3: Уменьшить данные ДО Shuffle (Pushdown) ──────────────────
# Вместо groupBy всей таблицы → агрегируй только нужное
# Без оптимизации:
df.groupBy("country", "category", "year", "month", "day").sum("revenue")
# С оптимизацией через Project Pushdown:
df.select("country", "category", "revenue") \
.groupBy("country", "category") \
.sum("revenue")
# Меньше колонок → меньше данных в Shuffle → меньше Spill
# ── Решение 4: spark.local.dir на NVMe ────────────────────────────────
# Если Spill неизбежен (очень большие данные) — хотя бы сделать его быстрым
# Spill на NVMe в 10-50x быстрее чем на HDD
spark.conf.set("spark.local.dir", "/mnt/nvme0/spark,/mnt/nvme1/spark")
# ── Решение 5: Включить компрессию Shuffle для уменьшения I/O ──────────
spark.conf.set("spark.shuffle.compress", "true")
spark.conf.set("spark.shuffle.spill.compress", "true")
spark.conf.set("spark.io.compression.codec", "lz4") # lz4 быстрее чем snappy
5. GC Time: скрытый убийца производительности¶
GC Time (Garbage Collection Time) — это время, которое JVM потратила на очистку неиспользуемой памяти вместо выполнения полезной работы. Эта метрика отображается на вкладке Stages и часто игнорируется начинающими инженерами, хотя высокий GC Time может сделать задачу в 5-10x медленнее.
Пороговые значения GC Time¶
GC Time % от Task Duration:
< 5% → Норма. JVM работает эффективно
5-10% → Допустимо, но стоит обратить внимание
10-20% → Проблема. Много краткоживущих объектов или нехватка памяти
> 20% → Критично! Кластер тратит 1/5 времени на "уборку мусора"
Причины высокого GC Time¶
Причина 1: Тяжёлые Python UDF создают миллионы Java-объектов
# ❌ Причина высокого GC: классический Python UDF
from pyspark.sql.functions import udf
from pyspark.sql.types import FloatType
@udf(returnType=FloatType())
def calculate_score(amount, category, region):
# Python ↔ JVM сериализация для КАЖДОЙ строки
# Spark создаёт Java Row объект, сериализует, передаёт в Python,
# Python обрабатывает, сериализует результат обратно в JVM
# 1M строк = 2M сериализаций = 2M+ Java объектов → GC storm!
score = amount * 0.5
if category == "Premium":
score *= 1.5
return float(score)
df.withColumn("score", calculate_score("amount", "category", "region"))
# ✅ Решение 1: Vectorized UDF (Pandas UDF) — batch передача через Apache Arrow
from pyspark.sql.functions import pandas_udf
import pandas as pd
@pandas_udf(FloatType())
def calculate_score_vectorized(amount: pd.Series,
category: pd.Series,
region: pd.Series) -> pd.Series:
# Данные приходят как pandas Series (колоночный батч)
# Один Arrow batch = тысячи строк → минимум Java объектов
score = amount * 0.5
score = score.where(category != "Premium", score * 1.5)
return score.astype(float)
df.withColumn("score", calculate_score_vectorized("amount", "category", "region"))
# ✅ Решение 2: Заменить на нативные Spark SQL функции (нет Python вообще!)
from pyspark.sql import functions as F
df.withColumn("score",
F.when(F.col("category") == "Premium",
F.col("amount") * 0.75)
.otherwise(F.col("amount") * 0.5)
)
# Этот код работает только на JVM в нативном формате Tungsten → нет GC!
Причина 2: Кешированные данные в неправильном StorageLevel
from pyspark import StorageLevel
# ❌ MEMORY_ONLY — данные хранятся как Java-объекты в Heap
# Куча JVM забита сложными объектами → GC должен обходить весь граф объектов
df.persist(StorageLevel.MEMORY_ONLY)
# ✅ MEMORY_ONLY_SER — данные сериализованы как byte[]
# byte[] = простые примитивные массивы → GC обходит их мгновенно
# Минус: небольшой overhead на де/сериализацию при чтении из кеша
# Но GC overhead значительно меньше
df.persist(StorageLevel.MEMORY_ONLY_SER)
# ✅ OFF_HEAP — данные вообще вне JVM Heap!
# GC вообще не видит эти данные → нулевой GC overhead от кеша
spark.conf.set("spark.memory.offHeap.enabled", "true")
spark.conf.set("spark.memory.offHeap.size", "10g")
df.persist(StorageLevel.OFF_HEAP)
Причина 3: Слишком большой Executor heap с G1GC
G1GC — стандартный GC в JVM 8-17. Он хорошо работает для heap до 32 GB. Больше — и время Full GC паузы начинает расти нелинейно.
spark = SparkSession.builder \
# Если executor.memory > 32 GB — переключайтесь на ZGC
.config("spark.executor.extraJavaOptions",
"-XX:+UseZGC " # ZGC: лучше для больших heap
"-XX:+UnlockExperimentalVMOptions "
"-Xms4g " # начальный heap = конечный (избегаем resize)
"-XX:+PrintGCDetails " # для логирования GC
"-XX:+PrintGCDateStamps") \
# Или G1GC с тюнингом для среднего heap (8-32 GB):
# .config("spark.executor.extraJavaOptions",
# "-XX:+UseG1GC "
# "-XX:G1HeapRegionSize=16m " # большие регионы для больших heap
# "-XX:+G1UseAdaptiveIHOP "
# "-XX:InitiatingHeapOccupancyPercent=35 " # GC начинает раньше
# "-XX:MaxGCPauseMillis=200") # целевая пауза 200 мс
.getOrCreate()
Диагностика GC через Spark UI¶
В вкладке Stages смотрите на Summary Metrics:
GC Time — 3.2 sec | Median
Task Duration — 12 sec | Median
Расчёт: GC% = 3.2 / 12 = 26.7% → КРИТИЧНО!
Если GC% > 10%, переходим к анализу:
def analyze_gc_issue(gc_time_ms: float, task_duration_ms: float,
executor_memory_gb: float) -> str:
"""
Диагностирует причину высокого GC и предлагает решение.
"""
gc_pct = gc_time_ms / task_duration_ms * 100
if gc_pct < 5:
return "GC нормальный"
elif gc_pct < 10:
return "Умеренный GC. Наблюдайте за трендом."
elif gc_pct < 20:
if executor_memory_gb > 32:
return (f"Высокий GC ({gc_pct:.0f}%). Executor {executor_memory_gb}GB > 32GB. "
"Попробуйте ZGC или уменьшите executor.memory.")
else:
return (f"Высокий GC ({gc_pct:.0f}%). Вероятно: Python UDF без Vectorization "
"или MEMORY_ONLY кеш. Переходите на Pandas UDF и MEMORY_ONLY_SER.")
else:
return (f"КРИТИЧЕСКИЙ GC ({gc_pct:.0f}%)! Немедленные действия: "
"1) Отключите все Python UDF, замените на SQL функции. "
"2) Переключите кеш на MEMORY_ONLY_SER или OFF_HEAP. "
"3) Уменьшите executor.memory (меньше heap = быстрее GC). "
"4) Включите spark.memory.offHeap.enabled=true.")
6. Декомпозиция Task Duration: Executor Computing Time vs Overhead¶
Общее время выполнения Task складывается из нескольких слагаемых. Понимание их соотношения помогает найти правильное направление оптимизации.
Анатомия Task Duration¶
Схема показывает что только 45% времени Task тратит на реальную работу. Оставшиеся 55% — инфраструктурный overhead. Оптимизация должна быть направлена на снижение overhead, а не на ускорение алгоритма вычислений (его в данном случае уже нет смысла оптимизировать).
Что означает каждое слагаемое¶
# Диагностика через Spark History Server API
import requests
def analyze_task_breakdown(app_id: str, stage_id: int,
history_server: str) -> dict:
"""
Декомпозирует время Task'ов для нахождения главного источника overhead.
Возвращает процентное распределение компонент Duration.
"""
url = (f"{history_server}/api/v1/applications/{app_id}"
f"/stages/{stage_id}/0/taskList?numTasks=100")
tasks = requests.get(url, timeout=30).json()
# Агрегируем по медиане
def median(values):
sorted_v = sorted(values)
n = len(sorted_v)
return sorted_v[n // 2] if n else 0
durations = [t.get("duration", 0) for t in tasks]
sched_delays = [t.get("schedulerDelay", 0) for t in tasks]
deser_times = [t.get("taskDeserializationTime", 0) for t in tasks]
gc_times = [t.get("jvmGcTime", 0) for t in tasks]
compute_times = [t.get("executorRunTime", 0) for t in tasks]
fetch_waits = [
t.get("taskMetrics", {})
.get("shuffleReadMetrics", {})
.get("fetchWaitTime", 0)
for t in tasks
]
median_dur = max(median(durations), 1)
breakdown = {
"median_duration_ms": median_dur,
"scheduler_delay_pct": median(sched_delays) / median_dur * 100,
"deserialization_pct": median(deser_times) / median_dur * 100,
"gc_time_pct": median(gc_times) / median_dur * 100,
"fetch_wait_pct": median(fetch_waits) / median_dur * 100,
"compute_pct": median(compute_times) / median_dur * 100,
}
# Определяем главный источник проблем
issues = []
if breakdown["gc_time_pct"] > 10:
issues.append(f"GC ({breakdown['gc_time_pct']:.0f}%) → Memory pressure")
if breakdown["fetch_wait_pct"] > 20:
issues.append(f"Fetch Wait ({breakdown['fetch_wait_pct']:.0f}%) → Network/Disk bottleneck")
if breakdown["scheduler_delay_pct"] > 5:
issues.append(f"Scheduler Delay ({breakdown['scheduler_delay_pct']:.0f}%) → Too many tasks")
if breakdown["deserialization_pct"] > 10:
issues.append(f"Deserialization ({breakdown['deserialization_pct']:.0f}%) → Heavy closure")
breakdown["main_issues"] = issues if issues else ["Нет явных проблем"]
return breakdown
7. Анализ Task Table: диагностика по отдельным Task'ам¶
В самом низу страницы Stage находится таблица всех Task'ов. Это место для "судебной экспертизы" — когда агрегированные метрики уже указали на проблему и нужно найти конкретный виновник.
Как работать с Task Table¶
Колонки Task Table (полный список с кнопкой "Show Additional Metrics"):
| Колонка | Что показывает | Когда важна |
|---|---|---|
Duration |
Общее время Task | Ищем выбросы (Stragglers) |
GC Time |
Время GC для этого Task | Если Stage медленный без видимых причин |
Input Size |
Размер входных данных | Подтверждение Data Skew |
Shuffle Read Size |
Данные прочитанные из Shuffle | Главный виновник при join/groupBy |
Shuffle Write Size |
Данные записанные в Shuffle | Map-сторона — до передачи по сети |
Spill (Memory) |
Spill для конкретного Task | Какой именно Task страдает |
Locality Level |
Где выполнялась Task | Эффективность Data Locality |
Executor ID |
На каком Executor | Выявить "плохой" Executor |
Host |
На каком физическом сервере | Аппаратные проблемы |
Errors |
Ошибки и retry | Нестабильность кластера |
Практика: поиск "битого" Executor по Task Table¶
# Симуляция: кластер с одним медленным Executor'ом
from pyspark.sql import SparkSession, functions as F
spark = SparkSession.builder \
.master("local[4]") \
.appName("stages-demo") \
.getOrCreate()
# Создаём данные для демонстрации разных паттернов
df = spark.range(2_000_000).select(
F.col("id"),
# Имитация Data Skew: 40% строк с key=0
F.when(F.rand() < 0.4, F.lit(0))
.otherwise((F.col("id") % 50).cast("int"))
.alias("join_key"),
(F.rand() * 1000).alias("revenue"),
)
# Stage с Data Skew (без AQE)
spark.conf.set("spark.sql.adaptive.enabled", "false")
spark.conf.set("spark.sql.shuffle.partitions", "50")
# groupBy с skew → один Task будет обрабатывать 40% данных!
result = df.groupBy("join_key").agg(F.sum("revenue"))
result.write.mode("overwrite").parquet("/tmp/result/")
# Проверьте в Spark UI → Stages:
# 1. Summary Metrics → Duration Max >> Median (в 10-20x)
# 2. Task Table → сортируйте по Duration desc
# 3. Первый Task: Input Size >> остальные
# 4. Этот Task обрабатывает join_key=0 (40% данных)
print("""
ЗАДАНИЕ:
1. Откройте http://localhost:4040 → Stages
2. Нажмите на Stage с groupBy
3. Проверьте Summary Metrics:
- Duration: Max vs Median — есть ли 10x разрыв?
- Input Size: Max vs Median — какой Task получил больше всего данных?
4. Включите 'Show Additional Metrics' в Task Table
5. Отсортируйте по Duration (desc) — нашли Straggler?
6. Посмотрите Locality Level — есть ли ANY (данные с удалённого узла)?
""")
spark.stop()
Шаблоны поиска проблем через Task Table¶
Симптом: один Task с Duration в 50x больше медианы
→ Сортируйте по Duration desc
→ Смотрите Shuffle Read Size этого Task
→ Если размер в 50x больше → Data Skew по ключу Shuffle
→ Смотрите Input Size → если огромный → неравные партиции исходных данных
Симптом: несколько Task'ов с Failed Status
→ Фильтруйте по Status = Failed
→ Смотрите Error Message
→ Все на одном Executor/Host → проблема с узлом кластера
→ Разные узлы → проблема в коде или данных
Симптом: все Task'ы медленные (равномерно)
→ Сортируйте по GC Time desc
→ Если GC Time > 10% от Duration → Memory Pressure
→ Сортируйте по Shuffle Fetch Wait → если высокий → сетевой/дисковый bottleneck
Симптом: Task'ы на удалённых Executor'ах работают медленнее
→ Смотрите Locality Level column
→ Много "ANY" или "RACK_LOCAL" → данные не локальны
→ Решение: правильная конфигурация Data Locality (spark.locality.wait)
Итоги: главные выводы о вкладке Stages¶
Таблица квантилей — первый инструмент диагностики. Max/Median ratio > 10 = проблема требующая немедленного внимания.
Три главных виновника медленных Stage'ов:
- Data Skew: Max Duration >> Median + Max Input/Shuffle >> Median
- Memory Pressure: Spill (Disk) > 0 + GC Time > 10%
- Network/Disk bottleneck: Shuffle Fetch Wait > 20% Task Duration
Декомпозиция Task Duration: найдите что реально тратит время — Fetch Wait, GC, Compute или Scheduler. Оптимизируйте самую большую долю.
Task Table для точной локализации: отсортируйте по Duration, найдите Straggler, проверьте его Input Size и Shuffle Read для подтверждения диагноза.
Алгоритм диагностики Stage за 5 минут:
1. Открыть Stage → Summary Metrics
2. Max/Median Duration ratio?
> 10x → Data Skew → AQE + Salting
< 3x → не Data Skew, смотрим дальше
3. Spill (Disk) > 0?
Да → нехватка Execution Memory → больше partitions или executor.memory
4. GC Time % > 10%?
Да → Memory pressure → MEMORY_ONLY_SER, Pandas UDF, ZGC
5. Shuffle Fetch Wait % > 20%?
Да → Network/Disk bottleneck → NVMe диски, меньше partition size
6. Task Table → найти конкретный Task → подтвердить диагноз