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.

optimization

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. Сразу видите:

  1. Summary Metrics → Duration: Max >> Median (в 10-1000+ раз)
  2. Summary Metrics → Shuffle Read Size: у одного Task объём в сотни раз больше
  3. Кнопка "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'ов:

  1. Data Skew: Max Duration >> Median + Max Input/Shuffle >> Median
  2. Memory Pressure: Spill (Disk) > 0 + GC Time > 10%
  3. 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 → подтвердить диагноз