Lesson 4 / 25
RDDs, DataFrames and Spark SQL
Choose between the low-level RDD API and the optimised DataFrame API.
Prefer DataFrames
The original abstraction is the RDD (resilient distributed dataset): a distributed collection of arbitrary objects, processed with functions like map and reduceByKey. A DataFrame is a distributed table with named, typed columns, like a pandas DataFrame or SQL table, and Spark SQL's Catalyst optimizer can analyse it, reorder operations, prune columns and choose efficient join strategies, which it cannot do for opaque Python functions on RDDs. Use DataFrames and SQL for almost everything; reach for RDDs only for low-level control or unstructured data.
The classic word count on an RDD, run
I ran this on Apache Spark 4.0.0 (PySpark, local mode, in the official Docker image). reduceByKey combines counts per word across partitions. The result is sorted here for a stable display.
lines = spark.sparkContext.parallelize(["to be or not to be", "to see or not to see"])
counts = lines.flatMap(lambda l: l.split()).map(lambda w: (w, 1)).reduceByKey(lambda a, b: a + b)
print(sorted(counts.collect()))
Output:
[('be', 2), ('not', 2), ('or', 2), ('see', 2), ('to', 4)]The same data as a DataFrame schema, run
I ran this on Apache Spark 4.0.0 (PySpark, local mode, in the official Docker image). Columns have names and types, including a proper date after the conversion, and every field is nullable.
orders.printSchema()
Output:
root |-- order_id: long (nullable = true) |-- customer: string (nullable = true) |-- city: string (nullable = true) |-- amount: double (nullable = true) |-- order_date: date (nullable = true)
Quick check: Why prefer DataFrames over raw RDDs?
- RDDs do not exist anymore
- The Catalyst optimizer can plan DataFrame operations efficiently
- DataFrames never use memory
- RDDs cannot run in parallel
Answer
The Catalyst optimizer can plan DataFrame operations efficiently — Named columns and declarative operations let Spark optimise the whole query.