Shuffle — перераспределение данных между партициями/экзекьюторами (по ключу join/groupBy/repartition): запись на диск, сеть, чтение. Дорогая операция.
Разбор
- Wide transformation обычно вызывает shuffle.
- Узкие (map/filter) — без обмена между partition.
- Смотрите Spill и Shuffle Read/Write в UI.
- Уменьшение shuffle — главный рычаг оптимизации.
Итог
Shuffle = сетевой обмен данными между стадиями; избегайте лишнего.