Cluster architecture
Every confusing Spark error and every bad performance decision traces back to not knowing which process is doing what. This lesson gives you the complete map: four components, seven vocabulary words, two deploy modes — and then you will never be lost again.
After this lesson you can…
- Draw the driver / cluster manager / executor picture from memory.
- Say precisely what an application, job, stage, task, slot, core and partition each are.
- Compute how many tasks your cluster can run at the same time.
- Explain why
print()inside a transformation shows nothing, and whycollect()can kill the driver. - Choose client vs cluster deploy mode, and know what each cluster manager gives you.
1. The mental picture
The construction site
A site manager (the driver) holds the blueprint. She never lifts a brick. She breaks the build into work packets, hands them to whoever is free, tracks what is finished, and assembles the final report. The workers (executors) each own a patch of the site, hold their own tools and materials (memory), and do the actual lifting. Neither of them owns the land — a site-leasing company (the cluster manager) decides which plots and crews the project gets. If a worker walks off, the manager notices and re-issues that packet to somebody else.
2. The four components
Who does what — the responsibility split
| Driver (1 per app) | Executor (many per app) |
|---|---|
Runs your main() / notebook cell | Runs tasks — the actual row-by-row work |
| Builds and optimises the query plan | Reads and writes data from storage |
| Splits the plan into stages and tasks | Holds cached partitions in memory |
| Decides which executor gets which task | Writes shuffle files to local disk |
| Tracks progress, retries failed tasks | Sends heartbeats and metrics back |
Receives collect() / show() results | Runs the Python worker for UDFs |
| Single point of failure — it dies, the app dies | Replaceable — Spark re-runs lost tasks elsewhere |
1. print() inside a UDF or foreach shows nothing. That code runs on an
executor, on another machine. Its stdout goes to that executor's log, not your terminal.
2. df.collect() crashes the driver with an OOM. collect() pulls every row from
every executor into the driver's single heap. A 200 GB DataFrame does not fit in a 4 GB driver.
Use df.show(20), df.limit(1000).collect(), or write to storage instead.
3. Seven words you must not mix up
Spark's vocabulary is small but people use the words loosely and then talk past each other. Here is the exact hierarchy, from biggest to smallest.
| Term | Exact meaning | Created by |
|---|---|---|
| Application | One SparkSession and its driver + executors. Your whole script. | SparkSession.builder.getOrCreate() |
| Job | All the work needed to satisfy one action. | Every count(), show(), write, collect() |
| Stage | A set of tasks that can run without moving data between machines. | The DAG scheduler, cutting at each shuffle |
| Task | The work for one partition of one stage. The smallest unit Spark schedules. | One per partition per stage |
| Partition | A chunk of the data, living on one executor. Typically 128 MB. | The data source, or repartition() |
| Core | A CPU thread given to an executor via --executor-cores. | Your submit configuration |
| Slot | A place to run one task. slots = executors × cores. | Derived — this is your true parallelism |
Work out your real parallelism
10 executors × 4 cores each = 40 slots
(40 tasks run at the same time)
Your DataFrame has 400 partitions = 400 tasks in the stage
400 tasks / 40 slots = 10 sequential "waves"
If each task takes 30s -> stage takes ~5 minutes (10 waves × 30s)
This single calculation explains most tuning decisions. More partitions than slots is good — it keeps every core busy and lets Spark rebalance. Fewer partitions than slots wastes cores you are paying for: 20 partitions on 40 slots leaves half your cluster idle no matter how big it is.
Aim for 2–4× as many partitions as total slots, with each partition around 100–200 MB. With 40 slots, 80–160 partitions is a healthy target. Lesson 19 covers how to get there.
4. What actually happens when you submit an application
Eight steps, in order. When something hangs, knowing which step it hung on tells you where to look.
- You run
spark-submit(or a notebook cell). A driver process starts and your Python code begins executing. SparkSession.builder.getOrCreate()runs. The driver starts its JVM, boots a SparkContext, and opens a web UI on port 4040.- The driver registers with the cluster manager and requests executors — "give me 10 containers of 4 cores and 8 GB".
- The cluster manager finds capacity on worker nodes and launches executor JVMs there. Each executor registers back with the driver. If the cluster is full, you wait here — this is the "job is stuck in ACCEPTED" state.
- Your transformations build a plan. Nothing runs yet. This is pure driver-side bookkeeping (lesson 06).
- You call an action. Catalyst optimises the plan, the DAG scheduler cuts it into stages, and the task scheduler queues the tasks of the first stage.
- Tasks are serialized and shipped to executors, preferring executors that already hold the relevant data (data locality). Executors run them, write shuffle files, and report back.
- All stages complete, the result returns to the driver or lands in storage, and on
spark.stop()the executors are released back to the cluster manager.
5. Deploy modes: where does the driver live?
This is a genuinely important choice and interviewers ask it constantly. The only difference is the physical location of the driver process.
pyspark shell, development — because you
need to see output in your terminal. Cluster mode is for production: the driver is inside the
datacentre, network latency is low, and your laptop going to sleep does not kill a 6-hour job.
# Interactive / development — driver on your machine, output in your terminal
spark-submit --master yarn --deploy-mode client my_job.py
# Production — driver launched inside the cluster, survives your laptop
spark-submit --master yarn --deploy-mode cluster \
--num-executors 10 --executor-cores 4 --executor-memory 8g \
--driver-memory 4g \
my_job.py
# Kubernetes only supports cluster mode for production submissions
spark-submit --master k8s://https://api.my-cluster:6443 --deploy-mode cluster my_job.py
6. The cluster managers compared
| Manager | What it is | Use it when | Watch out for |
|---|---|---|---|
| local[*] | Not a cluster. Driver and "executors" are threads in one JVM on your machine. | Learning, unit tests, development. This is what the whole course uses. | No real network, so shuffle and locality behave differently from production. |
| Standalone | Spark's own built-in manager. A master process plus worker processes. | A dedicated Spark-only cluster with minimal moving parts. | Basic scheduling; no multi-framework sharing or rich queues. |
| YARN | Hadoop's resource manager. Spark asks the ResourceManager for containers. | You already run Hadoop — on-prem, EMR, Cloudera. | Queue limits cause silent "stuck in ACCEPTED" waits; memory overhead config matters. |
| Kubernetes | Driver and executors run as pods. Now the mainstream cloud choice. | Containerised platforms, mixed workloads, autoscaling infrastructure. | Pod startup latency; image and dependency management is on you. |
| Databricks / EMR / Dataproc | Managed services wrapping the above. | You want the cluster to be someone else's problem. | Cost. Also they add non-standard extensions you can become locked into. |
Static vs dynamic allocation — how many executors do you actually get?
By default the driver requests a fixed number of executors and holds them for the whole application, idle or not. Dynamic allocation lets Spark add executors when tasks are queueing and release them when they have been idle, which matters a lot for cost on shared or cloud clusters.
spark-submit \ --conf spark.dynamicAllocation.enabled=true \ --conf spark.dynamicAllocation.minExecutors=2 \ --conf spark.dynamicAllocation.maxExecutors=50 \ --conf spark.dynamicAllocation.executorIdleTimeout=60s \ --conf spark.shuffle.service.enabled=true \ my_job.py
The catch: an executor holds shuffle files that other executors may still need to read. Removing it would lose them. That is what the external shuffle service is for — it keeps shuffle files available on the worker node after the executor is gone. Without it (or without Spark 3's shuffle tracking), dynamic allocation can cause expensive stage recomputation.
Recap
- Driver = your code + the plan + the schedulers. One per application. Single point of failure. Never holds the data.
- Cluster manager = the landlord. Grants containers; knows nothing about your query.
- Executor = a JVM that runs tasks, caches partitions and writes shuffle files. Replaceable.
- Application ⊃ Job ⊃ Stage ⊃ Task, and one task always processes exactly one partition.
- Slots = executors × cores — that is your real concurrency. Aim for 2–4× more partitions than slots.
- Client mode for interactive, cluster mode for production.
- Code inside transformations runs on executors;
print()there goes to their logs, andcollect()brings everything back to a single driver heap.
Checkpoint
print() statement inside a UDF "isn't working" — nothing appears in the terminal. What is actually happening?
df.limit(20).collect() and print there.local[*] is not a cluster at all.