Stage Boundary: почему shuffle создаёт новый Stage

Глубокий разбор Stage Boundary: pipelining внутри Stage, механика shuffle (write → sort → spill → network → read), narrow vs wide зависимости, квадратичный сетевой трафик, Spark UI метрики и способы минимизации shuffle.

core internals optimization

Каждый раз, когда вы видите в Spark UI надпись Exchange или видите что job разбился на несколько Stage - это след Stage Boundary. Именно здесь Spark вынужден остановить потоковую обработку данных, записать промежуточные результаты на диск и перетасовать данные по сети. Понимание того, почему это происходит и какова физическая стоимость, - ключ к написанию эффективных pipeline.


Pipelining: идеальный мир без границ

Начнём с того, как работает Spark в идеальном случае - внутри одного Stage, без shuffle.

Представьте цепочку трансформаций:

df.filter(col("age") > 18) \
  .withColumn("country_upper", upper(col("country"))) \
  .select("user_id", "country_upper")

Наивно можно представить это как три отдельных прохода по данным:

Проход 1: читать все данные → filter → записать во временный файл
Проход 2: читать временный файл → withColumn → записать
Проход 3: читать → select → результат

Но Spark так не работает. Внутри Stage все narrow-трансформации объединяются в единый пайплайн:

Каждая запись проходит через весь стек операторов, прежде чем обрабатывается следующая. Это называется volcano model (или iterator model): каждый оператор - это iterator, который тянет записи из нижестоящего оператора по одной. Нет промежуточных материализаций, нет записи на диск между операциями.

Почему narrow-трансформации можно пайплайнить

Условие для pipeline - данные одной выходной партиции зависят только от одной входной партиции. Это называется narrow dependency (узкая зависимость):

Каждый Executor может обрабатывать свою партицию независимо. Не нужно ждать других Executor-ов, не нужна сеть, не нужна синхронизация. Данные текут через операторы как через трубу - отсюда и называется pipelining.


Stage Boundary: стена, в которую упирается конвейер

Stage Boundary возникает тогда, когда выходная партиция зависит от нескольких входных партиций. Это называется wide dependency (широкая зависимость, или shuffle dependency).

Простой пример: groupBy("city").count()

Данные о пользователях из Москвы могут быть в партициях 0, 3, 7, 12... После groupBy все записи с одинаковым city должны оказаться в одной выходной партиции - иначе нельзя посчитать корректный count.

Фундаментальная причина границы: Stage 1 не может начать обрабатывать ни одну запись, пока Stage 0 не завершил запись всех своих shuffle-файлов. Это обязательная точка глобальной синхронизации.

Почему нельзя начать Stage 1 раньше? Потому что Partition A в Stage 1 должна получить данные из всех партиций Stage 0. Если хотя бы одна задача Stage 0 ещё не завершена - данные неполны, и агрегация даст неверный результат.


Физика shuffle: что происходит за кулисами

Shuffle - самая дорогая операция в Spark. Разберём каждую фазу.

Shuffle Write: подготовка данных

1. Хеширование ключей

Для каждой записи Spark вычисляет partition_id = hash(key) % numOutputPartitions. Это определяет, в какую выходную партицию попадёт запись. Все записи с одинаковым city = "Moscow" получат одинаковый partition_id.

2. Сортировка (Sort-based Shuffle)

Начиная с Spark 1.2, используется Sort-based Shuffle Manager (заменил Hash-based). Записи внутри каждой входной партиции сортируются по partition_id (и по ключу, если нужно). Это позволяет создавать один индексированный файл вместо отдельного файла на каждую выходную партицию.

3. Shuffle Spill

Если сортируемые данные не помещаются в память Executor-а - Spark сбрасывает их на диск в промежуточный spill-файл. Финальный шаг - merge-sort всех spill-файлов в один shuffle block file.

