Spark UI: вкладка Jobs — DAG визуализация и поиск узкого места

Полный разбор вкладки Jobs в Spark UI: иерархия Job→Stage→Task, DAG визуализация и narrow/wide зависимости, Event Timeline и Dynamic Allocation, механизм Skipped стадий и кеширование, Speculative Execution и Stragglers, алгоритм поиска узкого места за 3 минуты.

platform

1. Архитектура Job в Spark и анатомия вкладки Jobs

Spark UI — это не просто красивый дашборд. Это главный инструмент диагностики производительности, который отличает инженера, умеющего «гуглить ошибки», от эксперта, способного за три минуты найти причину деградации production-пайплайна.

Вкладка Jobs является отправной точкой любого расследования. Именно здесь формируется глобальная картина выполнения приложения, и именно здесь нужно начинать — а не сразу погружаться в низкоуровневые метрики вкладок Stages или SQL.

Как Spark превращает код в Jobs

Spark работает по принципу ленивых вычислений (lazy evaluation). Каждый вызов read(), filter(), groupBy(), join() и других трансформаций только строит логический план — ничего не вычисляется. Вычисление начинается только при вызове Actions — операций, требующих материализации результата.

Схема показывает: три Action создают три независимых Job. Каждый Job запустит полную цепочку вычислений заново — если DataFrame не был кеширован.

Иерархия выполнения: Job → Stage → Task

Это фундаментальная концепция, без которой нельзя правильно читать Spark UI:

Job (1 Action)
  ├── Stage 0 (чтение файлов + фильтрация)
  │     ├── Task 0 (Partition 0 → один файл или его часть)
  │     ├── Task 1 (Partition 1)
  │     └── ... Task N
  │
  ├── Stage 1 (после Shuffle: groupBy aggregation)
  │     ├── Task 0 (hash partition 0)
  │     ├── Task 1 (hash partition 1)
  │     └── ... Task N
  │
  └── Stage 2 (после Shuffle: join)
        ├── Task 0
        └── ...

Job — это единица работы, соответствующая одному Action. В Spark UI вкладка Jobs показывает список всех Job.

Stage — это группа Task'ов, которые могут выполняться без промежуточного Shuffle. Граница между Stage'ами — это всегда Shuffle Exchange (перераспределение данных по сети). Stage'ов в одном Job может быть много — по одному на каждую «широкую» (wide) зависимость.

Task — минимальная единица параллелизма. Каждый Task обрабатывает одну Partition данных. Количество Task'ов в Stage = количество Partition входных данных этого Stage.

Интерфейс вкладки Jobs: что видно

На главной странице вкладки Jobs отображается таблица всех Job:

ID  | Description        | Submitted | Duration | Stages | Tasks
----|--------------------|-----------|----------|--------|------
0   | count at Main.py:42| 10:01:15  | 45 s     | 3/3    | 1200/1200
1   | save at Main.py:57 | 10:02:05  | 2m 30s   | 4/4    | 2400/2400
2   | show at Main.py:63 | 10:04:40  | 3 s      | 2/2    | 200/200

Колонка Duration — ваш первый маяк. Найдите самый долгий Job и начните с него. Остальные Job в большинстве случаев быстрые и не требуют внимания.


2. Чтение DAG на уровне Job: narrow vs wide зависимости

Нажав на любой Job, вы попадаете на страницу с детальным описанием всех Stage'ов и визуализацией DAG. Это самое информативное место Spark UI для понимания того, что именно делает ваш код.

Narrow и Wide зависимости: основа формирования Stage'ов

Spark разделяет все зависимости между RDD/DataFrame на два типа, и это напрямую определяет структуру DAG:

Narrow зависимость: каждая партиция входных данных используется максимум одной партицией выходных данных. Примеры: filter(), map(), select(), withColumn(). Такие операции объединяются в одну Stage — это называется pipelining.

Wide зависимость: одна партиция выходных данных может зависеть от многих партиций входных данных. Это всегда требует Shuffle — физической передачи данных по сети. Примеры: groupBy(), join(), distinct(), repartition(). Wide зависимость создаёт границу между Stage'ами.

Типичная структура DAG для join-запроса

Рассмотрим запрос, который часто встречается в ETL:

df = spark.read.parquet("s3://bucket/events/")     # Stage 0 begins
filtered = df.filter(df.amount > 100)              # ↑ Narrow
enriched = filtered.join(F.broadcast(products),    # BHJ: остаётся в Stage 0
                         "product_id")
