Shuffle: что это, когда происходит и почему стоит так дорого
Shuffle — перераспределение данных между исполнителями так, чтобы все строки с одинаковым ключом оказались в одной партиции. Он неизбежен для groupBy, join, distinct, repartition, оконных функций с partitionBy и сортировки. Дорог он потому, что задействует всё сразу: данные сериализуются, пишутся на диск каждого исполнителя, передаются по сети и читаются принимающей стороной. Кроме того, shuffle — это барьер: следующая стадия не начнётся, пока не завершится вся предыдущая, поэтому одна медленная задача задерживает весь джоб. Отсюда стратегия оптимизации: сначала уменьшить объём данных до shuffle (фильтры и select нужных колонок как можно раньше, предагрегация), затем избежать самого shuffle там, где можно (broadcast-джойн вместо sort-merge), и лишь потом настраивать его параметры. Полезный ориентир при чтении Spark UI: узел Exchange в плане — это и есть shuffle.