Spark SQL exposes hundreds of configuration properties, and most tuning advice is a list of them with a suggested value. That advice ages badly. Defaults change between releases, Adaptive Query Execution (AQE) now rewrites plans at runtime, and a value that fixed one job can slow down the next twenty. A setting is only useful once you know where it is read, what it trades off, and how to confirm the running query actually saw it.

This article treats configuration as an engineering surface rather than a list. It covers where settings live and which one wins, how to read the effective value, the dozen properties that cover most real problems, the arithmetic behind scan and shuffle partition sizes, a worked tuning pass on a nightly job, and the failure modes that make configuration drift dangerous. Defaults quoted here come from the Spark 4.x documentation; check them against the version you run, because several have moved over time.

Advertisement

Where a setting lives and which one wins

A Spark setting can be supplied in four places, and later layers override earlier ones: spark-defaults.conf on the cluster, --conf flags on spark-submit, SparkSession.builder.config(...) in code before the session starts, and spark.conf.set or a SQL SET statement while the session runs. Managed platforms add their own layer, such as a cluster policy or a job definition, that renders down to one of these.

The last layer only works for runtime SQL configurations. Core settings like spark.executor.memory and spark.executor.cores size the JVMs, so they are fixed once executors launch. Static SQL configurations such as spark.sql.extensions and spark.sql.warehouse.dir are fixed when the session is created. Setting one of those at runtime either raises an error or is silently ignored by older code paths, which is why the first question when a setting appears to do nothing is which layer it was set in.

spark-defaults.confcluster-wide baselinespark-submit --confper applicationSparkSession.builder.configin code, before startspark.conf.set / SETruntime SQL confs onlyoverridden byoverridden byoverridden bySQLConf of the sessioneffective valuesCatalyst plannerbroadcast, file splitsAQE re-plannercoalesce, skew splitExecutorstasks, shuffle filesread at plan timeread per stagetask sizesStatic and core settings stop at session start; runtime SQL confs can change between queries
Where a Spark SQL setting comes from and where it is consumed. Later layers override earlier ones, and the value a query sees is the one in the session SQLConf at the moment the planner reads it.

Runtime SQL configurations are read by the planner when a query is planned, and AQE reads several of them again at each stage boundary. Changing spark.sql.shuffle.partitions between two actions in the same notebook therefore affects the second action and not the first. That is useful: you can scope a setting to one heavy step and restore it afterwards instead of raising it for the whole job.

Verify the value the query actually saw

Never assume a value took effect. Read it back from the session the query runs in, and check the plan for the behaviour the setting is supposed to cause:

# PySpark: read effective values in the running session
for key in ["spark.sql.shuffle.partitions",
            "spark.sql.adaptive.enabled",
            "spark.sql.adaptive.advisoryPartitionSizeInBytes",
            "spark.sql.autoBroadcastJoinThreshold",
            "spark.sql.files.maxPartitionBytes"]:
    print(key, "=", spark.conf.get(key))

spark.sql("SET -v").filter("key LIKE 'spark.sql.adaptive%'").show(truncate=False)

df = orders.join(customers, "customer_id").groupBy("region").sum("amount")
df.explain(mode="formatted")   # look for BroadcastHashJoin vs SortMergeJoin
df.collect()
df.explain(mode="formatted")   # after execution: AdaptiveSparkPlan isFinalPlan=true

The Spark UI confirms the rest. The Environment tab lists every property explicitly set for the application. The SQL tab shows the final adaptive plan, including AQEShuffleRead nodes with the coalesced partition count and any skew splits. If a setting is in the Environment tab but the plan shows no change, either the planner never consults it on this code path, or a later layer or hint overrides it. A /*+ BROADCAST(t) */ hint, for example, wins over the threshold.

Advertisement

The properties that matter

Most production problems come down to partition sizes, join strategy, skew and SQL semantics. These properties cover them:

PropertyDefault (4.x docs)What it controls
spark.sql.shuffle.partitions200Initial number of shuffle partitions; AQE can coalesce them down
spark.sql.adaptive.enabledtrueTurns on runtime re-planning at stage boundaries
spark.sql.adaptive.advisoryPartitionSizeInBytes64 MBTarget size when coalescing or splitting shuffle partitions
spark.sql.adaptive.coalescePartitions.parallelismFirsttrueIf true, ignore the advisory size and only respect minPartitionSize
spark.sql.adaptive.coalescePartitions.minPartitionSize1 MBFloor for coalesced partitions
spark.sql.autoBroadcastJoinThreshold10 MBLargest estimated side to broadcast; -1 disables
spark.sql.adaptive.autoBroadcastJoinThresholdunset (falls back to the above)Separate threshold for AQE decisions made on runtime sizes
spark.sql.broadcastTimeout300 sHow long to wait for a broadcast to build
spark.sql.files.maxPartitionBytes128 MBUpper bound on bytes per scan partition
spark.sql.files.openCostInBytes4 MBPadding added per file when packing scan splits
spark.sql.adaptive.skewJoin.skewedPartitionFactor5.0Partition is skewed if this many times the median...
spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes256 MB...and also larger than this
spark.sql.ansi.enabledtrue since Spark 4.0ANSI semantics: overflow and bad casts raise errors instead of returning NULL

