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")