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
- Start with the Spark UI: find the longest stages, then look at task-time distribution, shuffle read/write, and spill.
- Fix data and code first: prune columns and rows early, avoid Python UDFs, fix skew, and compact small files.
- Joins: broadcast small tables, keep sort-merge for large-large joins, and let AQE switch strategies at runtime.
- Executors: ~4–5 cores each, with memory sized to avoid spills. Use dynamic allocation for variable workloads.
- Cache only DataFrames reused multiple times, and unpersist them afterward.
- Measure every change: runtime, cost, and shuffle bytes, on realistic data volumes.
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:
- ~4–5 cores per executor: fewer wastes JVM overhead, more causes HDFS/S3 client and GC contention.
- Memory per core of roughly 4–8 GB for shuffle-heavy work. Increase it if you see spills.
- Leave headroom for the OS and cluster managers. On Kubernetes, requests must fit node allocatable memory.
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
- Columnar formats (Parquet, ORC, table formats) enable column pruning and predicate pushdown using file and row-group statistics.
- Partition pruning and data skipping (clustering, Z-order, Iceberg metadata) avoid reading irrelevant files.
- Compact small files: planning and opening thousands of tiny files dominates runtime for many jobs.
- Cloud storage: use optimized committers (S3A magic committer) and avoid listing-heavy layouts.
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
- Reproduce with realistic data volume, and record the baseline runtime, cost, and shuffle bytes.
- Find the dominant stage(s) in the Spark UI.
- Diagnose: skew, spill, huge shuffle, full scan, small files, UDF, or a bad join.
- Fix in code and data first: pruning, broadcast, pre-aggregation, salting, compaction, and replacing UDFs.
- Then tune resources: executor size, memory overhead, partition sizes, dynamic allocation.
- 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
- Apache Spark — Pillar overview
- Spark Partitioning — Shuffles, skew, and file layout
- Spark DataFrames — Writing optimizable code
- PySpark — Python UDF and Arrow performance
- Performance Optimization — General principles
- Data Lakehouse — Storage layout that speeds reads