DAG Scheduler: как трансформации превращаются в задачи
Полный путь от строки кода до Task на Executor: как DAG Scheduler делит граф на Stage по shuffle-границам, превращает Stage в TaskSet, учитывает Data Locality и восстанавливается после сбоев через Lineage.
Когда вы пишете df.groupBy("city").agg(count("*")) - Spark не выполняет это немедленно. Вместо этого он строит граф операторов и ждёт Action. В момент вызова .show() или .write этот граф попадает в DAG Scheduler - компонент, который разбивает его на этапы и задачи для Executor-ов. Разобраться в DAG Scheduler - значит понять, почему один join дешёвый, а другой убивает кластер.
От кода к DAG: что происходит до DAG Scheduler¶
Spark обрабатывает ваш код в два этапа ещё до DAG Scheduler:
Logical Plan - дерево реляционных операторов без привязки к кластеру. Spark накапливает трансформации (filter, join, groupBy) как узлы дерева. Ничего не выполняется.
Catalyst Optimizer - переписывает Logical Plan: переставляет фильтры выше джойнов, убирает лишние колонки, применяет десятки правил оптимизации.
SparkPlanner - конвертирует оптимизированный план в Physical Plan: выбирает алгоритм join (SortMerge vs Broadcast), планирует exchange (shuffle) между операторами.
Только когда вызван Action - show(), count(), write, collect() - Physical Plan передаётся в DAG Scheduler.
Что такое DAG¶
DAG - Directed Acyclic Graph (ориентированный ациклический граф):
- Directed (ориентированный) - данные текут в одном направлении: от источника к результату. Каждое ребро показывает зависимость.
- Acyclic (ациклический) - нет циклов. Если бы циклы были, Spark не мог бы определить порядок выполнения и не смог восстановиться после сбоев.
Единственный критерий деления на Stage: Shuffle¶
DAG Scheduler делит граф на Stage по одному принципу: Shuffle = граница Stage.
Shuffle возникает при wide transformations - операциях, где данные одной партиции нужны в нескольких выходных партициях:
Внутри одного Stage данные текут через операторы без записи на диск - это называется pipelining. Записи на диск происходят только между Stage в точках Shuffle.
Какие трансформации создают Shuffle (Wide)¶
| Трансформация | Почему Shuffle |
|---|---|
groupBy / groupByKey |
Все записи одного ключа → одна партиция |
reduceByKey / aggregateByKey |
То же + map-side combine |
join / leftJoin (SortMerge) |
Совмещение записей по ключу из двух партиций |
repartition(N) |
Явное перераспределение по всем N партициям |
coalesce(N) |
Только если N < текущего кол-ва (и full shuffle) |
distinct |
Дедупликация требует глобальной видимости |
sortBy / orderBy |
Глобальная сортировка |
intersection / subtract |
Сравнение записей из разных источников |
Что НЕ создаёт Shuffle (Narrow)¶
| Трансформация | Почему Narrow |
|---|---|
map / flatMap |
1 входная партиция → 1 выходная |
filter |
Фильтрует внутри партиции |
mapPartitions |
Обрабатывает каждую партицию независимо |
union |
Просто конкатенирует партиции |
sample |
Берёт подмножество в каждой партиции |
withColumn |
Вычисление новой колонки внутри записи |
Два типа Stage¶
ShuffleMapStage - промежуточный этап. Пишет результат в shuffle-файлы на диск, разбитые по ключам будущих партиций. Может быть несколько в цепочке (join → groupBy → orderBy → три ShuffleMapStage).
ResultStage - финальный этап. Читает shuffle-файлы (или исходные данные), выполняет оставшиеся операции и передаёт результат Action. Всегда один на Job.
От Stage к Task: связь с партициями¶
Каждый Stage разбивается на Tasks - одна Task на одну партицию:
Если у вас 200 партиций в Stage - будет 200 Tasks. Executor-ы выполняют их параллельно: Executor с 4 cores выполняет 4 Tasks одновременно.
Количество партиций определяется:
- При чтении файлов: количество HDFS/S3 блоков (обычно 128MB)
- После Shuffle: значением
spark.sql.shuffle.partitions(default: 200) - При
repartition(N): явно заданным N
Data Locality: задача идёт к данным¶
DAG Scheduler назначает Tasks не случайно - он старается запустить Task на Executor, который физически ближе к данным:
Spark ждёт spark.locality.wait (default: 3 сек) на каждом уровне. Если локальный Executor не освободился - понижает уровень и назначает Task на более далёкий узел.
Это объясняет поведение: первые Tasks часто PROCESS_LOCAL (кэш из предыдущего Stage), последние - ANY (все ближние Executor-ы заняты).
Pipeline execution внутри Stage¶
Spark не выполняет каждую операцию отдельным проходом по данным. Внутри Stage цепочка narrow-трансформаций объединяется в единый пайплайн:
df.filter(col("age") > 18) \
.withColumn("name_upper", upper(col("name"))) \
.select("id", "name_upper", "city")
Без pipelining (гипотетически):
1. Читать все данные → filter → записать во временный файл
2. Читать временный файл → withColumn → записать
3. Читать → select → результат
С pipelining (как работает Spark):
Для каждой записи: читать → filter → withColumn → select → результат
Одни проход по данным, нет временных записей
Это называется volcano model или operator pipelining: каждая запись проходит через весь стек операторов Stage, прежде чем начинается следующая. Нет материализации промежуточных результатов.
Жизненный цикл Task¶
Сериализация Task - Task содержит не только данные, но и весь ваш код в виде байткода JVM (или CloudPickle-сериализацию для Python UDF). Именно поэтому в PySpark все объекты, попадающие в UDF, должны быть сериализуемы.
Отказоустойчивость: Lineage и Retry¶
RDD Lineage¶
Каждый RDD хранит ссылку на родительский RDD и трансформацию, которая его породила. Это и есть Lineage - «родословная».
Если Executor упал и потерял партицию mappedRDD[3] - Spark не перечитывает весь датасет. Он видит в Lineage: «mappedRDD[3] = filteredRDD[3].map()» → «filteredRDD[3] = rawRDD[3].filter()». Пересчитывает только эту цепочку для одной партиции.
Retry-механизм¶
| Уровень | Параметр | Default | Поведение |
|---|---|---|---|
| Task | spark.task.maxFailures |
4 | Retry Task до N раз (на разных Executor-ах) |
| Stage | - | - | Если Task превысила лимит → Stage Failed |
| Job | spark.stage.maxConsecutiveAttempts |
4 | Retry Stage до N раз |
| Executor | spark.executor.maxNumFailures |
- | Executor блэклистится после N провалов |
Если Stage не удалось после всех ретраев - Job помечается как Failed и бросает исключение в Driver.
Speculative Execution¶
Один Executor может "застрять" - медленный диск, GC, шумный сосед на узле. Остальные 99 Tasks давно завершились, а эта одна тормозит весь Job.
spark.conf.set("spark.speculation", "true")
spark.conf.set("spark.speculation.multiplier", "1.5") # 1.5× медленнее медианы
spark.conf.set("spark.speculation.quantile", "0.75") # стартует после 75% Tasks
При spark.speculation=true Spark запускает копию медленной Task на другом Executor. Первая завершившаяся копия побеждает, вторая отменяется. Это не гарантирует idempotency - убедитесь, что ваши Task-и не имеют побочных эффектов (запись в БД, внешние системы).
Пример: разбор реального pipeline¶
# Код
orders = spark.read.parquet("s3a://lake/orders/") # [1]
users = spark.read.parquet("s3a://lake/users/") # [2]
active_orders = orders.filter(col("status") == "active") # [3] narrow
enriched = active_orders.join(users, "user_id") # [4] wide → SHUFFLE
by_city = enriched.groupBy("city").agg( # [5] wide → SHUFFLE
count("*").alias("cnt"),
sum("amount").alias("total")
)
by_city.write.mode("overwrite").parquet("output/") # [6] Action
Что видит DAG Scheduler:
4 Stage, 2 Shuffle. Tasks в каждом Stage:
- Stage 0: по числу блоков файла
orders(например, 50) - Stage 1: по числу блоков файла
users(например, 10) - Stage 2:
spark.sql.shuffle.partitions= 200 (default) - Stage 3: 200 (те же shuffle-партиции)
Итого: 460 Tasks на 4 Stage. Stage 0 и Stage 1 выполняются параллельно.
Spark UI: как читать DAG-граф¶
В Spark UI вкладка Stages показывает каждый Stage с метриками:
- Shuffle Write (MB) - сколько данных записано в Shuffle. Большое значение → дорогой Shuffle
- Shuffle Read (MB) - сколько считано. Если сильно отличается по Tasks - подозрение на Data Skew
- Duration - время Stage. Самый долгий Stage - узкое место pipeline
- Tasks Failed - число ретраев Task
Вкладка SQL → Details → "Details" → "DAG Visualization" показывает полный граф с цветовым обозначением Exchange (Shuffle) узлов.
Что искать:
- Exchange блоки - каждый это Shuffle и граница Stage
- BroadcastExchange - broadcast join, без Shuffle (хорошо)
- Толстые стрелки с большими MB между Exchange - кандидаты на оптимизацию
Оптимизация через понимание DAG¶
Минимизировать число Shuffle
# Плохо: два отдельных groupBy → два Shuffle
counts = df.groupBy("city").count()
totals = df.groupBy("city").sum("amount")
result = counts.join(totals, "city") # ещё один Shuffle
# Хорошо: один groupBy → один Shuffle
result = df.groupBy("city").agg(
count("*").alias("cnt"),
sum("amount").alias("total")
)
Использовать Broadcast Join для маленьких таблиц
from pyspark.sql.functions import broadcast
# SortMergeJoin → 2 Shuffle (оба датасета сортируются и шаффлятся)
enriched = large_df.join(small_df, "id")
# BroadcastJoin → 0 Shuffle (small_df копируется на каждый Executor)
enriched = large_df.join(broadcast(small_df), "id")
Правильное число партиций
# Default 200 - часто неоптимально
# Мало данных → 200 мелких Task → overhead > полезной работы
spark.conf.set("spark.sql.shuffle.partitions", "20")
# Много данных → 200 огромных Task → OOM
spark.conf.set("spark.sql.shuffle.partitions", "2000")
# AQE автоматически подбирает - включите и не трогайте вручную
spark.conf.set("spark.sql.adaptive.enabled", "true")
Кэшировать перед fork
# df используется в двух разных Stage-ах → будет перечитан дважды
df = spark.read.parquet("big_table/").filter(...)
counts = df.groupBy("city").count()
totals = df.groupBy("region").sum("amount")
# Правильно: кэшировать точку разветвления
df.cache()
df.count() # материализовать кэш
counts = df.groupBy("city").count()
totals = df.groupBy("region").sum("amount")
Итог¶
DAG Scheduler - центральный диспетчер Spark. Его работа:
- Получить Physical Plan от Catalyst при вызове Action
- Разбить DAG на Stage по границам Shuffle (wide transformations)
- Для каждого Stage создать TaskSet: одна Task на одну партицию
- Назначить Tasks на Executor-ы с учётом Data Locality
- Отслеживать выполнение, перезапускать при сбоях через Lineage
Понимание этой механики объясняет:
- Почему
groupByпослеfilterбыстрее чем до - Почему
repartition(200)в начале pipeline дорог - Почему
broadcastубирает целый Stage - Почему Data Skew убивает Stage: одна Task из 200 выполняется 100× дольше остальных
Следующий урок - Task Scheduler: как Tasks назначаются на конкретные Executor-ы и что происходит внутри Executor при выполнении Task.