Enable adaptation in an isolated session
Adaptive query execution can revise supported physical-plan choices using runtime statistics. This example creates a separate SQL session, then rebuilds the tiny product DataFrame in it so the settings apply to that query. It shares the SparkContext but does not change the main session settings. collect finishes the query before explain; the six rows are not a benchmark.
Lesson reference: PySpark and Scala
Enable adaptation in an isolated session
Adaptive query execution can revise supported physical-plan choices using runtime statistics. This example creates a separate SQL session, then rebuilds the tiny product DataFrame in it so the settings apply to that query. It shares the SparkContext but does not change the main session settings. collect finishes the query before explain; the six rows are not a benchmark.
PySpark example
from pyspark.sql.functions import col, count
# Separate SQL settings; keep the shared SparkContext running.
session = spark.newSession()
session.conf.set("spark.sql.adaptive.enabled", "false")
session.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
local_products = session.createDataFrame(products.collect(), products.schema)
result = local_products.groupBy("category").agg(count("*").alias("product_count")).orderBy("category")
result.collect() # Finish the query before inspecting its plan.
result.explain("formatted")
Scala example
import org.apache.spark.sql.functions._
// Separate SQL settings; do not stop the shared SparkContext.
val session = spark.newSession()
session.conf.set("spark.sql.adaptive.enabled", "false")
session.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
val local_products = session.createDataFrame(products.rdd, products.schema)
val result = local_products.groupBy("category").agg(count("*").alias("product_count")).orderBy("category")
result.collect() // Finish the query before inspecting its plan.
result.explain("formatted")
Inspect shuffle coalescing
With AQE enabled, shuffle coalescing can combine small post-shuffle partitions based on runtime sizes. Disabling that feature does not disable every adaptive optimization. Neither setting guarantees a particular final partition count. Compare the printed plans; source partitioning, query shape, and runtime statistics can affect what changes.
PySpark example
from pyspark.sql.functions import col, count
# Separate SQL settings; keep the shared SparkContext running.
session = spark.newSession()
session.conf.set("spark.sql.adaptive.enabled", "true")
session.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
local_products = session.createDataFrame(products.collect(), products.schema)
result = local_products.groupBy("category").agg(count("*").alias("product_count")).orderBy("category")
result.collect() # Finish the query before inspecting its plan.
result.explain("formatted")
Scala example
import org.apache.spark.sql.functions._
// Separate SQL settings; do not stop the shared SparkContext.
val session = spark.newSession()
session.conf.set("spark.sql.adaptive.enabled", "true")
session.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
val local_products = session.createDataFrame(products.rdd, products.schema)
val result = local_products.groupBy("category").agg(count("*").alias("product_count")).orderBy("category")
result.collect() // Finish the query before inspecting its plan.
result.explain("formatted")
Observe a join after filtering
AQE can revise eligible join strategies and handle certain skewed joins, but it is not guaranteed to do so for every query. Here a quantity filter shrinks order_items before joining products. The tiny lookup may already be broadcast by initial planning. Inspect the completed plan rather than claiming that a particular conversion happened.
PySpark example
from pyspark.sql.functions import col, count
# Separate SQL settings; keep the shared SparkContext running.
session = spark.newSession()
session.conf.set("spark.sql.adaptive.enabled", "true")
session.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
local_products = session.createDataFrame(products.collect(), products.schema)
local_items = session.createDataFrame(order_items.collect(), order_items.schema)
result = local_items.filter(col("quantity") >= 2).join(local_products, "product_id").select("order_item_id", "product_name", "quantity").orderBy("order_item_id")
result.collect() # Finish the query before inspecting its plan.
result.explain("formatted")
Scala example
import org.apache.spark.sql.functions._
// Separate SQL settings; do not stop the shared SparkContext.
val session = spark.newSession()
session.conf.set("spark.sql.adaptive.enabled", "true")
session.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
val local_products = session.createDataFrame(products.rdd, products.schema)
val local_items = session.createDataFrame(order_items.rdd, order_items.schema)
val result = local_items.filter(col("quantity") >= 2).join(local_products, Seq("product_id")).select("order_item_id", "product_name", "quantity").orderBy("order_item_id")
result.collect() // Finish the query before inspecting its plan.
result.explain("formatted")