skip to content

In cluster mode on YARN, where do you find the Spark driver's stack trace after a failed job?

level: seniorimportance: should knowfreq 48%

answer

  1. the driver is not where you are standing
  2. find out which process owns the exception
  3. one container holds both the AM and the driver
  4. an application ID unlocks everything
  5. aggregation must have been switched on

basics

~20 s

In YARN cluster mode the driver runs inside the ApplicationMaster container, so its stdout and stderr are container logs. Fetch them with yarn logs -applicationId <appId>, or open the ApplicationMaster log link in the ResourceManager UI.

solid answer

~40 s

In cluster mode the driver is the ApplicationMaster container, so `spark-submit` only ever prints the final status and YARN's short diagnostics string — the exception itself is in that container's log. While the application is running, follow the ResourceManager's tracking URL, which proxies the live Spark UI on the driver, and read the AM's stdout/stderr from the node. Once it has finished, `yarn logs -applicationId application_1724_0042` returns the aggregated logs for every container; `-am 1` narrows it to the ApplicationMaster (the driver) and `-containerId` to one executor. Aggregation only exists if `yarn.log-aggregation-enable` is on; otherwise the files sit on each NodeManager under `yarn.nodemanager.log-dirs` and expire. For the Spark-side view after the fact, the History Server serves the event log written when `spark.eventLog.enabled=true`. On Kubernetes the equivalent is `kubectl logs <driver-pod>`.

code

bash · 9 lines
bash
# -- driver (ApplicationMaster) log only
yarn logs -applicationId application_1724_0042 -am 1

# -- one executor container named in the stage failure
yarn logs -applicationId application_1724_0042 \
  -containerId container_e12_1724_0042_01_000007

# -- kubernetes equivalent
kubectl logs daily-rollup-1724-driver

go deeper

for a junior

Know that in cluster mode the driver's output is not on your screen, and that yarn logs -applicationId is how you retrieve it. Keep the application ID that spark-submit prints.

for a middle

Explain why the log is in the ApplicationMaster container, the role of log aggregation, and the difference between the live Spark UI reached through the tracking URL and the History Server's replay of the event log.

for a senior

Demonstrate a post-mortem method: driver log first, then the executor named in the stage failure, then the NodeManager log for container kills. Know what disappears when aggregation is off or retention has passed.

for a principal

Own observability as platform policy: aggregation and event logging on by default, retention long enough for a Monday post-mortem, application IDs captured by the orchestrator, and log volume budgeted so a chatty driver cannot fill the cluster.

## Why the terminal is empty The usual complaint — "the job failed silently" — is a direct consequence of deploy mode. In `--deploy-mode cluster` the driver JVM runs inside the YARN ApplicationMaster container, so its stdout, stderr and log4j output are that container's log files on whichever node YARN chose. The `spark-submit` process is only a launcher; when the application ends it prints the tracking URL, the final status (`SUCCEEDED`/`FAILED`/`KILLED`) and YARN's `Diagnostics:` string, which is usually a truncated one-liner such as `User class threw exception` with no stack trace. Nothing is broken; you are simply reading the wrong process's output. ## While the application is running - **ResourceManager UI → Application → Tracking URL.** For a live application the tracking URL proxies the driver's Spark UI (port 4040 on the driver, reached through the YARN proxy), giving you stages, tasks, SQL plans and the executors tab. - **The AM container log.** The same application page links to "logs" for each attempt's ApplicationMaster container — that is your driver's live stdout/stderr. - **Executors tab.** Each executor row links to its container's stdout/stderr on the owning NodeManager, which is how you reach task-side exceptions without knowing node names. ## After it has finished The Spark UI dies with the driver, so two different systems take over: 1. **YARN log aggregation.** If `yarn.log-aggregation-enable` is `true`, when containers exit the NodeManagers upload their logs into HDFS and you retrieve them with: ```bash yarn logs -applicationId application_1724_0042 -am 1 yarn logs -applicationId application_1724_0042 -containerId container_..._000007 ``` Without aggregation, logs remain in the local `yarn.nodemanager.log-dirs` on each node and are deleted after `yarn.nodemanager.log.retain-seconds`, so you must reach the node before they vanish. This is the single most common reason a post-mortem fails: aggregation was off, or the retention window had passed. 2. **The Spark History Server.** Set `spark.eventLog.enabled=true` and `spark.eventLog.dir` and the driver writes a JSON event log; the History Server (default port 18080) replays it into the familiar Spark UI long after the application is gone. Note what it does *not* contain: it holds Spark's structured events — jobs, stages, tasks, SQL plans, executor add/remove — not your `println`s or a thrown exception's stack trace. For that you still need container logs. ## Reading the failure correctly Work outward from the driver log. A stage that dies from an executor failure will show a driver-side `SparkException: Job aborted due to stage failure` whose message names the *first* failing task and its executor, plus a truncated cause. The real reason often lives in the executor's container log: a `FetchFailedException` on the driver side may be an out-of-memory kill on the executor side, and YARN's own `Container killed by YARN for exceeding memory limits` message appears in the NodeManager log for that container rather than in the Spark log. Get the application ID, pull the AM log first, then pull the named executor's container. ## Practical hygiene - **Print the application ID at submit time and keep it.** Orchestrators should capture it into the task's own log; every later lookup is keyed on it. - **Configure driver logging deliberately.** `spark.driver.extraJavaOptions` can point log4j at a config that keeps a readable level; a driver spamming INFO makes an aggregated log unreadable at 400 MB. - **Turn event logging on everywhere.** It is cheap, and it is the only way to answer "why was this slow" once the job is over. - **On Kubernetes**, the equivalents are `kubectl logs <driver-pod>` (and `kubectl describe pod` for scheduling/pull failures). Note that executor pods are deleted when they terminate unless `spark.kubernetes.executor.deleteOnTermination` is set to `false`, so an executor-side stack trace can disappear within seconds — ship logs off the node if you care about them. - **In client mode none of this applies**: the driver log is your terminal, which is exactly why interactive debugging is easier there and why teams sometimes reproduce a cluster-mode failure in client mode on a smaller input. ## Streaming and long-running jobs For a job that never ends, the History Server shows nothing until it stops. Rely on the live tracking URL, and if the job runs for weeks, consider the incremental event-log settings and a log shipper rather than trying to `yarn logs` a container that is still writing.

  • Why does the Spark History Server not show your application's stack trace?
    The History Server replays the event log, which records Spark's structured events — jobs, stages, tasks, SQL plans, executor lifecycle — not the driver's console output. An exception thrown by your code appears in the ApplicationMaster's container log instead. Use the History Server for why a job was slow, and container logs for why it threw.
  • An executor died and the driver only reports a stage failure. How do you get the executor's own log?
    Take the executor id or container id named in the driver's stage-failure message, then fetch that container: `yarn logs -applicationId <id> -containerId <containerId>`, or click through the Spark UI executors tab while the app is alive. Container kills for exceeding memory are logged by the NodeManager, so check that log too.
  • What changes about log access when the same job runs on Kubernetes instead of YARN?
    The driver is a pod, so `kubectl logs <driver-pod>` replaces `yarn logs -am`, and `kubectl describe pod` explains image-pull or scheduling failures that never produce Spark output. Executor pods are removed on termination by default, so their logs vanish quickly unless you disable that or ship logs to a central store.

saying these in an interview costs you the question

  • Expects driver output in the terminal in cluster mode
  • Says the History Server stores application stack traces
  • Cannot name an application ID as the lookup key
  • Assumes YARN log aggregation is always enabled
  • Looks only at the driver log when an executor was killed

context