Sort-Merge Join vs Shuffle Hash Join — когда каждый применяется

Глубокое сравнение двух алгоритмов Join для больших данных: физика SMJ (Shuffle + Sort + Merge), физика SHJ (Hash Table Build + Probe), стойкость SMJ к OOM и Disk Spill, риски OOM в SHJ, конфигурация preferSortMergeJoin, Hints, влияние Data Skew, Bucket Join для устранения Shuffle, AQE runtime switching и матрица выбора.

optimization

1. Анатомия Shuffle: общий знаменатель для SMJ и SHJ

Когда обе соединяемые таблицы слишком велики для Broadcast Join — они не умещаются в памяти Driver'а и Executor'ов — Spark вынужден выполнить распределённый Shuffle. Именно здесь начинается история Sort-Merge Join и Shuffle Hash Join.

Что такое Shuffle Exchange и почему он дорого стоит

Shuffle Exchange — это операция перераспределения данных между Executor'ами так, чтобы строки с одинаковыми Join Key'ями гарантированно оказались на одном Executor'е. Без этого невозможно корректно выполнить JOIN: строки с ключом customer_id=42 из левой и правой таблиц могут находиться на разных серверах.

Схема показывает результат Shuffle: Executor 4 получает все строки с key=100 из обеих таблиц. Если key=100 встречается часто (Data Skew) — это будет катастрофически большая партиция.

Стоимость Shuffle Exchange:

  • Shuffle Write: сериализация строк в бинарный формат, запись на локальный диск каждого Executor'а (~всё объём обоих датасетов в байтах)
  • Shuffle Read: чтение нужных партиций по сети с чужих Executor'ов (~аналогичный объём)
  • Накладные расходы: TCP соединения, буферы передачи, разрыв pipeline'а (барьер между Map и Reduce стадиями)

И SMJ, и SHJ начинают с одинакового Shuffle. Принципиальная разница — что каждый алгоритм делает с данными после того как они оказались на Executor'е.


2. Sort-Merge Join: физика алгоритма

Sort-Merge Join — это дефолтная стратегия для соединения больших таблиц в Spark начиная с версии 2.3. Алгоритм надёжен и масштабируем, хотя и не самый быстрый.

Полный жизненный цикл SMJ

Диаграмма показывает ключевую особенность SMJ: даже если данные не умещаются в Execution Memory — они уходят на диск и всё равно обрабатываются корректно. Спилл добавляет I/O, но не вызывает OOM.

Детали каждой фазы

Фаза Shuffle Read (общая для SMJ и SHJ): Executor читает партиции от других Executor'ов по сети. Данные поступают в произвольном порядке.

Фаза Sort: каждый Executor выполняет ExternalSort для своих партиций обоих датасетов по ключу JOIN. Слово «External» означает, что сортировка умеет работать с данными, не умещающимися в RAM — избыточные данные пишутся на диск как «Spill Runs», которые затем сливаются в один отсортированный поток. Сложность: O(N log N) с константой зависящей от числа Spill Run'ов.

Фаза Merge: после того как обе партиции отсортированы, Spark итерируется по ним одновременно двумя указателями. Поскольку оба потока отсортированы, поиск совпадений — просто линейное движение вперёд. Сложность: O(N + M) — линейная по числу строк в обоих датасетах.

from pyspark.sql import SparkSession, functions as F

spark = SparkSession.builder \
    .config("spark.sql.join.preferSortMergeJoin", "true") \  # дефолт: SMJ
    .getOrCreate()

# Большой JOIN: оба датасета > autoBroadcastJoinThreshold
orders = spark.table("bronze.orders")     # 500 GB
customers = spark.table("silver.customers")  # 100 GB

# Catalyst автоматически выберет SMJ
result = orders.join(customers, "customer_id", "inner")

# Проверяем что выбран SMJ:
result.explain("formatted")
# == Physical Plan ==
# SortMergeJoin [customer_id#1], [customer_id#15], Inner
# :- Sort [customer_id#1 ASC NULLS FIRST], false, 0
# :  +- Exchange hashpartitioning(customer_id#1, 200), ENSURE_REQUIREMENTS
# :     +- Scan parquet bronze.orders
# +- Sort [customer_id#15 ASC NULLS FIRST], false, 0
#    +- Exchange hashpartitioning(customer_id#15, 200), ENSURE_REQUIREMENTS
#       +- Scan parquet silver.customers
#
# Видим: ДВА Sort узла + ДВА Exchange (Shuffle) узла

