Skip to content
Reliable Data Engineering
Lesson
Open in the interactive app

Spark Performance Tuning Deep Dive

This module is the hands-on companion to Spark internals and shuffle, spill & salting. It covers the twelve topics that come up most in senior Spark interviews, each with what it is, how to see it and how to decide.


1. Reading a query plan

df.explain("formatted")    # also: "simple", "extended" (parsed → analyzed → optimized → physical), "cost", "codegen"

Spark turns your code into a logical plan → optimised logical plan (Catalyst rules: predicate pushdown, column pruning, constant folding, join reordering with CBO) → physical plan (concrete operators and join algorithms). Read physical plans bottom-up.

== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=false
+- HashAggregate(keys=[country], functions=[sum(amount)])
   +- Exchange hashpartitioning(country, 200)                                  ← shuffle #2
      +- HashAggregate(keys=[country], functions=[partial_sum(amount)])
         +- Project [country, amount]
            +- BroadcastHashJoin [customer_id], [id], Inner, BuildRight       ← no shuffle for the join
               :- Filter isnotnull(customer_id)
               :  +- FileScan parquet orders[customer_id, amount, order_date]
               :       PartitionFilters: [order_date >= 2024-05-01]           ← partition pruning ✓
               :       PushedFilters: [IsNotNull(customer_id)]                ← pushed to the reader ✓
               :       ReadSchema: struct<customer_id:bigint,amount:double>   ← column pruning ✓
               +- BroadcastExchange HashedRelationBroadcastMode(...)           ← small side broadcast
                  +- Filter (segment = 'enterprise')
                     +- FileScan parquet customers[id, country, segment]

What to look for:

Node / fieldMeaningWorry if…
Exchange hashpartitioning(k, n)Shuffle by keyMore Exchanges than expected (e.g. repeated on the same key)
Exchange rangepartitioningGlobal sort (orderBy)You didn’t need a global order
BroadcastHashJoin … BuildRightRight side broadcastThe broadcast side is large (driver/executor memory risk)
SortMergeJoinBoth sides shuffled and sortedOne side is small enough to broadcast
BroadcastNestedLoopJoin / CartesianProductNon-equi or missing join conditionAlmost always a bug or a range join needing a rewrite
PartitionFilters: [] on a partitioned tableNo partition pruningYou filtered on a derived expression (e.g. year(ts)) instead of the partition column
PushedFiltersFilters evaluated by the file reader (row-group skipping)A filter you expected isn’t there (UDFs and casts block pushdown)
*(n) prefixesWhole-stage codegen stage nMissing on hot operators (often caused by Python UDFs)
AdaptiveSparkPlan isFinalPlan=true (after running)AQE’s final, re-optimised planCompare with the initial plan to see what AQE changed

2. Jobs, stages, tasks and the DAG

Reading the DAG in the UI: a stage with a huge input and a short duration is fine; a stage with many retries, a long tail of slow tasks, or a much larger shuffle write than input is where to dig.

Lazy evaluation trap: reusing a DataFrame in two actions recomputes its whole lineage twice unless it’s cached or written out. Seeing the same scan twice in the SQL tab is the giveaway.


3. Executor sizing: the arithmetic

Interviewers often hand you a cluster and ask for executor settings.

Cluster: 10 worker nodes, each 16 cores and 64 GB RAM (YARN or a standalone cluster).

  1. Leave room for the OS and node daemons: 1 core and about 1 GB per node → 15 cores, 63 GB usable.
  2. Cores per executor: 4-5 is the sweet spot (enough parallelism per JVM; above about 5, HDFS/S3 client throughput and GC suffer). Choose 5 → 15 / 5 = 3 executors per node.
  3. Memory per executor: 63 GB / 3 = 21 GB, minus overhead (max(384 MB, 10%), more for PySpark) → spark.executor.memory ≈ 18-19g, memoryOverhead ≈ 2-3g.
  4. Executor count: 10 nodes × 3 = 30, minus 1 for the YARN ApplicationMaster/driver → 29 executors, 145 cores in total.
  5. Per-task memory: about 19 GB × 0.6 (unified fraction) / 5 cores ≈ 2.3 GB of execution+storage per concurrent task. If shuffle partitions are ~200 MB compressed (maybe 1 GB deserialised), that fits without spill.
Anti-patternWhy it’s bad
Tiny executors (1 core each)No sharing of broadcast variables or cache across tasks; per-JVM overhead multiplied; more executors to coordinate
Fat executors (all 16 cores, 60 GB)Long GC pauses, poor I/O throughput per core, one failure loses lots of work
Ignoring overhead in PySparkPython workers and Arrow buffers exceed the overhead → containers killed