The last row is not a performance knob, but it is the configuration change most likely to break a job during an upgrade. Spark 4.0 turned ANSI mode on by default, so a cast that used to produce NULL now fails the query. Set it explicitly in every job so the semantics do not change underneath you, and use try_cast where NULL-on-failure is what you want.

Scan partitions: the split arithmetic

Scan partitions are not simply files. For file sources, Spark packs file chunks into splits of at most maxSplitBytes, computed roughly as follows:

total     = sum(file.size + openCostInBytes for file in files)
perCore   = total / defaultParallelism          # total executor cores
maxSplit  = min(maxPartitionBytes, max(openCostInBytes, perCore))
# files larger than maxSplit are cut into chunks; small files are packed
# together until a split reaches maxSplit, each file charged openCost extra

Two consequences follow. On a large table, maxPartitionBytes dominates: 1 TB at 128 MB is about 8,000 tasks, and each task holds one split's decompressed rows in memory. Raise it to 256 MB to halve task overhead if tasks are short and memory is comfortable; lower it if scan tasks spill or run out of memory on wide rows. On a directory of tiny files, openCostInBytes decides how many files share a task. Raising it spreads small files across more tasks; lowering it packs more into each. Neither fixes the real problem of tiny files, which is a write-side issue to solve with compaction or maxRecordsPerFile and repartitioning before the write.

Shuffle partitions under AQE

With AQE on, spark.sql.shuffle.partitions is the starting partition count for every shuffle, and AQE merges adjacent small partitions after the map side finishes and real sizes are known. That changes how to set it. Too low, and AQE cannot help, because it only coalesces and never splits ordinary partitions: 200 partitions for a 2 TB shuffle means 10 GB per reducer and heavy spilling. Too high, and the map side writes many tiny blocks. AQE only merges partitions smaller than the advisory size, so a practical rule is to overprovision: size the initial count so partitions start at or below the advisory size, set advisoryPartitionSizeInBytes to the size you want each reducer to handle, such as 128 MB, and let AQE merge up to it. spark.sql.adaptive.coalescePartitions.initialPartitionNum lets you set this starting count for AQE separately.

The setting most people miss is parallelismFirst. With its default of true, AQE coalesces only as far as minPartitionSize (1 MB) allows while keeping the cluster's parallelism, and ignores the 64 MB advisory target. The documentation recommends setting it to false on a busy cluster so the advisory size is respected. In practice that means fewer, larger reduce tasks, less scheduling overhead and fewer output files. On a dedicated cluster where only this job runs, leaving it true can finish faster because every core stays busy.

Join strategy and skew settings

Broadcast joins avoid shuffling the large side, which is why the threshold is the most often raised setting. The risk is in estimates. At planning time Spark uses table statistics or file sizes, and a filter or aggregate in between can make the estimate wrong in either direction. If the real build side is far larger than estimated, the driver collects it, serialises it and ships it to every executor. That can exhaust driver memory, hit broadcastTimeout, or put several gigabytes of hash table in every executor.

Two settings separate the planning-time and runtime decisions. Keep spark.sql.autoBroadcastJoinThreshold conservative, because it acts on estimates. Set spark.sql.adaptive.autoBroadcastJoinThreshold higher if you want AQE to switch a sort-merge join to a broadcast join after a stage has run and the real size is known. AQE then uses a local shuffle reader to avoid a second shuffle. Raising broadcastTimeout is rarely the fix; a broadcast that takes five minutes to build is usually one that should not be a broadcast.

Skew handling works only for sort-merge joins whose partitions cross both skew conditions: larger than five times the median and larger than 256 MB by default. If one key holds 40 GB, AQE can split that partition into advisory-size pieces and replicate the matching other side, but a skewed aggregation or a skewed window is not covered and still needs salting or a two-phase aggregate.

Worked example: tuning a nightly join

