All concepts

The Shuffle

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.

Distributed Processing · Advanced · ~6 min

In plain English

Everyone in a hall regrouping by birthday month. Everybody moves at once, and until the last person sits down nothing else can happen.

Why it's worth your time

It is the dominant cost in almost every distributed job, and the number of them is set by how you wrote the query.

If you remember three things

  • Every row is serialised, written to disk and fetched across the network
  • Target ~128–200 MB per shuffle task; both extremes are slow
  • Bucketing pays the exchange once, at write time

Overview

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.

In an interview

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.

Production defaults

Partition sizing
check spill metrics rather than guessing a number
Bucketing
worth it for tables joined on the same key daily
Diagnosis
max task duration ÷ median > 5 means skew, not partition count

What breaks

  • Massive disk spill in one stage — Too few shuffle partitions. Increase until per-task input is a few hundred MB.
  • Thousands of tiny fast tasks — Over-partitioned; scheduling overhead dominates. Let adaptive execution coalesce them.

Watch it explained

Spark Basics | Partitions — Palantir Developers, 5:12

Related