Menu

Apache Spark course · Lesson 1 of 5

Apache Spark Architecture: Driver, Executors and Cluster Managers

How a Spark application is built: the driver and executors, standalone, YARN and Kubernetes cluster managers, deploy modes, and SparkSession versus SparkContext.

  • Beginner
  • Pillar guide
  • 14 min read
  • Updated Oct 2026
On this page
  1. Following along
  2. Spark architecture: the driver and executors
  3. What it is
  4. How the pieces talk
  5. PySpark adds Python worker processes
  6. Failure handling
  7. Pitfalls
  8. In interviews
  9. Cluster managers: standalone, YARN and Kubernetes
  10. What it is
  11. Client and cluster deploy mode
  12. How each one differs in practice
  13. Pitfalls
  14. In interviews
  15. SparkSession and SparkContext
  16. What they are
  17. How getOrCreate and newSession behave
  18. Runtime versus static configuration
  19. Spark Connect
  20. Pitfalls
  21. In interviews
  22. How the rest of the course fits
  23. Practice questions
  24. Key takeaways

Every Spark application, whether it runs on a laptop or on a thousand-node cluster, has the same shape: one driver that plans the work, many executors that do it, and a cluster manager that hands out the machines. Knowing which process does what explains most of Spark’s behaviour, from why collect() can crash a job to why a lost executor only costs you a few retried tasks. This lesson is the map for the rest of the Spark course.

Following along

The examples run in local mode, where the driver and a single executor share one JVM on your machine and local[2] gives that executor two task slots. The architecture is the same as on a cluster, only smaller.

from pyspark.sql import SparkSession

spark = (SparkSession.builder
         .master("local[2]")
         .appName("architecture-demo")
         .config("spark.sql.shuffle.partitions", "8")
         .getOrCreate())
sc = spark.sparkContext

print("Spark version:", spark.version)
print("master:", sc.master)
print("app name:", sc.appName)
print("default parallelism:", sc.defaultParallelism)
Spark version: 4.2.0
master: local[2]
app name: architecture-demo
default parallelism: 2

Spark architecture: the driver and executors

What it is

A Spark application is one driver process plus a set of executor processes that belong only to it. The driver is the brain; executors are the hands. Nothing is shared between two applications except the cluster’s hardware, so one application’s cached data is invisible to another.

Component Where it runs What it does
Driver Your spark-submit machine (client mode) or a container on the cluster (cluster mode) Runs your main program, creates the SparkContext, turns code into a plan, splits it into jobs, stages and tasks, schedules tasks, tracks where cached blocks and shuffle files live, receives results of collect()
Executor A JVM on a worker node (a YARN container or a Kubernetes pod) Runs tasks in threads, one per task slot; stores cached partitions and shuffle files in its block manager; reports status and metrics back to the driver
Cluster manager Separate service (standalone master, YARN ResourceManager, Kubernetes API server) Starts executor processes when the driver asks and gives them CPU and memory; knows nothing about your query

How the pieces talk

Inside the driver, three schedulers cooperate. The DAG scheduler cuts the plan into stages at shuffle boundaries. The task scheduler turns each stage into one task per partition and sends tasks to free executor slots, preferring an executor that already holds the data (data locality). The scheduler backend is the adapter for the chosen cluster manager. Executors run each task in a thread, so an executor with spark.executor.cores=4 runs up to four tasks at once.

Executors talk to each other directly during a shuffle: a reduce task fetches map output files from whichever executors wrote them. The driver is not in the data path for shuffles, only for results you bring back with collect(), take() or toPandas().

PySpark adds Python worker processes

Spark’s engine is written in Scala and runs on the JVM. When you use DataFrame functions, all work stays in the JVM. When you run Python code on the data (an RDD lambda, a Python UDF), each executor starts Python worker processes and streams rows to them. You can see this: the code below records the operating-system process ID inside each task and compares it with the driver’s.

import os
from pyspark import TaskContext

driver_pid = os.getpid()

def where_am_i(rows):
    ctx = TaskContext.get()
    yield (ctx.stageId(), ctx.partitionId(), os.getpid() != driver_pid)

for row in sc.parallelize(range(8), 4).mapPartitions(where_am_i).collect():
    print(row)
(0, 0, True)
(0, 1, True)
(0, 2, True)
(0, 3, True)

