Skip to content
Everyone knows 4 min read

Driver, executors and the cluster manager

Who plans the work, who does it, and who hands out the machines in a Spark application.

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.

TL;DR A Spark application has one driver that plans the work and many executors that do it. A cluster manager (YARN, Kubernetes or Spark standalone) starts the executors. Your code runs on the driver; the functions inside DataFrame operations run on executors.

Three roles

Driver

Runs your main program and the SparkSession. Turns DataFrame code into a plan, splits it into stages and tasks, sends tasks to executors, and collects results for actions like collect().

Executors

JVM processes on worker machines. Each has a number of cores (task slots) and memory. They run tasks, cache data and write shuffle files.

Cluster manager

Hands out machines. Spark asks it for executors; it starts them. YARN, Kubernetes and Spark's own standalone manager are supported. Mesos was deprecated in Spark 3.2.

What happens when you run a job

  1. 1
    You start an application. The driver process starts and creates a SparkSession.
  2. 2
    The driver asks the cluster manager for executors, for example 10 executors with 4 cores and 16 GB each.
  3. 3
    Your code builds DataFrames. Nothing runs yet: the driver only records a plan.
  4. 4
    An action such as write or count triggers a job. The driver optimises the plan, cuts it into stages and creates one task per partition.
  5. 5
    Tasks are scheduled onto free executor cores. 10 executors with 4 cores run 40 tasks at a time.
  6. 6
    Executors 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.

CodeRuns on
df = spark.read.parquet(...), df.filter(...)Driver: builds the plan only
The filter condition, F.upper(...), joins, aggregationsExecutors, compiled to JVM code
A Python UDF bodyExecutors, 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 listDriver
Watch out: 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

The driver runs where you launched the application: a laptop, an edge node, a notebook server. Good for interactive work. If that machine disconnects, the job dies.

Cluster mode

The driver runs inside the cluster, on a node the cluster manager picks. Good for scheduled production jobs: nothing depends on the submitting machine.

Sizing an application

SettingMeaning
spark.executor.instancesNumber of executors (when dynamic allocation is off)
spark.executor.coresConcurrent tasks per executor
spark.executor.memoryJVM heap per executor
spark.executor.memoryOverheadOff-heap memory per executor (default: the larger of 384 MB and 10% of executor memory)
spark.driver.memoryDriver heap
spark.dynamicAllocation.enabledLet Spark add and remove executors with the workload
PySpark · Setting resources at submit time
# 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

Driver out-of-memory. Write to storage instead.

Giving the driver tiny memory while broadcasting

Broadcast tables are built on the driver first. A 1 GB broadcast needs a driver with room for it.

One huge executor per machine

Long GC pauses and poor throughput. Several medium executors usually work better.

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 questions

1. 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