Каждый Executor Stage 0 создаёт один файл данных и один файл индекса:

  • shuffle_{shuffleId}_{mapId}_0.data - данные, отсортированные по partition_id
  • shuffle_{shuffleId}_{mapId}_0.index - байтовые смещения каждой партиции в файле данных

Shuffle Read: сбор данных

4. Fetching

Executor-ы Stage 1 запрашивают у MapOutputTracker (Driver) адреса shuffle-файлов для своей партиции. Затем скачивают нужные блоки от всех Executor-ов Stage 0 по HTTP/Netty.

Вот где возникает квадратичная нагрузка на сеть: если Stage 0 имеет M partition и Stage 1 имеет N partition, то число сетевых соединений - M × N. При M=200 и N=200 это 40 000 HTTP-запросов.

5. Десериализация и merge

Полученные блоки десериализуются, декомпрессируются и сортируются/мержатся для передачи в оператор агрегации или join.

Стоимость shuffle: три вида ресурсов

Ресурс Что происходит Почему дорого
Disk I/O Shuffle Write: данные → диск Диск в 10–1000× медленнее RAM
Network I/O Shuffle Read: данные между узлами Пропускная способность сети ограничена
CPU Сериализация, сортировка, merge GC-давление при десериализации

Правило: один shuffle ≈ один полный цикл чтения/записи всего датасета на диск + передача по сети.


Narrow vs Wide: таксономия зависимостей

Narrow Dependency

Каждая выходная партиция зависит от фиксированного подмножества входных партиций (обычно от одной). Данные не покидают Executor.

Примеры: map, filter, withColumn, flatMap, sample, union, coalesce (если без полного перераспределения).

Wide Dependency (Shuffle Dependency)

Каждая выходная партиция потенциально зависит от всех входных партиций. Данные должны физически переместиться между Executor-ами.

Примеры: groupBy, reduceByKey, aggregateByKey, distinct, join (Sort-Merge Join), repartition, orderBy, intersection, subtract.

Исключение: операции без shuffle там, где ожидается

Broadcast Join - если один DataFrame достаточно мал (по умолчанию < 10 MB, параметр spark.sql.autoBroadcastJoinThreshold), Spark копирует его целиком на каждый Executor. Join происходит локально, shuffle не нужен:

Bucket Join - если оба DataFrame заранее разбиты (bucketed) на одинаковое число бакетов по одному ключу и сохранены как Hive-таблицы. Spark знает, что данные уже правильно распределены, и пропускает shuffle:

# Запись с бакетированием
df_orders.write.bucketBy(50, "user_id").sortBy("user_id") \
    .saveAsTable("orders_bucketed")
df_users.write.bucketBy(50, "user_id").sortBy("user_id") \
    .saveAsTable("users_bucketed")

# Join - без shuffle (оба датасета уже партиционированы по user_id)
spark.table("orders_bucketed").join(
    spark.table("users_bucketed"), "user_id"
).explain()
# → SortMergeJoin без Exchange оператора

Детерминизм и материализация

Ключевое следствие Stage Boundary - полная материализация:

Stage N не может начаться, пока все задачи Stage N-1 не завершены и записали shuffle-файлы.

Это фундаментальное ограничение - Stage являются барьерами синхронизации. В отличие от потоковых движков (Apache Flink, Kafka Streams), Spark batch гарантирует, что каждый Stage полностью завершён до начала следующего.

Именно поэтому одна медленная задача в Stage может задержать весь Job: пока она не завершится, следующий Stage не стартует.

Именно для борьбы с этим существует Speculative Execution (spark.speculation=true) - дублирование медленных задач.


Map-side combine: оптимизация до границы

Некоторые операции поддерживают предварительную агрегацию (pre-aggregation, map-side combine) до shuffle. Это значительно сокращает объём данных, который нужно передать по сети:

groupByKey - не делает map-side combine. Все сырые записи идут в shuffle. При большом числе дублирующихся ключей это катастрофа.