Consider a nightly job that joins a 1.2 TB Parquet fact table (after partition pruning) with a 3 GB customer dimension, aggregates by region and day, and writes the result. It runs on 50 executors with 8 cores each, 400 cores in total, under Spark 4.x defaults. The figures below are illustrative, but the reasoning is the same on any job.

  1. Measure first. The SQL tab shows a sort-merge join with two shuffles. The fact shuffle writes about 900 GB into 200 partitions, about 4.5 GB each. Reduce tasks spill to disk and the slowest takes 38 minutes. One partition is 31 GB because a single default customer ID holds a large share of rows.
  2. Fix shuffle sizing. Start partitions at or below 64 MB: 900 GB / 64 MB is about 14,000. Set spark.sql.shuffle.partitions=14000, an advisory size of 128 MB and parallelismFirst=false so AQE merges toward 128 MB, leaving roughly 7,000 reducers. Spilling disappears from all but the skewed partition.
  3. Let skew handling see the hot key. The 31 GB partition is now hundreds of times the median and above 256 MB, so AQE splits it. Confirm a skew-join marker in the final plan. If it is missing, check whether the join is actually an outer join on the side that cannot be split.
  4. Consider the broadcast. 3 GB is far too large for the 10 MB planning threshold, and broadcasting it to 400 cores would cost about 3 GB of memory per executor once deserialised. Leave it as a sort-merge join. If a filter shrinks the dimension to 200 MB at runtime, an AQE-only threshold of 256 MB would let the runtime planner switch to a broadcast.
  5. Fix the output. The aggregate by region and day is small. With parallelismFirst=false, AQE merges its shuffle into a few partitions instead of thousands of 1 MB ones, so the write produces a few files rather than thousands. Add maxRecordsPerFile if downstream readers need bounded file sizes.
# Scope the heavy settings to this job's write step, then restore.
heavy = {
    "spark.sql.shuffle.partitions": "14000",
    "spark.sql.adaptive.advisoryPartitionSizeInBytes": "128m",
    "spark.sql.adaptive.coalescePartitions.parallelismFirst": "false",
    "spark.sql.adaptive.autoBroadcastJoinThreshold": str(256 * 1024 * 1024),
    "spark.sql.ansi.enabled": "true",
}
saved = {k: spark.conf.get(k, None) for k in heavy}
for k, v in heavy.items():
    spark.conf.set(k, v)
try:
    result.write.mode("overwrite").partitionBy("day").parquet(out_path)
finally:
    for k, v in saved.items():
        spark.conf.unset(k) if v is None else spark.conf.set(k, v)

The job's runtime is now dominated by the scan, which is where it should be.

A tuning process that survives upgrades

Treat configuration like code. Keep each job's settings in version control with the job, not in a shared cluster default that silently changes forty jobs at once. Change one setting at a time and compare the same input. Run-to-run noise is several percent. Prefer stage-level metrics such as shuffle bytes, spill and maximum task time over total wall-clock time, because they point at the cause.

Pin settings whose defaults have changed or may change: ANSI mode, AQE, and any threshold you depend on. On upgrade, diff the effective configuration of a representative job on the old and new versions using SET -v output, and read the SQL migration guide.

Failure modes and trade-offs

  • Global defaults set for one job. Raising the cluster-wide broadcast threshold to 1 GB to fix one join causes driver out-of-memory errors in an unrelated job whose estimate was wrong. Scope settings per job or per step.
  • Setting a static conf at runtime. spark.conf.set on a static SQL conf fails or has no effect, and the job runs with the old value. Verify with spark.conf.get and the Environment tab.
  • Tuning around stale statistics. If table statistics are old, the planner's estimates are wrong and no threshold is correct. Refresh statistics or rely on AQE's runtime sizes.
  • Hiding a data problem. Bigger partitions and longer timeouts can mask a duplicated join key or an exploding join. Check row counts before and after each join before you tune.

The underlying trade-off is between parallelism and per-task overhead. Many small tasks keep every core busy but pay scheduling, metadata and small-file costs. Few large tasks are cheap to schedule but spill, straggle and hold memory longer. The right point depends on cluster size, cluster load and what reads the output.

For the mechanics behind the settings, see Adaptive Query Execution, broadcast joins, data skew, partition pruning and coalesce versus repartition.

What to do next

  1. Pick your three most expensive jobs and record, from the SQL tab, shuffle bytes, spill, maximum task time and output file count per stage.
  2. Print the effective values of the properties in the table above from inside each job, and note which layer set each one.
  3. Move job-specific settings out of cluster defaults and into the job's code or submit configuration, under version control.
  4. Size spark.sql.shuffle.partitions so initial partitions start at or below the advisory size, set the advisory size to your target reducer size, and set parallelismFirst to false on shared clusters.
  5. Keep the planning-time broadcast threshold conservative and use the AQE threshold for runtime switches.
  6. Pin spark.sql.ansi.enabled explicitly and test the job under the value you will run in production.
  7. Change one setting at a time, compare stage metrics on the same input, and keep the results with the job.
Key takeaway: Spark SQL settings come from four layers, and only runtime SQL configurations can change while a session runs. Read the effective value back and confirm the plan changed before trusting any setting. Overprovision initial shuffle partitions, set the advisory size to the reducer size you want, and let AQE merge up to it, with parallelismFirst off on shared clusters. Keep the planning-time broadcast threshold conservative, pin ANSI mode explicitly, scope heavy settings to the job that needs them, and change one setting at a time.