3. Главное преимущество SMJ: стойкость к OOM

Самое важное свойство SMJ — out-of-core processing: алгоритм продолжает работать даже когда данные не умещаются в RAM.

Механизм External Sort и Disk Spill

Схема объясняет почему SMJ «никогда» не падает по OOM: при нехватке памяти данные просто уходят на диск. Цена — дополнительные I/O операции (Spill to Disk), которые замедляют работу, но не прерывают её.

Практическая стоимость Spill

# Диагностика Spill в Spark UI → Stages
# Метрики Task'ов показывают:
#   Spill (Memory): 4.2 GB   ← данные ожидающие записи на диск
#   Spill (Disk):   3.8 GB   ← данные на диске (после сжатия)
#
# Disk Spill overhead:
#   Without spill: 1x time
#   With spill:    2-5x time (из-за write + read диска)
#
# Правило: если Spill (Disk) > 20% от Input Size → нужно больше памяти
# или больше партиций (меньше данных на Task)

# Как минимизировать Spill без изменения алгоритма:
spark.conf.set("spark.sql.shuffle.partitions", "400")  # больше партиций → меньше данных на Task

# Или через AQE (динамически):
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
# AQE сам подберёт число партиций чтобы каждая ~advisoryPartitionSizeInBytes

4. Shuffle Hash Join: физика алгоритма

Shuffle Hash Join — альтернатива SMJ, которая устраняет фазу Sort, заменяя её построением хэш-таблицы в памяти. Это делает SHJ быстрее при правильных условиях — и опасным при неправильных.

Полный жизненный цикл SHJ

Принципиальное отличие от SMJ: SHJ строит Hash Table только для меньшей партиции (Build Side), а большую партицию стримит — читает строку за строкой, не загружая в память. Поэтому ограничение на память применяется только к одной стороне, а не к обеим.

Но именно это и является источником риска: если Build Side не помещается в Execution Memory — Spill невозможен (нельзя частично заполнить Hash Table). Spark упадёт с OOM.

Детали SHJ: почему нет Spill

В SMJ External Sort может писать на диск «Sorted Runs» и потом их мержить. Это работает потому что к результату сортировки можно обращаться последовательно.

В SHJ Hash Table не поддерживает частичное заполнение. Hash lookup требует весь словарь в памяти одновременно. Если записать половину Hash Table на диск — нельзя выполнить lookup, не загрузив всё обратно. Это принципиальное алгоритмическое ограничение.

# Принудительный SHJ через Hint
from pyspark.sql import functions as F

# Способ 1: DataFrame Hint
result = orders.join(
    customers.hint("SHUFFLE_HASH"),
    "customer_id"
)

# Способ 2: через preferSortMergeJoin = false
spark.conf.set("spark.sql.join.preferSortMergeJoin", "false")
result = orders.join(customers, "customer_id")

# Способ 3: SQL Hint
result = spark.sql("""
    SELECT /*+ SHUFFLE_HASH(c) */ o.*, c.name
    FROM orders o
    JOIN customers c ON o.customer_id = c.id
""")

# Проверяем план:
result.explain("formatted")
# == Physical Plan ==
# ShuffledHashJoin [customer_id#1], [id#15], Inner, BuildRight
# :- Exchange hashpartitioning(customer_id#1, 200)
# :  +- Scan parquet bronze.orders
# +- Exchange hashpartitioning(id#15, 200)
#    +- Scan parquet silver.customers
#
# Видим: ShuffledHashJoin + BuildRight (hash table из customers)
# НЕТ Sort узлов! Это и есть преимущество SHJ.

5. Риски SHJ: когда Hash Table не помещается в память

Это самый критичный раздел урока. SHJ — мощный инструмент, но без понимания ограничений он превращается в бомбу замедленного действия в production.

Почему размер партиции непредсказуем

Размер Hash Table на Executor'е зависит от:

  • Числа партиций (shuffle.partitions)
  • Равномерности распределения данных по ключам
  • Реального объёма данных (который может отличаться от оценки Catalyst)
# Оценка безопасности SHJ для конкретных данных