grouped = enriched.groupBy("category").sum("amount")  # Wide! Stage 1 begins
grouped.write.parquet("s3://bucket/result/")

В Spark UI этот код создаст DAG примерно такого вида:

Stage 0: Scan parquet → Filter → BroadcastHashJoin → Partial Aggregation
Stage 1: Exchange (hashpartitioning) → Sort → Aggregation → Write

Обратите внимание: BroadcastHashJoin остаётся в Stage 0, потому что broadcast не требует Shuffle — небольшая таблица рассылается всем Executor'ам заранее. А groupBy вызывает Exchange (Shuffle), разрывая DAG.

Что искать в DAG визуализации

Открыв страницу Job и увидев DAG, смотрите на:

  1. Количество Stage'ов: много Stage = много Shuffle = много потенциальных узких мест
  2. Симметрию: Stage'ы должны иметь похожее число Task'ов. Резкий дисбаланс указывает на проблему
  3. Цвет Stage'ов: зелёный = завершён, синий = выполняется, серый = Skipped
  4. Узлы с иконкой обмена (Exchange): каждый такой узел — это Shuffle, который дорого стоит

3. Поиск «зависших» Jobs и выявление Stragglers

В реальных production-пайплайнах часто встречается ситуация: 99% задач завершились за минуту, но Job продолжает висеть в статусе Running уже час. Виновник — одна-единственная Task, которую называют Straggler (отстающая).

Почему возникают Stragglers

На вкладке Jobs прогресс-бар Stage'а показывает 199/200 completed и останавливается. Task 2 содержит 50 миллионов строк вместо нормальных 200 тысяч — классический Data Skew.

Как обнаружить Straggler на вкладке Jobs

  1. Прогресс-бар показывает что Stage «почти завершён», но не двигается
  2. Duration у Job растёт, хотя большинство Stage'ов зелёные
  3. Нажмите на Stage → посмотрите Summary Metrics → колонка Duration: если Max Task Time в десятки раз больше Median Task Time — это Straggler

Speculative Execution: попытка автоматического лечения

Spark имеет механизм спекулятивного выполнения, который автоматически перезапускает медленные задачи на других узлах:

spark = SparkSession.builder \
    # Включает спекулятивное выполнение
    .config("spark.speculation", "true") \

    # Задача считается "медленной" если выполняется в multiplier раз
    # дольше медианы всех задач данного Stage
    .config("spark.speculation.multiplier", "1.5") \

    # Запускать спекуляцию только если завершено >= quantile задач
    .config("spark.speculation.quantile", "0.75") \

    .getOrCreate()

Когда спекуляция помогает: медленный узел (hardware issue), сетевой сбой, временная перегрузка DataNode.

Когда спекуляция не поможет: Data Skew (одна задача реально имеет больше данных — запуск на другом узле займёт столько же времени).

На вкладке Jobs спекулятивные задачи видны по суффиксу в ID: обычная Task имеет ID 42, её спекулятивная копия — 42.1.


4. Временная шкала событий: Event Timeline

Нажав кнопку Show Additional MetricsEvent Timeline на странице Job, вы увидите горизонтальную временную шкалу. Это мощный инструмент для понимания жизни вашего приложения во времени.

Что показывает Event Timeline на уровне Jobs

Что можно прочитать с Event Timeline:

  • Паузы между Job'ами — если между окончанием Job N и началом Job N+1 есть заметная пауза, это время тратит Driver: запросы к Hive Metastore, Python-обработка данных на Driver, ожидание внешних API
  • Циклы добавления/удаления Executor'ов — частое добавление и немедленное удаление указывает на неправильную настройку Dynamic Allocation
  • Executor gaps — периоды когда Executor'ы есть, но не работают (синий цвет = idle) — потенциальное место для оптимизации

Dynamic Allocation: паттерны в Event Timeline

spark = SparkSession.builder \
    # Dynamic Allocation: Spark сам добавляет/удаляет Executor'ы
    .config("spark.dynamicAllocation.enabled", "true") \

    # Минимальное число Executor'ов (всегда доступно)
    .config("spark.dynamicAllocation.minExecutors", "2") \

    # Максимальное число Executor'ов
    .config("spark.dynamicAllocation.maxExecutors", "50") \

    # Удалять Executor если он простаивал N секунд
    # ПРОБЛЕМА: если установить слишком маленькое значение →
    # кластер постоянно убивает и создаёт Executor'ы
    # → overhead на JVM startup, потеря кешированных данных!
    .config("spark.dynamicAllocation.executorIdleTimeout", "120s") \

    # Executor с кешированными данными — дольше ждём
    .config("spark.dynamicAllocation.cachedExecutorIdleTimeout", "3600s") \

    .getOrCreate()