Four partitions became four tasks of stage 0, and every task ran in a Python worker, not in the driver’s Python process. That is why a variable you change inside a task does not change on the driver (the RDD lesson covers closures in detail).

Failure handling

  • A task fails: the driver retries it, by default up to spark.task.maxFailures (4) attempts before the stage and job fail.
  • An executor dies: its running tasks are rescheduled elsewhere. Cached partitions and shuffle files it held are lost, and Spark recomputes them from the lineage (the recorded chain of transformations).
  • The driver dies: the application ends. All executors are released and cached data is gone. On YARN in cluster mode, the ResourceManager can start a new application attempt (spark.yarn.maxAppAttempts), but that attempt starts the job again from the beginning.

Pitfalls

  • Treating the driver as a worker. collect() and toPandas() pull the whole result into driver memory. Limit, aggregate or write to storage instead. spark.driver.maxResultSize caps how much a single action can return.
  • Forgetting parallelism is slot-bound. Tasks running at once = executors × cores per executor (divided by spark.task.cpus, normally 1). 1,000 partitions on 20 slots run in 50 waves.
  • Assuming executors share state. Each executor has its own memory; a Python global set in one task is not visible in another.

In interviews

“Walk me through what happens when you submit a Spark job” is the classic opener. A strong answer goes: the driver starts and creates a SparkContext; it asks the cluster manager for executors; your code builds a logical plan; an action triggers optimisation and the DAG scheduler splits the job into stages at shuffles; the task scheduler sends one task per partition to executors; executors run tasks, write shuffle files, and send results or status back; the driver reports success when the last stage finishes.

Cluster managers: standalone, YARN and Kubernetes

What it is

The cluster manager is the resource broker. The driver says “I need ten executors with four cores and 8 GB each”, and the cluster manager decides which machines run them. Spark supports three cluster managers. In Spark 4.0, support for Apache Mesos was removed (it had been deprecated since Spark 3.2), so standalone, YARN and Kubernetes are the full list.

Cluster manager Master URL How executors run Typical use
Local (not a cluster) local, local[4], local[*] Driver and one executor in one JVM; the number in brackets is the number of task threads Development, tests, this course
Standalone spark://host:7077 Spark’s own master and worker daemons start executor JVMs Small dedicated Spark clusters, simple setups
YARN yarn (cluster location comes from HADOOP_CONF_DIR or YARN_CONF_DIR) Executors run in YARN containers on NodeManagers; an ApplicationMaster negotiates them Hadoop clusters, Amazon EMR on EC2, on-premises data lakes
Kubernetes k8s://https://<api-server>:<port> Driver and executors are pods built from a container image Cloud-native platforms, shared clusters running other workloads

Managed platforms such as Databricks run Spark on their own resource layer. The driver and executor model is unchanged; you just do not choose the cluster manager.

Client and cluster deploy mode

--deploy-mode decides where the driver runs.

  • Client mode (the default): the driver runs in the spark-submit process on the machine you launched from. Interactive shells and notebooks use client mode. The driver must stay reachable by executors for the whole job, so launching from a laptop over a VPN is fragile.
  • Cluster mode: the cluster manager starts the driver inside the cluster (a YARN container or a Kubernetes pod). spark-submit can exit, and the driver sits close to the executors. Use it for scheduled production jobs.
# Standalone cluster, cluster deploy mode
spark-submit --master spark://spark-master:7077 --deploy-mode cluster \
  --executor-memory 8g --executor-cores 4 --total-executor-cores 40 \
  jobs/daily_sales.py

# YARN, cluster deploy mode
spark-submit --master yarn --deploy-mode cluster \
  --num-executors 10 --executor-memory 8g --executor-cores 4 \
  --conf spark.yarn.maxAppAttempts=2 \
  jobs/daily_sales.py

# Kubernetes, cluster deploy mode (the script is inside the image, hence local://)
spark-submit --master k8s://https://k8s-api.example.internal:6443 --deploy-mode cluster \
  --name daily-sales \
  --conf spark.kubernetes.container.image=registry.example.internal/spark-py:4.2.0 \
  --conf spark.kubernetes.namespace=data-jobs \
  --conf spark.kubernetes.authenticate.driver.serviceAccountName=spark \
  --conf spark.executor.instances=10 \
  local:///opt/app/jobs/daily_sales.py