def is_shj_safe(
    build_side_gb: float,
    spark_executor_memory_gb: float,
    executor_cores: int,
    memory_fraction: float = 0.6,
    storage_fraction: float = 0.5,
    n_shuffle_partitions: int = 200,
) -> dict:
    """
    Оценивает безопасность применения SHJ.

    Build Side должен помещаться в Execution Memory одной Task'и.
    Execution Memory = spark_memory × (1 - storageFraction)
    При N Tasks на Executor каждая получает 1/N долю.
    """
    # Доступная Execution Memory на Executor
    usable_gb = spark_executor_memory_gb - 0.3  # reserved
    spark_memory_gb = usable_gb * memory_fraction
    max_execution_gb = spark_memory_gb  # весь Spark Memory при пустом Storage

    # При N Task'ах одновременно:
    execution_per_task_gb = max_execution_gb / executor_cores

    # Build Side размер на одну партицию:
    build_partition_gb = build_side_gb / n_shuffle_partitions

    # Hash Table overhead ~2x в памяти vs на диске
    hash_table_in_memory_gb = build_partition_gb * 2

    is_safe = hash_table_in_memory_gb < execution_per_task_gb * 0.8  # 20% запас

    return {
        "build_side_total_gb": build_side_gb,
        "build_partition_avg_gb": build_partition_gb,
        "hash_table_in_memory_gb": hash_table_in_memory_gb,
        "execution_per_task_gb": execution_per_task_gb,
        "is_safe": is_safe,
        "recommendation": (
            "SHJ БЕЗОПАСЕН" if is_safe
            else "SHJ ОПАСЕН: используйте SMJ или увеличьте executor.memory/shuffle.partitions"
        )
    }

# Пример:
result = is_shj_safe(
    build_side_gb=50,           # Build Side = 50 GB
    spark_executor_memory_gb=20, # 20 GB на Executor
    executor_cores=5,
    n_shuffle_partitions=200,
)
for k, v in result.items():
    print(f"  {k}: {v}")
# build_partition_avg_gb: 0.25 GB
# hash_table_in_memory_gb: 0.5 GB  ← в 2x из-за overhead
# execution_per_task_gb: 1.42 GB
# is_safe: True ✓
# 0.5 GB << 1.42 GB → БЕЗОПАСНО

OOM при Data Skew + SHJ: катастрофический сценарий

Data Skew делает SHJ особенно опасным. Если key=NULL встречается в 50% строк:

Build Side (customers) = 50 GB, 200 партиций
Средняя партиция = 50/200 = 0.25 GB → OK
НО партиция с key=NULL = 50 × 50% = 25 GB (все NULL на одном Executor!)
Hash Table для 25 GB → OOM!
# Типичный OOM от SHJ в production:
ERROR Executor: Exception in task 47.0 in stage 3.0 (TID 1299)
java.lang.OutOfMemoryError: Java heap space
    at org.apache.spark.sql.execution.joins.ShuffledHashJoinExec.buildHashedRelation(...)
    at org.apache.spark.sql.execution.joins.ShuffledHashJoinExec.doExecute(...)

# Или при YARN OOM:
Container killed by YARN for exceeding memory limits.
  12.4 GB of 8 GB physical memory used.
  Consider boosting spark.executor.memoryOverhead.

6. Конфигурационные параметры и триггеры выбора алгоритма

Главный параметр: spark.sql.join.preferSortMergeJoin

# По умолчанию: Spark предпочитает SMJ
spark.conf.get("spark.sql.join.preferSortMergeJoin")
# 'true'

# Когда true: SMJ используется для всех больших JOIN'ов
# Catalyst выберет SHJ только если он ЯВНО безопаснее (редко)

# Когда false: Catalyst рассматривает SHJ как опцию
# Условия для автоматического выбора SHJ (при preferSortMergeJoin=false):
# 1. Одна сторона значительно меньше другой
# 2. Меньшая сторона делённая на num_partitions умещается в Execution Memory
# 3. Нет явного риска OOM (по оценке статистики)
spark.conf.set("spark.sql.join.preferSortMergeJoin", "false")

Полная конфигурация JOIN стратегий

