ETL CompareETL Compare

University · Running it · Lesson 7

Adaptive query execution and the Spark settings that shape shuffle

ETL Compare staff · Published 15 September 2026 · 3 min read

In brief

Adaptive query execution (AQE) re-plans a Spark SQL query while it runs, using statistics from finished stages to merge small shuffle partitions, split skewed ones and switch join strategies. It is on by default. Check these settings, and record them, before any accelerator trial, because a baseline with poor settings makes any product look better than it will on a tuned job.

What does adaptive query execution change at runtime?

Without AQE, Spark fixes the physical plan before the job starts, using estimates. With AQE, each shuffle becomes a checkpoint: once a stage finishes, Spark knows the real size of its output partitions and can re-plan the rest of the query. The Spark SQL performance tuning guide describes three main features: coalescing post-shuffle partitions, splitting skewed partitions in sort-merge joins, and converting sort-merge joins to broadcast joins when one side turns out to be small.

Which settings matter most, and what are the defaults?

Configs to check

  • spark.sql.adaptive.enabled: true by default. Turns adaptive query execution on.
  • spark.sql.adaptive.coalescePartitions.enabled: true by default. Merges small post-shuffle partitions so Spark does not launch thousands of tiny tasks.
  • spark.sql.adaptive.skewJoin.enabled: true by default. Splits skewed partitions in sort-merge joins into smaller tasks.
  • spark.sql.adaptive.skewJoin.skewedPartitionFactor: 5.0 by default. A partition counts as skewed if it is larger than this factor times the median partition size...
  • spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes: 256MB by default. ...and also larger than this threshold.
  • spark.sql.shuffle.partitions: 200 by default. The starting number of partitions for shuffles in joins and aggregations.
  • spark.sql.files.maxPartitionBytes: 128 MB by default. The most data packed into one partition when reading files.
  • spark.sql.autoBroadcastJoinThreshold: 10 MB by default. Tables smaller than this are sent to every executor instead of being shuffled.

Defaults are from the Spark SQL performance tuning guide for the current release line, Spark 4.2, which was released on 14 July 2026.

What does AQE not fix?

AQE works at shuffle boundaries, so it cannot help a stage before its first shuffle, such as a scan over too many small files. Its skew handling applies to sort-merge joins, not to every skewed aggregation. And it cannot give a task more memory than its executor has: a partition that is still too large after splitting will spill. Those limits are where accelerator claims start.

How do managed platforms change the defaults?

Managed platforms ship their own Spark runtimes. AWS documents Spark performance optimizations in the Amazon EMR runtime, including adaptive query execution, dynamic partition pruning and Bloom filter joins, and states that the EMR runtime is up to 5.4x faster than open-source Apache Spark. On Databricks, the runtime and Photon carry their own behavior. Read the settings from the Environment tab of the Spark UI for the job you are testing rather than assuming open-source defaults.

Why set these before a trial?

An accelerator trial compares two runs. If the baseline runs with 200 shuffle partitions on a job that needs thousands, or with AQE switched off, the accelerator is measured against a handicapped job. Set the baseline to sensible values, record them with the event log, and keep them identical across runs. The proof of concept guide covers the rest of the method.

Related

Sources