Антипаттерн: если на Event Timeline видно что Executor'ы добавляются и удаляются каждые 30-60 секунд — executorIdleTimeout слишком маленький. Каждый перезапуск JVM занимает 10-30 секунд, плюс теряется LocalState (кешированные данные, shuffle файлы).


5. Понимание Skipped и Failed стадий

Skipped Stage: маркер правильного кеширования

В DAG визуализации серые Stage'ы с надписью «Skipped» — это хорошая новость. Они означают, что Spark повторно использует уже материализованные данные.

Когда Stage становится Skipped:

df = spark.read.parquet("s3://bucket/large_dataset/")  # 500 GB
filtered = df.filter(df.year == 2024)

# Кешируем дорогостоящий промежуточный результат
filtered.cache()
filtered.count()  # Job 0: материализует кеш
# ↑ Stage 0 (Scan + Filter) → выполнен, данные в памяти

# Повторное использование
result1 = filtered.groupBy("region").count()
result1.show()  # Job 1: Stage 0 → SKIPPED! читаем из кеша

result2 = filtered.join(another_df, "user_id").count()
result2.count()  # Job 2: Stage 0 → SKIPPED! снова из кеша

В Spark UI Job 1 и Job 2 будут показывать «Stage 0: Skipped» — визуальное подтверждение что .cache() работает.

Когда отсутствие Skipped — проблема

Если вы добавили .cache(), но Stage всё равно не Skipped — ищите причину:

# Проверяем что кеш реально сохранился
df.cache()
df.count()  # материализуем

# Проверяем в каталоге:
print(spark.catalog.isCached("my_table"))  # для таблиц

# Через Storage вкладку: смотрим на RDD Blocks
# Если видим "0 Blocks Cached" → кеш не занял место

# Типичные причины "кеш не работает":
# 1. Слишком большой DataFrame, не помещается в Storage Memory
# 2. Executor был убит, кеш потерян
# 3. .unpersist() вызван раньше времени
# 4. .cache() вызван но Action не запущен

Failed Stage: читаем ошибки правильно

Failed Stage — худшее что может показать вкладка Jobs. Красный Stage означает что хотя бы одна Task завершилась с ошибкой и все retry исчерпаны.

Вкладка Jobs агрегирует ошибки на верхний уровень. Нажав на Failed Stage, вы увидите сообщение об ошибке. Типичные паттерны:

Ошибка в UI Вероятная причина
java.lang.OutOfMemoryError: Java heap space Executor heap слишком мал для данной задачи
Container killed by YARN for exceeding memory limits memoryOverhead недостаточен (особенно PySpark)
org.apache.spark.shuffle.FetchFailedException DataNode с shuffle данными стал недоступен
java.net.SocketTimeoutException Timeout при чтении данных с удалённого DataNode
BlockMissingException Блок HDFS недоступен (DataNode упал)

6. Влияние кеширования на структуру DAG

DAG до и после .cache()

Кеширование кардинально меняет визуальный вид DAG. Понимание этого изменения позволяет быстро верифицировать что кеш работает правильно.

На левой схеме (без кеша): Job 2 полностью пересчитывает всё от чтения файлов. На правой схеме (с кешем): DAG начинается с зелёного узла InMemoryTableScan, а прошлые Stage'ы пропускаются.

Антипаттерн: избыточное кеширование

# Антипаттерн: кешируем всё подряд
df1 = spark.read.parquet("s3://bucket/orders/").cache()
df2 = spark.read.parquet("s3://bucket/customers/").cache()
df3 = df1.join(df2, "customer_id").cache()  # уже кешированы df1 и df2 → тройное занятие памяти!
df4 = df3.groupBy("region").count().cache()  # и это тоже...

# Последствия:
# 1. Storage Memory переполнена
# 2. Execution Memory вытеснена → больше Spill на диск
# 3. GC давление растёт из-за большого heap
# 4. В Spark UI → Storage вкладка: видим что % кеша в памяти уменьшается

# Правильный подход: кешировать только то что используется 2+ раз
# и занимает дорогостоящие вычисления для пересчёта

Кеш имеет смысл только если:

  1. DataFrame будет использоваться минимум 2 раза
  2. Вычисление DataFrame дороже чем десериализация из кеша
  3. У вас достаточно Storage Memory (смотрите на вкладку Storage)

7. Практический алгоритм поиска узкого места

