Все sql-джойны внутри превращаются в один из 4х типов:
➡️ BroadcastHashJoin
➡️ SortMergeJoin
➡️ BroadcastNestedLoopJoin
➡️ CartesianProduct
Какой из них используется - смотрим через
df.explain()Сейчас пройдемся по каждому и посмотрим на условия срабатывания.
1️⃣ SortMergeJoin - дефолтный
- условие на равенство
F.col('df1.id') == F.col('df2.id')- ключи должны быть сортируемы
Примеры несортируемых ключей: бинарный формат, сложные структуры
2️⃣ BroadcastHashJoin
Что такое broadcast?
⬇️⬇️Как мы знаем, спарк нужен для параллельной обработки данных. При джойнах экзекьюторам приходится обмениваться данными (чтобы сопоставить одинаковые ключи), что приводит к операции
shuffle - группируются одинаковые ключи из разных экзекьюторов на одном => большие расходы на сеть. Но если у нас есть маленькая табличка (десятки мегабайт, но по умолчанию просто 10), то мы можем скопировать ее на все экзекьюторы, поджойнить там же и избежать шафла. Лимит на размер - память самого экзекьютора.- условие на равенство
- включен broadcast в настройках сессии
# 2й аргумент - размер в МБ
.config('spark.sql.autoBroadcastJoinThreshold', 100)
# так отключается broadcast
.config('spark.sql.autoBroadcastJoinThreshold', -1)
# broadcast join
df1.join(F.broadcast(df2), condition, join_type)
3️⃣ BroadcastNestedLoopJoin
- условие на неравенство
- включен broadcast
NestedLoop - потому что мы итерируемся по маленькому датасету и проверяем каждую строчку
4️⃣ CartesianProduct
- условие на неравенство
Самый медленный, может приводить к OOM-ошибкам (Out of Memory).
#spark