Lesson 14 / 25

Reading Plans with explain() and Adaptive Query Execution

Use explain() to see how Spark will run a query and know what AQE does at run time.

Look at the plan before you tune

df.explain() prints the physical plan: how Spark will actually run the query. Read it from the bottom (data sources) to the top (final result). Look for Exchange (shuffles), join types (BroadcastHashJoin, SortMergeJoin), scan details (PushedFilters, PartitionFilters) and AdaptiveSparkPlan. In recent Spark versions Adaptive Query Execution (AQE) is on by default: it re-optimises the plan at run time using real statistics, for example coalescing small shuffle partitions, switching a sort-merge join to a broadcast join when one side turns out small, and splitting skewed partitions. explain("formatted") gives a more readable layout, and the Spark UI shows the same plan with metrics.

A broadcast join plan, run

I ran this on Apache Spark 4.0.0 (PySpark, local mode, in the official Docker image). Wrapping the small cust DataFrame in F.broadcast() makes Spark ship it to every executor (BroadcastExchange) and join with BroadcastHashJoin, avoiding a shuffle of the large side. AdaptiveSparkPlan shows AQE is active; isFinalPlan=false because nothing has executed yet.

orders.join(F.broadcast(cust), "customer").explain()

Output:

== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=false
+- Project [customer#1, order_id#0L, city#2, amount#3, order_date#5, tier#72]
   +- BroadcastHashJoin [customer#1], [customer#71], Inner, BuildRight, false
      :- Project [order_id#0L, customer#1, city#2, amount#3, cast(order_date#4 as date) AS order_date#5]
      :  +- Filter isnotnull(customer#1)
      :     +- Scan ExistingRDD[order_id#0L,customer#1,city#2,amount#3,order_date#4]
      +- BroadcastExchange HashedRelationBroadcastMode(List(input[0, string, false]),false), [plan_id=587]
         +- Filter isnotnull(customer#71)
            +- Scan ExistingRDD[customer#71,tier#72]

Watch for SortMergeJoin on huge tables

A SortMergeJoin shuffles and sorts both sides. If one side is small enough, a broadcast join is much cheaper. AQE can switch automatically, but hints (F.broadcast) make your intent explicit.

Quick check: Which plan operator indicates a shuffle?

  • Scan
  • Filter
  • Project
  • Exchange
Answer

Exchange — Exchange is where data is redistributed across partitions or machines.