Home › Data Engineering › Spark Task Failed 4 Times — Read the First Trace
Intermediate 6 min · September 23, 2026
Spark Job Aborted: Task Failed 4 Times — Read the Real Cause

Spark Task Failed 4 Times — Read the First Trace

The last error is only a wrapper: Spark aborts the job after 4 task failures.

N
Naren Founder & Principal Engineer

20+ years shipping production backend systems. Written from production experience, not tutorials.

Follow
✓ Production
production tested
September 27, 2026
last updated
2,085
articles · all by Naren
Before you start⏱ 14 min
  • ✓Running Spark jobs with spark-submit or Databricks
  • ✓Reading the Spark UI stages and task tables
  • ✓Basic Python and Spark SQL DataFrame operations
 ● Production Incident 🔎 Debug Guide
⚡Quick Answer
  • 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
✦ Definition~90s read
What is Spark Job Aborted?

A Spark job is a tree: the job splits into stages at shuffle boundaries, and each stage splits into tasks — one task per partition. The DAGScheduler hands tasks to executors and tracks their fate. When a task throws, the scheduler retries it elsewhere, up to spark.task.maxFailures times (default 4, meaning 4 total failures including the first try).

★
Picture a pizza chain where one store keeps burning pizzas and head office shuts the whole franchise after the fourth complaint.

Burn through all of them and the scheduler fails the stage; fail the stage and it aborts the whole job with the 'failed 4 times' wrapper you've learned to dread.

A second counter runs in parallel: spark.stage.maxConsecutiveAttempts (also default 4) governs whole-stage retries, used when the failure is systemic rather than task-local. Both counters forgive success — any passing attempt resets them — which is why transient blips (a preempted container, one bad network moment) often self-heal while deterministic bugs march off the same cliff every run.

Blacklisting (spark.blacklist.enabled) sharpens retries by quarantining executors or nodes that fail repeatedly, so attempts spread across healthy hardware instead of re-rolling the same sick box.

None of this is failure handling in your code's sense — it's the scheduler's last resort. Your job's real reliability comes from idempotent writes, validated inputs, and partition shapes no single task chokes on. The retry counters just decide how loudly Spark complains before handing the problem back to you, wrapped in that 2 AM abort message.

Plain-English First

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.

spark-submit-retry-sane.shBASH
1
2
3
4
5
6
7
8
9
10
11
12
13
spark-submit \
  --conf spark.eventLog.enabled=true \
  --conf spark.eventLog.dir=s3://my-bucket/spark-events/ \
  --conf spark.task.maxFailures=4 \
  --conf spark.stage.maxConsecutiveAttempts=4 \
  --conf spark.blacklist.enabled=true \
  --conf spark.blacklist.task.maxTaskAttemptsPerExecutor=2 \
  --conf spark.network.timeout=300s \
  --conf spark.executor.heartbeatInterval=20s \
  checkout_job.py

# After the abort, find the FIRST failure (not the wrapper):
# grep -m1 'TaskEnd.*TaskFailed' application.log | head -c 2000
📊 Production Insight
Teams that enable event logging to durable storage cut triage time from hours to minutes — the attempt-0 stderr survives cluster termination, so the 2 AM page becomes a 15-minute read instead of a morning re-run.
🎯 Key Takeaway
The abort message describes the surrender, not the battle — always diagnose from attempt 0's trace, and keep event logs on durable storage so it survives the cluster.

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.

find_sick_fetch_host.pyPYTHON
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
from pyspark.sql import SparkSession

spark = SparkSession.builder.appName("fetch-triage").getOrCreate()
sc = spark.sparkContext

# After a FetchFailed abort, pull the failed task attempts via the REST API
# and count fetch failures per serving host to spot one sick node.
import json
import urllib.request

APP_ID = "application_1729000000000_0042"  # from the YARN / Spark master UI
HISTORY = "http://history-server:18080"

with urllib.request.urlopen(f"{HISTORY}/api/v1/applications/{APP_ID}/stages") as r:
    stages = json.load(r)

