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 и матрица выбора.
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 всегда.