Glossary
What is a Spark shuffle?
In brief
A shuffle is the step where Spark redistributes data across executors so that rows with the same key end up in the same partition, which joins, aggregations, distinct and repartition operations need. Each executor writes its output to local disk, and the next stage fetches it over the network, so a large shuffle costs disk, network and serialization time.
How do you spot a shuffle problem?
In the Spark UI stage view: Shuffle Write Size and Shuffle Write Time on the stage that produces the data, and Shuffle Read Size, Shuffle Remote Reads and Shuffle Read Fetch Wait Time on the stage that consumes it. High fetch wait time means tasks sat idle waiting for data.
How do you reduce shuffle cost?
Avoid the shuffle where you can (broadcast the small side of a join; spark.sql.autoBroadcastJoinThreshold defaults to 10 MB), filter and project before the wide operation, size spark.sql.shuffle.partitions (default 200) to your data, and keep adaptive query execution on so small partitions are coalesced.
Which accelerators address shuffle?
Documented
- Gluten + Velox: ColumnarShuffleManager
- Comet: native and columnar shuffle
- RAPIDS Accelerator: RAPIDS Shuffle Manager
- DualBird: “Reduce shuffle data”, per its product page
- Auron: compacted shuffle formats
Related terms
Sources
- Web UI - Spark documentation · Fetched 27 Sep 2026
- Performance Tuning - Spark SQL documentation · Fetched 27 Sep 2026
- Apache Gluten: Velox backend getting started · Fetched 27 Sep 2026
- Apache DataFusion Comet tuning guide · Fetched 27 Sep 2026
- RAPIDS Accelerator for Apache Spark FAQ · Fetched 27 Sep 2026
- DualBird product page · Fetched 27 Sep 2026
- Apache Auron on GitHub · Fetched 27 Sep 2026