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