How each one differs in practice

  • Standalone is the simplest: start a master (sbin/start-master.sh) and workers, and Spark manages everything. It has no multi-tenant queues beyond basic limits, so it suits clusters dedicated to Spark.
  • YARN adds queues, capacity limits and Kerberos-secured HDFS access. Each executor is a container whose size is executor memory plus memory overhead; if the process exceeds the container limit, YARN kills it, which shows up as a lost executor. The memory and tuning lesson works through the arithmetic.
  • Kubernetes packages Spark and your dependencies in a container image. The driver pod needs a service account allowed to create executor pods. Shared clusters let Spark sit next to other services, with namespaces and resource quotas as the isolation tools. The deployment lesson covers it in depth.

Pitfalls

  • Mixing up --master and --deploy-mode. The master picks the cluster manager; deploy mode picks where the driver runs. local has no cluster mode.
  • Expecting spark.executor.cores to mean the same thing everywhere. On YARN and Kubernetes it defaults to 1 core per executor; on standalone an executor takes all free cores on its worker unless you set it. Always set it explicitly in production.
  • Client mode behind a firewall. Executors must open connections back to the driver; blocked ports produce executors that register and then time out.

In interviews

Expect “Which cluster managers have you used and how do they differ?” and “Client versus cluster mode?”. Name the three managers, note that Mesos was removed in Spark 4.0, and explain deploy mode in terms of where the driver lives and what happens if the submitting machine disconnects.

SparkSession and SparkContext

What they are

SparkContext is the original entry point (Spark 1.x). It represents the connection to the cluster: it owns the schedulers, talks to the cluster manager, and creates RDDs, broadcast variables and accumulators. There is one active SparkContext per JVM.

SparkSession (Spark 2.0 onward) is the unified entry point for DataFrames and SQL. It wraps a SparkContext and adds the SQL pieces: the catalog of tables and views, the SQL parser, the Catalyst optimiser, and runtime SQL configuration. You build a SparkSession; you reach the SparkContext through spark.sparkContext when you need RDD-level features.

How getOrCreate and newSession behave

getOrCreate() returns the existing session if there is one. newSession() creates a second session that shares the SparkContext (and therefore the executors and cached data) but has its own SQL configuration and temporary views. Notebooks and servers use this to isolate users.

same = SparkSession.builder.getOrCreate()
print("getOrCreate returns same session:", same is spark)

other = spark.newSession()
print("newSession is a different session:", other is not spark)
print("...but shares the SparkContext:", other.sparkContext is sc)

other.conf.set("spark.sql.shuffle.partitions", "3")
print(spark.conf.get("spark.sql.shuffle.partitions"), other.conf.get("spark.sql.shuffle.partitions"))

spark.range(5).createOrReplaceTempView("nums")
print("nums visible in other session:", other.catalog.tableExists("nums"))
getOrCreate returns same session: True
newSession is a different session: True
...but shares the SparkContext: True
8 3
nums visible in other session: False

A builder .config() call on getOrCreate() when a session already exists does not restart anything. Static settings are fixed when the SparkContext starts.

Runtime versus static configuration

SQL settings such as spark.sql.shuffle.partitions can change at any time with spark.conf.set. Settings that shape processes (executor memory, cores, the cluster manager) are read once at start-up. Trying to change one at runtime raises an error in Spark 4:

try:
    spark.conf.set("spark.executor.memory", "4g")
except Exception as e:
    print(type(e).__name__, str(e).splitlines()[0])

print(spark.conf.isModifiable("spark.sql.shuffle.partitions"),
      spark.conf.isModifiable("spark.executor.memory"))
AnalysisException [CANNOT_MODIFY_CONFIG] Cannot modify the value of the Spark config: "spark.executor.memory".
True False

Set static values with spark-submit --conf, the builder before the first session is created, or spark-defaults.conf. Precedence, highest first: values set in code on the builder, then spark-submit flags, then spark-defaults.conf.

Spark Connect

Spark Connect (introduced in 3.4) splits the client from the server: your program uses a thin client that sends unresolved plans to a remote Spark driver over gRPC, using a sc://host:port address with SparkSession.builder.remote(). The client has a SparkSession but no SparkContext, so RDDs, sparkContext.broadcast and accumulators are unavailable. If code calls spark.sparkContext, it fails under Spark Connect. Prefer DataFrame APIs in new code so it runs both ways.

