Stage Boundary: почему shuffle создаёт новый Stage
Глубокий разбор Stage Boundary: pipelining внутри Stage, механика shuffle (write → sort → spill → network → read), narrow vs wide зависимости, квадратичный сетевой трафик, Spark UI метрики и способы минимизации shuffle.
Каждый раз, когда вы видите в 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_idshuffle_{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. Все они работают по одному принципу: уменьшить либо объём данных на границе, либо само число границ.