The Logical Plan

RDDs, DataFrames and the Logical Plan

A DataFrame is not data; it is a logical plan, a tree of relational operators (scan, filter, join, aggregate) that describes the result. Each transformation returns a new DataFrame with a bigger tree, and nothing is read until an action runs. explain(mode="extended") prints the tree at each step of its journey.

BookNest revenue by genre, and its plans
from pyspark.sql import functions as F
lines = spark.read.parquet("data/order_lines.parquet")
books = spark.read.parquet("data/books.parquet")
genre_revenue = (lines.where(F.col("status") == "delivered")
                 .join(books, lines.book_id == books.id)
                 .groupBy("genre")
                 .agg(F.sum(F.col("qty") * F.col("unit_price")).alias("revenue"))
                 .orderBy(F.desc("revenue")))
genre_revenue.explain(mode="extended")
Output
== Parsed Logical Plan ==
'Sort ['revenue DESC NULLS LAST], true
...
== Optimized Logical Plan ==
Sort [revenue#20 DESC NULLS LAST], true
+- Aggregate [genre#9], [genre#9, sum((cast(qty#6 as decimal(10,0)) * unit_price#7)) AS
  revenue#20]
   +- Project [qty#6, unit_price#7, genre#9]
      +- Join Inner, (cast(book_id#5 as bigint) = id#10L)
         :- Project [book_id#5, qty#6, unit_price#7]
         :  +- Filter ((isnotnull(status#4) AND (status#4 = delivered)) AND
           isnotnull(book_id#5))
         :     +- Relation [order_id#0L,customer_id#1,order_ts#2,channel#3,status#4,...]
           parquet
         +- Project [genre#9, id#10L]
...

The leading quote in 'Sort marks an unresolved plan; analysis binds every name to a typed column with an ID (genre#9). Optimization then added Project nodes that keep only three of the eighteen columns and isnotnull filters that make the join safe to push down. Underneath, a DataFrame still executes as RDDs of binary rows.