Catalyst cannot know how many rows a filter keeps. Adaptive Query Execution (AQE), on by default since Spark 3.2 129 , runs the plan one shuffle stage at a time and re-optimizes the rest from real sizes. It coalesces small shuffle partitions toward spark.sql.adaptive.advisoryPartitionSizeInBytes (64 MB), switches a sort-merge join to a broadcast join when one side turns out small, and splits skewed partitions in joins. The sample data is small, so the listing scales the 10 MB broadcast threshold down to 1 MB to put the planner in its production-scale dilemma.
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "1m") # 10m scaled to sample data
returns = lines.where("status = 'returned'").groupBy("customer_id").count()
spend = (lines.where("status = 'delivered'").groupBy("customer_id")
.agg(F.sum(F.col("qty") * F.col("unit_price")).alias("spend")))
risky = spend.join(returns, "customer_id") # big spenders who also return books
risky.explain() # before running: AQE's initial plan
print("rows:", len(risky.collect()))
risky.explain() # after running: the final, re-optimized plan
spark.conf.unset("spark.sql.autoBroadcastJoinThreshold")== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=false
+- Project [customer_id#1, spend#85, count#75L]
+- SortMergeJoin [customer_id#1], [customer_id#96], Inner
...
rows: 22717
== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=true
+- == Final Plan ==
...
+- *(4) BroadcastHashJoin [customer_id#1], [customer_id#96], Inner, BuildRight, false,
false
...
: +- AQEShuffleRead coalesced
...Catalyst estimated both sides above 1 MB and chose a sort-merge join. Once the aggregation stages had run, the returns side was small, so AQE broadcast it, dropped both sorts and coalesced the 200 shuffle partitions.