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