ETL CompareETL Compare

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 Documented

Not stated Not stated

  • Photon: not stated on the pages we reviewed
  • Flarion: not stated on the pages we reviewed

Related terms

Sources