Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. What is a shuffle?

Pipelines & scenarios · Core Pipeline Concepts

What is a shuffle?

Mediumpipe-11
shuffleSparknetworkgroup byperformance

Question

What is a shuffle in distributed data processing?

Solution

A shuffle is when a distributed engine must redistribute data across workers so that records with the same key land on the same machine. It is network + disk heavy.

Before shuffle (data local to workers):
  Worker A: (user=1, ...), (user=7, ...)
  Worker B: (user=1, ...), (user=3, ...)

After shuffle for GROUP BY user_id:
  Worker A: all user=1 rows
  Worker B: all user=3,7 rows

Operations that trigger shuffles

  • GROUP BY, DISTINCT, many joins
  • ORDER BY / global sort
  • repartition(key)

Why shuffles hurt

Bytes cross the network, stages wait on the slowest worker, skew amplifies pain, and spills to disk if memory is tight.

How to reduce shuffled bytes

  • Filter early
  • Use columnar formats and only needed columns
  • Broadcast small tables
  • Prefer map-side combiners / partial aggregates
  • Avoid unnecessary repartition

Interview tip: "Shuffle = redistributing records by key across the cluster." Say it is often necessary for correctness; the skill is minimizing shuffled bytes and handling skew.

PreviousNext