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 spacein the driver, orTotal 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()ortoPandas()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-memoryis 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_listorexplodecreating 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.