In YARN, what happens when an ApplicationMaster container fails, and what caps the retries?
answer
- a new attempt, not a resumed process
- the default budget is very small
- a lifetime budget hurts long-running jobs
- old containers die unless asked not to
- losing a whole node is detected slowly
basics
~20 sThe ResourceManager notices the ApplicationMaster is gone and starts a fresh attempt in a new container elsewhere, up to yarn.resourcemanager.am-max-attempts (default 2). When attempts run out the whole application is marked FAILED, and by default its earlier containers are killed too.
solid answer
~40 sThe ResourceManager detects the loss — either from the NodeManager reporting the container's exit or from the ApplicationMaster missing its liveness heartbeats — and launches a **new attempt** with a new attempt id in a new container, usually on a different node. The cluster-wide ceiling is `yarn.resourcemanager.am-max-attempts`, **2** by default; an application may request a lower `maxAppAttempts` but not a higher one. When the ceiling is reached the application ends as FAILED. By default the previous attempt's containers are killed and the new attempt starts clean, unless the application sets `keepContainersAcrossApplicationAttempts`. Because attempts accumulate over a long-running job's lifetime, framework configs let you count failures over a sliding window instead — Spark on YARN exposes this as `spark.yarn.maxAppAttempts` and `spark.yarn.am.attemptFailuresValidityInterval`.
code
properties · 6 lines# cluster-wide ceiling on ApplicationMaster attempts
yarn.resourcemanager.am-max-attempts=2
# Spark on YARN: per-application budget over a sliding window
spark.yarn.maxAppAttempts=4
spark.yarn.am.attemptFailuresValidityInterval=1hgo deeper
Know that YARN retries an application's coordinator a small number of times and then gives up, marking the application FAILED.
Explain the attempt model: a new attempt id, a fresh container, the cluster-wide ceiling that clamps whatever the application asked for, and the fact that the old containers are discarded by default.
Show you read attempt diagnostics to separate deterministic code failures from infrastructure loss, and know why long-running jobs need failures counted over a window rather than a lifetime.
Own the policy: what attempt ceiling the platform sets for batch versus long-running workloads, and how recovery expectations differ between a framework that can resume completed work and one that restarts from zero.
## Attempts, not restarts A YARN application has an id like `application_1699034221_0007`; each try at running its coordinator is an **application attempt**, `appattempt_1699034221_0007_000001`, and the ApplicationMaster's own container is that attempt's container `_01_000001`. When the coordinator dies, YARN does not resume it — it starts attempt `000002` from scratch. Understanding that the unit of recovery is a whole attempt, not a checkpointed process, explains most of the behaviour below. ## How the failure is detected Two paths. If the process exits, the NodeManager hosting it reports the container's completion and exit status to the ResourceManager on its next heartbeat. If the whole machine or the network is lost, no exit status ever arrives; instead the ResourceManager's AMLivelinessMonitor notices the ApplicationMaster has stopped sending `allocate()` heartbeats within the expiry interval and declares the attempt failed. That second path is slower — the expiry interval is on the order of ten minutes by default — which is why a lost node can leave an application apparently frozen before recovery begins. ## The retry ceiling `yarn.resourcemanager.am-max-attempts` is a **cluster-wide maximum**, default **2**. An application may specify its own `maxAppAttempts` at submission, but the ResourceManager clamps it to the cluster value — an application cannot grant itself more retries than the administrator allows. When the last attempt fails, the application's final state is FAILED and its diagnostics record the last attempt's exit information. There is an important subtlety for long-running applications. A streaming job that runs for months will eventually lose an ApplicationMaster to an unrelated node failure; with a lifetime budget of two attempts it dies on the second such event, possibly weeks apart. The fix is the **attempt failures validity interval**: failures older than the interval are forgotten, so the budget becomes "two failures within an hour" rather than "two failures ever". Spark exposes this as `spark.yarn.am.attemptFailuresValidityInterval` alongside `spark.yarn.maxAppAttempts`, and has the analogous pair `spark.yarn.max.executor.failures` and `spark.yarn.executor.failuresValidityInterval` for worker containers. ## What happens to the running containers By default, when an attempt fails the ResourceManager kills the containers that attempt had been granted, and the new ApplicationMaster starts with nothing. That is the right default for batch: the new coordinator has no memory of what the old one dispatched. An application can instead set `keepContainersAcrossApplicationAttempts` at submission, which preserves the allocated containers so a new attempt can reconnect to them — the mechanism long-running services use to survive a coordinator restart without losing their workers. What the new attempt *recovers* is framework-specific and is the part that surprises people. A MapReduce ApplicationMaster can recover completed tasks from its job history so a retry does not redo finished map output. A Spark application in cluster deploy mode has its **driver** inside the ApplicationMaster, so a new attempt means a new driver: all in-memory state and cached data are gone and the job restarts from the beginning — except Structured Streaming, which resumes from its own checkpoint directory in HDFS or object storage. That is a checkpoint in Spark's sense, unrelated to anything YARN does. ## ResourceManager failure is a different thing Do not confuse an ApplicationMaster failure with a ResourceManager failure. With `yarn.resourcemanager.ha.enabled` and a ZooKeeper-backed state store, a standby ResourceManager takes over, reloads application state, and NodeManagers and ApplicationMasters re-sync with it. This is **work-preserving** restart: running containers keep running across the failover and applications are not restarted. An ApplicationMaster failure, by contrast, is *not* work-preserving by default. ## Operating implications A repeatedly failing ApplicationMaster is usually one of three things: the coordinator ran out of memory (a Spark driver collecting too much to the driver, or a MapReduce ApplicationMaster tracking a huge number of tasks) and was killed by the NodeManager for exceeding its container memory; the application code threw during startup, in which case both attempts fail identically and quickly; or the node hosting it was lost, in which case the attempts fail on different nodes with different diagnostics. Read them with `yarn logs -applicationId <id>`, which after log aggregation contains every attempt's container logs, and compare attempt 1 with attempt 2 — identical failures point at code, different ones at infrastructure. A final trap: burning attempts costs real time. Each retry re-queues for capacity, re-localises jars, and restarts the work. Raising the attempt ceiling to mask a deterministic bug simply multiplies the time to failure.
- Why do long-running streaming applications need an attempt failures validity interval?Because the default budget is counted over the application's whole lifetime. A job running for months will lose its ApplicationMaster to unrelated node failures weeks apart and exhaust two attempts without ever having a real bug. The validity interval makes failures expire, so the policy becomes 'N failures within this window', which is what you actually want for a service-shaped workload.
- If an application's two attempts fail with identical stack traces seconds apart, what does that tell you?That it is deterministic application-side failure, not infrastructure. A lost node or an out-of-memory kill produces different diagnostics, different hosts and usually different timings between attempts. Identical, immediate failures mean the coordinator is throwing during startup — bad configuration, a missing dependency, or an unreadable input path — and raising the attempt limit only wastes more time.
saying these in an interview costs you the question
- Thinks a new attempt resumes the old ApplicationMaster's state
- Assumes an application can raise its own attempt limit above the cluster ceiling
- Confuses ApplicationMaster failure with ResourceManager failover
- Says a Spark cluster-mode driver survives an ApplicationMaster retry
- Raises the attempt limit to mask a deterministic startup failure