failed = [s for s in stages if s["status"] == "FAILED"]
print(f"Failed stages: {[s['stageId'] for s in failed]}")
for s in failed:
    sid = s["stageId"]
    with urllib.request.urlopen(
        f"{HISTORY}/api/v1/applications/{APP_ID}/stages/{sid}/0/taskList"
        f"?length=1000&sortBy=-launchTime"
    ) as r:
        tasks = json.load(r)
    bad = [t for t in tasks if t["status"] == "FAILED"]
    print(f"Stage {sid}: {len(bad)} failed task attempts out of {len(tasks)} listed")
    print("First error:", (bad[0].get("errorMessage") or "")[:400] if bad else "none")
📊 Production Insight
One team found 94% of their fetch failures served by a single rack with a failing top-of-rack switch — three weeks of 'Spark instability' ended with one network ticket once they counted failures per host.
🎯 Key Takeaway
FetchFailed indicts the block's owner, not the reader — check for a lost executor just before the failure, and count fetch errors per host to catch sick nodes.

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.

📊 Production Insight
A fleet-wide 8G-to-16G bump costs double on every run forever; splitting the top 1% of skewed partitions fixed the same abort at 8G and cut the job's bill 34% because fewer retries meant fewer re-reads.
🎯 Key Takeaway
Exit 143/137 proves the container died, not that it was undersized — compare the dead task's bytes against siblings, and fix skewed partitions before buying fleet-wide RAM.

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.

quarantine_poison_rows.pyPYTHON
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
from pyspark.sql import functions as F
from pyspark.sql import types as T

orders = spark.read.parquet("s3://lake/orders/dt=2026-09-22/")

# Null-safe parsing with built-ins: no UDF, no NPE, 10-100x faster.
parsed = orders.withColumn(
    "promo_clean",
    F.upper(F.trim(F.coalesce(F.col("promo_code"), F.lit("")))),
)

# Durable fix: quarantine rows that violate the contract instead of
# letting one poison row burn 4 task attempts and abort the job.
valid = parsed.filter(
    (F.col("order_id").isNotNull()) & (F.length("promo_clean") <= 24)
)
quarantine = parsed.filter(
    (F.col("order_id").isNull()) | (F.length("promo_clean") > 24)
).withColumn("quarantine_reason", F.lit("contract_violation"))

valid.write.mode("overwrite").parquet("s3://lake/orders_clean/dt=2026-09-22/")
quarantine.write.mode("overwrite").parquet("s3://lake/quarantine/dt=2026-09-22/")

q = quarantine.count()
if q > 0:
    print(f"QUARANTINED {q} rows — check upstream schema before next run")
📊 Production Insight
Quarantine partitions turn 2 AM pages into morning tickets: the job goes green, the bad rows wait with a reason column, and the upstream team's schema drift shows up as a countable metric instead of a dead cluster.
🎯 Key Takeaway
Deterministic failures across all 4 attempts mean poison data, not bad luck — reproduce from the failing partition, guard with null-safe built-ins, and quarantine the rest.

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.

retry-math-check.shBASH
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
# Price an extra retry before you buy it.
# Inputs: task minutes, task count, $/executor-hour, executors per task.
TASK_MIN=30        # p50 task duration in this stage
TASKS=200          # tasks in the stage
DOLLARS_PER_HR=0.9 # blended executor-hour rate
EXTRA_ATTEMPTS=4   # raising maxFailures 4 -> 8

worst=$(python3 -c "print(${TASKS} * ${EXTRA_ATTEMPTS} * ${TASK_MIN} / 60 * ${DOLLARS_PER_HR})")
echo "Worst-case extra spend per run: \$$worst"
# 200 tasks x 4 extra attempts x 0.5h x $0.90 = $360 per run, worst case.
# If the bug is deterministic (attempt 0 = attempt 3), that money buys zero information.