Это системный подход, который позволяет за 3-5 минут локализовать проблему в любом производственном пайплайне.

Алгоритм: от общего к частному

Шаг 1: Найти самый долгий Job

# В Spark UI: Jobs → сортируем по Duration (клик на заголовок колонки)
# Или программно через History Server API:
import requests

def find_slowest_jobs(app_id: str, history_server_url: str) -> list[dict]:
    """
    Находит 5 самых медленных Job'ов в приложении.
    Полезно для автоматического мониторинга SLA.
    """
    url = f"{history_server_url}/api/v1/applications/{app_id}/jobs"
    jobs = requests.get(url, timeout=30).json()

    completed_jobs = [j for j in jobs if j.get("status") == "SUCCEEDED"]
    sorted_jobs = sorted(
        completed_jobs,
        key=lambda j: j.get("completionTime", 0) - j.get("submissionTime", 0),
        reverse=True
    )

    return [{
        "jobId": j["jobId"],
        "name": j.get("name", "")[:60],
        "duration_ms": j.get("completionTime", 0) - j.get("submissionTime", 0),
        "stages": j.get("numCompletedStages", 0),
        "tasks": j.get("numCompletedTasks", 0),
    } for j in sorted_jobs[:5]]

Шаг 2: Найти критический путь в DAG

Не все Stage влияют на итоговое время одинаково. Некоторые Stage'ы выполняются параллельно — их ускорение не уменьшит общее время Job'а.

Критический путь — самая длинная цепочка зависимых Stage'ов. Только ускорение Stage'ов в критическом пути реально ускоряет Job.

Пример DAG:
Stage 0: Scan facts (10 сек) ─────────────────────┐
Stage 1: Scan dim_1 (2 сек) → Join (5 сек) ────── join ─→ Stage 4: Write (3 сек)
Stage 2: Scan dim_2 (1 сек) → Join (3 сек) ───────┘     (критический путь: 10+5+3=18 сек)

Ускорение Stage 2 (1+3=4 сек) до 0 → итог: 18 сек (без изменений!)
Ускорение Stage 0 (10 сек) до 3 сек → итог: 11 сек (значительное улучшение!)

Шаг 3: Симптомы проблем видные из Jobs

Опытный инженер может поставить предварительный диагноз уже на уровне Jobs, не заходя в Stages:

Симптом → Вероятная причина → Следующий шаг

"Stage с 200 Tasks, но только 2 быстро, остальные висят"
→ Data Skew или ресурсов недостаточно
→ Смотреть Task distribution в Stage (медиана vs максимум)

"Jobs выполняются 10+ секунд с 0 Tasks"
→ Overhead планирования: HMS запросы, Python init, Driver computation
→ Смотреть Event Timeline, логи Driver'а

"Много Skipped Stages, но Job всё равно медленный"
→ Кеш неэффективен: данные не в Storage Memory (на диске или evicted)
→ Смотреть Storage вкладку: Fraction Cached, Storage Level

"Job Failed с Exit Code 137"
→ OOM Killer убил контейнер (превышен memory limit)
→ Увеличить executor.memoryOverhead или уменьшить executor.memory

"Stage с огромным Shuffle Read (> 100 GB)"
→ Слишком мало партиций → каждая партиция огромная
→ Увеличить spark.sql.shuffle.partitions или включить AQE

Полный чек-лист анализа вкладки Jobs

□ 1. Открыть Jobs → отсортировать по Duration
□ 2. Найти Job с максимальным Duration → нажать
□ 3. Проверить: есть ли Failed Stage'ы? (красный цвет)
□ 4. Если есть Failed → прочитать ошибку, найти Root Cause
□ 5. Посмотреть DAG: сколько Stage'ов? Есть ли широкие зависимости?
□ 6. Найти самый долгий Stage (по цвету/ширине полосы)
□ 7. Проверить: Stage Skipped? → верификация кеширования
□ 8. Нажать на проблемный Stage → перейти на вкладку Stages
□ 9. В Stage проверить: Task Distribution (Median vs Max)
□ 10. Если Max >> Median → Data Skew, применить AQE или salting
□ 11. Если равномерно медленно → Shuffle, Memory, I/O проблема
□ 12. Записать гипотезу и проверить через конкретные метрики

Практика: разбор реального Job

Создаём тестовый сценарий с проблемами

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

spark = SparkSession.builder \
    .master("local[4]") \
    .appName("spark-ui-demo") \
    .config("spark.sql.shuffle.partitions", "200") \
    .config("spark.ui.enabled", "true") \
    .getOrCreate()

