From RDDs to DataFrames

From RDDs to DataFrames: Spark's Evolution

A Python lambda is a black box: Spark 129 cannot see which fields it reads, so it cannot skip columns or push filters into the reader. Spark 1.3 (2015) added the DataFrame, a table of named, typed columns planned by the Catalyst optimizer; 2.0 (2016) added whole-stage code generation, 3.0 (2020) Adaptive Query Execution and 3.2 (2021) the pandas 16,086 API on Spark.

The same revenue per book with RDD lambdas and with a DataFrameJavaScript
import time
from operator import add
from pyspark.sql import functions as F
lines = spark.read.parquet("data/order_lines.parquet")      # one row per order line
rdd_way = lambda: (lines.rdd.filter(lambda r: r.status == "delivered")
                   .map(lambda r: (r.book_id, r.qty * r.unit_price))
                   .reduceByKey(add).collect())
df_way = lambda: (lines.where("status = 'delivered'").groupBy("book_id")
                  .agg(F.sum(F.col("qty") * F.col("unit_price"))).collect())
for name, fn in [("RDD", rdd_way), ("DataFrame", df_way)] * 2:   # round 2 is warm
    t0 = time.perf_counter()
    fn()
    print(f"{name:<9} {time.perf_counter() - t0:5.2f}s")
Output
RDD       39.03s
DataFrame  7.29s
RDD       26.82s
DataFrame  1.78s

Warm, the DataFrame is about fifteen times faster on this shared 4-CPU machine. The RDD version turns every row into a Python object; the DataFrame reads four of the eight Parquet 129 columns, filters inside the scan and aggregates in compiled JVM code. Write new code against DataFrames and SQL.