DAG Scheduler: как трансформации превращаются в задачи

Полный путь от строки кода до Task на Executor: как DAG Scheduler делит граф на Stage по shuffle-границам, превращает Stage в TaskSet, учитывает Data Locality и восстанавливается после сбоев через Lineage.

core internals

Когда вы пишете 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. Его работа:

  1. Получить Physical Plan от Catalyst при вызове Action
  2. Разбить DAG на Stage по границам Shuffle (wide transformations)
  3. Для каждого Stage создать TaskSet: одна Task на одну партицию
  4. Назначить Tasks на Executor-ы с учётом Data Locality
  5. Отслеживать выполнение, перезапускать при сбоях через Lineage

Понимание этой механики объясняет:

  • Почему groupBy после filter быстрее чем до
  • Почему repartition(200) в начале pipeline дорог
  • Почему broadcast убирает целый Stage
  • Почему Data Skew убивает Stage: одна Task из 200 выполняется 100× дольше остальных

Следующий урок - Task Scheduler: как Tasks назначаются на конкретные Executor-ы и что происходит внутри Executor при выполнении Task.