Skip to content
Reliable Data Engineering
Practice problem hard spillexecutor-tuningshuffle-partitionsspark-uimemory
Practise with timer, notes and rubric

Diagnose Spill From Spark UI Metrics and Size the Executors

Difficulty: Hard · Topics: spill, executor sizing, shuffle partitions, memory model · Asked at: Databricks, Amazon, LinkedIn, Pinterest

Scenario

A nightly PySpark aggregation over 1.2 TB of Parquet (8 columns read out of 60) runs for 2 hours 40 minutes on a YARN cluster of 12 nodes, each 32 cores and 128 GB RAM. Current settings:

spark.executor.instances=12
spark.executor.cores=31
spark.executor.memory=110g
spark.sql.shuffle.partitions=200
spark.sql.adaptive.enabled=true

The job reads the data, applies a pandas UDF that parses a JSON column, then does groupBy(account_id, day).agg(...) with 6 aggregates.

Evidence

Spark UI, the aggregation stage (200 tasks):

MetricMinMedian75th pctMax
Duration38 min41 min43 min52 min
Shuffle Read Size2.9 GB3.1 GB3.2 GB3.6 GB
Spill (Memory)19 GB21 GB22 GB25 GB
Spill (Disk)3.1 GB3.4 GB3.5 GB4.0 GB
GC Time9 min11 min12 min15 min

Also observed: several executors were lost with Container killed by YARN for exceeding memory limits. 118.4 GB of 115.5 GB physical memory used, and the whole stage retried once.

Your task

  1. Is this skew, under-partitioning or an executor-shape problem? Justify from the numbers.
  2. Explain the executor losses.
  3. Propose a new configuration with the arithmetic, and other changes to the job.
  4. Say what you’d check in the UI after the change.

Hints

Hint 1

Compare max vs median shuffle read. What does an even distribution with large values tell you?

Hint 2

Where do pandas UDF workers allocate memory: on the JVM heap or outside it? How many run concurrently in one 31-core executor?

Solution

1. Diagnosis: under-partitioning plus a bad executor shape, not skew.

2. Executor losses. YARN counts the whole container: heap (110 GB) + overhead. Pandas UDFs run in Python worker processes outside the JVM heap, one per concurrent task (up to 31 per executor), each holding Arrow batches and pandas DataFrames. The default overhead (10% ≈ 11 GB) can’t hold 31 Python workers, so the container exceeds its limit and is killed. Losing an executor also loses its shuffle output, so the stage retries.

3. New configuration.

Executors (per node: 32 cores, 128 GB):

Shuffle partitions:

spark.executor.instances=71
spark.executor.cores=5
spark.executor.memory=14g
spark.executor.memoryOverhead=6g
spark.sql.shuffle.partitions=4000
spark.sql.adaptive.coalescePartitions.initialPartitionNum=4000
spark.sql.adaptive.advisoryPartitionSizeInBytes=128m
spark.sql.execution.arrow.maxRecordsPerBatch=5000       # smaller Arrow batches → less Python memory

Job changes:

4. Verify after the change: spill ≈ 0, GC time < 10% of task time, task durations of seconds to a couple of minutes, max/median ratios still ~1-2×, no executor loss, and the stage’s total shuffle read unchanged (or lower after projection). Expected runtime: roughly 10-20 minutes.

What interviewers look for