Spark UI: рентген вашего приложения
Полное руководство по Spark UI: от навигации по вкладкам до диагностики skew, spill, GC и неправильных join-стратегий. Учимся читать execution plan и находить узкие места без чтения тысяч строк логов.
Почему Spark UI - главный инструмент Spark-инженера¶
Представьте, что вы написали PySpark-задание, запустили его и оно работает… 40 минут. Вы не знаете, в чём проблема: не хватает памяти? Данные перекошены? Join выбрал неправильную стратегию? Не хватает CPU? Сеть перегружена?
Читать тысячи строк логов Spark - это как ставить диагноз без инструментов. Spark UI - это рентген вашего приложения: он показывает, что происходит внутри кластера в реальном времени и после завершения задания.
Правильное чтение Spark UI отличает инженера, который меняет конфиги наугад, от инженера, который точно знает, какой конфиг изменить и почему. В production-окружении нельзя перезапускать задание 10 раз с разными настройками - Spark UI даёт ответ с первого прохода.
Как получить доступ к Spark UI¶
Spark UI становится доступен автоматически после создания SparkSession. В зависимости от окружения:
| Среда | Адрес | Примечание |
|---|---|---|
| Локальный запуск | http://localhost:4040 |
Активен только пока живёт приложение |
| Spark Standalone | http://<driver-host>:4040 |
Порт может сместиться на 4041, 4042 если занят |
| YARN | http://<resourcemanager>:8088 |
Через YARN Application Manager |
| Databricks | Кнопка «Spark UI» в кластере | Встроенный интерфейс |
| AWS EMR | Application UI через консоль EMR | После запуска job |
| History Server | http://<server>:18080 |
Для завершённых заданий |
History Server - особенно важен в production. Когда задание завершилось (успешно или с ошибкой), живой Spark UI исчезает. History Server читает event log файлы (обычно хранятся в HDFS или S3) и позволяет анализировать завершённые задания сколько угодно времени после. Включается через:
# spark-defaults.conf
spark.eventLog.enabled=true
spark.eventLog.dir=hdfs:///spark-logs
spark.history.fs.logDirectory=hdfs:///spark-logs
После запуска start-history-server.sh сервер доступен на порту 18080 и показывает все завершённые задания из указанной директории.
Архитектура Spark execution и её отражение в UI¶
Прежде чем читать Spark UI, нужно понять иерархию выполнения Spark. Это фундамент, без которого вкладки UI будут выглядеть как бессвязные числа.
Разберём каждый уровень детально, потому что это напрямую влияет на то, что вы видите в UI:
Application - это весь жизненный цикл вашего PySpark-скрипта от SparkSession.builder.getOrCreate() до завершения. Одно приложение может породить десятки и сотни Jobs.
Job - это единица работы, порождаемая каждым Action в вашем коде: count(), show(), collect(), write.parquet(), save(). Трансформации (filter, map, join) сами по себе Jobs не создают - это lazy-вычисления. Только Action «заставляет» Spark выполнить накопленный план.
Stage - это набор задач, которые можно выполнить без shuffle (перемешивания данных между узлами). Boundary (граница) между Stages возникает именно в точках, где данные должны перераспределиться: groupBy, join, repartition, distinct. Всё это - wide transformations. Narrow transformations (filter, map, select) не создают новый Stage - они выполняются в том же Stage, что и предыдущая операция.
Task - это минимальная единица параллельного выполнения. Одна Task обрабатывает одну партицию данных на одном ядре одного Executor. Количество Tasks в Stage = количество партиций данных в этом Stage.
Это фундаментальное соответствие: 1 партиция = 1 Task = 1 ядро.
Почему граница Stage возникает на Shuffle¶
Для агрегации (groupBy) все записи с одним ключом должны попасть на один Executor, чтобы тот мог посчитать агрегат. До Shuffle записи с одним user_id могут лежать в разных партициях на разных машинах. Shuffle - это процесс сортировки и передачи данных по сети так, чтобы каждый Executor получил «свои» ключи. После Shuffle можно делать агрегацию локально.
Именно поэтому Shuffle - самая дорогая операция в Spark. Это не просто вычисление - это сетевая передача, которая может занимать гигабайты данных.
Обзор вкладок Spark UI¶
Порядок чтения при диагностике не совпадает с порядком вкладок в интерфейсе. Рекомендуемый порядок: SQL → Jobs → Stages → Executors → Storage.
Вкладка SQL: архитектура запроса¶
SQL-вкладка - это место, где Spark показывает физический план выполнения вашего запроса. Если вы хотите понять, почему запрос медленный, нужно начинать именно здесь.
Каждый раз, когда вы запускаете DataFrame-операцию с Action (или явный spark.sql()), в SQL-вкладке появляется новая запись. Нажав на неё, вы видите DAG физического плана.
Как читать граф физического плана¶
Важное правило: граф читается снизу вверх. Данные «текут» снизу (источник) наверх (результат). Самый нижний узел - это откуда берутся данные (FileScan, InMemoryTableScan), самый верхний - это результат.
Каждый узел в графе показывает:
- Название оператора - что делается (FileScan, Exchange, Join, Aggregate)
- Метрики - сколько строк прошло, сколько байт, сколько времени
Ключевые операторы и что они означают¶
FileScan parquet - чтение файлов с диска или объектного хранилища. Рядом с оператором в UI вы увидите:
number of files read: 128- сколько файлов прочитаноsize of files read: 2.3 GB- общий объём- Хорошо видно, применился ли predicate pushdown (фильтр передан в читалку файлов) или нет. Если
rows outputзначительно меньшеrows read- pushdown не сработал и Spark читает лишнее.
Exchange - это Shuffle. Каждый Exchange в графе - это точка, где данные перемещаются по сети. Рядом будет написано: hashpartitioning(key, 200) - это значит, данные хешируются по указанному ключу и раскладываются по 200 партициям. Много Exchange-узлов в плане = много Shuffle = медленный запрос. Задача оптимизатора - минимизировать их количество.
BroadcastExchange - специальный Shuffle, при котором вся таблица рассылается на каждый узел. Работает только для маленьких таблиц (по умолчанию < 10 МБ, настраивается через spark.sql.autoBroadcastJoinThreshold). Дешевле, чем двусторонний Shuffle в SMJ.
BroadcastHashJoin - Join, при котором маленькая таблица уже на каждом Executor, и Join делается локально. Не требует Shuffle для большой таблицы. Один из самых быстрых видов Join в Spark.
SortMergeJoin - Join по умолчанию для двух больших таблиц. Требует Shuffle обеих таблиц по ключу Join + сортировку. Более надёжный, чем BHJ, но дороже.
WholeStageCodegen (часто обозначается числом в кружке) - это «упаковка» нескольких операторов в один компилируемый Java-метод. Spark не интерпретирует операторы по одному, а генерирует bytecode, выполняющий несколько шагов за один проход. Это ускоряет выполнение в разы. Операторы внутри одного WholeStageCodegen выполняются без материализации промежуточных данных.
HashAggregate (partial) + HashAggregate (final) - это двухфазная агрегация. Сначала каждый Executor вычисляет частичные агрегаты локально (partial), потом после Shuffle вычисляется финальный результат (final). Это значительно эффективнее, чем собирать все данные в одном месте.
Как найти узкое место через SQL-вкладку¶
Hover (наведение мыши) на любой оператор в графе показывает дополнительные метрики. Особенно важны:
number of output rows- сколько строк прошло через оператор. Если после Filter прошло 95% строк из 100% прочитанных - фильтр не эффективен или pushdown не работает.shuffle bytes writtenрядом сExchange- сколько гигабайт перемещается по сети. Если это число в десятки GB - это главный кандидат на оптимизацию.spill (memory)иspill (disk)рядом с HashAggregate или Sort - означает, что оператору не хватило памяти и он начал писать временные данные на диск. Это критически замедляет выполнение.
Вкладка Jobs: декомпозиция Actions¶
Вкладка Jobs показывает список всех Actions, выполненных вашим приложением. Для каждого Job отображается:
- Описание (обычно имя метода и строка кода)
- Статус: Running, Succeeded, Failed
- Число Stages и Tasks
- Продолжительность
- Прогресс-бар
Что такое Job DAG¶
Нажав на конкретный Job, вы видите DAG Visualization - граф зависимостей между Stages этого Job. Граф показывает, как Stages зависят друг от друга: один Stage может начаться только после завершения другого (если он потребляет его shuffle-вывод).
Stages, которые не зависят друг от друга (например, чтение двух источников перед Join), могут выполняться параллельно - Spark запускает их одновременно, если есть свободные ресурсы. Это видно в Event Timeline.
Event Timeline¶
Event Timeline - это горизонтальная диаграмма, показывающая, что делал каждый Executor в каждый момент времени. Каждый цветной отрезок - это Task на конкретном Executor.
На схеме видно, что Task 1 на Executor 2 (красная) занимает в 3-4 раза больше времени, чем остальные Tasks. Это типичный симптом Data Skew - одна партиция содержит непропорционально много данных.
Что нужно искать на Event Timeline:
- Равномерная загрузка - хорошо: все ядра заняты примерно одинаково долго
- Один длинный отрезок - плохо: это straggler task, скорее всего из-за skew
- Большие серые зазоры между Tasks - это Scheduler Delay (время ожидания выделения ресурсов)
- Параллельные короткие Tasks - если Tasks занимают миллисекунды, возможно партиций слишком много и overhead на создание Task значителен
Вкладка Stages: детальный анализ¶
Вкладка Stages - это место, где можно найти детали о каждом Stage: сколько Tasks, сколько данных было прочитано и записано, каковы метрики отдельных Tasks.
Summary Statistics и Task-уровень¶
Самая важная часть страницы Stage - таблица Summary Statistics с метриками Tasks:
| Метрика | Min | 25th %ile | Median | 75th %ile | Max |
|---|---|---|---|---|---|
| Duration | 1s | 2s | 2.3s | 2.8s | 47s |
| GC Time | 50ms | 80ms | 100ms | 150ms | 2.1s |
| Input Size | 120MB | 150MB | 160MB | 170MB | 3.2GB |
| Shuffle Write | 10MB | 12MB | 13MB | 15MB | 890MB |
Здесь Max >> Median по Duration и Input Size - это классический признак Data Skew. Одна Task обрабатывала 3.2 GB входных данных вместо типичных 160 MB. Пока эта Task не завершится - весь Stage ждёт.
Scheduler Delay и его причины¶
В детальных метриках каждой Task есть поле Scheduler Delay. Это время между моментом, когда Task была запланирована, и моментом, когда Executor начал её выполнять. Большой Scheduler Delay означает:
- Не хватает ресурсов - нет свободных слотов на Executor'ах, Task ждёт в очереди
- Проблемы с Data Locality - Spark ждёт
spark.locality.wait(3 сек по умолчанию) прежде чем понизить уровень локальности и запустить Task на другом узле - Overhead Driver'а - при очень большом количестве Tasks (десятки тысяч) Driver тратит время на само планирование
Если суммарный Scheduler Delay в Stage занимает 20-30% времени - это сигнал пересмотреть число партиций или конфигурацию Executor'ов.
Shuffle Read Time и его составляющие¶
В метриках Task есть Shuffle Read Blocked Time - время, которое Task провела в ожидании данных от предыдущего Stage. Это чистое время сетевой передачи данных. Большое значение указывает на:
- Узкое место в сети кластера
- Слишком большой объём Shuffle (нужно оптимизировать запрос)
- Проблему с Shuffle Service на конкретном узле
Диагностика по симптомам¶
Симптом 1: Data Skew (перекос данных)¶
Data Skew - одна из самых частых причин медленных Spark-заданий. Он возникает, когда данные неравномерно распределены по партициям: большинство партиций маленькие, а одна или несколько содержат огромную долю данных.
Откуда берётся skew? Типичные причины:
- NULL-значения в ключе Join или GroupBy: все строки с
NULLпопадают в одну партицию - Доминирующий ключ: если 80% заказов принадлежит одному пользователю (боту, тестовому аккаунту), партиция этого пользователя будет огромной
- Неравномерное распределение в источнике: исходные файлы имеют разный размер
Как увидеть в Spark UI:
- Зайти в Stages → нажать на Stage с Join или GroupBy
- Посмотреть Summary Statistics → Duration
- Если
Max Durationв 10+ раз большеMedian Duration- это skew - Также смотреть на
Input Size- если Max >> Median по размеру входных данных - В таблице Tasks найти конкретную Task с максимальной Duration - она и есть «горячая» партиция
Дополнительный признак: в Event Timeline все Tasks завершились, но Stage всё ещё выполняется - это один straggler task держит весь Stage.
Симптом 2: Избыточный Shuffle¶
Shuffle - это «дорога» между двумя Stages. Чем больше данных пересылается по этой дороге, тем дольше Stage.
Как измерить в Spark UI:
В Summary Stats Stage смотрим на Shuffle Write Size (что Stage записал для следующего) и Shuffle Read Size (что Stage прочитал из предыдущего). Если эти значения составляют гигабайты - стоит подумать об оптимизации.
Обычные приёмы снижения объёма Shuffle:
- Фильтровать данные до Join, а не после
- Использовать Broadcast Join для маленькой таблицы
- Выбирать только нужные колонки (
select) до операции, требующей Shuffle - При повторном использовании данных - кешировать после первого Shuffle
Симптом 3: Spill (сброс на диск)¶
Spill происходит, когда Executor начинает обработку данных, но обнаруживает, что данных больше, чем помещается в отведённую память. В этом случае Spark вынужден:
- Сериализовать часть данных
- Записать их во временные файлы на диск
- По мере освобождения памяти - читать данные обратно с диска
Spill критически замедляет выполнение, потому что диск на 100-1000x медленнее памяти. Задание, которое в памяти выполнялось бы 2 минуты, со Spill может работать 20-30 минут.
Как увидеть в Spark UI:
В Summary Stats Stage ищем колонки Spill (Memory) и Spill (Disk). Если они не равны 0 - есть Spill. Spill (Memory) - это объём данных до сериализации (до сжатия), Spill (Disk) - после сжатия. Разница показывает коэффициент сжатия.
Причины Spill:
- Executor memory слишком мал для объёма данных в партиции
- Data Skew: одна партиция намного больше, чем Executor может держать в памяти
spark.sql.shuffle.partitionsслишком мал → партиции слишком большие- Слишком высокое
spark.memory.storageFraction→ мало памяти для execution
Лечение:
- Увеличить число shuffle partitions:
spark.sql.shuffle.partitions=500 - Увеличить память Executor:
spark.executor.memory=8g - Исправить Data Skew (salting, или включить AQE Skew Join)
Симптом 4: Высокое GC Time¶
GC (Garbage Collection) - это процесс освобождения памяти JVM от объектов, которые больше не нужны. В нормальном приложении GC занимает 1-5% времени выполнения. Если GC Time превышает 10-20% - это сигнал проблем с памятью.
Stop-the-World GC - особенно опасен: во время Major GC весь Executor полностью останавливается (все потоки приостанавливаются). Если эта пауза длится 5+ секунд, Driver может решить, что Executor умер, и пометить его Tasks как failed. Tasks будут перезапущены - это напрасно потраченное время.
Как увидеть в Spark UI:
В вкладке Executors есть колонка GC Time. Если GC Time / Total Time > 10% - плохо. Если > 20% - критично.
В Summary Stats Stage на странице конкретного Stage есть также GC Time по Tasks - можно найти конкретные Tasks с очень высоким GC.
Причины высокого GC:
- Executor memory слишком мал для обрабатываемых данных
- Использование Python UDF без Pandas UDF/Arrow: каждая строка сериализуется в Python-объект и обратно → огромное давление на GC
- Слишком много данных закешировано в
MEMORY_AND_DISK- Storage съедает память, оставляя мало для Execution
Лечение:
- Увеличить
spark.executor.memory - Заменить Python UDF на Pandas UDF или встроенные функции
- Использовать G1GC:
spark.executor.extraJavaOptions=-XX:+UseG1GC - Уменьшить
spark.memory.storageFractionесли Execution страдает
Вкладка Executors¶
Вкладка Executors даёт сводку по каждому воркеру: что он сделал за время жизни приложения.
Ключевые колонки¶
Storage Memory - сколько памяти занимают кешированные данные на этом Executor. Если у одних Executor'ов значение 100%, а у других 20% - кеш распределился неравномерно (обычно из-за Data Skew).
GC Time - суммарное время, потраченное JVM на сборку мусора. Мы уже обсудили, почему это важно.
Shuffle Read / Write - сколько данных этот Executor отправил и получил через Shuffle. Если один Executor отправил на порядок больше других - он обрабатывал «горячую» партицию.
Tasks: Active / Failed / Complete - Failed Tasks - плохой знак. Spark перезапускает упавшие Tasks автоматически (до spark.task.maxFailures=4), но каждый перезапуск тратит время. Причины:
- OOM (OutOfMemoryError) - не хватило памяти Executor'у
- ExecutorLostFailure - Executor убит кластером (Preemption в YARN/K8s, SPOT instance terminated)
- Task timeout
Blacklisted - если Executor слишком часто вызывает Failed Tasks, Spark может внести его в blacklist и перестать посылать ему новые Tasks. Это видно в UI.
Executor Dead (умерший Executor)¶
Умерший Executor продолжает отображаться в UI (серым цветом) с пометкой Dead. Его Tasks были перераспределены на других Executor'ов. Если Executors часто умирают - это серьёзная проблема:
ExecutorLostFailure+Exit code: 137→ OOM Kill (Linux убил процесс по нехватке памяти)ExecutorLostFailure+Container killed by YARN→ превышен memory overhead- Много Dead Executors → задание будет работать значительно дольше из-за повторных запусков
Вкладка Storage: анализ кеша¶
Вкладка Storage показывает все DataFrame'ы и RDD, которые были явно закешированы через cache() или persist().
Что важно видеть¶
Fraction Cached - какой процент партиций DataFrame реально попал в кеш. Если вы закешировали DataFrame с расчётом, что он войдёт в память, но видите Fraction Cached: 45% - кеш неполный. При следующем обращении Spark будет пересчитывать некешированные партиции с нуля.
Size in Memory vs Size on Disk - если StorageLevel = MEMORY_AND_DISK, то часть кеша может оказаться на диске. Это значительно медленнее, чем из памяти.
Partitions Cached - количество партиций в кеше. Если кеш неполный - видно по разнице между Total Partitions и Partitions Cached.
Когда кеш вредит¶
Кеш занимает Storage Memory, конкурируя с Execution Memory. Если вы кешировали очень большой DataFrame, который занял 80% памяти Executor'ов, то для выполнения JOIN-операций остаётся мало места → Spill, медленная работа.
Признак этой проблемы в UI: высокий Spill при наличии большого кеша в Storage-вкладке. Решение: df.unpersist() после использования или уменьшение spark.memory.storageFraction.
Physical Plan: как связать df.explain() с UI¶
Команда df.explain(mode="formatted") в консоли показывает тот же физический план, что и SQL-вкладка UI, но в текстовом виде. Умение читать explain() ускоряет диагностику, потому что не нужно запускать задание - достаточно посмотреть план.
# Посмотреть физический план до запуска
df_result = (
orders
.join(users, "user_id")
.groupBy("region")
.agg(F.sum("amount").alias("total"))
)
df_result.explain(mode="formatted")
Вывод в консоли:
== Physical Plan ==
AdaptiveSparkPlan (1) ← AQE включён
+- HashAggregate (2) ← Финальная агрегация
+- Exchange (3) ← SHUFFLE для groupBy
+- HashAggregate (4) ← Partial агрегация (до shuffle)
+- BroadcastHashJoin (5) ← Join через broadcast
:- Filter (6) ← Фильтры
: +- Scan parquet (7) ← Чтение orders
+- BroadcastExchange (8) ← Users отправлен как broadcast
+- Scan parquet (9) ← Чтение users
Читать нужно снизу вверх: сначала читаются orders (7) и users (9), users отправляется как broadcast (8), делается Join (5), применяются фильтры (6), частичная агрегация (4), Shuffle (3), финальная агрегация (2).
В UI этот же план будет представлен в виде графа, где каждый оператор - это узел. Числа в скобках соответствуют ID операторов в графе UI.
AQE: как план меняется в рантайме¶
При включённом AQE (spark.sql.adaptive.enabled=true) физический план может измениться во время выполнения. Это видно в SQL-вкладке как «AQE Plan» vs начальный план.
В SQL-вкладке при AQE вы увидите два раздела: Initial Plan и Final Plan. Если они отличаются - AQE что-то изменил. Это нормально и хорошо. Если Initial Plan был SortMergeJoin, а Final Plan стал BroadcastHashJoin - AQE обнаружил, что одна таблица маленькая и переключился на более быструю стратегию.
Чек-лист диагностики при открытии Spark UI¶
Практика: анализ медленного задания по симптомам¶
Кейс 1: JOIN занимает 40 минут¶
Представим ситуацию: у вас есть задание, которое делает JOIN двух таблиц и агрегирует результат. Оно работает 40 минут, хотя по объёму данных должно укладываться в 5.
Шаг 1 - SQL-вкладка: открываем план. Видим SortMergeJoin с двумя Exchange-узлами. Рядом с одним Exchange написано Shuffle Write: 45 GB. Это много - половина времени, скорее всего, уходит на это.
Шаг 2 - смотрим правую таблицу Join. В SQL-план она тоже проходит через Exchange. Нажимаем на FileScan правой таблицы - там написано rows: 50,000, size: 8 MB. Таблица маленькая!
Вывод: Catalyst не выбрал BroadcastHashJoin, потому что не смог оценить размер таблицы (возможно, статистики нет). Решение:
from pyspark.sql.functions import broadcast
result = big_table.join(broadcast(small_table), "user_id")
После добавления hint broadcast - Exchange для маленькой таблицы исчезает, BroadcastHashJoin вместо SortMergeJoin, Shuffle Write падает с 45 GB до 2 GB, время выполнения - с 40 минут до 4 минут.
Кейс 2: Stage завис на 99% готовности¶
Задание выполняется нормально, почти все Tasks завершились, но Stage на 99% стоит ещё 30 минут.
Шаг 1 - Stages-вкладка: открываем зависший Stage. Видим, что 199 Tasks из 200 завершились за 2-3 минуты каждая, но одна Task всё ещё running.
Шаг 2 - смотрим Summary Stats: Max Duration = 45 min, Median Duration = 2.5 min. Разница в 18 раз.
Шаг 3 - смотрим Input Size: Max = 28 GB, Median = 145 MB. Разница в 200 раз. Это явный skew.
Шаг 4 - проверяем, есть ли у нас NULL-ключи в данных Join:
# Проверка NULL-ключей
orders.filter(F.col("user_id").isNull()).count() # → 5,000,000 строк
Оказывается, 5 миллионов строк имеют user_id = NULL. При хеш-партиционировании NULL % 200 = 0 (или фиксированная партиция), поэтому все NULL попали в партицию 0.
Решение: фильтруем NULL до Join, если они не нужны в результате, или включаем AQE Skew Join:
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")
spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionFactor", "5")
spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes", "256m")
Кейс 3: Spill и медленная агрегация¶
Задание делает groupBy("city").agg(collect_list("event")). Для некоторых городов (Москва, Санкт-Петербург) в collect_list попадают миллионы событий.
Шаг 1: В Stage, где происходит агрегация, видим Spill (Memory): 18 GB, Spill (Disk): 4 GB. Executor записывал 18 GB данных на диск (сжатых до 4 GB).
Шаг 2: Max Input Size = 8 GB на одну Task - явный skew по городу.
Решение: collect_list на большие партиции - это антипаттерн. Нужно ограничить размер списка или использовать другую агрегацию. Если collect_list необходим - нужно предварительно salting:
# Salting: добавляем суффикс к ключу, чтобы разбить горячий город на N частей
N = 10
df_salted = df.withColumn(
"city_salted",
F.concat(F.col("city"), F.lit("_"), (F.rand() * N).cast("int"))
)
# Первая агрегация по salted ключу
partial = df_salted.groupBy("city_salted").agg(
F.collect_list("event").alias("events_partial")
)
# Восстанавливаем оригинальный ключ
partial = partial.withColumn("city", F.split("city_salted", "_")[0])
# Вторая агрегация - соединяем частичные списки
from pyspark.sql.functions import flatten
result = partial.groupBy("city").agg(
flatten(F.collect_list("events_partial")).alias("all_events")
)
Советы эксперта¶
Первым делом смотреть на SQL-вкладку, а не Jobs или Stages. Физический план сразу покажет, правильный ли Join выбран, есть ли лишние Exchange-узлы, применился ли predicate pushdown.
Сравнивать Max и Median в Summary Statistics всегда. Если Max в 10+ раз больше Median - есть skew. Это самая быстрая диагностика Data Skew.
Следить за Shuffle Write объёмом. Умножьте Shuffle Write одного Stage на количество Stages - это суммарный сетевой трафик вашего задания. Если это терабайты - пора оптимизировать.
Spill ≠ ошибка, но spill = предупреждение. Задание не падёт из-за Spill (Spark умеет работать с диском), но оно будет работать в 10-100 раз медленнее. Нулевой Spill - это цель.
GC Time > 10% - красный флаг. Особенно если это происходит из-за Python UDF: каждый вызов UDF требует передачи строки из JVM в Python-процесс и обратно. Переключение на Pandas UDF (vectorized UDF) снижает GC до нормального уровня.
History Server - обязательно в production. Живой Spark UI доступен только пока работает приложение. В production задания запускаются ночью, завершаются к утру. Без History Server вы не сможете проанализировать завершённые задания. Включите event logging с первого дня.
# Включить event logging при создании SparkSession
spark = SparkSession.builder \
.appName("MyJob") \
.config("spark.eventLog.enabled", "true") \
.config("spark.eventLog.dir", "s3a://bucket/spark-logs/") \
.getOrCreate()
Не оптимизировать вслепую: любое изменение конфигурации должно быть обоснованием из Spark UI. «Увеличим executor.memory с 4g до 16g» без знания текущего GC time и Spill - это потеря денег на кластере. Сначала UI, потом изменения.