Shuffles and Spilling

On the map side of a shuffle, each task sorts its output by target partition and writes one data file plus an index; on the reduce side, each next-stage task fetches its slice from every map output, over the network on a real cluster. When a sort or aggregation buffer outgrows its share of execution memory, Spark 129 spills the sorted buffer to local disk and merges the runs later. The listing gives the driver a small 512 MB heap and sorts the million orders, items and all.

Measuring shuffle and spill through the UI's REST APIJavaScript
import json, urllib.request
from pyspark.sql import SparkSession
spark = (SparkSession.builder.master("local[4]").appName("spill-demo")
         .config("spark.driver.memory", "512m")              # a deliberately small heap
         .config("spark.sql.shuffle.partitions", "4")
         .config("spark.ui.port", "33040").getOrCreate())
(spark.read.parquet("data/orders.parquet")                   # 1M orders with item arrays
 .repartition(4, "customer_id").sortWithinPartitions("customer_id", "order_ts")
 .write.mode("overwrite").parquet("out/orders_by_customer"))
sc = spark.sparkContext                                      # the UI's REST API, same port
url = f"{sc.uiWebUrl}/api/v1/applications/{sc.applicationId}/stages?status=complete"
for s in sorted(json.load(urllib.request.urlopen(url)), key=lambda s: s["stageId"]):
    print(f"stage {s['stageId']}: {s['numTasks']} tasks, shuffle write/read "
          f"{s['shuffleWriteBytes'] >> 20}/{s['shuffleReadBytes'] >> 20} MiB, spilled "
          f"{s['memoryBytesSpilled'] >> 20}/{s['diskBytesSpilled'] >> 20} MiB (memory/disk)")
Output
stage 0: 1 tasks, shuffle write/read 0/0 MiB, spilled 0/0 MiB (memory/disk)
stage 1: 4 tasks, shuffle write/read 45/0 MiB, spilled 0/0 MiB (memory/disk)
stage 3: 4 tasks, shuffle write/read 0/45 MiB, spilled 107/25 MiB (memory/disk)

Stage 3 read 45 MiB of shuffle data and spilled 107 MiB of in-memory rows (25 MiB compressed on disk). With spark.driver.memory at 2g the spill drops to zero. Spill is not an error, but it multiplies disk I/O (Tuning and Memory).