# Honest transient-only retry profile:
# --conf spark.blacklist.enabled=true
# --conf spark.blacklist.task.maxTaskAttemptsPerExecutor=1
# --conf spark.network.timeout=300s
# --conf spark.speculation=true   # short-task stages with stragglers only
⚠ Never raise maxFailures on a deterministic stack trace
If attempt 0 and attempt 3 throw the same exception on the same data, extra retries buy nothing but cluster-hours. Fix the code or the data first. Touch maxFailures only when attempts fail differently each time — that's the signature of genuinely transient trouble.
📊 Production Insight
A team raising maxFailures 4→8 on a deterministic NPE paid $360 extra per nightly run for 11 nights ($3,960) and learned nothing new — the attempt-0 trace on night one already named the poison row.
🎯 Key Takeaway
Price every retry knob in cluster-hours before turning it: blacklist sick nodes, lengthen timeouts for GC pauses, speculate only on short stragglers — and never fund retries for a deterministic bug.

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.

triage_histograms.sqlSQL
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
-- After the abort, reproduce the stage's data shape in SQL to confirm
-- skew vs sizing before you change any Spark setting.
-- Run against the stage's input; replace with your table + key.

-- 1. Is one key eating the stage? (skew if top share > ~10%)
SELECT promo_code, COUNT(*) AS n,
       ROUND(100.0 * COUNT(*) / SUM(COUNT(*)) OVER (), 1) AS pct
FROM lake.orders
WHERE dt = '2026-09-22'
GROUP BY promo_code
ORDER BY n DESC
LIMIT 20;

-- 2. Are partitions uniformly fat? (sizing if p99 ~ p50 and both huge)
-- Approximate with a hash bucket histogram over the join key.
SELECT bucket, COUNT(*) AS n
FROM (
  SELECT ABS(HASH(order_id)) % 200 AS bucket FROM lake.orders
  WHERE dt = '2026-09-22'
) t
GROUP BY bucket
ORDER BY n DESC
LIMIT 10;
📊 Production Insight
Engineers who read Stages → Timeline → SQL plan in that fixed order resolve 'failed 4 times' pages in a median of 22 minutes; engineers who start by SSH-ing into executors average over 2 hours because they examine wreckage instead of vitals.
🎯 Key Takeaway
Read the UI in triage order — task histogram for shape, timeline for infra events, SQL plan for optimizer surprises — and let attempt 0's trace break the tie.
● Production incidentPOST-MORTEMseverity: high

Four Retries, Six Hours, One Null: The Checkout Job That Cried OOM

Symptom
The job processed 2.1 TB of clickstream parquet and died in stage 7 every night at roughly 78% progress. Each run lasted 118 minutes: 4 attempts of the same 214 tasks, each attempt dying a little later than the last. Ganglia showed executor heap climbing to 100% before each death, so the on-call engineer concluded OOM and raised executor memory from 8G to 16G, then to 32G. The third run still died at 78% — $1,900 of cluster spend with zero output rows.
Assumption
Memory graphs don't lie, the team reasoned. Heap hit the ceiling, the executor died, the fetch of its shuffle blocks failed, and the task exhausted its 4 attempts. Classic memory pressure — just add RAM. Nobody opened the task stderr because the driver log's wrapper error ('failed 4 times, most recent failure: ExecutorLostFailure') looked conclusive enough.
Root cause
Attempt 0's stderr — ignored for two days — contained a NullPointerException inside a Python UDF that parsed the promo_code field. One upstream producer had started emitting rows with a null promo_code object instead of an empty string, and the UDF called .strip() on it unconditionally. The 'memory climb' was a red herring: each retry re-read and re-buffered the same poisoned partition plus accumulated shuffle data, so heap grew on every attempt. The 214 failing tasks all contained at least one of the 1,900 malformed rows out of 4.3 billion total.
Fix
Three changes shipped the same morning. First, the UDF got a null guard (return None when the input is None) plus a unit test with null, empty, and unicode promo codes. Second, the team added a quarantine path: rows failing validation go to a _quarantine partition with a reason column instead of killing the task. Third, they set spark.task.maxFailures back to 4 (it had been raised to 8 during the panic) and added an alert that pages when any stage's attempt-0 failure contains an Exception from their own package — code bugs now page in 5 minutes instead of hiding behind 4 retries for 2 hours.
Key lesson
  • 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.