spark = SparkSession.builder \

    # ── Основной выбор алгоритма ─────────────────────────────────────
    # true (дефолт): предпочтение SMJ — надёжнее, но медленнее
    # false: Catalyst рассматривает SHJ если условия подходят
    .config("spark.sql.join.preferSortMergeJoin", "true") \

    # ── Порог для автоматического Broadcast ──────────────────────────
    # Таблицы меньше этого порога → BHJ (полностью минуют Shuffle!)
    # Дефолт: 10 MB. Повышайте если уверены в памяти Driver'а
    .config("spark.sql.autoBroadcastJoinThreshold", "100MB") \

    # ── AQE: динамический выбор алгоритма в runtime ───────────────────
    # AQE может конвертировать SMJ → BHJ или SMJ → SHJ
    # на основе реальных размеров после предыдущих Stage'ов
    .config("spark.sql.adaptive.enabled", "true") \

    # Порог для AQE-conversion SMJ → BHJ
    .config("spark.sql.adaptive.autoBroadcastJoinThreshold", "50MB") \

    # ── Число Shuffle партиций ────────────────────────────────────────
    # Больше партиций = меньше данных на Task = безопаснее для SHJ
    # AQE может уменьшить это число если данные небольшие
    .config("spark.sql.shuffle.partitions", "400") \

    .getOrCreate()

7. Явные Hints для SHJ в PySpark и SQL

Когда Catalyst выбирает SMJ, но вы знаете что ваши данные позволяют SHJ — используйте Hints.

PySpark Hints

from pyspark.sql import SparkSession, functions as F

spark = SparkSession.builder \
    .config("spark.sql.join.preferSortMergeJoin", "true") \  # дефолт
    .getOrCreate()

orders = spark.table("bronze.orders")      # 500 GB
products = spark.table("silver.products")  # 5 GB (меньше в 100 раз!)

# ── Вариант 1: SHUFFLE_HASH Hint ─────────────────────────────────────
# Принудительно используем SHJ с построением Hash Table из products
result_shj = orders.join(
    products.hint("SHUFFLE_HASH"),  # Hash Table строится из products
    "product_id",
    "inner"
)
result_shj.explain("formatted")
# ShuffledHashJoin [product_id#1], [id#15], Inner, BuildRight
# (BuildRight = строим hash table из products)

# ── Вариант 2: SHUFFLE_MERGE Hint (принудительный SMJ) ───────────────
# Полезно когда AQE пытается применить SHJ, а вы не хотите
result_smj = orders.join(
    products.hint("SHUFFLE_MERGE"),  # принудительный SMJ
    "product_id",
    "inner"
)
result_smj.explain("formatted")
# SortMergeJoin [product_id#1], [id#15], Inner

# ── Вариант 3: SHUFFLE_REPLICATE_NL Hint (особый случай) ────────────
# Nested Loop Join: для non-equi join (>, <, BETWEEN)
# ОЧЕНЬ медленный, только для специальных случаев
result_nl = orders.join(
    products.hint("SHUFFLE_REPLICATE_NL"),
    orders.price > products.min_price,  # non-equi!
    "inner"
)

SQL Hints

-- Принудительный SHJ для правой таблицы
SELECT /*+ SHUFFLE_HASH(p) */
    o.order_id,
    o.amount,
    p.product_name
FROM orders o
JOIN products p ON o.product_id = p.id;

-- Принудительный SMJ (когда AQE пытается делать что-то другое)
SELECT /*+ SHUFFLE_MERGE(p) */
    o.order_id,
    p.product_name
FROM orders o
JOIN products p ON o.product_id = p.id;

