Driver, executors and the cluster manager
Who plans the work, who does it, and who hands out the machines in a Spark application.
On this page
You will learn
- What the driver does and what executors do
- What a cluster manager is, and the ones Spark supports
- The difference between client and cluster deploy mode
- Which settings size an application, and how to reason about them
Read first
Nothing. This lesson starts from scratch.
Three roles
Driver
collect().Executors
Cluster manager
What happens when you run a job
- 1You start an application. The driver process starts and creates a SparkSession.
- 2The driver asks the cluster manager for executors, for example 10 executors with 4 cores and 16 GB each.
- 3Your code builds DataFrames. Nothing runs yet: the driver only records a plan.
- 4An action such as
writeorcounttriggers a job. The driver optimises the plan, cuts it into stages and creates one task per partition. - 5Tasks are scheduled onto free executor cores. 10 executors with 4 cores run 40 tasks at a time.
- 6Executors read data, run the tasks and write shuffle files or output. The driver tracks progress and retries failed tasks.
Where does my code run?
This is the question behind many confusing errors.
| Code | Runs on |
|---|---|
df = spark.read.parquet(...), df.filter(...) | Driver: builds the plan only |
The filter condition, F.upper(...), joins, aggregations | Executors, compiled to JVM code |
| A Python UDF body | Executors, in Python worker processes next to each executor |
df.collect(), toPandas() | Executors compute; all rows are sent to the driver |
| A plain Python loop over a list | Driver |
collect() on a large DataFrame pulls every row into the driver's memory and is the classic way to crash a driver. Write results to storage, or limit before collecting. spark.driver.maxResultSize (1 GB by default) is the safety net.Client vs cluster mode
Client mode
Cluster mode
Sizing an application
| Setting | Meaning |
|---|---|
spark.executor.instances | Number of executors (when dynamic allocation is off) |
spark.executor.cores | Concurrent tasks per executor |
spark.executor.memory | JVM heap per executor |
spark.executor.memoryOverhead | Off-heap memory per executor (default: the larger of 384 MB and 10% of executor memory) |
spark.driver.memory | Driver heap |
spark.dynamicAllocation.enabled | Let Spark add and remove executors with the workload |
# spark-submit --master yarn --deploy-mode cluster \ # --num-executors 10 --executor-cores 4 --executor-memory 16g app.py spark = (SparkSession.builder .config("spark.executor.cores", "4") .config("spark.executor.memory", "16g") .getOrCreate())
A common rule of thumb is 4 to 5 cores per executor: fewer wastes memory on per-JVM overhead, more causes contention and long garbage-collection pauses. Total parallelism is executors × cores; aim for a few tasks per core per stage so that slow tasks do not leave cores idle.
Common mistakes
Collecting big results to the driver
Giving the driver tiny memory while broadcasting
One huge executor per machine
Key takeaways
- The driver plans and schedules; executors run tasks; the cluster manager provides machines.
- DataFrame code builds a plan on the driver; the work runs on executors.
- Client mode keeps the driver on your machine; cluster mode puts it in the cluster.
- Parallelism is executors × cores per executor.
Check yourself
3 questions1. Where does the body of a Python UDF run?
Show the answer
In Python worker processes on the executors. Executors start Python workers and send them rows for Python UDFs.
2. 8 executors with 5 cores each. How many tasks can run at the same time?
Show the answer
40. Each core is a task slot: 8 × 5 = 40.
3. Which deploy mode suits a nightly production job?
Show the answer
Cluster mode. In cluster mode the driver runs in the cluster, so the job does not depend on the submitting machine.
Go deeper
Primary sources: Cluster mode overview · Configuration