В спарке существуют 2 вида трансформаций: узкие и широкие.
💃Узкие не требуют перемещений данных и на любом маленьком кусочке могут выполняться параллельно и независимо:
where(), withColumn(), union(). Например, чтобы отфильтровать строки, нам не нужно знать весь датасет. Мы берем одну строку, применяем условие - готово.🍊Широкие же требуют шафла:
join(), groupBy(), sort(), distinct(). Здесь же нам нужен весь датасет. Допустим, мы хотим сделать дистинкт по полю color: на первом экзекьюторе лежат red, blue, green, на втором yellow, violet, blue. Если брать отдельно каждый экзекьютор, то цвета уникальны, но если мы возьмем все, то будут дубликаты. То есть нам сначала надо одинаковые значения собрать (это и есть шафл) и только потом почистить.Аналогично работают и джойны, поэтому нужно уменьшить количество данных на стадии шафла.
Есть несколько советов:
1️⃣Все фильтры до джойнов
2️⃣Использовать equi-джойны (SortMergeJoin, BroadcastHashJoin)
3️⃣Если можно увеличить данные, но вместо non-equi (NestedLoop, Cartesian) использовать equi, то делать именно так
4️⃣Если правый датасет помещается в память экзекьютора, использовать broadcast
5️⃣Избегать cross-join'ов
6️⃣Перепроверять в плане запросов
#spark
