Spark Performance Tuning

Apache Spark jobs that are correct but slow waste cluster hours and delay pipelines. Tuning Spark isn't about flipping dozens of configuration flags. It's a diagnostic process: find where time actually goes (usually a few expensive stages), understand why (shuffle volume, skew, spills, bad joins, small files, Python UDFs), and fix the root cause in the code or data layout before touching cluster settings.

Modern Spark (3.x and later) does a lot automatically through the Catalyst optimizer and Adaptive Query Execution (AQE), so the biggest wins usually come from writing optimizable code, laying out data well, and sizing executors sensibly.

TL;DR

Quick Example

A tuned session configuration and join hints:

Core Concepts

Reading the Spark UI

Spill means a task ran out of execution memory and wrote intermediate data to disk. It's a strong signal of partitions that are too big, or too little memory per core.

Executor Sizing and Memory

Each executor JVM has a heap (spark.executor.memory) split between execution memory (shuffles, joins, sorts, aggregations) and storage memory (cache), which share a unified region (spark.memory.fraction). Off-heap overhead (memoryOverhead) covers native memory, Python workers, and Arrow buffers, so raise it for PySpark and pandas UDF-heavy jobs.

Rules of thumb:

Join Strategies

Hints: F.broadcast(df), or /*+ BROADCAST(t) /, /+ MERGE /, /+ SHUFFLE_HASH */ in SQL. Broadcasting a table that's too large can OOM executors or the driver, so size-check it first. Dynamic partition pruning filters large partitioned fact tables using dimension filters at runtime.

Adaptive Query Execution

AQE re-plans after each shuffle stage, using real statistics: it coalesces partitions, converts sort-merge to broadcast joins when a side turns out small, and splits skewed partitions. Keep it enabled, and tune the advisory size and skew thresholds rather than hand-setting partition counts for every job. See Spark partitioning.

I/O Efficiency

Caching

cache() or persist() is worthwhile when a DataFrame is expensive to compute and reused by multiple actions (iterative ML, multiple outputs from one intermediate). It's wasteful when used once, since caching itself costs memory and time. Prefer MEMORY_AND_DISK, check the Storage tab, and unpersist() when done. For cross-job reuse, write an intermediate table instead.

Python-Specific Costs

Python UDFs serialize every row between the JVM and Python workers. Replace them with built-in functions, or with vectorized pandas UDFs using Arrow. See PySpark.

Tuning Workflow

  1. Reproduce with realistic data volume, and record the baseline runtime, cost, and shuffle bytes.
  2. Find the dominant stage(s) in the Spark UI.
  3. Diagnose: skew, spill, huge shuffle, full scan, small files, UDF, or a bad join.
  4. Fix in code and data first: pruning, broadcast, pre-aggregation, salting, compaction, and replacing UDFs.
  5. Then tune resources: executor size, memory overhead, partition sizes, dynamic allocation.
  6. Re-measure, and keep the changes that help. Document the reasons.

Best Practices

Optimize the Biggest Stage First

Typically 1–2 stages account for most of the runtime. Tuning anything else yields little.

Right-Size Instead of Over-Provisioning

Doubling executors rarely halves runtime when the bottleneck is skew or a single huge task. Use dynamic allocation, and measure cost per run, not just speed.

Keep Plans Optimizable

Built-in functions, early filters, and explicit schemas let Catalyst do its job. UDFs and complex Python logic block pushdown and code generation.

Monitor Over Time

Data grows. Track job duration, shuffle size, and input size per run, and alert on regressions before SLAs slip. See observability.

Common Mistakes

Tuning Configs Before Reading the UI

Changing shuffle.partitions or memory blindly often masks the real problem (skew, a bad join). Diagnose first.

Caching Everything

Caching many DataFrames that are used once fills memory, causes evictions and spills, and slows jobs down. Cache selectively.

Broadcasting Large Tables

A forced broadcast of a multi-gigabyte table can crash the driver, which collects it first, or the executors. Check sizes, and let AQE choose when unsure.

FAQ

How do I find out why my Spark job is slow?

Open the Spark UI, find the stages consuming the most time, and inspect task-duration distribution, shuffle sizes, and spill. Large gaps between median and max task times indicate skew, spill indicates memory pressure, and large shuffle sizes indicate expensive joins or aggregations.

How many cores and how much memory should each executor have?

A common starting point is 4–5 cores and 4–8 GB of memory per core, plus overhead (more for PySpark). Adjust based on spill, GC time, and failures observed in the UI, and on node sizes in your cluster.

When should I use broadcast joins?

When one side of the join is small enough to fit comfortably in each executor's memory (typically tens to a few hundred megabytes). Broadcasting avoids shuffling the large side, which often speeds joins dramatically.

Does caching always make Spark faster?

No. Caching helps only when the same computed DataFrame is reused by several actions. Otherwise it adds memory pressure and serialization cost. Measure, and unpersist cached data when it's no longer needed.

Related Topics

References