What you will be able to do
- Explain how executor CPU cores become task slots and set an application's maximum parallelism
- Describe what executor memory, local disk and driver memory are each used for, and recognise the signs that one is too small
- Compare clusters by total executor cores and memory instead of by worker count
1.CPU cores become task slots
In Spark, work arrives at an executor as tasks, and a task is a unit of work sent to one executor. Each executor runs its tasks in multiple threads. How many tasks one executor can run at the same time depends on the CPU cores it has.
Spark uses a 1:1 mapping between task slots and available cores. Add up the slots across the cluster and you get the ceiling on parallelism. Databricks calls this total executor cores: the total number of cores across all executors, which determines the maximum parallelism of a compute.
The word *executor* in that definition matters. Cores on the driver node do not add task slots, because tasks are sent to executors. Slots are also shared by everything running on the cluster. When several queries run on one cluster, they all share the task slots on the executors, and stages from one query can occupy the available slots and delay or starve the others. Databricks' advice is to make sure enough cores are available if queries need to run concurrently.
Checkpoint 1 of 4· Check yourself
A cluster has a 16-core driver and 4 workers with 8 cores each, one executor per worker. What is the maximum number of tasks that can run at the same time?
Maximum parallelism comes from total executor cores: 4 × 8 = 32. The driver's cores are not task slots.
“Total executor cores (compute): The total number of cores across all executors. This determines the maximum parallelism of a compute.”Source: docs.databricks.com
2.Memory: executors hold the data, the driver holds the coordination
An executor does two things: it runs tasks and it keeps data, in memory or on disk, across them. How much of that data stays in memory is limited by the RAM you give the executors. Databricks defines total executor memory as the total amount of RAM across all executors. It determines how much data can be stored in memory before Spark spills it to disk.
When data does not fit, it goes to the executor's local storage. Local disk is mainly used for spills during shuffles and for caching. Spilling to disk is slower than staying in memory. Running out of memory altogether is a failure. For compute-intensive work Databricks gives a direct rule: if you see significant spill to disk or OOM errors, increase the memory available on your instances.
Executors usually carry the heavier memory load. In general, executors might perform more memory-intensive operations than the driver node, so you may need to tune the executor JVM and off-heap memory allocation parameters for your application's load. On Databricks you can pass JVM options to the driver and executors through spark.driver.extraJavaOptions and spark.executor.extraJavaOptions. The API docs show this with GC logging switched on for the driver:
{"spark.driver.extraJavaOptions": "-verbose:gc -XX:+PrintGCDetails"}Worker memory and local disk can also be set through environment variables exported when the driver and workers launch. The Jobs API gives this example of spark_env_vars:
Checkpoint 2 of 4· Fill the gap
Which environment variable in this spark_env_vars example sets the worker memory to 28000m?
{" ? ": "28000m", "SPARK_LOCAL_DIRS": "/local_disk0"}SPARK_WORKER_MEMORY carries the memory value. SPARK_LOCAL_DIRS points at the local disk that spills and caching use.
Source: docs.databricks.comThe driver is sized separately. In the Jobs API, node_type_id encodes the resources available to each Spark node, such as memory-optimized or compute-optimized hardware. driver_node_type_id is optional: if it is unset, the driver gets the same node type as node_type_id. A driver struggling with collect() results or many concurrent queries can be given a bigger node without changing the workers.
| Resource | What it determines | Symptom when too small |
|---|---|---|
| Total executor cores | Maximum parallelism (1 task slot per core) | Queries wait for or starve each other of task slots |
| Total executor memory | How much data stays in memory before spilling | Significant spill to disk or OOM errors |
| Executor local storage | Room for spills during shuffles and caching | Workers running low on disk space |
| Driver node (driver_node_type_id) | Capacity for scheduling and driver-side results | High driver CPU/memory in the Spark UI; driver OOM |
Checkpoint 3 of 4· Check yourself
A complex ETL job shows significant spill to disk in the Spark UI. Which change targets the cause most directly?
Spill happens when the data does not fit in executor memory. The executors run on the workers, so worker memory is the resource to increase.
“If you observe significant spill to disk or OOM errors, increase the amount of memory available on your instances.”Source: docs.databricks.com
3.Sizing a cluster by cores and memory, not by worker count
With cores setting parallelism and memory setting how much data stays in memory, the number of workers turns out to be the wrong measure of cluster size. Databricks says people often think of compute size in terms of the number of workers, but other factors matter as well. Its example: two workers with 16 cores and 128 GB of RAM each have the same compute and memory as 8 workers with 4 cores and 32 GB each. Both give 32 executor cores and 256 GB of executor RAM.
How the totals are split across nodes affects data movement. Databricks notes that fewer, larger nodes reduce the network and disk I/O needed for shuffles, which is why it recommends them for analyst and complex ETL workloads.
Databricks' workload guidance follows from this. A compute resource with fewer, larger nodes can reduce the network and disk I/O needed for shuffles. For complex ETL, Databricks recommends fewer workers with larger instances to make up for them. When executors run short, the advice is to size executor nodes appropriately in CPU, memory and disk, and to scale vertically first. If vertical scaling is not possible, add more worker nodes.
Checkpoint 4 of 4· Exam question
A job repeatedly fails with `java.lang.OutOfMemoryError: Java heap space` on the executors. The cluster has 10 worker nodes, each configured with `--executor-cores 8` and `--executor-memory 4g`, with one executor per worker. Increasing `--executor-cores` to 16 without changing `--executor-memory` makes the failures worse. Why?
Correct answer: A — More cores let the executor run more concurrent tasks inside the same JVM heap, dividing the fixed 4g memory pool across more simultaneous tasks and shrinking each one's share.
- A. Correct. All task slots in an executor share the same 4g JVM heap; doubling the core count doubles the number of tasks that can run at once, so each task gets a smaller effective share of that fixed heap.
- B. Incorrect. Changing `--executor-cores` does not automatically adjust `spark.executor.memory`; the two settings are configured independently.
- C. Incorrect. Additional cores run additional task threads inside the same executor JVM; they do not spawn separate JVM processes each with their own heap.
- D. Incorrect. There is no cluster-manager-enforced heap cap tied to worker identity; the 4g figure comes directly from the `--executor-memory` setting applied to that executor.
- E. Incorrect. Partition count is driven by the input data and any repartitioning calls, not by the number of cores assigned to an executor, and partitions do not each reserve a fixed heap buffer.
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.A cluster with more workers always has more processing power than one with fewer workers.Why is that wrong?
Capacity comes from total executor cores and total executor memory, not from how many workers there are. Two 16-core/128 GB workers match eight 4-core/32 GB workers.
Covered in Sizing a cluster by cores and memory, not by worker count
Sources
Every claim above is drawn from one of these pages, quoted as it was written on the date shown.
- 1.
“Spark uses a 1:1 mapping between task slots and available cores.”
↩︎ CPU cores become task slots“With multiple queries running on the same cluster, all queries share task slots on the executors.”
↩︎ CPU cores become task slots“In general, executors might perform more memory-intensive operations than the driver node.”
↩︎ Memory: executors hold the data, the driver holds the coordination“Tune the executor JVM and off-heap memory allocation parameters if needed to handle your application load.”
↩︎ Memory: executors hold the data, the driver holds the coordination“Ensure that executor nodes are sized appropriately in terms of CPU, memory, and disk space and scale vertically if needed.”
↩︎ Sizing a cluster by cores and memory, not by worker count“If vertical scaling is not possible, consider adding additional worker nodes to the cluster.”
↩︎ Sizing a cluster by cores and memory, not by worker count - 2.
“Total executor cores (compute): The total number of cores across all executors. This determines the maximum parallelism of a compute.”
↩︎ CPU cores become task slots“The total amount of RAM across all executors. This determines how much data can be stored in memory before spilling it to disk.”
↩︎ Memory: executors hold the data, the driver holds the coordination“Local disk is primarily used in the case of spills during shuffles and caching.”
↩︎ Memory: executors hold the data, the driver holds the coordination“People often think of compute size in terms of the number of workers, but there are other important factors to consider”
↩︎ Sizing a cluster by cores and memory, not by worker count“A compute resource with a smaller number of larger nodes can reduce the network and disk I/O needed to perform these shuffles.”
↩︎ Sizing a cluster by cores and memory, not by worker count“has the same compute and memory as configuring compute with 8 workers, each with 4 cores and 32 GB of RAM.”
↩︎ Exam trap 1“If you observe significant spill to disk or OOM errors, increase the amount of memory available on your instances.”
↩︎ Checkpoint - 3.https://spark.apache.org/docs/latest/cluster-overview.htmlSecondary source
“A unit of work that will be sent to one executor”
↩︎ CPU cores become task slots - 4.
“You can also pass in a string of extra JVM options to the driver and the executors via spark.driver.extraJavaOptions and spark.executor.extraJavaOptions respectively.”
↩︎ Memory: executors hold the data, the driver holds the coordination“if unset, the driver node type will be set as the same value as node_type_id defined above.”
↩︎ Memory: executors hold the data, the driver holds the coordination“This field encodes, through a single value, the resources available to each of the Spark nodes in this cluster.”
↩︎ Memory: executors hold the data, the driver holds the coordination