Catalyst Optimizer

The Catalyst Optimizer's Phases

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

The physical plan in formatted mode
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.