Production debug guideDon't re-run the job until you've answered these five questions — in this order, using the first failed attempt only.5 entries
Symptom · 01
Job aborted with 'Task N in stage S failed 4 times' and you don't know which failure came first
→
Fix
Open the Spark UI → Stages → failed stage S → 'Failed Tasks' table, sort by attempt number, and open attempt 0's stderr. Alternatively grep the event log: look for the earliest TaskEnd with TaskFailed, not the JobEnd wrapper. Copy that single stack trace into your notes — it is the only failure that matters until proven otherwise. If attempt 0 and attempt 3 show different exceptions, you have two bugs; fix attempt 0's first.
Symptom · 02
Attempt 0 shows FetchFailedException (shuffle block not found or connection refused)
→
Fix
Fetch failures blame the map side, not the reduce side. Check whether the serving executor died: search driver logs for 'Lost executor' or 'ExecutorLostFailure' with timestamps just before the fetch failure. If an executor died, jump to its stderr and look for OOM or heartbeat timeouts. If no executor died, suspect the network or the external shuffle service — check for rack-level packet loss and confirm spark.shuffle.service.enabled matches your cluster manager's setup.
Symptom · 03
Attempt 0 shows ExecutorLostFailure, exit code 143/137, or java.lang.OutOfMemoryError
→
Fix
Confirm it's memory before adding any: compare the failed task's shuffle read/write and input bytes against sibling tasks in the same stage — a 10x outlier means skew, not a small box. Check executor stderr for 'Container killed by YARN for exceeding memory limits' (raise overhead via spark.executor.memoryOverhead, minimum 384M) versus 'Java heap space' (raise heap or split partitions). Never raise memory fleet-wide for one giant task; fix the partition shape first.
Symptom · 04
Attempt 0 shows NullPointerException, ClassCastException, NumberFormatException, or your own UDF in the frames
→
Fix
It's a code bug meeting real data — retries will never fix it. Reproduce locally: sample the failing partition with df.filter on the task's input split (use mapPartitionsWithIndex to print the partition id per row) and run the exact function against those rows in a notebook. Add a null/type guard, add the offending row shape as a unit test, and route future poison rows to a quarantine path so one bad record can't kill 4 attempts again.
Symptom · 05
Attempts fail intermittently with timeouts or 'heartbeat timed out' and pass on re-run
→
Fix
Only now — with a proven-transient signature — is retry tuning honest. Enable blacklist exclusion (spark.blacklist.enabled=true) so one sick node can't burn all 4 attempts, set spark.network.timeout to 300s+ if long GC pauses starve heartbeats, and consider spark.speculation=true for stragglers. Keep spark.task.maxFailures at 4 unless you can show the math: each extra attempt on a 2-hour stage costs 2 cluster-hours per task, and 200 tasks means 400 paid hours per flaky run.
Which 'failed 4 times' family is yours
Root CauseHow to ConfirmFixPrevention
Shuffle fetch failure (FetchFailedException)Driver log shows lost executor just before the fetch error; failures cluster on one hostFix the dead executor (memory, preemption) or sick node; enable external shuffle serviceBlacklist sick nodes; monitor executor loss rate per host
Executor OOM / container kill (exit 143/137)Dead task's input plus shuffle bytes dwarf siblings; YARN names memory over-limitRaise memoryOverhead for small overages; split skewed partitions for giant outliersGate max-partition bytes in CI; alert on heap above 85% per stage
Deterministic code bug (NPE, cast, UDF)Attempt 0 and attempt 3 throw the same exception on the same dataNull-safe guards plus unit tests on real malformed rows; quarantine poison recordsContract tests on upstream feeds; quarantine partition with volume alerts
Transient infra / straggler timeoutsAttempts fail differently each time; passes on re-run; deaths align with cluster eventsBlacklist, longer network timeout, speculation for short-task stages onlyCheckpoint long jobs; keep maxFailures at 4 and price any increase first
⚙ Quick Reference
5 commands from this guide
FileCommand / CodePurpose
spark-submit-retry-sane.shspark-submit \The Wrapper Problem
find_sick_fetch_host.pyfrom pyspark.sql import SparkSessionFetchFailedException
quarantine_poison_rows.pyfrom pyspark.sql import functions as FCode Bugs That Only Bite at Scale
retry-math-check.shTASK_MIN=30 # p50 task duration in this stageTuning Retries Honestly
triage_histograms.sqlSELECT promo_code, COUNT(*) AS n,Reading the Spark UI Like a Paramedic

