Project Tungsten (Spark 1.4 129 to 2.0) keeps rows in a compact binary UnsafeRow format, partly off-heap, which cuts garbage collection and makes hashing and sorting cache-friendly. Whole-stage code generation then fuses every operator in a pipelined chain into one Java method, compiled at run time by the Janino compiler, so a row passes from scan to aggregate without virtual calls between operators. The *(2) prefix in a plan marks operators fused into generated stage 2.
genre_revenue.collect() # under AQE the final plan, and its code, exist after a run
genre_revenue.explain(mode="codegen")Found 4 WholeStageCodegen subtrees. ... /* 087 */ filter_value_3 = columnartorow_value_0.binaryEquals(((UTF8String) references[8] /* li... /* 088 */ if (!filter_value_3) continue; ... /* 106 */ // find matches from HashedRelation /* 107 */ UnsafeRow bhj_buildRow_0 = bhj_isNull_0 ? null: (UnsafeRow)bhj_relation_0.getValue(bh... ... /* 122 */ hashAgg_doConsume_0(columnartorow_value_2, columnartorow_isNull_2, columnartorow_valu...
Filter, hash-join lookup and partial aggregate sit a few lines apart in one loop. One trap: run explain(mode="codegen") before any action and, with adaptive execution on, it printed Found 0 WholeStageCodegen subtrees here, because the final plan does not exist until the query has run.