Every row with the same key has to end up on the same machine, and getting it there means writing the whole dataset to disk and reading it back across the network.
Everyone in a hall regrouping by birthday month. Everybody moves at once, and until the last person sits down nothing else can happen.
It is the dominant cost in almost every distributed job, and the number of them is set by how you wrote the query.
A shuffle is the redistribution of data so that rows sharing a key land together. Each task writes its output into buckets by hash of the key, those files are read by the tasks that own each bucket, and the whole dataset therefore makes a round trip through serialisation, local disk and the network. It is the single most expensive thing a distributed job does, which makes the two questions worth asking about any slow pipeline: how many shuffles does it perform, and is the data spread evenly across the buckets it creates.
A shuffle redistributes rows so equal keys are co-located: map-side tasks write hash-partitioned files, reduce-side tasks fetch their bucket from every writer. That means serialisation, disk and network for the whole dataset — the dominant cost in most jobs. Fewer shuffles and even partition sizes are what make a job fast; partition count controls how big each task's slice is.
Spark Basics | Partitions — Palantir Developers, 5:12