What you will be able to do
- Explain what it means for Spark to be a unified engine, with many APIs on one execution engine
- Identify the RDD abstraction and shared variables as the features of Spark's core layer
- Describe what Spark SQL adds over the basic RDD API and the ways you can interact with it
- Distinguish DataFrames from Datasets, including which languages support each
- Identify MLlib as Spark's machine learning library and name the kinds of algorithms and utilities it provides
Key concept
One engine, many APIs — Spark exposes several modules (SQL, DataFrames, pandas-style APIs, streaming, machine learning), but they all run on the same execution engine. The module you pick only changes how you express the work, not what runs it.
1.Spark as a unified analytics engine
Databricks describes Apache Spark as "a unified analytics engine for big data and machine learning." "Unified" is the key word for this objective. Spark is not a set of separate products glued together. It is one engine with several modules on top, and each module is a different way of expressing work. The PySpark documentation lists them: relational queries with Spark SQL and DataFrames, stream processing with Structured Streaming, pandas-style analysis with Pandas API on Spark, machine learning with MLlib, and graph computation with GraphX (which this exam objective does not list).
The Spark SQL guide spells out what follows from this: because the engine is shared, developers "can easily switch back and forth between different APIs based on which provides the most natural way to express a given transformation." Keep that in mind for the rest of this lesson. Every module below is one way into the same engine, and an exam question that asks what a module is "for" is really asking which way in suits a given task.
The machine learning module is MLlib (the exam guide spells it "MLib"). Databricks describes it as "the Apache Spark machine learning library", made up of "common learning algorithms and utilities, including classification, regression, clustering, collaborative filtering, dimensionality reduction, and underlying optimization primitives." The PySpark documentation adds that "MLlib is a scalable machine learning library built on Spark that provides a uniform set of APIs" that help users create and tune practical machine learning pipelines. In Python you reach it through the pyspark.ml package, and the Databricks example notebooks build classification and regression applications with the MLlib Pipelines API. Because MLlib is built on Spark, it is one more way into the same engine: it is what you choose when the task is training a model rather than querying or streaming data.
Checkpoint 1 of 5· Check yourself
A data team needs to run collaborative filtering and dimensionality reduction on a large dataset in Spark. Which module provides these?
MLlib is Spark's machine learning library. Its common learning algorithms and utilities include classification, regression, clustering, collaborative filtering and dimensionality reduction.
“common learning algorithms and utilities, including classification, regression, clustering, collaborative filtering, dimensionality reduction, and underlying optimization primitives”Source: docs.databricks.com
2.Spark Core: RDDs and shared variables
The sources here do not describe a separate "Core" module page. What they do describe, in the RDD Programming Guide, is the foundational abstraction every other layer is built on. Every Spark application has a driver program that runs the user's main function and executes parallel operations on a cluster. The main abstraction is the resilient distributed dataset (RDD): "a collection of elements partitioned across the nodes of the cluster" that can be operated on in parallel.
The guide names three RDD features. You can create an RDD from a file in a Hadoop-supported file system or from an existing collection in the driver. You can ask Spark to persist an RDD in memory so it can be reused across operations. And RDDs "automatically recover from node failures." A second core abstraction is shared variables. Normally Spark ships each task its own copy of every variable a function uses, but two kinds can be shared: broadcast variables, which cache a value in memory on all nodes, and accumulators, which can only be "added" to, such as counters and sums.
Checkpoint 2 of 5· Check yourself
Which two types of shared variables does Spark's core programming model support?
The RDD guide names exactly two shared-variable types: broadcast variables, which cache a value on every node, and accumulators, which can only be added to.
“Spark supports two types of shared variables: broadcast variables, which can be used to cache a value in memory on all nodes, and accumulators”Source: spark.apache.org
RDDs are still there underneath the higher-level APIs. A DataFrame's rdd property "Returns the content as an RDD of Row." so you can always drop down to the core layer from a DataFrame:
df = spark.range(1)
type(df.rdd)
# <class 'pyspark.core.rdd.RDD'>3.Spark SQL: structured data with extra optimization
"Spark SQL is a Spark module for structured data processing." What sets it apart from the basic RDD API is information. Spark SQL's interfaces tell Spark about the structure of both the data and the computation, and "Internally, Spark SQL uses this extra information to perform extra optimizations." An RDD is a collection of elements that Spark cannot see inside. A Spark SQL query tells Spark which columns and operations are involved, so Spark can plan around them.
There are several ways to use Spark SQL. You can run SQL queries directly. You can read data from an existing Hive installation. You can use the command line or connect over JDBC/ODBC. And you can embed SQL inside a program: "Spark SQL allows you to mix SQL queries with Spark programs." When SQL runs from inside another programming language, "the results will be returned as a Dataset/DataFrame", so a query's output feeds straight into DataFrame code.
Checkpoint 3 of 5· Check yourself
A Python program runs a SQL statement through Spark SQL. What type comes back?
SQL run from within another programming language returns its results as a Dataset/DataFrame, which is why SQL and DataFrame code mix so easily.
“When running SQL from within another programming language the results will be returned as a Dataset/DataFrame.”Source: spark.apache.org
4.DataFrames and Datasets
"DataFrames are the primary objects in Apache Spark." A DataFrame is a dataset organized into named columns. You can think of it as a spreadsheet or SQL table. It has a schema, which "defines the column names and types of a DataFrame"; its records are Row objects; and its columns can hold simple types or complex ones like arrays and maps. The entry point to all of this is SparkSession: "The entry point to programming Spark with the Dataset and DataFrame API."
A Dataset is the typed relative of the DataFrame. Added in Spark 1.6, it combines the benefits of RDDs (strong typing, powerful lambda functions) with Spark SQL's optimized execution engine. The Scala and Java references describe it as a strongly typed collection of domain-specific objects, and add that each Dataset "also has an untyped view called a DataFrame, which is a Dataset of Row." In Scala, DataFrame is simply a type alias of Dataset[Row]; in Java you write Dataset<Row>.
| API | What it is | Languages |
|---|---|---|
| DataFrame | A Dataset organized into named columns | Python, Scala, Java and R |
| Dataset | A distributed collection of data with strong typing and lambda functions | Scala and Java only |
Python does not have the Dataset API. The guide notes that because Python is dynamic, many of the benefits are already there anyway. For example, you can read a row's field by name as row.columnName.
Checkpoint 4 of 5· Match them up
Match each term to its description
Tap a term, then the definition that fits it.
The Spark API reference defines SparkSession as the entry point, a Dataset as strongly typed, and a DataFrame as a Dataset's untyped view of Rows. The RDD is the core partitioned collection.
“Each Dataset also has an untyped view called a DataFrame, which is a Dataset of Row.”Source: docs.databricks.com
Checkpoint 5 of 5· Exam question
A data science team has an existing pandas-based exploratory notebook that relies heavily on `DataFrame.plot()` for visualizing distributions. The dataset has grown too large to fit in a single machine's memory, so the team wants to keep writing pandas-style code while Spark executes it across the cluster, starting from a PySpark DataFrame named `sales_df`. Which line correctly converts `sales_df` into a `pyspark.pandas` DataFrame so plotting and other pandas methods can be used directly?
Correct answer: A — `sales_ps = sales_df.pandas_api()`, which wraps the existing distributed DataFrame in a pandas-compatible interface without collecting any data to the driver.
- A. `pandas_api()` is the current, documented way to view a native PySpark DataFrame through the pandas-on-Spark interface, and it keeps execution distributed across the cluster. This lets the team reuse pandas-style calls like `.plot()` without losing scale.
- B. `toPandas()` materializes every row on the driver as a local `pandas.DataFrame`, which is exactly the memory bottleneck the team is trying to avoid. It produces a pandas object, but not a distributed one.
- C. `ps.from_pandas()` is designed to convert an in-memory pandas object into a pandas-on-Spark object, not to wrap an already-distributed Spark DataFrame. Calling it on `sales_df` does not perform the intended conversion.
- D. `toDF()` on a DataFrame simply returns an equivalent DataFrame and adds nothing useful here, and PySpark's `DataFrame` class has no `.pandas()` method to chain afterward. This line would fail before any conversion occurred.
- E. Spark's DataFrameWriter supports formats like Parquet, JSON, CSV, and Delta, but "pandas" is not a registered output format. This line writes nothing useful and does not produce a pandas-on-Spark object.
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.PySpark gives you the typed Dataset API just like Scala does.Why is that wrong?
The Dataset API exists only in Scala and Java. Python works with DataFrames, and its dynamic nature already provides many of the Dataset benefits.
Covered in DataFrames and Datasets
2.SQL queries run on a different engine from DataFrame code, so you should pick one style for performance.Why is that wrong?
Spark SQL uses the same execution engine whichever API or language expresses the computation, so you can choose the style that reads most naturally.
Covered in Spark SQL: structured data with extra optimization
3.MLlib is a separate machine learning system that runs outside Spark and needs its own engine.Why is that wrong?
MLlib is a scalable machine learning library built on Spark itself. It is another module on the unified engine, used for model training the way Spark SQL is used for queries.
Covered in Spark as a unified analytics engine
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/pysparkOfficial docs
“Databricks is built on top of Apache Spark, a unified analytics engine for big data and machine learning.”
↩︎ Spark as a unified analytics engine“MLlib is a scalable machine learning library built on Spark that provides a uniform set of APIs”
↩︎ Spark as a unified analytics engine“Spark SQL allows you to mix SQL queries with Spark programs.”
↩︎ Spark SQL: structured data with extra optimization“A schema defines the column names and types of a DataFrame.”
↩︎ DataFrames and Datasets“MLlib is a scalable machine learning library built on Spark that provides a uniform set of APIs”
↩︎ Exam trap 3 - 2.https://spark.apache.org/docs/latest/sql-programming-guide.htmlSecondary source
“This unification means that developers can easily switch back and forth between different APIs”
↩︎ Spark as a unified analytics engine“Internally, Spark SQL uses this extra information to perform extra optimizations.”
↩︎ Spark SQL: structured data with extra optimization“In the Scala API, DataFrame is simply a type alias of Dataset[Row].”
↩︎ DataFrames and Datasets“When computing a result, the same execution engine is used, independent of which API/language you are using to express the computation.”
↩︎ Key concept“The Dataset API is available in Scala and Java. Python does not have the support for the Dataset API.”
↩︎ Exam trap 1“When computing a result, the same execution engine is used, independent of which API/language you are using to express the computation.”
↩︎ Exam trap 2“When computing a result, the same execution engine is used, independent of which API/language you are using to express the computation.”
↩︎ Prediction“When running SQL from within another programming language the results will be returned as a Dataset/DataFrame.”
↩︎ Checkpoint - 3.
“Apache Spark MLlib is the Apache Spark machine learning library”
↩︎ Spark as a unified analytics engine“common learning algorithms and utilities, including classification, regression, clustering, collaborative filtering, dimensionality reduction, and underlying optimization primitives”
↩︎ Spark as a unified analytics engine“The pyspark.ml package from Apache Spark MLlib is supported on serverless, standard, and dedicated compute.”
↩︎ Spark as a unified analytics engine - 4.
“Returns the content as an RDD of Row.”
↩︎ Spark Core: RDDs and shared variables - 5.https://spark.apache.org/docs/latest/rdd-programming-guide.htmlSecondary source
“The main abstraction Spark provides is a resilient distributed dataset (RDD), which is a collection of elements partitioned across the nodes of the cluster”
↩︎ Spark Core: RDDs and shared variables“Finally, RDDs automatically recover from node failures.”
↩︎ Spark Core: RDDs and shared variables“Spark supports two types of shared variables: broadcast variables, which can be used to cache a value in memory on all nodes, and accumulators”
↩︎ Checkpoint - 6.https://docs.databricks.com/aws/en/reference/sparkOfficial docs
“SparkSession - The entry point to programming Spark with the Dataset and DataFrame API.”
↩︎ DataFrames and Datasets“Each Dataset also has an untyped view called a DataFrame, which is a Dataset of Row.”
↩︎ Checkpoint