spark.sparkContext.setLogLevel("WARN")

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


# ── ДЕМОНСТРАЦИЯ 1: Лишние пересчёты без .cache() ──────────────────────
df = spark.range(10_000_000).select(
    F.col("id"),
    (F.rand() * 100).cast("int").alias("category"),
    (F.rand() * 1000).alias("revenue"),
)

# Job 0 (Action 1 без cache): пересчёт от нуля
count1 = df.count()  # Job 0: полный scan
# Job 1 (Action 2 без cache): ещё раз от нуля!
result1 = df.groupBy("category").sum("revenue")
result1.count()  # Job 1: полный scan + groupBy

# Смотрите в Spark UI → Jobs:
# Job 0 и Job 1 имеют почти одинаковые DAG → дублирование!

# ── ДЕМОНСТРАЦИЯ 2: С .cache() → Skipped Stage'ы ─────────────────────────
df_cached = spark.range(10_000_000).select(
    F.col("id"),
    (F.rand() * 100).cast("int").alias("category"),
    (F.rand() * 1000).alias("revenue"),
)
df_cached.cache()

count2 = df_cached.count()   # Job 2: материализация кеша
result2 = df_cached.groupBy("category").sum("revenue")
result2.count()  # Job 3: Stage 0 → SKIPPED! Намного быстрее!

# ── ДЕМОНСТРАЦИЯ 3: Data Skew → Straggler Task ───────────────────────────
from pyspark.sql.types import LongType

# Скошенные данные: 90% строк с key=0
skewed_data = spark.range(1_000_000).select(
    F.when(F.rand() < 0.9, F.lit(0)).otherwise(F.col("id") % 100).alias("join_key"),
    (F.rand() * 100).alias("value"),
)
lookup = spark.range(100).select(
    F.col("id").alias("join_key"),
    F.concat(F.lit("cat_"), F.col("id")).alias("category")
)

# Без AQE: один Task получит 90% данных → Straggler!
spark.conf.set("spark.sql.adaptive.enabled", "false")
joined = skewed_data.join(lookup, "join_key")
joined.count()  # Job 4: смотрите один медленный Task в Stages!

# С AQE: Bloom Filter + Skew Join → равномерно
spark.conf.set("spark.sql.adaptive.enabled", "true")
joined_aqe = skewed_data.join(lookup, "join_key")
joined_aqe.count()  # Job 5: намного быстрее!


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

Вкладка Jobs:
1. Сравните Duration Job 0 и Job 2 (с cache vs без)
2. Найдите Job 3 → убедитесь что есть Skipped Stage'ы
3. Сравните Job 4 (без AQE) и Job 5 (с AQE)
   - Job 4 → нажмите на Stage с Join → посмотрите Distribution
   - Должен быть один Task намного медленнее остальных

Вкладка DAG:
- Нажмите на любой Job → посмотрите граф Stage'ов
- Найдите Exchange (Shuffle) узлы на границах Stage'ов
- Найдите InMemoryTableScan в Job 3 (зелёный узел кеша)
""")

spark.stop()

Что вы увидите в Spark UI

После запуска этого кода в Spark UI появятся 6 Job'ов. Анализируя их вкладку Jobs, вы научитесь:

  • Job 0 vs Job 2: сравнение длительности — Job 2 значительно быстрее благодаря кешу
  • Job 3: Stage 0 помечен как Skipped — визуальное подтверждение работы кеша
  • Job 4: один Task занимает 90% времени Stage → классический Data Skew
  • Job 5: равномерное распределение благодаря AQE Skew Join

Итоги: главные выводы о вкладке Jobs

Вкладка Jobs — это диспетчерский пульт: она даёт общую картину и позволяет быстро определить где тратится время приложения.

Трёхминутный алгоритм диагностики: сортировка Jobs по Duration → вход в самый долгий Job → просмотр DAG → переход на проблемный Stage.

Skipped Stage — ваш друг: наличие серых Stage'ов подтверждает что кеширование работает правильно. Их отсутствие там где они должны быть — повод проверить Storage вкладку.

Straggler Task — главный симптом Data Skew: если прогресс-бар Stage останавливается на «199/200», виновник — Data Skew. Следующий шаг — AQE Skew Join или ручное сальтирование ключей.

DAG — перевод кода в план: научившись читать DAG, вы видите не строки Python-кода, а реальные операции над данными. Каждый узел Exchange — это дорогой Shuffle, каждый зелёный узел — сэкономленные вычисления из кеша.