What you will be able to do
- Tell expected executor removals (autoscaling, spot loss) apart from executors running out of memory
- Follow the documented path from a failed job to the logs of failed executors
- Confirm a suspected memory problem by changing the memory-per-core ratio, and list its common causes
- Spot an overloaded or hanging driver and use thread dumps to investigate it
- Use compute metrics and task placement to recognize an underutilized cluster
1.Failed jobs and removed executors
A removed executor is not always a problem. Databricks lists three common reasons: autoscaling (expected, not an error), spot instance losses (the cloud provider reclaiming VMs), and executors running out of memory. Only the last one points to your code or sizing. So the diagnosis works by elimination.
Start with the failing job. Click it, scroll to the failed stage and read its failure reason. If the reason is generic, click the link in the description and scroll down to see why each task failed. Then check the compute's Event log: it may show the cluster resizing or spot instances being lost. If the event log explains nothing, go back to the Spark UI, open the Executors tab and get the logs from the failed executors. If you have got that far, the guide says "the likeliest explanation is a memory issue."
Checkpoint 1 of 6· Put it in order
Put the documented steps for investigating failing executors in order
- 1.Check the compute's Event log for resizing or spot instance losses
- 2.Open the Spark UI Executors tab and get the logs from the failed executors
- 3.Click the description link to see why each task failed
- 4.Click the failing job and scroll to the failed stage and its failure reason
You work down from the job to the stage to the tasks. Next you rule out infrastructure causes in the Event log, and only then read the failed executors' logs.
“If you don't see any information in the event log, navigate back to the Spark UI then click the Executors tab:”Source: docs.databricks.com
Sources1
2.Confirming an out-of-memory problem
Memory problems often show up as an error like the one below. Note its last sentence: it sends you to the driver logs.
SparkException: Job aborted due to stage failure: Task 3 in stage 0.0 failed 4 times, most recent failure: Lost task 3.3 in stage 0.0 (TID 30) (10.139.64.114 executor 4): ExecutorLostFailure (executor 4 exited caused by one of the running tasks) Reason: Remote RPC client disassociated. Likely due to containers exceeding thresholds, or network issues. Check driver logs for WARN messages.Don't read too much into this message. The guide warns that such messages "are often generic and can be caused by other issues". The message itself mentions network issues as an alternative. To confirm a memory problem, run an experiment: double the memory per core and see whether the failure changes. What matters is the ratio of cores to memory, not the total memory.
| Worker type | Cores | Memory | Memory per core |
|---|---|---|---|
| Original | 4 | 16GB | 4GB |
| Test | 4 | 32GB | 8GB |
If the job takes longer to fail, or doesn't fail at all, you're on the right track. More memory may then be the whole fix. If it doesn't help, or the extra cost is too high, look further. The guide lists the usual causes: too few shuffle partitions, a large broadcast, UDFs, a window function without a PARTITION BY statement, skew, and streaming state. The compute metrics UI also helps here. Its JVM heap usage chart shows heap use against heap capacity and the configured maximum, and Container memory usage shows the Spark container's memory against its configured limit.
Checkpoint 2 of 6· Exam question
In a Databricks cluster, which statement correctly distinguishes the purpose of driver logs from executor logs when diagnosing a failed Spark application?
Correct answer: A — Driver logs capture the driver process's control-flow output — Python tracebacks, notebook `print` statements, and scheduler messages — while executor logs capture per-task standard output and JVM errors from the worker nodes running the computation.
- A. This matches how Spark instruments the two processes: the driver runs the user's program and scheduler, so its logs surface driver-side exceptions and print output, while each executor runs assigned tasks, so its logs surface task-level errors such as out-of-memory failures on that worker.
- B. Garbage collection pauses are recorded for both the driver and every executor JVM, not just the driver, and query plans generated by Catalyst are visible through `explain()` output and the Spark UI's SQL tab rather than being confined to the executor logs.
- C. Shuffle read and write metrics are collected per task and surfaced through the Stages tab of the Spark UI and the executor logs where those tasks ran, not through the driver logs, and DAG visualization data is UI rendering, not raw executor log content.
- D. Both driver and executor logs are produced in local mode and in cluster mode; local mode simply runs the driver and a single executor in the same JVM, so the distinction between the two log types still applies to whichever processes exist, not to the deployment mode itself.
- E. Cluster autoscaling and node provisioning events are surfaced through the cluster event log in the Databricks UI, and Delta Lake transaction log entries are stored as JSON files in the table's `_delta_log` directory, not written into either driver or executor application logs.
Checkpoint 3 of 6· Check yourself
You suspect an OOM. Your workers have 8 cores and 32GB. Which test change best follows the documented approach to confirm it?
The test doubles memory per core (4GB to 8GB). 16 cores/64GB, more workers of the same type, and 4 cores/16GB all keep 4GB per core.
“It's the ratio of cores to memory that matters here.”Source: docs.databricks.com
3.When the driver is the bottleneck or hangs
Executors are not the only place memory and CPU run out. The driver can be overloaded too, and the guide gives the most common reason: "there are too many concurrent things running on the cluster". That means too many streams, queries or Spark jobs, sometimes launched from threads. You have three options: increase the size of the driver, reduce the concurrency, or spread the load over multiple clusters. Databricks recommends trying the first one first, by doubling the driver size.
When the driver seems to hang (no Spark progress bars, or bars stuck at 100%), take a thread dump, which is a snapshot of a JVM's thread states. Open the Executors tab and, in the driver row, click the link in the Thread Dump column. A slow task works the same way. Note the task's Task ID and Executor ID from the stage's task list, open that executor's thread dump, and find the thread whose name contains TID followed by the Task ID. If the task has already finished, there is no matching thread.
Checkpoint 4 of 6· Check yourself
Several notebooks each launch many concurrent Spark jobs on one cluster, and the driver is overloaded. What does Databricks recommend trying first?
Concurrency overloads the driver, not the workers. The first recommended step is to double the driver size and measure the effect.
“Databricks recommends you first try doubling the size of the driver and see how that impacts your job.”Source: docs.databricks.com
4.Recognizing an underutilized cluster
The opposite of an overloaded cluster is one paying for workers that do nothing. One cause is non-Spark code running on the driver. If you see gaps in the timeline caused by that code, all the workers are idle during those gaps and probably wasting money. Rewriting that code to use Spark lets you fully utilize the cluster. Another cause is poor task spread. Make sure tasks run on multiple executors, because sometimes only one executor does all the work even though the compute has more than one (a single streaming receiver is the documented example). The Spark UI's task list shows which executor each task ran on.
For a view over time, open the compute's Metrics tab. It shows hardware metrics by default, and you switch the drop-down from Hardware to Spark for Spark metrics. Metrics are collected every minute, kept for 30 days, and averaged across all nodes, including the driver, unless you pick a node. Serverless compute uses query insights instead of this UI.
| Metric chart | What it shows | Underutilization signal |
|---|---|---|
| CPU utilization and active nodes | Active node count and time the CPU spent in each mode, including idle | High idle share while nodes are active |
| Server load distribution | CPU utilization over the past minute for each node | A few busy tiles while the rest are quiet |
| Active tasks | Total number of tasks executing at any given time | Few active tasks on a large compute |
| By node view | Metrics for all individual nodes on one page | Outlier nodes |
Checkpoint 5 of 6· Exam question
A PySpark job aggregates web clickstream data and consistently fails with executor `OutOfMemoryError` exceptions, even though the cluster has ample total memory. The Spark UI's Stages tab shows one task in the failing stage processing 40 GB of shuffle data while the other 199 tasks in the same stage each process under 200 MB: ```python clicks = spark.read.parquet("/mnt/clickstream/events") sessions = clicks.groupBy("session_id").agg(collect_list("page").alias("pages")) sessions.write.format("delta").mode("overwrite").save("/mnt/gold/sessions") ``` What is the most likely root cause, and which change addresses it?
Correct answer: A — A small number of `session_id` values account for a disproportionate share of the clicks, so one task inherits far more rows than the rest; salting the skewed keys or enabling adaptive skew join handling redistributes that task's work across more partitions.
- A. Uneven task sizes in the Stages tab, where one task's shuffle input dwarfs the rest, is the classic signature of data skew on the grouping key, and mitigations that spread the heavy key's rows across multiple partitions — key salting or Spark's adaptive skewed-join handling — directly relieve the executor that would otherwise run out of memory.
- B. The described symptom is skew in the shuffle stage that produces `sessions`, which happens before the write step even begins; changing how many records go into each output file changes the shape of the Delta write, not the distribution of work in the preceding `groupBy` aggregation.
- C. Adaptive query execution is enabled by default in current Spark releases and already targets this exact scenario through its skewed-join optimization, but that optimization applies to joins, and there is no configuration named `spark.sql.adaptiveJoinEnabled` that converts a `groupBy` aggregation into a broadcast operation.
- D. `collect_list` runs as an executor-side aggregation like any other grouped function and does not inherently ship its output to the driver; the failure here is an executor-level `OutOfMemoryError` tied to one task's data volume, which points to skew rather than to a driver-bound aggregation function.
- E. Row group layout in the source Parquet files affects how Spark splits input for reading, but the failure occurs later, in the shuffle stage produced by `groupBy("session_id")`, and the uneven task sizes reported there are keyed by `session_id` values, not by how the upstream files were originally written.
Checkpoint 6 of 6· Check yourself
A job's timeline shows long gaps between Spark jobs while a driver-side Python loop runs. What does this mean for the cluster?
Non-Spark code keeps only the driver busy, so the workers sit idle during those gaps.
“If you see gaps in your timeline caused by running non-Spark code, this means your workers are all idle”Source: docs.databricks.com
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.A removed executor always means the job hit an error such as out-of-memory.Why is that wrong?
Autoscaling removes executors as expected behaviour, and spot instance losses are infrastructure events. Check the compute's Event log before assuming a memory problem.
Covered in Failed jobs and removed executors
2.An ExecutorLostFailure message proves the executor ran out of memory.Why is that wrong?
These messages are often generic. Confirm a memory problem by doubling memory per core and seeing whether the failure changes.
Covered in Confirming an out-of-memory problem
3.You can filter the Spark metrics charts to a single node to find the idle one.Why is that wrong?
Only hardware (and GPU) metrics can be shown per node. Spark metrics are cluster-wide, so use Server load distribution or the By node view to find outlier nodes.
Covered in Recognizing an underutilized cluster
Sources
Every claim above is drawn from one of these pages, quoted as it was written on the date shown.
- 1.
“To find out why your executors are failing, you'll first want to check the compute's Event log”
↩︎ Failed jobs and removed executors“If you've gotten this far, the likeliest explanation is a memory issue.”
↩︎ Failed jobs and removed executors“Autoscaling: In this case it's expected and not an error.”
↩︎ Exam trap 1“Autoscaling: In this case it's expected and not an error.”
↩︎ Prediction“If you don't see any information in the event log, navigate back to the Spark UI then click the Executors tab:”
↩︎ Checkpoint - 2.
“you can verify the issue by doubling the memory per core to see if it impacts your problem.”
↩︎ Confirming an out-of-memory problem“If it takes longer to fail with the extra memory or doesn't fail at all, that's a good sign”
↩︎ Confirming an out-of-memory problem“These error messages, however, are often generic and can be caused by other issues.”
↩︎ Exam trap 2“It's the ratio of cores to memory that matters here.”
↩︎ Checkpoint - 3.
“JVM heap usage: The JVM heap memory usage, averaged across all applicable nodes.”
↩︎ Confirming an out-of-memory problem“Server load distribution: These tiles show the CPU utilization over the past minute for each node in the compute resource.”
↩︎ Recognizing an underutilized cluster“To help identify any outlier nodes within the cluster, you can also view metrics for all individual nodes on a single page.”
↩︎ Recognizing an underutilized cluster“Spark metrics are not available for individual nodes.”
↩︎ Exam trap 3 - 4.https://docs.databricks.com/aws/en/optimizations/spark-ui-guide/spark-driver-overloadedOfficial docs
“The most common reason for this is that there are too many concurrent things running on the cluster.”
↩︎ When the driver is the bottleneck or hangs“Databricks recommends you first try doubling the size of the driver and see how that impacts your job.”
↩︎ Checkpoint“If you see gaps in your timeline caused by running non-Spark code, this means your workers are all idle”
↩︎ Checkpoint - 5.
“In the Executors table, in the driver row, click the link in the Thread Dump column.”
↩︎ When the driver is the bottleneck or hangs“Thread dumps are useful in debugging a specific hanging or slow-running task.”
↩︎ When the driver is the bottleneck or hangs“Ensure that the tasks are executed on multiple executors (nodes) in your compute to have enough parallelism while processing.”
↩︎ Recognizing an underutilized cluster“sometimes only one executor might be doing all the work though you have more than one executor in your compute.”
↩︎ Recognizing an underutilized cluster