Когда мы пишем обычный df1.join(df2, "id"), Spark не просто сравнивает строки, как это делает классическая база данных. Под капотом он выбирает конкретную физическую стратегию соединения. И если одна из таблиц не подходит для broadcast, то чаще всего Spark использует Sort Merge Join. Это основной вариант по умолчанию для больших таблиц.
Этапы работы Sort Merge Join
🔀Shuffle
Строки обеих таблиц перегоняются по ключам так, чтобы одинаковые ключи оказались в одной партиции.
🤔Sort
Внутри каждой партиции данные сортируются по ключу.
🔥Join
Два отсортированных набора объединяются. Это работает по принципу merge-шага в алгоритме сортировки слиянием.
Плюсы и минусы
У SMJ есть сильная сторона — универсальность. Он подходит для очень больших датасетов и поддерживает все типы соединений, включая outer join. Но у такого подхода есть и минусы: shuffle перегружает сеть и диск, сортировка на больших объёмах требует ресурсов процессора и памяти, а при перекосе ключей одна партиция может стать узким местом и замедлить всю задачу.
Настройки и оптимизация
Spark по умолчанию выбирает Sort Merge Join:
spark.conf.set("spark.sql.join.preferSortMergeJoin", True)Количество shuffle-партиций задаётся параметром spark.sql.shuffle.partitions. Значение по умолчанию — 200, но для больших таблиц этого часто недостаточно, а для маленьких может быть слишком много:
spark.conf.set("spark.sql.shuffle.partitions", 400)Чтобы бороться с перекосами ключей, можно использовать hint
sales.hint("skew").join(customers, "customer_id")Важно учитывать сортировку данных. Если подготовить таблицы заранее, Spark сможет пропускать часть работы:
sales.repartition("customer_id") \
.sortWithinPartitions("customer_id") \
.write.mode("overwrite").parquet("/data/sales_sorted")В таком случае данные будут сразу распределены и отсортированы по ключу, и последующие join-ы выполнятся быстрее.
Как проверить тип join-а
Посмотреть физический план:
print(result.explain(True))
В нём можно увидеть строку:
*(3) SortMergeJoin [customer_id#...], [customer_id#...], Inner
Это означает, что используется именно Sort Merge Join.
Sort Merge Join — это мощная стратегия Spark для больших таблиц: он надёжен и универсален, но требует внимания к настройкам и подготовке данных, поэтому понимание как он работает, и умение читать физический план выполнения позволяет более гибко и эффективно его использовать
Рекомендую разбираться в этой теме — она регулярно встречается в продакшне и на собеседованиях. Полезно и в теории, и в практике.
Ставьте ❤️ если хотите ещё разборы технологий
Ставьте 🔥 если уже сталкивались с проблемами на Sort Merge Join
@bigdata_postupashki