Spill and GC

Spill, Garbage Collection and Memory Pressure

Before memory runs out, it gets tight: operators spill sorted runs to disk and the JVM spends more time collecting garbage. Both show in the Stages tab and its REST API. The listing numbers each customer's order lines with a window function, a sort of 2.8 million rows, at two heap sizes.

Spill and GC time for one sort-heavy jobPython
"""Spill and GC for one sort-heavy job: spark-submit --driver-memory 512m l5104_pressure.py"""
import json, time, urllib.request
from pyspark.sql import SparkSession, Window, functions as F
spark = SparkSession.builder.appName("pressure").getOrCreate()
sc = spark.sparkContext
lines = spark.read.parquet("data/order_lines.parquet").crossJoin(spark.range(2))   # 2.8 M rows
w = Window.partitionBy("customer_id").orderBy("order_ts", "id")
t0 = time.perf_counter()
(lines.withColumn("n", F.row_number().over(w))                      # sorts every customer
 .write.format("noop").mode("overwrite").save())
secs = time.perf_counter() - t0
api = f"{sc.uiWebUrl}/api/v1/applications/{sc.applicationId}/stages"
st = json.load(urllib.request.urlopen(api))
mib = lambda k: sum(s[k] for s in st) / 2**20
task, gc = sum(s["executorRunTime"] for s in st), sum(s["jvmGcTime"] for s in st)
print(f"driver {sc.getConf().get('spark.driver.memory'):<4} {secs:5.1f}s  "
      f"spill {mib('memoryBytesSpilled'):4.0f} MiB (disk {mib('diskBytesSpilled'):3.0f})  "
      f"GC {100 * gc / task:4.1f}% of task time")
Output
driver 512m  22.9s  spill  276 MiB (disk  47)  GC  4.4% of task time
driver 2g    21.5s  spill    0 MiB (disk   0)  GC  4.0% of task time

The small heap spilled 276 MiB (47 MiB compressed on disk) yet ran only a little slower: spilling to a local SSD is cheap and is Spark 129 working as designed. Spill hurts when it is large relative to the input or lands on slow disks; GC above roughly 10% of task time signals real pressure. Remedies: more partitions, fewer cores per executor, then more memory. G1 is the default collector; pass GC flags in spark.executor.extraJavaOptions.