Out-of-Memory Causes

Common Causes of Out-of-Memory Errors

The listing provokes the three classic failures on the order lines repeated four times (5.5 million rows). A shell loop ran each case as spark-submit --driver-memory 512m l5103_oom.py CASE and echoed its exit code.

Three ways to run out of memoryPython
"""One out-of-memory failure on purpose: spark-submit --driver-memory 512m l5103_oom.py CASE"""
import sys
from pyspark.sql import SparkSession, functions as F
spark = SparkSession.builder.appName("oom_" + sys.argv[1]).getOrCreate()
heap = spark.sparkContext._jvm.java.lang.Runtime.getRuntime().maxMemory() / 2**20
lines = spark.read.parquet("data/order_lines.parquet").crossJoin(spark.range(4))  # 5.5 M rows
cases = {"collect": lambda: lines.collect(),                       # everything to the driver
         "broadcast": lambda: lines.join(F.broadcast(lines.select("order_id", "id")),
                                         ["order_id", "id"]).count(),
         "huge-group": lambda: lines.groupBy("status")              # every row of a status
                                    .agg(F.size(F.collect_list(F.struct("*")))).collect()}
try:
    cases[sys.argv[1]]()
    print(f"{sys.argv[1]} (heap {heap:.0f} MiB): ok")
except Exception as e:
    cause = e.java_exception                                     # the JVM-side exception...
    while cause.getCause() is not None:
        cause = cause.getCause()                                 # ...and its root cause
    print(f"{sys.argv[1]} (heap {heap:.0f} MiB):\n  {str(cause).split('. ')[0]}")
Output
collect (heap 512 MiB):
  java.lang.OutOfMemoryError: Java heap space
  exit code 0
broadcast (heap 512 MiB):
  org.apache.spark.SparkException: Not enough memory to build and broadcast the table to all
    worker nodes
  exit code 0
  exit code 52

collect() pulled every row into the driver, but the error came back as a Python exception and the driver survived. The broadcast failed while the driver built the hash table; Spark 129 's message goes on to suggest disabling broadcasts or adding driver memory. The third case printed nothing: collect_list put 4.9 million delivered rows into one group, the OOM struck inside a task, and the uncaught-exception handler killed the JVM with exit code 52, Spark's out-of-memory code, before Python could report it (two of five runs printed the error first); only the log showed java.lang.OutOfMemoryError: Java heap space. On a cluster that is an ExecutorLostFailure. Fix in this order: keep big data off the driver, don't broadcast big tables, break up huge groups, and only then add memory.