reduceByKey, aggregateByKey, countByKey, .groupBy().agg(count()) - делают map-side combine. Данные частично агрегируются внутри каждой партиции Stage 0 ещё до записи в shuffle-файлы.


Цепочка shuffle: стоимость накапливается

Каждая wide-трансформация добавляет Stage Boundary. В сложных pipeline границ может быть несколько:

df \
  .filter(col("status") == "active")       # narrow
  .join(users, "user_id")                  # wide → Stage Boundary 1
  .groupBy("city")                         # wide → Stage Boundary 2
  .agg(count("*"), sum("revenue"))         # агрегация в Stage 3
  .orderBy(col("revenue").desc())          # wide → Stage Boundary 3
  .write.parquet("output/")               # Action

Итого: 4 Stage, 3 Shuffle. Каждый shuffle - полный цикл write-to-disk + network transfer.

Цепочка shuffle означает, что каждый Stage ждёт завершения предыдущего. Параллелизм между Stage ограничен: Stage 1 не начнётся, пока Stage 0 не запишет все shuffle-файлы.


Анализ Stage Boundary в Spark UI

Вкладка SQL → Details → DAG Visualization

В SQL-плане каждый shuffle виден как узел Exchange:

== Physical Plan ==
*(3) Sort
+- Exchange rangepartitioning(revenue DESC, 200)     ← Shuffle 3
   +- *(2) HashAggregate(keys=[city], functions=[count, sum])
      +- Exchange hashpartitioning(city, 200)         ← Shuffle 2
         +- *(1) Project [city, revenue]
            +- SortMergeJoin [user_id], [user_id]
               :- Exchange hashpartitioning(user_id, 200)  ← Shuffle 1
               :  +- *(1) Filter (status = active)
               +- Exchange hashpartitioning(user_id, 200)  ← Shuffle 1 (правый)

Каждый Exchange - это Stage Boundary.

Вкладка Stages: метрики shuffle

Метрика Что означает
Shuffle Write (MB) Объём данных записан в shuffle-файлы Stage N
Shuffle Read (MB) Объём данных считан из shuffle Stage N-1
Shuffle Spill (Memory) Данные, которые не поместились в RAM
Shuffle Spill (Disk) Данные, сброшенные на диск (если > 0 - проблема)
Duration Время Stage; включает и Write, и wait, и Read

Признаки Data Skew в метриках Stage:

  • Min/Median/Max time Tasks сильно расходятся (одна задача 5 мин, остальные 30 сек)
  • Shuffle Read Size резко различается по Tasks (одна Task читает 10 GB, остальные по 50 MB)

Skipped Stages

Если в Spark UI Stage помечен Skipped - это означает, что данные уже закэшированы и Stage не нужно перевычислять. Кэш работает как мемоизация: повторный Action на закэшированном DataFrame пропускает все Stage до точки кэширования.


Практика: наблюдаем границы на коде

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, count, sum as spark_sum, upper

spark = SparkSession.builder \
    .master("local[*]") \
    .appName("StageBoundaryDemo") \
    .config("spark.sql.shuffle.partitions", "4") \
    .getOrCreate()

# Датасет
orders = spark.createDataFrame([
    (1, "Moscow",  "electronics", 1200),
    (2, "Berlin",  "clothing",     300),
    (3, "Moscow",  "clothing",     450),
    (4, "Paris",   "electronics", 2100),
    (5, "Berlin",  "electronics",  900),
    (6, "Moscow",  "electronics",  750),
], ["order_id", "city", "category", "revenue"])

# ── Вариант A: три отдельных groupBy → три shuffle ──────────────────
city_count    = orders.groupBy("city").agg(count("*").alias("orders"))
cat_count     = orders.groupBy("category").agg(count("*").alias("orders"))
city_cat      = orders.groupBy("city", "category").agg(count("*"))

# Три Action → три Job → каждый с двумя Stage и одним Shuffle
city_count.show()
cat_count.show()
city_cat.show()