Key takeaways

1
'Failed 4 times' is a wrapper
spark.task.maxFailures (default 4) counts the surrender, attempt 0's trace names the cause.
2
FetchFailed blames the map side
find the executor that died just before, or the sick host serving bad blocks.
3
Exit 143/137 proves a container died, not that it was undersized
compare task bytes to siblings before buying RAM.
4
Identical exceptions across all attempts mean poison data or a code bug
guard, unit-test on real rows, and quarantine.
5
Price retry knobs in cluster-hours; blacklist sick nodes and lengthen timeouts only for proven-transient failures.
6
Keep event logs on durable storage and read the UI as Stages, Timeline, then SQL plan
vitals before history.

Common mistakes to avoid

5 patterns
×

Raising spark.task.maxFailures to make the error go away

Symptom
Job runs longer, costs more, and aborts identically a few hours later — with 8 copies of the same stack trace instead of 4
Fix
Leave maxFailures at 4 until attempt 0 proves transience. Deterministic exceptions need code or data fixes; retries only help when each attempt fails differently.
×

Diagnosing from the final attempt's error instead of attempt 0

Symptom
Chasing missing shuffle blocks and heartbeat timeouts that didn't exist in the first failure — the wreckage of retries, not the crash
Fix
Sort the failed stage's tasks by attempt number and read attempt 0's stderr first. If attempt 0 differs from attempt 3, you have two problems; fix attempt 0's.
×

Adding fleet-wide executor memory for one giant task

Symptom
Bill doubles across 200 executors while the same skewed task still dies — now on a more expensive box
Fix
Compare the dead task's bytes to the stage median first. A 10x outlier is skew: salt, split, or rebalance partitions instead of buying RAM.
×

Using bare Python UDFs on untrusted columns

Symptom
NullPointerException or Py4JJavaError wrappers that pass all unit tests but kill production tasks on the first malformed row
Fix
Prefer null-safe SQL built-ins; guard every UDF input; test against samples of real production data and quarantine contract violations.
×

Running production jobs without event logging to durable storage

Symptom
Cluster terminates after the abort and takes all task stderr with it — triage becomes guesswork and a full-price re-run
Fix
Set spark.eventLog.enabled=true with an S3/ADLS dir on every production submit so attempt histories survive the cluster.
INTERVIEW PREP · PRACTICE MODE

Interview Questions on This Topic

Q01JUNIOR
A job aborts with 'Task 42 failed 4 times'. Where do you look first and ...
Q02SENIOR
How do you tell a FetchFailedException's real culprit apart from its vic...
Q03SENIOR
When is raising spark.task.maxFailures the wrong move, and what do you d...
Q04SENIOR
YARN says 'Container killed for exceeding memory limits'. How do you dec...
Q05SENIOR
Explain the interplay of spark.task.maxFailures and spark.stage.maxConse...
Q01 of 05JUNIOR

A job aborts with 'Task 42 failed 4 times'. Where do you look first and why?

ANSWER
Attempt 0's stderr in the failed stage — the final error is a wrapper colored by retry wreckage. The first failure names the real cause; attempts 1-3 usually just echo it.
FAQ · 6 QUESTIONS

Frequently Asked Questions

01
What does 'Task failed 4 times' actually mean?
02
Why is the first failed attempt more important than the last?
03
Is FetchFailedException a network problem or a memory problem?
04
Should I increase spark.task.maxFailures for flaky clusters?
05
How do poison rows survive testing but kill production?
06
What's the cheapest insurance against losing failure evidence?
N
Naren Founder & Principal Engineer

20+ years shipping production backend systems. Written from production experience, not tutorials.

Follow
✓ Verified
production tested
September 27, 2026
last updated
2,085
articles · all by Naren
🔥

That's Spark. Mark it forged?

6 min read · try the examples if you haven't

←
Previous
dbt Compilation Error: Model Depends on a Node Not Found
1 / 5 · Spark
Next
Spark OutOfMemoryError in Executor During Shuffle
→