When Spark has to move rows between executors, as joins, sorts and aggregations require. Per stage, the Spark UI lists Shuffle Read and Shuffle Write; heavy shuffling frequently causes spill.
Read more: Microsoft Learn
In the Ultra Transcenders books
Each book explains Shuffle in context, with comparison tables and the common traps.
Terms in this definition
- Aggregations
Summary queries over a large DirectQuery table can be answered from memory instead of the source thanks to this semantic model feature: it keeps a concealed summary table cached and sends qualifying queries there, while anything needing fine detail still goes back to the source.
- Spark UI
Classic compute's Apache Spark web console. Databricks advises checking the jobs timeline first, then the slowest stage, skew or spill, I/O and lastly the SQL DAG. Serverless and SQL warehouses offer a query profile in its place.
- CRUD
Shorthand for create, read, update and delete, the four basic things you do with data. Data-plane roles in Azure Cosmos DB, for instance, authorise those operations on items.
- Spill
Spark writing data to disk because shuffles, joins, sorts or aggregations have exhausted execution memory, which is costly. It is reported in the stage details of the Spark UI.
Related terms
- AQE
Adaptive query execution: Spark re-plans a query while it runs, for instance turning a sort-merge join into a broadcast join, merging shuffle partitions, dealing with skew and propagating empty relations. Databricks advises leaving it enabled.
- Autotune
A preview option in Fabric Spark that studies previous runs to pick better values for shuffle partitions, broadcast join threshold and max partition bytes query by query. It's disabled by default (
spark.ms.autotune.enabled), limited to Runtime 1.2, and incompatible with private endpoints and high concurrency mode. - Broadcast join
A join method that copies the smaller table to every executor so that the larger table needs no shuffle. You can force it with the
BROADCASThint (orBROADCASTJOIN/MAPJOIN), or AQE may pick it while the query runs. - Join hints
Hints in SQL that recommend a particular join algorithm, such as
BROADCAST,MERGE,SHUFFLE_HASHorSHUFFLE_REPLICATE_NL; they are suggestions only, so the optimiser may choose otherwise.