Shuffle Explosion: COUNT DISTINCT, множественный GROUP BY и sort без причины
Три операции, которые незаметно генерируют терабайты shuffle: точный COUNT DISTINCT на больших колонках, несколько несвязанных GROUP BY в одном pipeline и orderBy без последующей window операции.
optimization
Что такое shuffle explosion¶
Shuffle explosion - ситуация когда объём shuffle данных в разы превышает объём исходных данных, или когда лишние shuffle-операции добавляются в план из-за неосторожного кода.
Анти-паттерн 1: точный COUNT DISTINCT на больших данных¶
# ❌ COUNT DISTINCT на колонке с высокой кардинальностью
df.select(countDistinct("user_id")).show()
# Catalyst создаёт:
# 1. Partial aggregation (на каждом Executor'е)
# 2. Full shuffle по всем данным
# 3. Final aggregation на одном Executor'е
# Для 1 TB данных с 100M уникальных user_id - всё это едет в одну точку
# ✅ Приближённый COUNT DISTINCT (HLL алгоритм, ошибка ~2%)
from pyspark.sql.functions import approx_count_distinct
df.select(approx_count_distinct("user_id", rsd=0.02))
# rsd = relative standard deviation (точность)
# В 10-100× быстрее, нет full shuffle в одну точку
# ✅ Если точность критична - сначала groupBy, потом count
df.select("user_id").distinct().count()
# Catalyst оптимизирует через partial agg + shuffle + merge
Анти-паттерн 2: несколько несвязанных GROUP BY в одном pipeline¶
# ❌ Каждый groupBy - отдельный Shuffle Exchange
df_cached = df.cache()
df_cached.count()
stats1 = df_cached.groupBy("region").agg(sum("revenue"))
stats2 = df_cached.groupBy("category").agg(count("*"))
stats3 = df_cached.groupBy("channel").agg(avg("amount"))
stats1.write.parquet("out/region/")
stats2.write.parquet("out/category/")
stats3.write.parquet("out/channel/")
# 3 отдельных Shuffle Exchange → 3 сетевых передачи одних и тех же данных
# ✅ CUBE / ROLLUP для связанных агрегатов (один shuffle)
from pyspark.sql.functions import grouping_id
df.rollup("region", "category", "channel") \
.agg(sum("revenue"), count("*"), avg("amount")) \
.write.parquet("out/stats/")
# ✅ Или сохранить результат одного groupBy и использовать повторно
base_agg = df.groupBy("region", "category", "channel") \
.agg(sum("revenue").alias("rev"), count("*").alias("cnt")) \
.cache()
region_stats = base_agg.groupBy("region").agg(sum("rev"))
category_stats = base_agg.groupBy("category").agg(sum("rev"))
Анти-паттерн 3: orderBy / sort без причины¶
# ❌ orderBy в середине pipeline (не перед write и не перед window)
df = spark.read.parquet("s3://events/") \
.filter("event_type = 'purchase'") \
.orderBy("user_id", "event_time") # ← ЗАЧЕМ? ради чего?
.groupBy("user_id").agg(sum("amount"))
# orderBy = full sort shuffle + ещё один shuffle от groupBy
# Spark оптимизирует orderBy → groupBy только если они рядом
# но в реальных DAG промежуточный orderBy часто остаётся
# ✅ orderBy только если нужен конечный порядок (limit, window, write)
# Правильно: перед window function
df.orderBy("event_time") \
.withColumn("rank", row_number().over(Window.partitionBy("user_id").orderBy("event_time")))
# Но! Window function сам обеспечивает сортировку - внешний orderBy избыточен:
df.withColumn("rank", row_number().over(
Window.partitionBy("user_id").orderBy("event_time")
)) # orderBy снаружи не нужен
# ✅ sortWithinPartitions - локальная сортировка без shuffle
df.sortWithinPartitions("event_time") # нет Exchange в плане!
Анти-паттерн 4: join + groupBy без оптимизации порядка¶
# ❌ join создаёт shuffle, groupBy - ещё один shuffle
result = orders.join(customers, "customer_id") \ # Shuffle 1
.groupBy("region").sum("amount") # Shuffle 2
# Два shuffle, хотя можно было объединить
# ✅ AQE автоматически объединяет shuffle если включён
spark.conf.set("spark.sql.adaptive.enabled", "true")
# AQE видит что после join идёт groupBy - переключает join strategy
# и оптимизирует порядок shuffle
Диагностика в EXPLAIN¶
df.explain(mode="formatted")
# Ищите в плане:
# Exchange hashpartitioning(..., 200) ← каждый такой узел = shuffle
# Sort [...] ← full sort shuffle
# Нормально: 1 Exchange для join или groupBy
# Подозрительно: 3+ Exchange в одном плане без очевидной необходимости
Метрики в Spark UI¶
Stage → Summary Metrics:
Shuffle Write: 450 GB ← данные записанные в shuffle files
Shuffle Read: 450 GB ← данные прочитанные из shuffle files
Если Shuffle Write >> Input Size × 2 - что-то взрывается
# Настройка числа shuffle партиций (дефолт 200 - часто неоптимально)
# Для 100 GB данных: ~200-500 партиций (цель: 200 MB - 1 GB на партицию)
spark.conf.set("spark.sql.shuffle.partitions", "400")
# Или автоматически через AQE:
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")