RDDs, DataFrames and Datasets
The three Spark APIs, what Catalyst can and cannot optimise, and when an RDD still makes sense.
On this page
Show code in
Every code block on the page follows this.
You will learn
- What RDDs, DataFrames and Datasets are
- Why DataFrames are usually much faster than RDDs
- What Catalyst and Tungsten can do with a DataFrame but not an RDD
- The few cases where an RDD still makes sense
The three APIs
| API | What it holds | Languages | Optimised by Catalyst |
|---|---|---|---|
| RDD (2011) | Any objects; Spark does not know their structure | Scala, Java, Python | No |
| DataFrame (2015) | Rows with a known schema of named, typed columns | Scala, Java, Python, R, SQL | Yes |
| Dataset (2016) | Typed objects with a schema; a DataFrame is Dataset[Row] | Scala and Java only | Yes, except inside lambdas |
PySpark has RDDs and DataFrames. Datasets need compile-time types, which Python does not have, so in Python the DataFrame is the structured API. Spark Connect (Spark 3.4+) and serverless platforms do not support RDDs at all.
The same job twice
# RDD: Spark sees only "apply this Python function" revenue = (sales.rdd .filter(lambda r: r.amount > 0) .map(lambda r: (r.country, r.amount)) .reduceByKey(lambda a, b: a + b)) # DataFrame: Spark sees filter, group and sum revenue = (sales .filter(F.col("amount") > 0) .groupBy("country") .agg(F.sum("amount").alias("revenue")))
SELECT country, SUM(amount) AS revenue FROM sales WHERE amount > 0 GROUP BY country
Why the DataFrame version is faster
CatalystCatalyst: Spark's query optimizer. It rewrites your DataFrame code into a faster equivalent plan before running it. Learn more → optimisation
amount > 0, so it pushes it into the ParquetParquet: A columnar file format: values of each column are stored together, so queries read only the columns they need. Learn more → reader and reads only the country and amount columns. The RDD version reads every column of every row, then calls Python.No Python in the loop
Tungsten memory format
Code generation
Partial aggregation
When an RDD still makes sense
- Low-level control of partitioning with a custom partitioner (Scala/Java).
- Data with no tabular structure at all, processed by arbitrary code (rare; mapInPandas or a Python data source usually covers it in PySpark).
- Reading legacy code:
sc.textFile,mapPartitions,reduceByKeystill appear in older jobs and interviews. df.rdd.getNumPartitions()is still the common way to check partitionpartition: A chunk of a DataFrame's rows. Spark processes each partition as one task, so partitions decide how much work runs in parallel. Learn more → counts.
What all three share
Underneath, every DataFrame query becomes RDD operations: the physical plan is executed as RDDs of binary rows. So the core ideas are the same for all three: lazy transformations, actions that trigger jobs, lineage used to recompute lost partitions, narrow and wide dependencies, and stages split at shuffles. reduceByKey vs groupByKey is the RDD version of "aggregate before the shuffle": reduceByKey combines values on each executorexecutor: A worker process on a cluster machine that runs tasks and holds data in memory. Learn more → first; groupByKey ships every value across the network.
Common mistakes
Converting to an RDD to do something a DataFrame can do
df.rdd.map(...) loses all optimisation and adds Python serialisation.groupByKey followed by a sum
Thinking Datasets exist in PySpark
Key takeaways
- RDDs are opaque objects; DataFrames are typed tables Spark understands.
- Catalyst, Tungsten and codegen make DataFrames faster.
- PySpark RDD lambdas run in Python, like UDFs.
- Datasets are Scala/Java only; prefer DataFrames in Python.
Check yourself
3 questions1. Why can Spark push a filter into the Parquet reader for a DataFrame but not for an RDD?
Show the answer
The DataFrame filter is a known expression; an RDD filter is an opaque function. Catalyst can only optimise what it can see.
2. Which RDD operation combines values on each executor before the shuffle?
Show the answer
reduceByKey. reduceByKey does map-side combining; groupByKey sends every value.
3. Which API is not available in PySpark?
Show the answer
Dataset. Typed Datasets need compile-time types, available in Scala and Java.
Practice it
Interview problems that use this: write the PySpark, run it, and get graded on hidden tests.
Keep going
Up next · lesson 4 of 24 · 4 min readPartitions
The unit of parallelism in Spark: where partitions come from and why their number matters.
Related lessons
Catalyst and physical plansSpark internals · 3 min read
Transformations vs actionsPySpark functions · 4 min read
UDFs and pandas UDFs
Previous: Transformations vs actions
Primary sources: RDD programming guide · SQL, DataFrames and Datasets