Pitfalls

  • Creating a SparkContext directly in new code. Use the SparkSession builder; it creates the context for you.
  • Calling spark.stop() in shared notebooks. It stops the SparkContext for every session that shares it.
  • Expecting a temporary view to cross sessions. Use a global temporary view (global_temp database) or a catalog table.

In interviews

“What is the difference between SparkSession and SparkContext?” Answer: SparkContext is the low-level connection to the cluster and the RDD entry point; SparkSession is the higher-level entry point for DataFrames and SQL that wraps a SparkContext; one SparkContext can serve many sessions, each with its own SQL config and temp views. Mentioning Spark Connect, where there is no SparkContext at all, shows current knowledge.

How the rest of the course fits

Each later lesson zooms into one part of this picture:

  1. RDD fundamentals: the low-level data abstraction the engine is built on.
  2. Execution model: how actions become jobs, stages and tasks.
  3. Partitions, shuffles and skew: where most cost comes from.
  4. Caching, broadcast and accumulators: sharing data between driver and executors.
  5. Catalyst and Tungsten: how the driver optimises a query.
  6. Adaptive Query Execution: re-planning at run time.
  7. Memory and executor tuning: sizing executors.
  8. Deployment and monitoring: Kubernetes and the History Server.

Practice questions

What happens, step by step, between spark-submit and the job finishing?

spark-submit starts the driver (locally in client mode, on the cluster in cluster mode). The driver creates a SparkContext, which registers with the cluster manager and requests executors. Your code builds a logical plan without running anything. An action triggers optimisation into a physical plan; the DAG scheduler splits it into stages at shuffle boundaries; the task scheduler launches one task per partition on free executor slots, preferring local data. Executors run tasks, write shuffle output for the next stage, and report back. When the final stage finishes, results go to the driver or to storage, and when the program ends the executors are released.

An executor is lost halfway through a job. What does Spark do?

The driver marks the executor’s running tasks as failed and reschedules them on other executors. Any cached partitions and shuffle files on that executor are gone, so stages whose output is missing are rerun for the lost partitions only, using lineage. The job continues unless a single task fails spark.task.maxFailures times (4 by default). Losing the driver, in contrast, ends the application.

When would you choose cluster deploy mode over client mode?

For scheduled or long-running production jobs: the driver runs inside the cluster, close to the executors, and the job does not depend on the submitting machine staying connected. Client mode suits interactive work (shells, notebooks) where you want results printed locally, and debugging where you want the driver logs in your terminal.

Your notebook calls spark.conf.set(“spark.executor.memory”, “16g”) and nothing changes. Why?

Executor memory is a static setting read when the SparkContext starts and executors are launched. In Spark 4 the call raises CANNOT_MODIFY_CONFIG, as shown above. Set it with spark-submit --conf, in spark-defaults.conf, or on the builder before the first session is created, and restart the application.

What is the relationship between SparkSession and SparkContext, and when do you still need SparkContext?

A SparkSession wraps one SparkContext and adds the SQL catalog, parser, optimiser and runtime SQL config. Several sessions can share one SparkContext. You still need the SparkContext for RDD operations, broadcast variables, accumulators and job groups. Under Spark Connect there is no SparkContext on the client, so such code must be rewritten with DataFrame APIs.

Which cluster managers does Spark 4 support, and what replaced Mesos?

Standalone, YARN and Kubernetes (plus local mode for development). Mesos support was deprecated in 3.2 and removed in 4.0; Kubernetes is the usual choice for container-based clusters that Mesos used to serve.

Key takeaways

  • The driver plans and schedules; executors run tasks, cache data and hold shuffle files; the cluster manager only allocates resources.
  • Parallelism equals total executor cores (divided by spark.task.cpus); each partition becomes one task.
  • Spark 4 supports standalone, YARN and Kubernetes; Mesos was removed in 4.0.
  • Deploy mode decides where the driver runs: client for interactive work, cluster for production.
  • SparkSession is the entry point for DataFrames and SQL and wraps a SparkContext; sessions share executors but not SQL config or temp views.
  • Static settings (memory, cores) are fixed at start-up; runtime SQL settings can change with spark.conf.set.

By Data Career Hub Editorial · Last reviewed Oct 2026 · Python examples run on PySpark 4.2.0 in local mode (local[2]). spark-submit commands for standalone, YARN and Kubernetes are not executed here (no cluster available); their options were checked against the Spark 4.2 documentation.

Progress is saved in this browser only. No account needed.

Search
Filter by type