Debugging Spark is harder than debugging a single-process program for three reasons. Evaluation is lazy, so an error in line 12 is reported at line 80, where the action runs. Work runs on many machines, so the useful stack trace is often in an executor log, not in front of you. And Spark retries, so a task can fail three times quietly before the fourth failure stops the job.
This article is about the bugs that stop jobs or corrupt results, not about making correct jobs faster. It shows how to read a Spark failure, the common failure types and what each means, how to find the one record that breaks a job, where the logs are, how to attach a debugger, and how to track down wrong answers. For reading the web UI itself, see the Spark UI in depth; for plans, see reading EXPLAIN plans.
Why the error appears in the wrong place
Transformations such as select, filter and join only build a plan. Nothing runs until an action such as count, collect or a write. At that point the scheduler cuts the plan into stages and tasks and sends them to executors. If a task throws, the executor logs the exception and reports it to the driver, which retries the task on another attempt, up to spark.task.maxFailures (4 by default). Only when the limit is reached does the driver abort the job and throw a SparkException from the line that called the action.
Two practical results follow. First, the line number in your notebook is the action, not the bug. Second, a flaky failure, such as a lost executor or a network blip, may be retried silently and never reach you, while a deterministic one, such as a bad record, fails all four attempts identically.
Reading a Spark failure
A job failure message has a fixed structure. Read it in this order.
org.apache.spark.SparkException: Job aborted due to stage failure: Task 17 in stage 4.0
failed 4 times, most recent failure: Lost task 17.3 in stage 4.0 (TID 212)
(10.0.3.7 executor 5): org.apache.spark.api.python.PythonException:
Traceback (most recent call last):
File "/jobs/enrich.py", line 41, in parse_amount
return Decimal(raw.replace(",", ""))
decimal.InvalidOperation: [<class 'decimal.ConversionSyntax'>]
...
Driver stacktrace:
at org.apache.spark.scheduler.DAGScheduler.failJobAndIndependentStages(...)
...- Which stage and task, how many attempts.
Task 17 in stage 4.0 failed 4 times. Four failures on the same task usually means a deterministic problem with that task's data. Failures spread across many tasks and executors point at the environment instead. - Where it ran. The executor ID and host. If every failure is on one host, suspect the host.
- The first real exception. Skip the scheduler frames and read the exception after the task line, then every
Caused bydown to the last one. The lastCaused byis usually the root. For Python UDFs, the Python traceback is embedded as aPythonException. - The stage in the UI. Stage 4 in the Stages tab tells you which part of your code it is, from the stage's description and its SQL plan node.
If the driver message is truncated or says only ExecutorLostFailure, the cause is in the executor's own log or in the cluster manager's record of why the container died.
A taxonomy of common failures
| Symptom | What it usually means | Where to look and what to try |
|---|---|---|
Task not serializable, NotSerializableException | A closure captured an object that cannot be shipped, often a client or connection created on the driver | Create the object inside mapPartitions on the executor, or make the captured field local |
OutOfMemoryError: Java heap space on an executor | One task holds too much: a skewed key, a huge row, an explode, a large collect_list | Stage task metrics for a max far above the median; salt or split the key, raise partitions |
| Container killed for exceeding memory limits, exit code 137 or 143 | The process as a whole used more than heap plus overhead (137 is 128 + SIGKILL); off-heap, Python workers or native buffers | Raise spark.executor.memoryOverhead (default 10% of executor memory, minimum 384 MiB) or reduce Python and Arrow batch memory |
| Driver OOM | collect or toPandas of a large result, broadcast of a big table, huge plans | Write results out instead of collecting; check broadcast sizes |
FetchFailedException | A reducer could not read shuffle output, usually because the executor holding it died | Find why that executor died first; the fetch failure is the second symptom. Stages are retried up to spark.stage.maxConsecutiveAttempts (4) |
AnalysisException | A plan error found before running: missing column, ambiguous reference, type mismatch | Read the message; it names the column and shows the plan |
DIVIDE_BY_ZERO, CAST_INVALID_INPUT | ANSI mode, on by default since Spark 4.0, raises errors where older versions returned NULL | Fix the data or use try_divide and try_cast where NULL is the intended answer |
| Python worker crashed, no traceback | A segfault in a native library called from a UDF | Enable spark.sql.execution.pyspark.udf.faulthandler.enabled (Spark 4.0 and later) to dump the Python stack |
Memory failures are often data-shape problems in disguise. Skew is covered in depth elsewhere; adaptive query execution handles some of it automatically, as described in adaptive query execution.
Shrink the problem: find the record
A deterministic task failure almost always comes down to a few records. Finding them is faster than reasoning about them. Three techniques cover most cases.
Keep malformed input. The JSON and CSV readers can keep rows they cannot parse in a corrupt-record column instead of failing or setting every field to NULL. The column must be declared in the schema. Spark does not allow a query that references only the corrupt-record column on raw files, so cache the DataFrame first.
Make a fragile function report instead of throw. Wrap the UDF so it returns a value and an error string. Now one run shows every failing input, with the file name and partition ID, instead of dying on the first one.
Reproduce small. Once you know the file and partition, copy that input locally and run the same code with local[2]. A 20-second loop on one file beats a 20-minute loop on a cluster.
from decimal import Decimal
from pyspark.sql import functions as F, types as T
# 1. Keep malformed input instead of failing or silently nulling it
schema = T.StructType([
T.StructField("order_id", T.StringType()),
T.StructField("amount", T.StringType()),
T.StructField("_corrupt", T.StringType()), # must be in the schema
])
raw = (spark.read.schema(schema)
.option("mode", "PERMISSIVE")
.option("columnNameOfCorruptRecord", "_corrupt")
.json("s3://bucket/orders/2026-09-29/")
.withColumn("file", F.input_file_name()) # capture before caching
.cache()) # needed to query _corrupt alone
raw.filter(F.col("_corrupt").isNotNull()).show(5, truncate=False)
# 2. Wrap a fragile UDF so a bad row is reported, not fatal
result_t = T.StructType([T.StructField("value", T.DecimalType(18, 2)),
T.StructField("error", T.StringType())])
@F.udf(result_t)
def parse_amount_safe(raw_amount):
try:
return (Decimal(raw_amount.replace(",", "")), None)
except Exception as e: # debugging only; narrow it later
return (None, f"{type(e).__name__}: {raw_amount!r}"[:300])
checked = raw.withColumn("p", parse_amount_safe("amount"))
bad = (checked.filter("p.error IS NOT NULL")
.select("order_id", "amount", "p.error", "file",
F.spark_partition_id().alias("partition")))
bad.show(20, truncate=False)Remove the broad except once you know the failure. Left in production, it converts a loud failure into silent NULLs.
Finding the logs
Print statements and log calls inside a UDF or mapPartitions run on executors, so their output goes to the executor's stdout and stderr, not your notebook. While the application runs, the Executors tab links to each executor's logs. After it ends, the location depends on the cluster manager: aggregated YARN logs, the executor pods' logs on Kubernetes (which disappear when pods are deleted unless shipped elsewhere), or your managed platform's log store. For finished applications, the event log gives you the UI back through the History Server, but not the executor logs; ship those separately.
Change the log level of a running session with spark.sparkContext.setLogLevel("DEBUG"). Do this briefly and on a small input: debug logging from every executor can produce gigabytes and slow the job enough to change its behaviour.
# Attach an IDE debugger: local mode runs driver and executors in one JVM,
# so breakpoints in Scala/Java UDFs and your own classes are hit.
spark-submit --master "local[2]" \
--driver-java-options "-agentlib:jdwp=transport=dt_socket,server=y,suspend=y,address=*:5005" \
--class com.example.Enrich target/enrich.jar --input sample/
# Logs of a finished YARN application, all containers
yarn logs -applicationId application_1727600000000_0042 > app.log
# Kubernetes: executor pods are named after the application
kubectl logs -n spark <executor-pod-name> --previousA debugger is most useful for JVM code in local mode, where the driver and executors share one JVM and a breakpoint in your UDF is hit. On a cluster, a live thread dump is more practical: the Executors tab has a thread dump link for each executor, which shows what a stuck task is doing, such as waiting on a lock, a remote call or a slow file system.
Debugging wrong results
A job that finishes with the wrong answer is worse than one that fails. The method is the same as for any data pipeline: check invariants at each step until the first step where they break.
def audit(df, name, keys=None):
"""Cheap-ish invariants at each step of a pipeline. Each call runs a job."""
n = df.count()
line = f"{name}: rows={n}"
if keys:
nulls = df.filter(" OR ".join(f"{k} IS NULL" for k in keys)).count()
dups = df.groupBy(*keys).count().filter("count > 1").count()
line += f" null_keys={nulls} duplicate_keys={dups}"
print(line)
return df
orders = audit(spark.table("orders").filter("dt = '2026-09-29'"), "orders", ["customer_id"])
customers = audit(spark.table("customers"), "customers", ["customer_id"])
joined = audit(orders.join(customers, "customer_id", "left"), "joined")
# joined rows > orders rows means duplicate keys on the right side fanned the join out- Row counts that grow after a join mean duplicate keys on one side. A left join of 10 million orders that returns 10.4 million rows has a customer table with duplicates.
- NULL join keys never match, so inner joins drop those rows silently. Use
<=>when NULL should match NULL. - Non-determinism.
first()after agroupBy,rand()without a seed,monotonically_increasing_idandlimitwithout an ordering can give different results on each run or on a task retry. If two runs of the same job differ, look for these first. - Time zones. Timestamp parsing and date truncation use
spark.sql.session.timeZone. A cluster and a laptop with different defaults produce different daily totals. - Accumulators used as counters. Updates made inside transformations can be applied more than once when tasks are retried or data is recomputed. Only updates inside actions are counted exactly once, so use accumulators as hints, not audit numbers.
When invariants hold but the numbers are still wrong, compare explain("formatted") output with what you intended: a filter placed after a join instead of before, or a join type different from the one you wrote.
Worked example: exit 137 and a short count
A nightly PySpark job enriches 60 million orders and writes a partitioned table. After a library upgrade it starts failing about half the time with ExecutorLostFailure, and on nights it succeeds the finance team reports 0.3% fewer orders than the source.
The failures. The YARN diagnostics for the lost executors say the container exceeded its memory limit, exit code 137. Heap usage in the Executors tab is modest, so the JVM heap is not the problem. The killed containers were running a stage with a pandas UDF; the same change had raised the job's Arrow batch size from 10,000 to 100,000 rows. Python workers run outside the JVM heap, inside the overhead allowance, which was the default 10%. Setting spark.sql.execution.arrow.maxRecordsPerBatch back to 10,000 and memory overhead to 2 GiB stops the kills.
The missing rows. The audit helper shows 60,000,000 rows after the read and 59,820,000 after the join with currencies. It is an inner join, and null_keys on the orders side is 180,000: rows with a missing currency code, which an inner join drops. The fix is a left join with a default currency and a data-quality alert on NULL codes. Both problems were visible within an hour once the evidence was read in the right place.
Pitfalls when debugging
- Debugging changes the job. Adding
count()orshow()runs extra jobs and can recompute everything upstream;cache()can hide non-determinism by fixing one result. Remove probes when done. - Collecting to look.
collect()on a large DataFrame moves the failure to the driver. Useshow,limitor write a sample. - Fixing the second symptom. Fetch failures and lost executors are usually consequences. Find the first executor that died and why.
- Raising memory for everything. More memory makes skew and bad records rarer, not gone. Look at task-level metrics first.
- Testing only on clean samples. Keep a library of real bad records and run every change against it.
What to do next
- Read the next failure in order: stage and task, attempts, host, then the last
Caused by. - Make sure executor logs are kept after the application ends, and practise finding one on your platform.
- Add
columnNameOfCorruptRecordhandling to readers of untrusted JSON and CSV. - Keep an audit helper that counts rows, NULL keys and duplicate keys, and run it at each step when a result looks wrong.
- Search your jobs for
first()aftergroupBy, unseededrand()and unorderedlimit. - Before moving to Spark 4.0, run your jobs with
spark.sql.ansi.enabledset to true to find the errors ANSI mode will raise. - For Python-heavy jobs, review overhead memory and Arrow batch sizes together; see PySpark performance.