A Query Through the Engine

How a BookNest Query Moves Through the Engine

Put together, the revenue-by-genre query took this path through Spark 4.2 129 on the 4-core workstation:

The genre query's path from DataFrame call to result
Step Component What happened to revenue by genre
1 Python API Built a logical plan; no data read
2 Catalyst Resolved columns, pruned 18 to 6, pushed the status filter
3 Physical planner Broadcast join for the 6-row catalog, two-phase aggregate
4 AQE, job 2 Broadcast the catalog from one task
5 Job 3, stage 3 4 tasks: scan, filter, join, partial sum; 22 rows to shuffle files
6 Job 4, stages 4-5 Stage 4 skipped; 1 task merged 22 rows into 6 and sorted
7 Driver Received six rows and printed them

The SQL tab's plan graph for the same run reported the row counts at each edge: 1,380,694 rows scanned, 1,229,452 after the filter and the join, 22 partial-aggregate rows into the exchange and 6 rows out. Almost all the cost sat in one fused stage of four tasks and the shuffle moved about 2 KiB: the shape of a healthy query.