On Databricks you pick instance types instead of executor flags (one executor per worker node using all its cores). The same reasoning becomes memory per core: choose memory-optimised instances for heavy joins/aggregations, compute-optimised for CPU-bound transforms, storage-optimised (local SSD) for heavy disk cache use, and enable autoscaling with sensible bounds. Serverless removes most of this tuning.

Dynamic allocation scales executor count with the backlog of pending tasks; it needs shuffle data to survive executor removal (an external shuffle service, or shuffle tracking/decommissioning on Kubernetes).


4. Shuffle partitions (and AQE coalescing)


5. Data partitioning: in memory and on disk

In memory:

OperationShuffle?Use for
repartition(n)Yes (round-robin)Increase parallelism, even out skewed input partitions
repartition(n, "k") / repartition("k")Yes (hash)Co-locate keys before several key-based ops; control output files per partition value
repartitionByRange(n, "k")Yes (range, sampled)Sorted, non-overlapping output files (good min/max stats)
coalesce(n)NoReduce partitions cheaply before writing; can create uneven tasks and reduce upstream parallelism

Input partitions: spark.sql.files.maxPartitionBytes (128 MB) controls how files are split into read tasks; spark.sql.files.openCostInBytes makes Spark pack many small files into one task.

On disk (partitionBy on write):


6. Bucketing

Bucketing pre-shuffles a table on write into a fixed number of buckets by hash(key) mod B, and records that in the metastore.

(orders.write.bucketBy(64, "customer_id").sortBy("customer_id")
       .mode("overwrite").saveAsTable("orders_bucketed"))

7. Caching: when it helps and when it hurts

Cache when the same expensive intermediate result is used by multiple actions: iterative ML, several aggregations over one filtered dataset, interactive exploration.

Don’t cache when: it’s used once, it’s cheap to recompute, it’s bigger than cluster memory (it spills, then evicts other data), or the source is already fast (Delta with disk cache).

active = events.filter("event_date >= '2024-05-01' AND country = 'DE'").select("user_id", "event", "ts")
active.cache()
active.count()                           # materialise (cache is lazy)
daily = active.groupBy("event").count()
funnel = active.groupBy("user_id").agg(...)
active.unpersist()                       # free memory when done

8. Broadcast joins

The small side is collected to the driver, then shipped to every executor, which builds an in-memory hash table; the big side streams through without shuffling.


9. Adaptive Query Execution (AQE)

AQE re-plans the query between stages, using real statistics from completed shuffle map stages.

FeatureWhat it doesKey settings
Coalesce shuffle partitionsMerges small adjacent partitions to the advisory sizeadvisoryPartitionSizeInBytes, coalescePartitions.minPartitionSize
Switch join strategySort-merge → broadcast (or shuffled hash) when a side turns out smalladaptive.autoBroadcastJoinThreshold, maxShuffledHashJoinLocalMapThreshold
Skew joinSplits skewed partitions and replicates the matching sideskewJoin.skewedPartitionFactor (5), skewedPartitionThresholdInBytes (256 MB)
Local shuffle readerAfter switching to broadcast, reads shuffle files locally instead of over the networkon by default
Optimize skews in rebalanceSplits skewed partitions in REBALANCE hints (useful before writes)optimizeSkewsInRebalancePartitions.enabled

What AQE can’t do: fix skewed aggregations or windows, split a partition that’s skewed within a single map output, repair a bad file layout, or undo a UDF that blocks optimisation. It also only acts at shuffle boundaries, so a plan with no Exchange gets no adaptive benefit.


10. Dynamic Partition Pruning (DPP)

Problem: a fact table partitioned by date_key is joined to a filtered date dimension. The filter is on the dimension (d.is_holiday = true), so static pruning can’t know which fact partitions to read and scans all of them.

SELECT f.store_id, SUM(f.amount)
FROM sales f JOIN dim_date d ON f.date_key = d.date_key
WHERE d.is_holiday = true          -- filter is on the dimension, not the fact's partition column
GROUP BY f.store_id

