What you will be able to do
- Describe the roles of the driver, cluster manager, worker nodes and executors in a Spark application
- Trace how an action becomes a job, stages and tasks
- Explain what SparkSession.builder.getOrCreate() does and how sessions relate to the SparkContext
- Distinguish DataFrames from Datasets and say which languages support each
Key concept
Driver and executors — A Spark application has one driver process. It runs your main program, holds the SparkContext and schedules work. It also has a set of executor processes on worker nodes that run that work and keep data for the application. Everything else in Spark, from sessions to caching, sits on this split.
1.Who does the work: driver, cluster manager and executors
Every Spark application starts with a driver program. The driver runs your main() function and creates the SparkContext. To run on a cluster, the SparkContext connects to a cluster manager: Spark's own standalone manager, Hadoop YARN or Kubernetes. The cluster manager allocates resources. Spark then acquires executors on the cluster's nodes, sends them your application code and finally sends them tasks to run.
On Databricks the same split shows up as the driver node and worker nodes. The driver node keeps the state of attached notebooks, holds the SparkContext and interprets every command you run. It also runs the Apache Spark master that coordinates with the Spark executors.
| Term | Meaning |
|---|---|
| Driver program | The process running the main() function of the application and creating the SparkContext |
| Cluster manager | An external service for acquiring resources on the cluster (e.g. standalone manager, YARN, Kubernetes) |
| Worker node | Any node that can run application code in the cluster |
| Executor | A process launched for an application on a worker node, that runs tasks and keeps data in memory or disk storage across them |
| Deploy mode | In "cluster" mode the driver is launched inside the cluster; in "client" mode the submitter launches it outside the cluster |
Several consequences follow from this design. Each application gets its own executors. They stay up for the whole application and run tasks in multiple threads. This isolates applications from each other: each driver schedules its own tasks, and tasks from different applications run in different JVMs. The cost is that two Spark applications cannot share data unless one writes it to external storage.
The driver also has to accept connections from its executors for its whole lifetime, so it must be network-addressable from the worker nodes and should run close to them. Each driver serves a web UI, typically on port 4040, that shows running tasks, executors and storage usage.
Checkpoint 1 of 6· Match them up
Match each component to its role
Tap a term, then the definition that fits it.
The driver coordinates and the cluster manager hands out resources. Executors, which live on worker nodes, do the computing and hold the data.
“A process launched for an application on a worker node, that runs tasks and keeps data in memory or disk storage across them.”Source: spark.apache.org
The driver splits work using a fixed hierarchy. Calling an action such as save or collect spawns a job, which is a parallel computation made of many tasks. The job is divided into stages: sets of tasks that depend on each other, similar to the map and reduce stages in MapReduce. A task is the smallest piece, a unit of work sent to one executor. Job and stage names appear in the driver's logs.
Checkpoint 2 of 6· Check yourself
What causes Spark to spawn a new job?
A job is the parallel computation Spark launches in response to an action. It is then split into stages and tasks.
“A parallel computation consisting of multiple tasks that gets spawned in response to a Spark action (e.g. save, collect)”Source: spark.apache.org
Checkpoint 3 of 6· Exam question
Which statement correctly describes the relationship between `SparkSession` and `SparkContext` in current Spark releases?
Correct answer: A — `SparkSession` wraps a `SparkContext` internally and exposes it through `spark.sparkContext`, unifying the older `SQLContext` and `HiveContext` entry points into one object.
- A. This is correct: since Spark 2.0, `SparkSession` is the single entry point and internally holds a `SparkContext`, accessible via `spark.sparkContext`, replacing the need for separate `SQLContext` and `HiveContext` objects.
- B. This reverses the actual containment. `SparkContext` does not wrap or create a `SparkSession`; instead `SparkSession.builder.getOrCreate()` creates or reuses the `SparkContext` under the hood.
- C. Both objects live in the same driver process and share the same cluster manager connection. There is no scenario in the current architecture where they connect independently.
- D. `SparkContext` still exists and is reachable through `spark.sparkContext`; it was not removed, only wrapped so most applications no longer need to instantiate it directly.
- E. `SparkSession` does more than manage catalog metadata: it is the unified entry point for DataFrame, SQL, and (via its internal context) RDD execution, so no separate context needs manual configuration.
2.SparkSession: creating, configuring and stopping the entry point
With the RDD API, the first thing a program had to do was build a SparkConf and create a SparkContext. Only one SparkContext should be active per JVM, and you must stop() the active one before you create another. The DataFrame and Dataset APIs wrap this in a higher-level object, the SparkSession. You use it to create DataFrames, register them as tables, run SQL, cache tables and read Parquet. It still exposes the SparkContext underneath through its sparkContext property.
spark = (
SparkSession.builder
.master("local")
.appName("Word Count")
.config("spark.some.config.option", "some-value")
.getOrCreate()
)getOrCreate() first checks for a valid global default SparkSession and returns it if one exists. If none exists, it creates a new session from the builder's options and makes it the global default. If you need a new session without that check, the builder also has create(). Options set with config() go to both SparkConf and the session's own configuration. You can also change them later at runtime through spark.conf.
On a cluster, the Spark docs advise against hardcoding master in the program. Launch the application with spark-submit and receive the master value there instead. Passing local is meant for local testing and unit tests.
Checkpoint 4 of 6· Fill the gap
Which builder method makes s2 the same session as s1, so that k1 still reads v1_new?
s2 = SparkSession.builder.config("k2", "v2"). ? ()
s1.conf.get("k1") == s2.conf.get("k1") == "v1_new"
# TruegetOrCreate() returns the existing global default session, so s1 and s2 share configuration. create() would skip the check and build a new session.
Source: docs.databricks.com| Method | What it does |
|---|---|
| getActiveSession() | Returns the active SparkSession for the current thread |
| newSession() | Returns a new SparkSession with separate SQLConf, registered temporary views, and UDFs, but shared SparkContext and table cache |
| stop() | Stops the underlying SparkContext |
The newSession() row explains two facts that can look contradictory. Several sessions inside one application share a single SparkContext and table cache, so they can see the same cached data. Two separate applications each have their own SparkContext and executors, so they cannot share data without external storage. Also note that stop() stops the underlying SparkContext, not just one session.
Checkpoint 5 of 6· Exam question
A developer new to PySpark asks why their code only ever produces `DataFrame` objects and never a typed `Dataset[Row]` the way equivalent Scala examples do. Which explanation is accurate?
Correct answer: A — The typed Dataset API is a compile-time JVM feature; PySpark is dynamically typed, so it exposes only the untyped `DataFrame`, itself `Dataset[Row]` under the hood.
- A. This is correct: the compile-time type safety that Datasets provide relies on JVM generics and encoders, which Python's dynamic typing cannot support, so PySpark only surfaces the untyped `DataFrame`, itself internally represented as `Dataset[Row]`.
- B. No such `.as[T]` mechanism with Python dataclasses exists in PySpark; the typed Dataset API is not reachable from Python code by any method call.
- C. Datasets were not deprecated; they remain fully supported in Scala and Java, where compile-time typing makes the encoder-based API practical to implement.
- D. The objects are not identical across languages: Scala's `Dataset[T]` carries a compile-time type parameter and encoder, while Python's `DataFrame` carries neither, so the difference is more than naming.
- E. No such configuration flag exists in Spark or Databricks Runtime; the absence of typed Datasets in PySpark is a language-level limitation, not a toggle.
3.DataFrames and Datasets: structured APIs on one engine
The SparkSession is the entry point to programming Spark with the Dataset and DataFrame API. Both APIs belong to Spark SQL, the module for structured data processing. Unlike the basic RDD API, they tell Spark about the structure of both the data and the computation, and Spark SQL uses that information to optimize further. Whether you write SQL, Python or Scala, the same execution engine computes the result. That lets you switch between APIs freely.
A Dataset is a distributed collection of data, added in Spark 1.6. It combines the benefits of RDDs, namely strong typing and powerful lambda functions, with Spark SQL's optimized engine. You build one from JVM objects and transform it with map, flatMap, filter and similar functions. The Dataset API exists only in Scala and Java. Python has no Dataset API, but because Python is dynamic it already gets many of the same benefits, such as reading a field by name with row.columnName.
A DataFrame is a Dataset organized into named columns. Conceptually it is like a relational table or an R/Python data frame, with richer optimizations underneath. The DataFrame API is available in Python, Scala, Java and R. In Scala, DataFrame is simply a type alias of Dataset[Row]. In Java you write Dataset<Row>.
| Aspect | Dataset | DataFrame |
|---|---|---|
| What it is | A distributed collection of data, strongly typed | A Dataset organized into named columns |
| Languages | Scala and Java | Python, Scala, Java and R |
| Scala representation | Dataset[T] | Type alias of Dataset[Row] |
| Built from | JVM objects | Structured data files, Hive tables, external databases, or existing RDDs |
Underneath both APIs sit RDDs. An RDD is a collection of elements partitioned across the cluster's nodes and processed in parallel. You can ask Spark to persist an RDD in memory so it can be reused quickly, and RDDs recover automatically from node failures. Keeping data in memory while staying fault-tolerant is the property that makes caching worth studying next.
Checkpoint 6 of 6· Check yourself
Which statement about DataFrames and Datasets is correct?
In Scala a DataFrame is a Dataset of Rows. Python has no Dataset API, and DataFrame code and SQL share one engine.
“In the Scala API, DataFrame is simply a type alias of Dataset[Row].”Source: spark.apache.org
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.PySpark supports the typed Dataset API, just as Scala does.Why is that wrong?
The Dataset API is available only in Scala and Java. Python works with DataFrames.
Covered in DataFrames and Datasets: structured APIs on one engine
2.getOrCreate() always builds a fresh session, or ignores the builder's options when a session already exists.Why is that wrong?
It returns the existing global default session if one exists, and applies the builder's config options to that session.
Covered in SparkSession: creating, configuring and stopping the entry point
Sources
Every claim above is drawn from one of these pages, quoted as it was written on the date shown.
- 1.https://docs.databricks.com/aws/en/compute/configureOfficial docs
“runs the Apache Spark master that coordinates with the Spark executors”
↩︎ Who does the work: driver, cluster manager and executors“The driver node maintains state information of all notebooks attached to the compute resource.”
↩︎ Who does the work: driver, cluster manager and executors - 2.https://spark.apache.org/docs/latest/cluster-overview.htmlSecondary source
“Each application gets its own executor processes, which stay up for the duration of the whole application and run tasks in multiple threads.”
↩︎ Who does the work: driver, cluster manager and executors“However, it also means that data cannot be shared across different Spark applications (instances of SparkContext) without writing it to an external storage system.”
↩︎ Who does the work: driver, cluster manager and executors“Each job gets divided into smaller sets of tasks called stages that depend on each other”
↩︎ Who does the work: driver, cluster manager and executors“Spark applications run as independent sets of processes on a cluster, coordinated by the SparkContext object in your main program (called the driver program).”
↩︎ Key concept“A process launched for an application on a worker node, that runs tasks and keeps data in memory or disk storage across them.”
↩︎ Checkpoint“A parallel computation consisting of multiple tasks that gets spawned in response to a Spark action (e.g. save, collect)”
↩︎ Checkpoint - 3.https://docs.databricks.com/aws/en/pyspark/reference/classes/sparksession/builder/getOrCreateOfficial docs
“Gets an existing SparkSession or, if there is no existing one, creates a new one based on the options set in this builder.”
↩︎ SparkSession: creating, configuring and stopping the entry point“To create a new session without checking, use create.”
↩︎ SparkSession: creating, configuring and stopping the entry point“When an existing SparkSession is returned, the config options specified in this builder are applied to it.”
↩︎ Exam trap 2“When an existing SparkSession is returned, the config options specified in this builder are applied to it.”
↩︎ Prediction - 4.
“Options are automatically propagated to both SparkConf and SparkSession's own configuration.”
↩︎ SparkSession: creating, configuring and stopping the entry point - 5.https://spark.apache.org/docs/latest/rdd-programming-guide.htmlSecondary source
“Only one SparkContext should be active per JVM. You must stop() the active SparkContext before creating a new one.”
↩︎ SparkSession: creating, configuring and stopping the entry point“you will not want to hardcode master in the program, but rather launch the application with spark-submit”
↩︎ SparkSession: creating, configuring and stopping the entry point“Finally, RDDs automatically recover from node failures.”
↩︎ DataFrames and Datasets: structured APIs on one engine - 6.
“Stops the underlying SparkContext.”
↩︎ SparkSession: creating, configuring and stopping the entry point“The entry point to programming Spark with the Dataset and DataFrame API.”
↩︎ DataFrames and Datasets: structured APIs on one engine - 7.https://spark.apache.org/docs/latest/sql-programming-guide.htmlSecondary source
“provides the benefits of RDDs (strong typing, ability to use powerful lambda functions) with the benefits of Spark SQL’s optimized execution engine”
↩︎ DataFrames and Datasets: structured APIs on one engine“When computing a result, the same execution engine is used, independent of which API/language you are using to express the computation.”
↩︎ DataFrames and Datasets: structured APIs on one engine“The DataFrame API is available in Python, Scala, Java and R.”
↩︎ DataFrames and Datasets: structured APIs on one engine“Python does not have the support for the Dataset API.”
↩︎ Exam trap 1“In the Scala API, DataFrame is simply a type alias of Dataset[Row].”
↩︎ Checkpoint