ML Engineer MasterClass (October) | 4 seats left

Spark · Adaptive execution
AmazonAmazon Analytics
Amazon · Tune Spark Workloads

Adaptive execution

Inspect runtime plan adaptation without relying on fixed join strategies or partition counts.

Step 1 of 6 · Learn

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")
example.pyPySpark
1. Prepare isolated SQL settingsRebuild the teaching inputs in a new session.
2. Execute the queryRuntime statistics may enable supported adaptive changes.
3. Read the completed planActual operators can vary; result values must not.
Source products6 rows
product_idproduct_namecategory
201Wireless keyboardElectronics
202Laptop standOffice
203Desk lampHome
204USB-C hubElectronics
205Travel backpackTravel
206Notebook setOffice
Enable adaptation in an isolated session
Result4 rows
categoryproduct_count
Electronics2
Home1
Office2
Travel1
The result remains two Electronics, one Home, two Office, and one Travel product. Inspect the live plan for any adaptive changes.

The example is loaded in the editor. Run it as written, then try a small change.

Runs on the Spark backend. First startup may take a moment.

Run your code to see the result.