Spark Task Failed 4 Times — Read the First Trace
The last error is only a wrapper: Spark aborts the job after 4 task failures.
20+ years shipping production backend systems. Written from production experience, not tutorials.
- ✓Running Spark jobs with spark-submit or Databricks
- ✓Reading the Spark UI stages and task tables
- ✓Basic Python and Spark SQL DataFrame operations
- Spark retries each failed task until spark.task.maxFailures (default 4) trips, then aborts the whole job — the 'failed 4 times' message is a wrapper, not the cause
- Open the failed stage in the Spark UI, click the first failed attempt (attempt 0), and read its stderr stack trace — attempts 1-3 almost always repeat the same root failure
- FetchFailedException points at shuffle trouble (a lost executor or bad network link), ExecutorLostFailure with exit 143 points at memory, and NullPointerException or cast errors point at your code meeting real data
- Don't raise maxFailures to hide the pain: fix the root cause first, then tune retries and blacklisting only for genuinely transient faults
Picture a pizza chain where one store keeps burning pizzas and head office shuts the whole franchise after the fourth complaint. The shutdown notice doesn't say why pizzas burned — just 'four failures, we're done.' That's Spark: after 4 task failures it aborts the job with a wrapper. The real story — a broken oven (dead executor), a lost delivery route (fetch failure), or a bad recipe (your code bug) — hides in the first log. Read the notice alone and you'll fix the wrong thing every time.
It's 2 AM, your Spark job has been running for three hours, and the driver log ends with the dreaded line: 'Job aborted due to stage failure: Task 42 in stage 7 failed 4 times.' Your first instinct is to restart the job and hope. Your second instinct is to bump spark.task.maxFailures and hope harder. Both are wrong, and both cost another three hours.
That final error message is a wrapper. It tells you Spark gave up, not why it gave up. The why lives somewhere else entirely — in the stderr of the first failed attempt, usually 500 lines above where you're looking. Engineers who don't know this burn entire nights re-running jobs that fail identically, because the root cause never changed between attempt 1 and attempt 4.
Three failure families cause nearly every 'failed 4 times' abort. Shuffle fetch failures mean a task couldn't read another executor's output — typically because that executor died or the network dropped. Executor losses with OOM signatures mean the box was too small for the partition it chewed on. And code bugs — nulls, bad casts, exploding UDFs — mean your logic met production data for the first time and lost.
This article teaches you to triage in minutes. You'll learn where Spark hides the first attempt's trace, how to read the signature of each failure family, and how to tune retries honestly instead of using them as a rug to sweep bugs under.
The Wrapper Problem: Why the Last Error Lies
When Spark aborts a job, the DAGScheduler throws whatever the fourth attempt reported — and by then the cluster is usually in a degraded state that colors the evidence. The first attempt died on your real bug; the fourth attempt often dies on the wreckage (missing shuffle blocks from executors killed during attempts 1-3, timed-out heartbeats on an overloaded driver). Reading only the final error is like diagnosing a car crash from the tow-truck receipt.
The mechanics are simple. spark.task.maxFailures (default 4) counts consecutive task failures; each failure is retried on (usually) a different executor. Separately, a stage is retried up to spark.stage.maxConsecutiveAttempts (default 4) when it fails wholesale. So 'failed 4 times' describes one unlucky task, while the job-level abort is the scheduler deciding the stage can't proceed. Both counters reset the moment anything succeeds, which is why flaky infrastructure sometimes squeaks through and real bugs never do.
Your first move is always the same: isolate attempt 0. In the Spark UI, the failed stage lists every task attempt with its own stderr link. In logs, search backward from the abort for the first TaskEnd/TaskFailed event. Enable event logging (spark.eventLog.enabled=true with a durable dir) on every production job so this archaeology is possible after the cluster is gone — without it, you're debugging from memory the next morning.
FetchFailedException: Blame the Map Side, Not the Task
FetchFailedException means a reduce task asked for shuffle blocks and couldn't get them. Beginners blame the reduce task and its memory; veterans check the executor that was supposed to serve those blocks. Nine times out of ten, that executor is already dead — killed by OOM, preempted by the cluster manager, or partitioned off the network — and the fetch failure is just the obituary.
The tell is timing. Open the driver log and line up timestamps: a 'Lost executor 7' message seconds before your FetchFailed means the map side died first. The shuffle data it held died with it (unless you run the external shuffle service, which keeps blocks alive across executor loss). Spark then resubmits the map tasks and retries your reduce — which is exactly why one dead executor can burn through all 4 attempts of hundreds of downstream tasks.
If no executor died, you're looking at genuine shuffle-service or network trouble: an overloaded external shuffle service, a bad NIC dropping packets on one rack, or shuffle files cleaned up too aggressively. Check spark.shuffle.service.enabled against your deploy mode (YARN needs it for decommissioning safety), and look for fetch failures clustering on one host — a single sick node producing most fetch errors is a hardware ticket, not a Spark ticket.
ExecutorLostFailure and OOM: Prove It's Memory First
ExecutorLostFailure with exit code 143 (SIGTERM, usually YARN killing an over-limit container) or 137 (SIGKILL, OOM-killer) looks like a memory verdict, but it only proves the container died — not that your sizing was wrong. A skewed partition 50x the median will kill any reasonably sized executor, and no amount of fleet-wide memory fixes a single giant task. The proof step takes two minutes: compare the dead task's input + shuffle-read bytes against its siblings.
YARN's 'exceeding memory limits' message deserves a close read because it names two numbers: physical memory used versus requested. If the overage is small (a few hundred MB over), bump spark.executor.memoryOverhead — the default max(384M, 10% of heap) is too thin for shuffle-heavy or native-code workloads. If the message says 'Java heap space' instead, the JVM heap itself filled: your partitions are too fat, your cache is retained too long, or a UDF is materializing whole partitions in memory.
Off-heap and GC logs turn guesses into measurements. Enable spark.memory.offHeap.enabled with spark.memory.offHeap.size for shuffle-heavy stages, and read the executor stderr for 'GC overhead limit exceeded' (heap churn from object-heavy code — switch to SQL functions or Kryo) versus a flat climb to the ceiling (genuinely too much data per task — split partitions). Fix the shape before the size: one repartition that halves max-partition bytes beats doubling memory across 200 executors on both cost and reliability.
Code Bugs That Only Bite at Scale: Nulls, Casts, UDFs
The cruelest 'failed 4 times' aborts come from code that passed every test. Tests run on clean samples; production runs on 4.3 billion rows including the one malformed record an upstream producer started emitting Tuesday. Nulls where your UDF assumes strings, 'N/A' where you cast to int, timestamps in a format nobody documented — each kills the task deterministically on every attempt, and retries just re-read the same poison row four times.
Python UDFs are the classic carrier because they run outside the Catalyst optimizer's safety nets: no null-safe codegen, whole-row serialization, and exceptions that surface as inscrutable Py4JJavaError wrappers. A scalar UDF calling .strip() on None fails exactly the way our incident did. Prefer built-in SQL functions (they're null-safe by design and run 10-100x faster), and when you must use a UDF, guard every input and PNL-test it against a sample of real production data, not hand-written fixtures.
The durable fix is a quarantine path, not just a guard. Route rows that fail validation into a separate output with a reason column, and alert on quarantine volume. Guards stop the crash; quarantine tells you the upstream schema drifted before the next field breaks. Pair it with mapPartitionsWithIndex during triage to identify exactly which partition — and therefore which input files — harbor the poison rows.
Tuning Retries Honestly: Blacklists, Timeouts, Speculation
Retry knobs exist for genuinely transient failures — a preempted container, a blipped rack, a straggler disk — and each knob has a price tag. Raising spark.task.maxFailures from 4 to 8 on a stage whose tasks run 30 minutes each means one sick task can burn 4 extra cluster-hours before aborting; across 200 tasks that's a potential 800 paid hours per run. Do that math in the incident channel before touching the knob, because 'just bump retries' is how a 2-hour outage becomes a $4,000 invoice.
Blacklisting is the cheapest honest retry: with spark.blacklist.enabled=true, executors (or whole nodes) that fail repeatedly are excluded from future attempts, so your 4 tries spread across healthy hardware instead of hammering the same sick box. Pair it with spark.blacklist.task.maxTaskAttemptsPerExecutor=1 for fetch-failure storms — one strike per executor forces the scheduler to move on immediately.
Timeouts and speculation handle the slow rather than the dead. If long GC pauses starve heartbeats, raise spark.network.timeout (which also lifts heartbeat and fetch timeouts) to 300s or more rather than watching healthy executors get declared dead. Enable spark.speculation=true only for short-task stages with stragglers — it launches duplicate copies of slow tasks and keeps whichever finishes first, which wastes work on long tasks but saves whole jobs when one disk goes sluggish.
Reading the Spark UI Like a Paramedic
The Spark UI rewards a fixed reading order the way an ER rewards triage order: vitals first, history second, labs last. Start at Stages — the failed stage is red, and its description names the operation (often a wholeStage codegen range like 'parquet at X'). Click through to the task table and sort by duration or shuffle-read: one task 20x slower or fatter than the median is skew; all tasks uniformly slow is sizing; a handful of instant failures is poison data or a missing file.
Next, open the Timeline view. Tasks that die in waves right after a cluster event (autoscaling, decommissioning, a deploy) point at infrastructure, not code. The Event Timeline also shows when executors were added and removed — if your failures start exactly when spot nodes were reclaimed, you've found the killer without reading a single stack trace.
Finally, read the SQL tab for the failed stage's query plan. It shows the exact join order, broadcast decisions, and partition counts the optimizer chose. A SortMergeJoin on two terabyte tables where you expected a broadcast tells you autoBroadcastJoinThreshold was missed (check for type mismatches killing the broadcast). The plan plus the task histogram together answer 80% of 'failed 4 times' pages before you ever SSH anywhere.
Four Retries, Six Hours, One Null: The Checkout Job That Cried OOM
- Attempt 0 is the crime scene; attempts 1-3 are usually echoes. Build the habit (and the alerting) that reads the first failure's stack before touching any memory setting — two days of 32G executors cost $1,900 and fixed nothing.
- Memory graphs can mislead when retries accumulate state. Heap climbing across attempts doesn't prove the first failure was memory-related; always correlate the heap curve with the exception type in the earliest trace.
- Poison rows need a quarantine path, not just a null guard. The 1,900 malformed rows out of 4.3 billion would have been invisible forever without a _quarantine partition — the next malformed field would have caused the same 2 AM page.
| File | Command / Code | Purpose |
|---|---|---|
| spark-submit-retry-sane.sh | spark-submit \ | The Wrapper Problem |
| find_sick_fetch_host.py | from pyspark.sql import SparkSession | FetchFailedException |
| quarantine_poison_rows.py | from pyspark.sql import functions as F | Code Bugs That Only Bite at Scale |
| retry-math-check.sh | TASK_MIN=30 # p50 task duration in this stage | Tuning Retries Honestly |
| triage_histograms.sql | SELECT promo_code, COUNT(*) AS n, | Reading the Spark UI Like a Paramedic |
Key takeaways
Common mistakes to avoid
5 patternsRaising spark.task.maxFailures to make the error go away
Diagnosing from the final attempt's error instead of attempt 0
Adding fleet-wide executor memory for one giant task
Using bare Python UDFs on untrusted columns
Running production jobs without event logging to durable storage
Interview Questions on This Topic
A job aborts with 'Task 42 failed 4 times'. Where do you look first and why?
Frequently Asked Questions
20+ years shipping production backend systems. Written from production experience, not tutorials.
That's Spark. Mark it forged?
6 min read · try the examples if you haven't