Apache Spark job failures almost always trace back to one of a small set of root causes: memory pressure on the driver or executors, data skew, shuffle instability, lost nodes, too many small files, serialization errors, dependency conflicts, resource misconfiguration, or schema drift, and each one leaves a distinct error string behind that points straight at the fix.
This guide works through ten of those causes in the rough order we run into them in production: driver and executor out-of-memory errors, skew, shuffle failures, lost executors, small files, serialization bugs, classpath conflicts, resource misconfiguration and schema drift. For each one we give the literal error string, the config keys involved with their actual current defaults, how to confirm the cause in the Spark UI or History Server, and the first fix worth trying. Apache Spark job failures rarely need guesswork once you know which tab to open.
Memory Errors Are the Most Common Apache Spark Job Failure Causes
Out-of-memory failures outnumber everything else on this list, and they split cleanly into two categories: the driver and the executors. They look similar in the logs but need opposite fixes, so conflating them wastes a debugging session.
1.Why does the driver run out of memory on collect() and toPandas()?
The driver has a hard ceiling on how much result data it will accept back from the cluster. We see this most often when a notebook calls .toPandas() on a result that was never meant to leave the cluster. spark.driver.maxResultSize defaults to 1g, and the Databricks KB entry for this failure documents the templated form of the message: Total size of serialized results of XXXX tasks (X.0 GB) is bigger than spark.driver.maxResultSize (X.0 GB). It is not limited to an explicit .collect() call. toPandas() triggers it, and so does writing a large file to driver-local storage during a save.
A related and less obvious trigger is a broadcast join against a table that turns out to be bigger than expected. Spark materializes the “small” side of a broadcast join on the driver before shipping it to executors, so a 200MB broadcast table can exhaust a driver that has little headroom even though 200MB sounds small. spark.sql.autoBroadcastJoinThreshold defaults to 10485760 bytes, 10MiB, and setting it to -1 disables broadcast joins outright, forcing a sort-merge join instead. spark.driver.memory defaults to 1g as well, which is frequently the actual gap: a driver sized for orchestration, not for collecting results.
// Stop driver OOM from an oversized broadcast join
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1")
df1.join(df2.hint("merge"), "key")
// Or raise the safety limit explicitly instead of letting collect()
// fail opaquely with a bare OutOfMemoryError
spark.conf.set("spark.driver.maxResultSize", "4g")Raising maxResultSize is a mitigation, not a cure for this class of Apache Spark job failures. The Databricks KB is blunt about it: setting the value very high “can lead to OOM errors,” because you have just moved the failure point further out rather than removing the underlying data volume problem.
2.Why do executors get killed with container exit codes?
Executor memory is really three budgets stacked together: heap, overhead, and (for PySpark) a separate worker process. spark.executor.memory defaults to 1g with a 450m floor, and spark.executor.memoryOverhead defaults to executorMemory * spark.executor.memoryOverheadFactor, where that factor defaults to 0.10, or 0.40 for non-JVM Kubernetes jobs, with a minimum floor applied on top. Overhead covers JVM metaspace, native allocations, and the interpreter process for PySpark UDFs, not the data itself, so a job that looks fine on heap usage can still be killed because the Python worker crept past the overhead allowance, one of the quieter Apache Spark job failures to diagnose from the stack trace alone.
On YARN or Kubernetes, that kill shows up as a container exit rather than a Java stack trace, which is why it reads differently from a driver OOM. Exit code 137 is the convention practitioners use for a SIGKILL from an out-of-memory condition, 143 for SIGTERM; treat those as community convention rather than an Apache-documented table, since no single primary source maps every exit code. The practical read is the same either way: give the executor more overhead headroom, or shrink how much data a single task holds in memory at once by increasing partition count.
Data Skew and Shuffle Failures Stall Otherwise Healthy Jobs
A job can be correctly sized on paper and still crawl, because the work was never spread evenly across tasks in the first place.
3.Why does one straggler task hold up an entire stage?
We treat the Stages tab as the first stop, not config changes, when a job runs long for no obvious reason. Open the Task Duration Distribution table; if the max duration is several times the median while 75th-percentile tasks finished minutes ago, that is skew, not a slow cluster. One partition, usually tied to a hot join or group-by key, is doing most of the work while every other executor sits idle, a pattern behind a good share of Apache Spark job failures that never throw an error at all, just a deadline.
Since Spark 3.0, Adaptive Query Execution can catch this automatically, and since Spark 3.2.0 (SPARK-33679) spark.sql.adaptive.enabled is on by default, so you no longer have to opt in. AQE’s skew-join handling flags a partition as skewed when it exceeds spark.sql.adaptive.skewJoin.skewedPartitionFactor, default 5.0, times the median partition size, and also exceeds spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes, default 256MB. A guide that still describes skew as a problem requiring manual salting on every job is usually describing pre-3.2 Spark; if AQE is on, which it is unless someone disabled it, the engine splits the skewed partition on its own. When AQE genuinely isn’t enough, typically on a key with extreme cardinality concentration, manual salting still works:
from pyspark.sql import functions as F
SALT = 16
facts_salted = facts.withColumn(
"salt", (F.rand(seed=42) * SALT).cast("int"))
dim_exploded = dim.withColumn(
"salt", F.explode(F.array([F.lit(i) for i in range(SALT)])))
joined = facts_salted.join(dim_exploded, ["join_key", "salt"]).drop("salt") 4.Why do jobs fail with FetchFailedException at the shuffle stage?
org.apache.spark.shuffle.FetchFailedException means an executor went looking for shuffle output that another executor was supposed to be holding, and it was gone. The Databricks KB on this error names three concrete causes: the cluster autoscaled down before the shuffle data was read, a spot instance was reclaimed mid-job, or an executor became unresponsive from memory pressure. Under the hood, SPARK-32003 documented a related mechanism in Spark 2.4.6/3.0.0: when the DAGScheduler marked an executor lost, it removed it from the BlockManagerMaster without unregistering that executor’s shuffle files, so tasks fetching from it failed even though the files might still physically exist elsewhere; that specific gap was fixed in 2.4.7/3.0.1/3.1.0, so treat it as a documented historical mechanism rather than current default behavior.
The wrapper message usually reads something like Job aborted due to stage failure: ShuffleMapStage 149 ... has failed the maximum allowable number of times: 4, which is spark.stage.maxConsecutiveAttempts doing its job and giving up. We’ve found the three knobs below worth setting together rather than one at a time: raise spark.shuffle.io.maxRetries above the Spark default of 3, set spark.network.timeout to 800s and spark.rpc.timeout to 600s, and revisit spark.sql.shuffle.partitions, which defaults to 200.
spark.sql.adaptive.enabled true
spark.sql.adaptive.skewJoin.enabled true
spark.sql.adaptive.skewJoin.skewedPartitionFactor 5
spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes 256MB
spark.shuffle.service.enabled true
spark.shuffle.io.maxRetries 10
spark.sql.shuffle.partitions 400 Shuffle spill is a related but distinct symptom worth telling apart from a fetch failure. The Stages tab’s Task Duration Distribution table includes explicit “Shuffle spill (memory)” and “Shuffle spill (disk)” columns; large numbers there mean tasks held more data than executor memory could handle and spilled to disk, which slows a stage without necessarily failing it outright.
Lost Executors and Small Files Are Common Reasons Spark Jobs Fail on Cloud Infrastructure
Two failure classes have nothing to do with the SQL plan and everything to do with the infrastructure underneath it.
5.Spot and Preemptible Node Loss
ExecutorLostFailure and “Executor heartbeat timed out” both point at a node that disappeared out from under a running task. spark.network.timeout defaults to 120s and spark.executor.heartbeatInterval defaults to 10s; the docs are explicit that the heartbeat interval “should be significantly less than” the network timeout, and getting that ratio wrong produces spurious timeouts under ordinary GC pressure, not just real node loss.
We rarely recommend disabling autoscaling to fix this; graceful decommissioning covers most of what autoscaling costs you without giving up the elasticity. It originated as SPARK-20624, filed in 2017 and explicitly aimed at “YARN over-commit, EC2 spot instances, and GCE preemptible instances,” and shipped as an experimental feature in Spark 3.1.1. Since Spark 3.4, setting both spark.decommission.enabled and spark.storage.decommission.enabled to true makes Spark try to migrate cached RDD and shuffle blocks off an executor before it is reclaimed, rather than losing that data and forcing a recompute. spark.storage.decommission.rddBlocks.enabled and spark.storage.decommission.shuffleBlocks.enabled both default to true already once decommissioning itself is on.
6.The Small-File and Too-Many-Partitions Problem
A directory with tens of thousands of small part files produces a Spark job that looks slow before it has done any real work: most of the elapsed time is the driver listing files and building the partition plan, and then the scheduler pays per-task overhead on thousands of tiny tasks that each do almost nothing. We treat a job that “succeeds” in six hours as a failure too, just a quieter one. This is one of the common reasons Spark jobs fail to meet an SLA without throwing any error at all; the job completes and is simply too slow to be useful.
With AQE on, spark.sql.adaptive.coalescePartitions.enabled, which defaults to true since Spark 3.0, merges small post-shuffle partitions toward a configured advisory size automatically, which reduces the need to hand-tune spark.sql.shuffle.partitions for every job. It addresses shuffle output, though, not an already-small-filed source table; for that, a periodic compaction job that rewrites a directory into fewer, larger files remains the standard fix, same as it would be on Hive.
Serialization, UDF and Dependency Failures Break the Job Before It Runs
Some Spark job errors happen before a single row of data moves, because the code shipped to executors was never valid to begin with.
7.Why does Spark throw Task not serializable?
java.io.NotSerializableException: Task not serializable happens when a closure sent to executors references something that cannot cross the wire, commonly a database connection or helper object initialized on the driver and then captured by a lambda. We see this most often on a first PySpark-to-Scala port, not on mature jobs that have already paid this tax. The Databricks Spark Knowledge Base gives four concrete fixes: make the offending class implement Serializable; instantiate the non-serializable object inside the lambda itself so it is created fresh on the worker; use a static singleton; or move initialization into foreachPartition/mapPartitions so a resource like a connection is built once per partition, on the executor, never shipped from the driver at all.
Kryo gets blamed for a related but different error. spark.kryoserializer.buffer defaults to 64k and spark.kryoserializer.buffer.max defaults to 64m with a hard ceiling of 2048m; exceeding it throws org.apache.spark.serializer.KryoException: Buffer overflow. SPARK-20071, filed against Spark 2.1.0 and closed as Incomplete rather than fixed or formally confirmed, describes one reporter’s case of a high-cardinality string column, for example going through StringIndexer, serializing to an object bigger than the buffer allows; treat it as an anecdotal data point, not project-documented behavior. The error message names its own remedy regardless: raise spark.kryoserializer.buffer.max, up to the 2048m ceiling, sized to the largest single object you expect to serialize. Kryo is not Spark’s default serializer. org.apache.spark.serializer.JavaSerializer still is, unless spark.serializer is explicitly set to Kryo, so most “Task not serializable” failures happen under the slower default, not under Kryo at all.
PySpark UDFs add a third variant: a PythonException wrapping a Python-side traceback, often from something as simple as passing a scalar where a Spark SQL function expected a Column. It reads differently from a JVM stack trace and is worth recognizing on sight so you don’t go looking for a Scala bug that isn’t there.
8.Classpath Conflicts, NoSuchMethodError and Mismatched Scala Versions
spark-submit --jars transfers jars to the cluster but performs no dependency resolution at all; every transitive dependency has to be listed by hand or it silently isn’t there at runtime. --packages is the better default for anything with transitive dependencies: it takes comma-delimited Maven coordinates and resolves the full dependency tree through Ivy, and --exclude-packages takes a comma-separated list of groupId:artifactId coordinates to drop conflicting transitive dependencies without hand-managing the rest.
# --jars: no dependency resolution, list every jar yourself
spark-submit --class com.example.MyApp \
--jars /path/dep1.jar,/path/dep2.jar myapp.jar
# --packages: resolves Maven coords + transitive deps via Ivy
spark-submit --class com.example.MyApp \
--packages com.google.guava:guava:31.1-jre \
--exclude-packages com.google.code.findbugs:jsr305 \
myapp.jar We’ve chased more NoSuchMethodError bugs back to a Scala version mismatch than to a genuinely broken build. Spark 3.2.0 added Scala 2.13 support alongside the existing 2.12 build, and mixing a 2.12-compiled jar with a 2.13 runtime, or the reverse, is a binary-incompatibility problem that compiles fine and fails only when the mismatched method signature is actually called. The other usual source is classpath ordering on managed platforms: Hadoop injects its own dependencies ahead of the application’s, so a Guava or Jackson version your build pinned deliberately can lose to an older one Hadoop supplied. Maven Shade relocation fixes it at build time; spark.driver.userClassPathFirst and spark.executor.userClassPathFirst fix it at submit time by telling Spark to prefer the application’s jars.
9.Misconfigured Resources Waste a Cluster Without Ever Crashing It
Not every item on this list ends in a stack trace. Under-provisioning shows up as a job that finishes, just slower and more expensively than it should.
The Spark tuning guide’s own recommendation is 2-3 tasks per CPU core, a different and more conservative figure than the “5 cores per executor” rule often repeated online, which actually traces back to a Cloudera blog post rather than Apache documentation. We default new jobs to the tuning-guide figure for that reason: it comes from the primary source, not a secondary one. spark.dynamicAllocation.enabled defaults to false, with minExecutors at 0 and executorIdleTimeout at 60s when it is turned on, so a cluster sized for peak load by default stays at peak load the whole job even during a mostly-idle tail stage, unless dynamic allocation is explicitly enabled. Retry behavior is tunable too: spark.task.maxFailures defaults to 4, meaning three retries before a task is given up on, and spark.stage.maxConsecutiveAttempts, since Spark 2.2.0, defaults to the same number for a stage as a whole. Spark job errors that look like flakiness are sometimes just these defaults doing exactly what they’re configured to do on a genuinely unstable cluster.
10.Schema Drift and Null Handling Cause Silent Failures
Not every failure on this list is loud. Some of the worst ones are quiet: a job finishes, the numbers are wrong, and nobody notices until a downstream report does. We treat a silently dropped column as a worse outcome than a thrown exception, because nobody triages a report that merely finished on time.
Why does AnalysisException appear after an upstream schema change?
AnalysisException covers a family of planning-time failures, two of the most common being a literal Path does not exist: <uri> and cannot resolve '<col>' given input columns, both typically the result of upstream schema drift or a path typo rather than a Spark defect. spark.sql.parquet.mergeSchema defaults to false, and has since Spark 1.5.0, specifically because schema merging across files is expensive to compute automatically. That means a Parquet dataset with column drift across files will not silently reconcile; it will either throw or silently drop a column, depending on how the read was structured, unless mergeSchema is explicitly set to true.
JSON and CSV readers default to PERMISSIVE mode: malformed records get their fields nulled out and, if the schema declares a columnNameOfCorruptRecord field, the raw offending string is stashed there for inspection; without that field declared, corrupt records are simply dropped with no error at all. FAILFAST throws on the first bad record instead, and DROPMALFORMED discards it silently, close to what PERMISSIVE already does by default when no corrupt-record column is declared.
Null handling causes a specific, recurring class of wrong-answer bug rather than a crash. Spark SQL uses three-valued logic, so WHERE col = NULL matches nothing, ever, because the comparison evaluates to unknown rather than true or false. The null-safe equality operator, <=>, is the one that actually tests for NULL correctly, and it’s an easy one-character fix once you know a filter is silently dropping every row you expected it to match.
Confirming Root Cause With the Spark UI and History Server
Most Apache Spark job failure causes are cheap to fix once confirmed, and expensive to fix blind. Confirming the actual cause before changing a config key is what separates a five-minute fix from an afternoon of trial and error.
Reading the Stages Tab and Task Duration Distribution
The Stages tab’s Task Duration Distribution table reports min, 25th percentile, median, 75th percentile and max for Duration, GC time, Shuffle Read Size, and both shuffle spill columns, one of the quickest ways to confirm skew or spill without reading a line of code. The Executors tab adds a per-executor Thread Dump link for catching an executor that is stuck rather than merely busy, plus direct stderr/stdout log links per executor. The Environment tab’s Classpath Entries sub-tab lists the driver classpath by source, a quick way to confirm which jar version actually won a conflict instead of guessing from a NoSuchMethodError message alone.
Event Logs and the History Server After a Job Exits
The web UI dies with the application. Spark’s own docs state plainly that once the application exits, the UI is no longer reachable. Enabling spark.eventLog.enabled with spark.eventLog.dir set, alongside spark.history.fs.logDirectory on the History Server side (default port 18080), keeps that data queryable after the fact instead of losing it the moment a job finishes or crashes. On YARN specifically, setting yarn.log-aggregation-enable=true means a single yarn logs -applicationId <app ID> pulls every driver and executor container log for a finished application in one shot, rather than hunting per-node log directories for the host that happened to run a given container. For a classpath problem that only reproduces under YARN, setting yarn.nodemanager.delete.debug-delay-sec to a large value keeps the container’s launch script, jars and environment on disk long enough to inspect before YARN cleans it up.
A Troubleshooting Table for Apache Spark Job Failures
The table below maps the error strings covered above to a likely cause and a first fix worth trying before anything else.
| Symptom / error string | Likely cause | First fix |
|---|---|---|
| Total size of serialized results ... bigger than spark.driver.maxResultSize | collect()/toPandas() or a large broadcast pulling too much data to the driver | Raise spark.driver.maxResultSize, or disable the broadcast via autoBroadcastJoinThreshold=-1 |
| Executor killed, container exit code (commonly cited as 137) | Executor memory + overhead undersized for off-heap or PySpark worker use | Raise spark.executor.memoryOverhead, or spark.executor.memory |
| One task far longer than the rest on the Stages tab | Data skew on a join or group-by key | Confirm AQE skew-join is enabled (default since 3.2); salt the key if not enough |
| org.apache.spark.shuffle.FetchFailedException | Executor lost (autoscale-down, spot reclaim, OOM) before shuffle data was read | Raise spark.shuffle.io.maxRetries, network.timeout, rpc.timeout |
| ExecutorLostFailure / “heartbeat timed out” | Spot/preemptible termination or GC pause past heartbeatInterval | Enable graceful decommissioning; widen the heartbeat/timeout ratio |
| Long pre-stage “listing” time, thousands of tiny part files | Upstream job wrote many small files, or shuffle partitions oversized | Compact the source directory; let coalescePartitions merge shuffle output |
| java.io.NotSerializableException: Task not serializable | Non-serializable object captured in a closure sent to executors | Init the object inside foreachPartition/mapPartitions, not on the driver |
| KryoException: Buffer overflow | Serialized object larger than spark.kryoserializer.buffer.max | Raise buffer.max, up to the 2048m ceiling |
| NoSuchMethodError despite a clean build | Scala 2.12/2.13 mismatch, or Hadoop classpath precedence | Align Scala versions; use --exclude-packages or userClassPathFirst |
| AnalysisException: cannot resolve '<col>' given input columns | Upstream schema drift, mergeSchema default of false | Set mergeSchema=true if intentional, otherwise fix the upstream schema |
How Yeedu Assistant Helps Debug Apache Spark Job Failures
Every section above ends with the same chore: open the driver and executor logs, find the run in the History Server, check what the cluster was doing at the time, and line the three up until the cause is obvious. Yeedu Assistant (Assistant X) is the debugging companion built into Yeedu that does that gathering for you. Point it at a failed or slow job and it fetches the job’s Spark logs, reads the Spark History Server event data for that run, and pulls the cluster monitoring details, so the evidence is collected in one place instead of across three tabs and a terminal.
Collecting the evidence is only half of the work. Yeedu Assistant is also trained with Spark troubleshooting skills, and it uses them to perform root-cause analysis: it explains why the job failed and recommends a fix. Take the maxResultSize failure from the first section. Yeedu Assistant can read the error from the driver log, connect it to the stage that sent the oversized result back, and recommend the change that fits, whether that is raising spark.driver.maxResultSize or keeping the data on the cluster instead of calling toPandas(). For a stage stalled on one straggler task, the History Server event data shows the uneven task durations that point to skew, and the recommendation follows from there.
Beyond a single failure, Yeedu Assistant also gives configuration recommendations and code assistance, which covers the resource-sizing and serialization problems on this list. It works as read-only analysis within a session, with no persistent storage of user data, so it can inspect a production job without changing it. The Spark UI and History Server habits described in this guide still apply; Yeedu Assistant runs the same checklist on every failed job, so you start from a diagnosis rather than a blank log window.
Matching the Fix to the Failure
Ten causes, ten different first moves, and almost all of them come down to a single config key once the root cause is confirmed rather than guessed. The Spark UI and History Server exist precisely so that confirmation step doesn’t require re-running a job three times to see if a hunch was right.
The pattern worth keeping is to read the error string literally before reaching for a config change. maxResultSize and FetchFailedException look unrelated on the surface, but both are the cluster telling you, correctly, where it hit a wall: one on the driver, one on a lost shuffle block. Treat the message as the starting point for the investigation, not an obstacle to silence with a bigger number.



