What you will be able to do
- Describe what a Spark application is made of and what the cluster manager does
- Explain what the driver node does and does not do, and why its placement and sizing matter
- Tell a worker node apart from an executor, and explain why executors are isolated per application
Key concept
Driver–executor architecture — A Spark application has one driver process and a set of executor processes. The driver runs your main program, holds the SparkContext and schedules work. The executors run on worker nodes, where they execute the tasks and keep the application's data.
1.The cluster, the application and the cluster manager
Databricks defines a cluster as a set of computation resources and configurations on which you run notebooks and jobs. Spark does not take over that cluster as a single program. Each application, meaning the user program built on Spark, runs on it as its own set of processes: one driver program and a number of executors.
Spark does not acquire machines itself. When an application starts, its SparkContext connects to a cluster manager, which the Spark docs describe as an external service for acquiring resources on the cluster. Spark can use its own Standalone manager, Hadoop YARN or Kubernetes. Once Spark is connected, it acquires executors on nodes in the cluster, sends them your application code (JAR or Python files) and then sends them tasks to run.
That order is the whole lifecycle of resource acquisition: connect to the manager, get executors, ship code, ship tasks.
Checkpoint 1 of 6· Put it in order
Put these steps in the order Spark follows when an application starts on a cluster
- 1.SparkContext sends tasks to the executors to run
- 2.Spark sends the application code to the executors
- 3.Spark acquires executors on nodes in the cluster
- 4.SparkContext connects to a cluster manager
The cluster overview gives this sequence. Resources come first, from the manager. Executors are acquired next, then the code is shipped, and only then are tasks sent.
“Once connected, Spark acquires executors on nodes in the cluster, which are processes that run computations and store data for your application.”Source: spark.apache.org
Spark is agnostic to the underlying cluster manager. As long as it can get executor processes and those processes can talk to each other, the rest of the architecture stays the same. On Databricks the shape of a cluster is fixed in its configuration. The Jobs API states that a cluster has one Spark driver and num_workers executors, for a total of num_workers + 1 Spark nodes. A cluster with num_workers set to 4 therefore runs on five machines.
It runs 11 nodes. num_workers counts only the worker (executor) nodes. Every cluster also has exactly one driver node on top of those.
2.The driver node: coordinator, not workhorse
The Spark glossary defines the driver program as the process that runs the application's main() function and creates the SparkContext. Databricks describes the same component from the node's side. The driver node maintains attached notebook state, maintains the SparkContext, interprets notebook and library commands, and runs the Spark master that coordinates with Spark executors.
The driver has three jobs: run your program, turn it into tasks, and coordinate the executors that run those tasks. It also hosts the application's web UI, typically on port 4040, which shows running tasks, executors and storage usage.
This coordination role puts constraints on where the driver lives. The driver program must listen for and accept incoming connections from its executors for its whole lifetime, so it must be network addressable from the worker nodes. Because it schedules tasks on the cluster, it should run close to the worker nodes, preferably on the same local area network. Deploy mode controls where it runs. In cluster mode the framework launches the driver inside the cluster. In client mode the submitter launches it outside the cluster.
Checkpoint 2 of 6· Check yourself
A team submits an application so that the driver runs on a laptop outside the cluster. Which deploy mode is this?
In client mode the submitter launches the driver outside the cluster. In cluster mode the framework launches it inside. Standalone is a cluster manager, not a deploy mode.
“In "cluster" mode, the framework launches the driver inside of the cluster. In "client" mode, the submitter launches the driver outside of the cluster.”Source: spark.apache.org
The driver is one process on one node, and everything in the application shares it. Databricks warns that the driver is a shared resource: multiple queries share its CPU, memory, DAG scheduler, task scheduler and driver-side UDF execution. The usual way to overload it is to bring data back to it. Databricks advises you to avoid calling collect() or other expensive DataFrame operations on the driver unless absolutely necessary, because materializing large result sets there causes out-of-memory (OOM) errors. If the Spark UI shows high driver CPU, memory or disk usage, you adjust driver sizing in the cluster compute settings.
Checkpoint 3 of 6· Exam question
A team submits `spark-submit --num-executors 5 --executor-cores 4 --executor-memory 16g job.py` on a standalone cluster. Each worker node has 8 CPU cores and 32 GB RAM and hosts exactly one executor. A colleague asks why only 4 tasks from the current stage run at once on each executor, even though every worker has 8 cores physically available. What should you tell them?
Correct answer: A — Each executor's task-slot count equals `--executor-cores`, so this executor can only run 4 tasks concurrently regardless of the worker's total core count.
- A. Correct. `--executor-cores` sets the number of task slots an executor can run in parallel; only 4 were requested, so the extra physical cores on the worker sit idle for this application.
- B. Incorrect. Spark's built-in shuffle service does not reserve a fixed pool of cores per executor, and this behavior is not how executor concurrency is limited.
- C. Incorrect. Executor memory governs how much data an executor can hold and process, not how many tasks run concurrently; that is controlled by the core count.
- D. Incorrect. Standalone mode does not impose a hidden half-core ceiling per executor; the executor simply runs as many concurrent tasks as `--executor-cores` specifies.
- E. Incorrect. The driver schedules tasks based on stage partitioning, but the number that can run simultaneously on a given executor is bounded by that executor's configured cores, not partition counts of unrelated RDDs.
3.Worker nodes and executors
Spark separates a machine from a process. A worker node is any node that can run application code in the cluster. An executor is a process launched for an application on a worker node. It runs tasks and keeps data in memory or disk storage across them. A task, the unit of work, is always sent to one executor.
| Term | Meaning |
|---|---|
| Driver program | The process running the main() function of the application and creating the SparkContext |
| Cluster manager | An external service for acquiring resources on the cluster (e.g. standalone manager, YARN, Kubernetes) |
| Worker node | Any node that can run application code in the cluster |
| Executor | A process launched for an application on a worker node, that runs tasks and keeps data in memory or disk storage across them |
| Task | A unit of work that will be sent to one executor |
Checkpoint 4 of 6· Match them up
Match each component to its role
Tap a term, then the definition that fits it.
These are the glossary definitions. The worker is the machine, and the executor is the process your application gets on that machine.
“A process launched for an application on a worker node, that runs tasks and keeps data in memory or disk storage across them.”Source: spark.apache.org
Databricks often uses the two words interchangeably. Its clusters consist of a Spark driver node and zero or more Spark worker (also known as executor) nodes, and worker nodes run one Spark executor per worker node. This one-to-one mapping is how Databricks is set up. Spark's own definitions still keep the node and the process separate.
The workers are where Spark work happens. If a Databricks compute resource has zero workers, you can run non-Spark commands on the driver node, but Spark commands fail. The exception is a single-node cluster, which has one driver node, no worker nodes, and runs Spark in local mode.
Executors also belong to a single application. Each application gets its own executor processes. Those processes stay up for the whole application and run tasks in multiple threads. Tasks from different applications run in different JVMs, which isolates the applications from each other. The cost is that data cannot be shared across different Spark applications without first writing it to an external storage system.
Checkpoint 5 of 6· Exam question
A cluster runs a single Spark application with one driver and three executors, each on a separate worker node. Partway through a long-running job, the driver process crashes due to an unrelated JVM issue on its host machine. What happens to the application?
Correct answer: A — The entire application terminates because the driver holds the SparkContext and coordinates all executors; losing it ends the job even though the executors are still healthy.
- A. Correct. The driver owns the SparkContext and is the single coordinator for scheduling and tracking the application's tasks; its loss ends the application regardless of executor health.
- B. Incorrect. Executors do not elect a replacement driver; there is no built-in leader-election mechanism among executor processes in standard Spark deployments.
- C. Incorrect. Losing the driver ends the whole application rather than triggering a selective stage restart, because the component that tracks stage state and dependencies is gone.
- D. Incorrect. Executors are not designed to continue as independent applications; they depend on the driver for task assignment and coordination and terminate along with it.
- E. Incorrect. A crashed driver process cannot be reattached to its original SparkContext; the application must be resubmitted rather than resumed from a paused state.
Checkpoint 6 of 6· Check yourself
Two separate Spark applications run on the same cluster. Application A has cached a DataFrame in its executors' memory. How can application B use that data?
Executors are per application and run in separate JVMs. Data cannot cross from one application to another without going through external storage.
“data cannot be shared across different Spark applications (instances of SparkContext) without writing it to an external storage system.”Source: spark.apache.org
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.A worker node and an executor are the same thing.Why is that wrong?
In Spark a worker node is a machine that can run application code, and an executor is a process launched on it for one application. Databricks happens to run one executor per worker node, which is why its docs use the terms interchangeably.
Covered in Worker nodes and executors
2.The driver does the heavy data processing, so pulling results back to it with collect() is harmless.Why is that wrong?
The driver coordinates and schedules, and executors run the tasks. The driver is a single shared process, and materializing large results on it with collect() can cause out-of-memory errors.
Covered in The driver node: coordinator, not workhorse
Sources
Every claim above is drawn from one of these pages, quoted as it was written on the date shown.
- 1.
“A set of computation resources and configurations on which you run notebooks and jobs.”
↩︎ The cluster, the application and the cluster manager - 2.
“A cluster has one Spark driver and num_workers executors for a total of num_workers + 1 Spark nodes.”
↩︎ The cluster, the application and the cluster manager - 3.https://spark.apache.org/docs/latest/cluster-overview.htmlSecondary source
“An external service for acquiring resources on the cluster (e.g. standalone manager, YARN, Kubernetes)”
↩︎ The cluster, the application and the cluster manager“Spark is agnostic to the underlying cluster manager.”
↩︎ The cluster, the application and the cluster manager“Finally, SparkContext sends tasks to the executors to run.”
↩︎ The cluster, the application and the cluster manager“the driver program must be network addressable from the worker nodes.”
↩︎ The driver node: coordinator, not workhorse“The process running the main() function of the application and creating the SparkContext”
↩︎ The driver node: coordinator, not workhorse“Each application gets its own executor processes, which stay up for the duration of the whole application and run tasks in multiple threads.”
↩︎ Worker nodes and executors“Spark applications run as independent sets of processes on a cluster, coordinated by the SparkContext object in your main program (called the driver program).”
↩︎ Key concept“Once connected, Spark acquires executors on nodes in the cluster, which are processes that run computations and store data for your application.”
↩︎ Checkpoint“Because the driver schedules tasks on the cluster, it should be run close to the worker nodes, preferably on the same local area network.”
↩︎ Prediction“In "cluster" mode, the framework launches the driver inside of the cluster. In "client" mode, the submitter launches the driver outside of the cluster.”
↩︎ Checkpoint“A process launched for an application on a worker node, that runs tasks and keeps data in memory or disk storage across them.”
↩︎ Checkpoint“data cannot be shared across different Spark applications (instances of SparkContext) without writing it to an external storage system.”
↩︎ Checkpoint - 4.https://docs.databricks.com/aws/en/sparkrOfficial docs
“runs the Spark master that coordinates with Spark executors”
↩︎ The driver node: coordinator, not workhorse“Databricks clusters consist of an Apache Spark driver node and zero or more Spark worker (also known as executor) nodes.”
↩︎ Worker nodes and executors“A single node cluster has one driver node and no worker nodes, with Spark running in local mode”
↩︎ Worker nodes and executors“Worker nodes run the Spark executors, one Spark executor per worker node.”
↩︎ Exam trap 1 - 5.
“The driver is a shared resource. Multiple queries share the same CPU, memory, DAG scheduler, task scheduler, and driver-side UDF execution”
↩︎ The driver node: coordinator, not workhorse“If you see high CPU, memory, or disk usage, adjust driver sizing in the cluster compute settings.”
↩︎ The driver node: coordinator, not workhorse“Avoid calling collect() or other expensive DataFrame operations on the driver unless absolutely necessary”
↩︎ Exam trap 2 - 6.https://docs.databricks.com/aws/en/compute/configureOfficial docs
“If the compute resource has zero workers, you can run non-Spark commands on the driver node, but Spark commands will fail.”
↩︎ Worker nodes and executors