# ── Вариант B: один groupBy → один shuffle ──────────────────────────
combined = orders.groupBy("city", "category").agg(
    count("*").alias("orders"),
    spark_sum("revenue").alias("total_revenue"),
)
combined.show()
# Один Job, два Stage, один Shuffle

# ── plan: смотрим Exchange ──────────────────────────────────────────
combined.explain(mode="formatted")

Откройте http://localhost:4040/SQL/ и сравните DAG-графы вариантов A и B - разница в числе Exchange-узлов хорошо видна.


Как минимизировать Stage Boundaries

1. Фильтровать до shuffle

# Плохо: shuffle всего датасета, затем filter
result = df.groupBy("city").agg(sum("revenue")) \
           .filter(col("sum(revenue)") > 1000)

# Хорошо: filter до groupBy - shuffle меньшего объёма данных
result = df.filter(col("status") == "active") \
           .groupBy("city").agg(sum("revenue")) \
           .filter(col("sum(revenue)") > 1000)

Catalyst обычно делает это автоматически (predicate pushdown), но явная фильтрация до shuffle полезна когда автоматика не срабатывает.

2. Объединять groupBy вместо нескольких

# Плохо: два shuffle
count_df = df.groupBy("city").count()
sum_df   = df.groupBy("city").sum("revenue")
result   = count_df.join(sum_df, "city")  # ещё один shuffle

# Хорошо: один shuffle
result = df.groupBy("city").agg(
    count("*").alias("cnt"),
    sum("revenue").alias("total"),
)

3. Broadcast join вместо Sort-Merge Join

from pyspark.sql.functions import broadcast

# Sort-Merge Join → 2 shuffle (оба датасета сортируются и шаффлятся)
result = large_df.join(medium_df, "id")

# Broadcast Join → 0 shuffle (medium_df копируется на каждый Executor)
result = large_df.join(broadcast(medium_df), "id")

Порог spark.sql.autoBroadcastJoinThreshold (default: 10 MB). При AQE включённом Spark может поднять порог динамически на основе реальной статистики.

4. Repartition один раз, заранее

# Плохо: shuffle перед каждой операцией
df.groupBy("city").count()
df.groupBy("city").sum("revenue")  # повторный shuffle

# Хорошо: один repartition → два groupBy используют один shuffle
df_city = df.repartition("city")
df_city.groupBy("city").count()
df_city.groupBy("city").sum("revenue")

5. Включить AQE

spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")

AQE (Adaptive Query Execution) меняет план во время выполнения: сливает мелкие shuffle-партиции, переключает SMJ на Broadcast Join если партнёр оказался меньше порога, разбивает skewed-партиции.

6. Правильное число shuffle-партиций

# Default 200 - часто неоптимально
# Слишком много мелких партиций → накладные расходы на scheduler
# Слишком мало крупных → OOM, медленная сортировка

# Правило: целевой размер партиции после shuffle - 100–200 MB
# total_shuffle_data / target_partition_size
spark.conf.set("spark.sql.shuffle.partitions", "400")  # если 80 GB данных

Итог

Stage Boundary - это точка, где Spark вынужден прервать pipeline и физически перераспределить данные. Её фундаментальная причина: агрегация или join требуют, чтобы все данные одного ключа находились в одной партиции, а в исходных данных они распределены по разным Executor-ам.

Стоимость каждой Stage Boundary складывается из трёх компонентов:

  • Disk I/O - shuffle-файлы записываются на диск каждого Executor
  • Network I/O - данные передаются между узлами (M×N соединений)
  • Глобальная синхронизация - следующий Stage ждёт завершения всех задач предыдущего

Способы борьбы - фильтрация до shuffle, объединение groupBy, Broadcast Join, Bucket Join, правильное число партиций и AQE. Все они работают по одному принципу: уменьшить либо объём данных на границе, либо само число границ.