Catalyst is the query compiler behind DataFrames and Spark 129 SQL. Its four phases are analysis (resolve names and types), logical optimization (rules such as predicate pushdown, column pruning and constant folding), physical planning (choose join and aggregate operators, insert Exchange nodes where data moves) and code generation (Tungsten and Code Generation).
genre_revenue.explain(mode="formatted")Output
== Physical Plan ==
AdaptiveSparkPlan (14)
...
+- Exchange (10)
+- HashAggregate (9)
...
(1) Scan parquet
Output [4]: [status#4, book_id#5, qty#6, unit_price#7]
PushedFilters: [IsNotNull(status), EqualTo(status,delivered), IsNotNull(book_id)]
ReadSchema: struct<status:string,book_id:int,qty:int,unit_price:decimal(6,2)>
...ReadSchema lists four of the eight columns, and PushedFilters hands status = 'delivered' to the Parquet 129 reader, which skips row groups whose statistics rule it out. The partial HashAggregate (9) before the shuffle meant only 22 rows crossed Exchange (10) instead of 1.2 million.