-- Несколько Hints одновременно (разные JOIN'ы в одном запросе)
SELECT /*+ SHUFFLE_HASH(p), BROADCAST(r) */
    o.order_id,
    p.product_name,
    r.region_name
FROM orders o
JOIN products p ON o.product_id = p.id
JOIN regions r ON o.region_id = r.id;

Когда использовать Hints: чек-лист

def recommend_join_strategy(
    left_gb: float,
    right_gb: float,
    driver_memory_gb: float,
    executor_memory_gb: float,
    is_data_skewed: bool,
    n_shuffle_partitions: int = 200,
) -> dict:
    """
    Рекомендует стратегию JOIN на основе параметров данных и кластера.
    """
    smaller_gb = min(left_gb, right_gb)
    larger_gb = max(left_gb, right_gb)
    ratio = larger_gb / smaller_gb

    # Execution Memory per Task
    usable_exec_mb = (executor_memory_gb - 0.3) * 0.6 * 1024  # ~60% heap
    exec_per_task_mb = usable_exec_mb / 5  # 5 cores per executor

    # Размер партиции Hash Build Side
    build_partition_mb = smaller_gb * 1024 / n_shuffle_partitions * 2  # 2x overhead

    if smaller_gb < 0.01:  # < 10 MB
        return {"strategy": "BroadcastHashJoin", "hint": "BROADCAST",
                "reason": "Таблица < autoBroadcastJoinThreshold"}

    if smaller_gb < 0.1 and driver_memory_gb > smaller_gb * 10:  # < 100 MB
        return {"strategy": "BroadcastHashJoin (raise threshold)",
                "hint": f"BROADCAST или config autoBroadcastJoinThreshold={int(smaller_gb*1.5*1024)}MB",
                "reason": f"Таблица {smaller_gb*1024:.0f} MB < Driver memory"}

    if not is_data_skewed and ratio > 5 and build_partition_mb < exec_per_task_mb * 0.8:
        return {"strategy": "ShuffleHashJoin", "hint": "SHUFFLE_HASH",
                "reason": f"Ratio {ratio:.0f}:1, партиция {build_partition_mb:.0f} MB < "
                          f"Exec Memory {exec_per_task_mb:.0f} MB"}

    if is_data_skewed:
        return {"strategy": "SortMergeJoin + AQE Skew",
                "hint": "SHUFFLE_MERGE + spark.sql.adaptive.skewJoin.enabled=true",
                "reason": "Data Skew: SHJ опасен OOM на горячих ключах"}

    return {"strategy": "SortMergeJoin (дефолт)", "hint": "Нет",
            "reason": "Безопасный выбор для неизвестных данных"}


# Примеры:
scenarios = [
    (500, 5, 32, 20, False),    # 500 GB + 5 GB: ratio 100:1, нет скоса
    (500, 100, 32, 20, False),  # 500 GB + 100 GB: ratio 5:1
    (500, 5, 32, 20, True),     # 500 GB + 5 GB: ЕСТЬ скос
    (500, 500, 32, 20, False),  # 500 GB + 500 GB: одинаковые
]

for left, right, drv, exc, skew in scenarios:
    rec = recommend_join_strategy(left, right, drv, exc, skew)
    print(f"\n{left} GB + {right} GB (skew={skew}):")
    print(f"  Стратегия: {rec['strategy']}")
    print(f"  Причина: {rec['reason']}")

8. Влияние Data Skew на SMJ и SHJ

Data Skew — неравномерное распределение данных по ключам — по-разному влияет на два алгоритма.

Data Skew в SMJ: Straggler Task, но не OOM

В SMJ перекошенная партиция попадает на один Executor и заставляет его сортировать и обрабатывать несоразмерно большой объём данных. Это создаёт Straggler Task — последнюю Task которая работает часами пока все остальные уже завершились.

SMJ при этом не падает. Данные уходят в Spill и обрабатываются медленно, но верно.

Data Skew в SHJ: мгновенный OOM

В SHJ горячая партиция попадает в Build Phase. Spark пытается построить Hash Table из 50 GB партиции в Execution Memory размером 2 GB → немедленный OOM. Никакого Spill, никакого восстановления — Task падает.

# Безопасная стратегия при подозрении на skew + SHJ:
# Проверьте распределение ДО применения SHJ

# 1. Быстрая проверка распределения ключей
key_dist = df.groupBy("join_key") \
    .count() \
    .orderBy(F.desc("count")) \
    .limit(10)
key_dist.show()

# 2. Считаем skew ratio
stats = df.groupBy("join_key").count() \
    .agg(
        F.max("count").alias("max_count"),
        F.percentile_approx("count", 0.5).alias("median_count")
    ).first()

skew_ratio = stats["max_count"] / max(stats["median_count"], 1)
print(f"Skew ratio: {skew_ratio:.0f}x (max/median)")
# Если > 10x и планируете SHJ → ОПАСНО

# 3. При skew + SMJ включайте AQE Skew Join
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")
spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionFactor", "5")

9. Bucket Join: устранение Shuffle и Sort полностью

В Data Lakehouse (Delta Lake, Iceberg) можно полностью избавиться от Shuffle и Sort этапов для SMJ через Bucketing. Это самая мощная оптимизация для повторяющихся JOIN'ов на production-кластерах.

Что такое Bucketing

Bucketing — это предварительное распределение данных по ключу при записи таблицы. Если обе таблицы сохранены с одинаковым числом bucket'ов по одному и тому же ключу, Spark при чтении может пропустить Shuffle Exchange.

# ── Запись таблицы с Bucketing ───────────────────────────────────────
# ОДИН РАЗ при создании таблицы

# Таблица фактов: orders
orders.write \
    .format("parquet") \
    .bucketBy(200, "customer_id") \     # 200 bucket'ов по customer_id
    .sortBy("customer_id") \             # предварительная сортировка!
    .saveAsTable("silver.orders_bucketed")

# Таблица измерений: customers
customers.write \
    .format("parquet") \
    .bucketBy(200, "customer_id") \     # ОДИНАКОВОЕ число bucket'ов!
    .sortBy("customer_id") \
    .saveAsTable("silver.customers_bucketed")

# ── Чтение с автоматическим Bucket Join ─────────────────────────────
# КАЖДЫЙ запрос: Shuffle и Sort исключены!
orders_b = spark.table("silver.orders_bucketed")
customers_b = spark.table("silver.customers_bucketed")

result = orders_b.join(customers_b, "customer_id")
result.explain("formatted")
# == Physical Plan ==
# SortMergeJoin [customer_id#1], [customer_id#15], Inner
# :- *(1) Project [...]                          ← НЕТ Exchange!
# :  +- *(1) FileScan parquet silver.orders_bucketed  ← НЕТ Sort!
# +- *(2) Project [...]                          ← НЕТ Exchange!
#    +- *(2) FileScan parquet silver.customers_bucketed
#
# Нет Exchange (Shuffle) узлов!
# Нет Sort узлов!
# Spark читает уже "перемешанные" и "отсортированные" данные!

Когда Bucket Join работает и когда нет

# Bucket Join РАБОТАЕТ если:
# 1. Обе таблицы bucketed по ОДНОМУ ключу JOIN
# 2. Обе таблицы имеют ОДИНАКОВОЕ число bucket'ов
# 3. Обе таблицы читаются из HMS (enableHiveSupport())

# Bucket Join НЕ РАБОТАЕТ если:
# 1. Разное число bucket'ов
# 2. Разные ключи bucketing
# 3. После фильтрации/агрегации (метаданные теряются)
# 4. Streaming DataFrame (нет статических bucket метаданных)

# Проверка что Bucket Join применяется:
# В EXPLAIN ищем: "selected buckets: [0, 1, 2, ...]"
# и отсутствие Exchange/Sort узлов

# Важно: для Delta Lake/Iceberg Bucket Join
# требует специальной конфигурации:
spark.conf.set("spark.sql.sources.bucketing.enabled", "true")
spark.conf.set("spark.sql.sources.bucketing.autoBucketedScan.enabled", "true")

10. AQE Runtime Switching: динамический выбор алгоритма

С Adaptive Query Execution (AQE) выбор между SMJ и другими стратегиями происходит не только статически, но и в процессе выполнения запроса.

Как AQE конвертирует SMJ → BHJ в runtime

# AQE конфигурация для максимальной эффективности:
spark = SparkSession.builder \
    .config("spark.sql.adaptive.enabled", "true") \

    # Порог для AQE-конвертации SMJ → BHJ (используется реальный размер!)
    # Можно поднять выше статического autoBroadcastJoinThreshold
    .config("spark.sql.adaptive.autoBroadcastJoinThreshold", "50MB") \

    # Уменьшение числа Shuffle партиций (убирает пустые партиции)
    .config("spark.sql.adaptive.coalescePartitions.enabled", "true") \
    .config("spark.sql.adaptive.advisoryPartitionSizeInBytes", "128MB") \

    # Skew Join: разбивает горячие партиции при SMJ
    .config("spark.sql.adaptive.skewJoin.enabled", "true") \

    .getOrCreate()

# Диагностика: что AQE сделал с JOIN
def check_aqe_join_conversion(df) -> str:
    """
    Проверяет финальный план после AQE конвертации.
    Полезно для аудита что AQE реально применил.
    """
    # Для финального плана нужен executed plan (после collect)
    plan = df._jdf.queryExecution().executedPlan().toString()

    if "BroadcastHashJoin" in plan:
        return "BHJ (AQE конвертировал из SMJ)"
    elif "ShuffledHashJoin" in plan:
        return "SHJ (AQE конвертировал или explicit hint)"
    elif "SortMergeJoin" in plan:
        # Проверяем на skew handling
        if "SkewedPartition" in plan or "CustomShuffleReader" in plan:
            return "SMJ + AQE Skew Join (горячие ключи разбиты)"
        return "SMJ (стандартный)"
    return "Unknown"

Полная матрица выбора алгоритма JOIN

Шпаргалка в одну строку для каждого алгоритма

Алгоритм Когда применять Главный риск
BroadcastHashJoin Одна сторона < 10-200 MB OOM на Driver/Executor
Bucket Join (SMJ) Обе таблицы заранее bucketed Только для повторяемых JOIN
ShuffleHashJoin ratio > 5:1, нет skew, достаточно памяти OOM при skew или большой Build Side
SortMergeJoin Все остальные случаи Straggler Task при skew, Disk Spill при OOM

Лабораторная работа: сравнение SMJ и SHJ

# lab_smj_vs_shj.py
from pyspark.sql import SparkSession, functions as F
import time

spark = SparkSession.builder \
    .master("local[8]") \
    .appName("smj-vs-shj") \
    .config("spark.sql.adaptive.enabled", "false") \
    .config("spark.sql.shuffle.partitions", "100") \
    .getOrCreate()

spark.sparkContext.setLogLevel("ERROR")

# Создаём тестовые данные с разными ratio
def create_test_data(spark, n_large=5_000_000, n_small=50_000):
    large = spark.range(n_large).select(
        F.col("id").alias("order_id"),
        (F.rand() * n_small).cast("long").alias("product_id"),
        (F.rand() * 1000).alias("amount")
    )
    small = spark.range(n_small).select(
        F.col("id").alias("product_id"),
        F.concat(F.lit("Product_"), F.col("id")).alias("name"),
        (F.rand() * 100).alias("margin")
    )
    return large, small

large, small = create_test_data(spark)
large.cache()
small.cache()
large.count()
small.count()

def benchmark_join(large, small, strategy: str, label: str) -> float:
    spark.conf.set("spark.sql.join.preferSortMergeJoin", "true")

    if strategy == "SMJ":
        result = large.join(small.hint("SHUFFLE_MERGE"), "product_id")
    elif strategy == "SHJ":
        result = large.join(small.hint("SHUFFLE_HASH"), "product_id")

    agg = result.groupBy("name").agg(F.sum("amount"))

    t0 = time.time()
    agg.count()
    elapsed = time.time() - t0
    print(f"  {label}: {elapsed:.1f}с")
    print(f"  Join strategy: {get_join_strategy(result)}")
    return elapsed

print(f"\n=== {large.count():,} строк JOIN {small.count():,} строк ===")
t_smj = benchmark_join(large, small, "SMJ", "SortMergeJoin")
t_shj = benchmark_join(large, small, "SHJ", "ShuffleHashJoin")

print(f"\nИтог:")
print(f"  SMJ: {t_smj:.1f}с")
print(f"  SHJ: {t_shj:.1f}с")
print(f"  SHJ {'быстрее' if t_shj < t_smj else 'медленнее'} SMJ в {abs(t_smj/t_shj):.1f}x")

spark.stop()

Итоги: ключевые правила выбора

SMJ (Sort-Merge Join) — надёжный выбор по умолчанию. Используйте когда не знаете точные размеры данных, когда данные могут быть скошены, когда стабильность важнее скорости. SMJ никогда не падает по OOM — худший сценарий это Disk Spill и замедление.

SHJ (Shuffle Hash Join) — быстрее SMJ при правильных условиях за счёт отсутствия Sort. Применяйте когда одна сторона значительно меньше другой (ratio > 5:1), данные равномерны, и вы уверены что Build Side партиция умещается в Execution Memory. В противном случае — OOM.

Bucket Join — полностью устраняет Shuffle и Sort для повторяемых JOIN'ов. Требует предварительной подготовки таблиц, но окупается многократно для production pipelineов.

AQE в Spark 3.2+ решает большинство случаев автоматически — конвертирует SMJ в BHJ когда одна сторона неожиданно оказывается маленькой после фильтрации. Включайте adaptive.enabled=true всегда.