Spark UI: вкладка Jobs — DAG визуализация и поиск узкого места
Полный разбор вкладки Jobs в Spark UI: иерархия Job→Stage→Task, DAG визуализация и narrow/wide зависимости, Event Timeline и Dynamic Allocation, механизм Skipped стадий и кеширование, Speculative Execution и Stragglers, алгоритм поиска узкого места за 3 минуты.
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, смотрите на:
- Количество Stage'ов: много Stage = много Shuffle = много потенциальных узких мест
- Симметрию: Stage'ы должны иметь похожее число Task'ов. Резкий дисбаланс указывает на проблему
- Цвет Stage'ов: зелёный = завершён, синий = выполняется, серый = Skipped
- Узлы с иконкой обмена (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¶
- Прогресс-бар показывает что Stage «почти завершён», но не двигается
- Duration у Job растёт, хотя большинство Stage'ов зелёные
- Нажмите на 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 Metrics → Event 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+ раз
# и занимает дорогостоящие вычисления для пересчёта
Кеш имеет смысл только если:
- DataFrame будет использоваться минимум 2 раза
- Вычисление DataFrame дороже чем десериализация из кеша
- У вас достаточно 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, каждый зелёный узел — сэкономленные вычисления из кеша.