Skip to content
Everyone knows 4 min read 1 practice problem ↓

RDDs, DataFrames and Datasets

The three Spark APIs, what Catalyst can and cannot optimise, and when an RDD still makes sense.

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

Read first

Comfortable with these? Read on.

TL;DR An RDD is a distributed collection of opaque objects transformed by your functions. A DataFrame is a distributed table with a schema, transformed by declarative operations Spark understands. Because Spark can see what a DataFrame query does, it optimises it (Catalyst) and stores rows compactly (Tungsten). In PySpark, use DataFrames; RDDs are a low-level escape hatch.

The three APIs

APIWhat it holdsLanguagesOptimised by Catalyst
RDD (2011)Any objects; Spark does not know their structureScala, Java, PythonNo
DataFrame (2015)Rows with a known schema of named, typed columnsScala, Java, Python, R, SQLYes
Dataset (2016)Typed objects with a schema; a DataFrame is Dataset[Row]Scala and Java onlyYes, 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

PySparkSpark SQL · Revenue per country
# 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

Spark knows the filter is 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

DataFrame operations run in the JVM. Every RDD lambda in PySpark ships rows to a Python worker and back, like a UDFUDF: User-defined function: your own Python function applied to each row. Flexible but much slower than built-in functions. Learn more →.

Tungsten memory format

DataFrame rows are stored in a compact binary format, often off-heap, instead of one Java or Python object per row. Less memory and less garbage collection.

Code generation

Whole-stagestage: A group of steps Spark can run without moving data between machines. A new stage starts at every shuffle. Learn more → codegen fuses a chain of operators into one tight compiled loop per stage.

Partial aggregation

Both versions combine values before the shuffleshuffle: Moving rows between machines so that all rows with the same key end up together. Needed by joins, groupBy and sorting, and usually the most expensive step of a job. Learn more →, but the DataFrame does it on binary rows with generated code.

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, reduceByKey still 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.

Tip: if an interviewer asks about RDDs, say what they are, then explain why you would use a DataFrame, mentioning Catalyst, Tungsten and the Python serialisation cost. That is the answer they are looking for.

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

Ships every value over the network. reduceByKey or a DataFrame groupBy combines locally first.

Thinking Datasets exist in PySpark

They are Scala/Java only; in Python, a DataFrame is the structured API.

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 questions

1. 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.

Solve: Word Count →

Keep going

Up next · lesson 4 of 24 · 4 min read
Partitions
The unit of parallelism in Spark: where partitions come from and why their number matters.

Related lessons

Previous: Transformations vs actions

Primary sources: RDD programming guide · SQL, DataFrames and Datasets