Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. Driver OOM vs executor OOM

PySpark · Memory, OOM & Tuning

Driver OOM vs executor OOM

Hardpyspark-55
scenariooomdriverexecutortroubleshooting

Question

Your job fails with OutOfMemoryError. How do you tell whether it's the driver or an executor, and what do you fix for each?

Solution

First read the error and find where it happened. The stack trace and the Spark UI tell you whether the driver or an executor ran out of memory. The fixes are different, so do not touch memory settings until you know which one it is.

Telling them apart

  • Driver OOM: the message appears in the driver log or your console, and the application dies at once. Typical text: java.lang.OutOfMemoryError: Java heap space in the driver, or Total size of serialized results ... is bigger than spark.driver.maxResultSize.
  • Executor OOM: the Spark UI shows a failed task with "ExecutorLostFailure" or an OutOfMemoryError on one executor, then retries. The Executors tab shows lost executors. If the message is "Container killed by YARN for exceeding memory limits" or the pod is OOMKilled, memory overhead was exceeded, not the heap.

Driver causes and fixes

  • collect() or toPandas() on a large result. Aggregate first, limit, or write to storage.
  • Broadcasting a huge table. Check the plan and the threshold.
  • A very large number of partitions or files. The driver keeps metadata for every task and every file. Reduce partition counts or compact the files.
  • Fix: stop pulling data to the driver. Raising --driver-memory is a last step.

Executor causes and fixes

  • Skewed partition: one task holds far more data than the others. Look at max versus median task input in the stage page. Fix the skew.
  • Partitions too big because there are too few of them. Increase shuffle partitions or repartition.
  • groupByKey, collect_list or explode creating huge groups or many rows.
  • Caching too much with a memory-only level.
  • Python UDFs or pandas conversions eating overhead (container killed). Raise memoryOverhead, and reduce batch size.

Process

Reproduce on a smaller slice if you can. Find the stage and the task that failed. Compare its size with the median. Then change one thing and rerun.

The usual wrong move is doubling executor memory. It sometimes works and hides the cause, and a skewed key will simply need more memory again next month. Show that you look for why the data does not fit before you add capacity.

🎯 Put this concept into practice

Solidify this answer with real hands-on interview drills in the browser studio.

Open related drill →
PreviousNext