DPP (on by default, spark.sql.optimizer.dynamicPartitionPruning.enabled) runs the dimension side first, collects the qualifying date_key values (reusing the broadcast when the join is a broadcast join) and injects them as a runtime partition filter on the fact scan. In the plan you see PartitionFilters: [dynamicpruningexpression(date_key IN dynamicpruning#…)] on the fact FileScan.

Conditions: the fact side must be partitioned on the join key, the join must be an equi-join, the dimension side must be filterable/selective, and the optimiser must estimate the pruning is worth it. On Delta with clustering instead of partitioning, Databricks applies a similar idea via dynamic file pruning using file-level min/max statistics.


11. Data skipping and layout (often the biggest win)

A 30-minute query that becomes 30 seconds usually came from reading 1% of the data, not from tuning executors.


12. Worked tuning scenarios

Scenario 1: A join of 1 TB orders with a 40 MB customers table takes 25 minutes. The plan shows SortMergeJoin. What do you do?

The small side is just above the 10 MB auto-broadcast threshold (estimates are often based on file size, and compressed Parquet can deserialize larger). Use a broadcast(customers) hint or raise autoBroadcastJoinThreshold to ~100 MB, select only the needed columns of customers first, and confirm that BroadcastHashJoin appears and the Exchange on the orders side disappears. Expected: the orders shuffle (about 1 TB of shuffle write) is eliminated.

Scenario 2: A daily aggregation shows 200 tasks, each spilling 3 GB, and runs for an hour.

The shuffle is about 600 GB of data in only 200 partitions (the default), so each task processes ~3 GB and spills. Increase shuffle partitions to ~4,000 (or set AQE’s initial partitions high with a 128-256 MB advisory size), and select only the needed columns before the aggregation. Check that partial aggregation is happening (partial_sum before the Exchange). Spill should drop to zero and tasks to seconds each.

Scenario 3: A query filters a 5-year partitioned fact table to last month via a join with a calendar table, but reads all partitions.

Check the plan for dynamicpruningexpression on the fact scan. If it’s missing: the fact may not be partitioned on the join key (e.g. partitioned by event_date but joined on date_key), the filter may be on a non-selective column, or DPP may be disabled. Fix by joining on the partition column, or filter the fact directly on its partition column (event_date >= add_months(current_date(), -1)), which gives static pruning.

Scenario 4: A PySpark job with a Python UDF normalising phone numbers runs 10× slower than the rest of the pipeline.

Rows are pickled to Python workers one by one, and the UDF blocks codegen and pushdown. Rewrite with built-in functions (regexp_replace, substring, when), or as a pandas UDF using vectorised string methods if the logic is complex. Verify the stage no longer shows BatchEvalPython/ArrowEvalPython (or that it shows ArrowEvalPython for the pandas UDF), and check memoryOverhead if using pandas UDFs.

Scenario 5: The same expensive filtered DataFrame feeds five different aggregations and the job reads the source five times.

Every action recomputes the lineage. Cache (or better, write to a temporary Delta table) the filtered, column-pruned DataFrame after selecting only the needed columns, materialise it once, run the five aggregations, then unpersist. Or restructure into one pass using GROUPING SETS/ROLLUP if the aggregations share the source.


Interview questions

How do you size executors for a 10-node cluster with 16 cores and 64 GB per node?

Reserve 1 core and about 1 GB per node for the OS/daemons → 15 cores, 63 GB. Use 5 cores per executor → 3 executors per node, 21 GB each; subtract overhead (~10%, more for PySpark) → ~18-19 GB heap. 30 executors minus 1 for the driver/AM = 29 executors, 145 cores. Then check the memory per concurrent task (~19 × 0.6 / 5 ≈ 2.3 GB) against the expected shuffle partition size.

repartition vs coalesce: when do you use each before writing?

coalesce(n) merges partitions without a shuffle, so it’s cheap, but it can produce uneven files and it also reduces the parallelism of the upstream stage (all work collapses into n tasks). repartition(n) (or by the partition column) shuffles but produces even partitions and keeps upstream parallelism. Before writing a partitioned table, repartition("date") (or optimized writes) controls the files per date.

Why doesn't bucketing help on Delta tables?

Delta Lake doesn’t support Spark/Hive bucketing metadata, so the bucket layout isn’t known to the planner and the shuffle still happens. On Delta, use liquid clustering/Z-order for data skipping, broadcast joins for small sides, and AQE; or keep a Hive-style bucketed Parquet table if eliminating that join shuffle is critical.

What does AQE change at runtime, and what can't it fix?

It coalesces small shuffle partitions, switches sort-merge joins to broadcast/shuffled-hash joins when a side is small after filtering, splits skewed join partitions, and uses local shuffle reads after a switch. It can’t fix skewed aggregations or windows, bad file layouts, missing data skipping, UDF overhead, or problems in stages with no shuffle boundary.

Explain dynamic partition pruning and when it doesn't kick in.

At runtime, Spark evaluates the filtered dimension side of a join, collects the qualifying join-key values and applies them as a partition filter on the partitioned fact table’s scan, so only matching partitions are read. It requires the fact table to be partitioned on the join key, an equi-join, a selective filter on the other side and the feature enabled. It does nothing for unpartitioned facts (where Databricks’ dynamic file pruning with file stats can help instead).

When is caching a bad idea?

When the data is used once, is cheap to recompute, or is larger than available memory (it spills and evicts other useful blocks); when the source is already fast (Delta with disk cache); or when you forget to unpersist in long-running sessions. Caching also hides lineage-based recovery costs: if executors die, cached partitions are recomputed anyway.