Module 1 · Foundations & Architecture

Cluster architecture

Beginner 20 min read ⭐ Most important lesson

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 why collect() 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

DRIVER PROCESS one per application · your code runs here SparkSession / SparkContext Catalyst — builds & optimises the plan DAG Scheduler Task Scheduler Holds: plan, lineage, block locations, accumulators, broadcast values, and the results of collect() CLUSTER MANAGER Standalone · YARN · Kubernetes Owns the machines. Grants containers to apps. ① the driver asks it for executors WORKER NODES physical/virtual machines Worker node 1 — 16 cores, 64 GB may host 1..n executors from 1..n apps ② launches EXECUTORS — JVM processes that do all the real work Executor A 4 cores · 8 GB TASK SLOTS (1 per core) task 1 task 2 task 3 task 4 Block manager — cached partitions live here Python worker processes (PySpark only) Executor B 4 cores · 8 GB TASK SLOTS (1 per core) task 5 task 6 idle idle Shuffle files written to local disk 2 wasted slots = 2 cores you paid for Executor C 4 cores · 8 GB TASK SLOTS (1 per core) task 7 task 8 task 9 task 10 Reads its partitions from storage, computes, and reports status back to the driver every heartbeat. ③ sends tasks (serialized closures) executors run inside workers ④ results, metrics and heartbeats flow back to the driver
Memorise this picture. The driver plans and coordinates but never touches a data partition. The cluster manager only hands out machines. Executors hold the data, run the tasks, cache blocks and write shuffle files. In the picture above the cluster can run 10 tasks at once (12 slots, 2 idle).

Who does what — the responsibility split

Driver (1 per app)Executor (many per app)
Runs your main() / notebook cellRuns tasks — the actual row-by-row work
Builds and optimises the query planReads and writes data from storage
Splits the plan into stages and tasksHolds cached partitions in memory
Decides which executor gets which taskWrites shuffle files to local disk
Tracks progress, retries failed tasksSends heartbeats and metrics back
Receives collect() / show() resultsRuns the Python worker for UDFs
Single point of failure — it dies, the app diesReplaceable — Spark re-runs lost tasks elsewhere
The two classic beginner bugs — both explained by this table

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.

APPLICATION — one SparkSession. Starts when you create it, ends at stop(). JOB 1 — triggered by df.count() one job per ACTION. No action, no job. STAGE 0 — read + filter stages split at shuffle boundaries t0 t1 t2 t3 … 200 tasks STAGE 1 — aggregate after shuffle t0 t1 … tasks run only after stage 0 finishes JOB 2 — triggered by df.write.parquet() a second action = a second job, re-reading the source STAGE 0 (again!) unless you cached — see lesson 20 this is why the same read appears twice in the UI TASK — the unit that actually runs 1 task processes exactly 1 partition 1 task occupies exactly 1 slot (= 1 core) tasks in a stage are identical code, different data
Application ⊃ Job ⊃ Stage ⊃ Task, and Task ↔ Partition is one-to-one. The count you see in the Spark UI ("Tasks: 200/200") is literally the number of partitions in that stage.
TermExact meaningCreated by
ApplicationOne SparkSession and its driver + executors. Your whole script.SparkSession.builder.getOrCreate()
JobAll the work needed to satisfy one action.Every count(), show(), write, collect()
StageA set of tasks that can run without moving data between machines.The DAG scheduler, cutting at each shuffle
TaskThe work for one partition of one stage. The smallest unit Spark schedules.One per partition per stage
PartitionA chunk of the data, living on one executor. Typically 128 MB.The data source, or repartition()
CoreA CPU thread given to an executor via --executor-cores.Your submit configuration
SlotA 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.

Rule of thumb

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.

  1. You run spark-submit (or a notebook cell). A driver process starts and your Python code begins executing.
  2. SparkSession.builder.getOrCreate() runs. The driver starts its JVM, boots a SparkContext, and opens a web UI on port 4040.
  3. The driver registers with the cluster manager and requests executors — "give me 10 containers of 4 cores and 8 GB".
  4. 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.
  5. Your transformations build a plan. Nothing runs yet. This is pure driver-side bookkeeping (lesson 06).
  6. 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.
  7. 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.
  8. 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.

CLIENT MODE — driver runs where you typed the command Your laptop / edge node DRIVER close laptop = job dies Cluster executor executor executor chatty control traffic over the WAN ↔ every task update crosses the network to your laptop CLUSTER MODE — driver runs inside the cluster Your laptop submits, then exits close laptop = job keeps running ✓ Cluster DRIVER executor executor executor all traffic is local ✓
Client mode is for interactive work — notebooks, 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

ManagerWhat it isUse it whenWatch 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.
StandaloneSpark'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.
YARNHadoop'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.
KubernetesDriver 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 / DataprocManaged 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, and collect() brings everything back to a single driver heap.

Checkpoint

1 · You have 8 executors with 5 cores each. Your DataFrame has 120 partitions. How many tasks run concurrently, and how many waves does the stage take?
Slots = executors × cores = 8 × 5 = 40, so 40 tasks run at once. With 120 partitions you get 120 tasks, which is 120 / 40 = 3 sequential waves. Note this is a healthy ratio — 3× more partitions than slots keeps every core busy.
2 · A teammate says their print() statement inside a UDF "isn't working" — nothing appears in the terminal. What is actually happening?
Transformations execute on executors. Anything they print goes to that executor's log (viewable through the Spark UI's Executors tab or your cluster manager's log aggregation). To debug, either inspect executor logs, or bring a small sample to the driver with df.limit(20).collect() and print there.
3 · Your production job must survive your laptop closing and must not send control traffic over the WAN. Which submit configuration is correct?
In client mode the driver is your laptop: it coordinates every task over the WAN and the job dies with your process. Cluster mode launches the driver as a container inside the cluster, so submission is fire-and-forget and all scheduling traffic stays on the fast internal network. local[*] is not a cluster at all.
4 · Which component is responsible for deciding that a query needs three stages?
Stage boundaries are a property of the query plan, not of the hardware. On the driver, Catalyst produces a physical plan and the DAG scheduler cuts it wherever data must be redistributed across the network — every shuffle creates a new stage. The cluster manager only supplies containers and never inspects your query.