GC Pressure in Spark: How to Find It, Read It, and Actually Fix It
GC pressure is the most misdiagnosed failure in Spark, and the one most often 'fixed' by adding heap, which usually makes it worse. Here is where to look in the UI and the logs, how to tell it apart from data skew, and the levers that actually move the number.
Your executors are spending a third of their life collecting garbage. You raised the heap. The job got slower.
GC pressure is the only Spark failure mode I have seen get worse after the “obvious” fix. That is because the obvious fix targets the wrong variable: people treat garbage collection as a memory shortage, when it is almost always a churn problem.
This is the field guide I wish I had the first time I watched a 4-hour job turn into a 6-hour job because someone doubled spark.executor.memory.
π§ GC is a tax on churn, not a memory shortage
A Spark executor is a JVM with a fixed heap. Every task running on it creates objects: rows, boxed primitives, hash-map entries, broadcast lookups, shuffle buffers. The JVM cannot free those objects individually β it waits, then stops the world and sweeps a whole generation at once.
The cost of that sweep is proportional to the number of live objects it has to walk. Spark’s own tuning guide puts it plainly: garbage collection cost is proportional to the number of Java objects, and Java objects carry enormous overhead β a small String is tens of bytes of header and indirection before it holds a single character.
So the two variables are:
- Allocation rate β how fast you are minting new objects (churn).
- Retention β how many of those objects are still alive when the collector runs (object lifetime).
A memory shortage is one possible cause of pressure. Churn is the far more common one, and it is invisible if you only look at heap size.
The number you actually care about is how much of your compute is being spent on bookkeeping rather than work:
$$\eta = 1 - \frac{t_{\text{GC}}}{t_{\text{task}}}$$When t_GC is 38% of t_task, you are paying for 100 cores to get the throughput of 62. That is a line item, not an anecdote.
πΊοΈ The four signals, and where they live
Before touching config, know your instruments. Everything you need is already in the Spark UI and the executor logs β you just have to read the right tab in the right order.
Spark UI β read it in this order
Executors β START HERE. GC Time vs Task Time, per executor
Stages β is one task 10x the median? (the skew detector)
Jobs β which stage actually dominates wall clock?
Storage β is your cache spilling to disk?
Environment β did your -XX flags even apply?
SQL β exchange sizes, shuffle partition counts
There are four independent signals. A real diagnosis needs at least two of them agreeing.
| # | Signal | Where it lives | What it proves |
|---|---|---|---|
| 1 | GC Time / Task Time ratio | Executors tab | That GC is expensive, and on which executors |
| 2 | GC log pause pattern | executor stderr | Why it’s expensive: churn vs retention vs full GC |
| 3 | Task-duration spread | Stages tab | Whether it’s really GC, or data skew wearing a GC costume |
| 4 | The OOM message | executor log / exit code | Which of three different failures you’re holding |
π Signal 1 β The Executors tab: the ratio that matters
Open Executors. You get one row per live executor. The two columns to read together are Task Time and GC Time. Simplified, it looks like this:
| Executor | Address | Cores | Complete tasks | Task Time | GC Time | GC / Task |
|---|---|---|---|---|---|---|
| 1 | 10.0.4.11:3301 | 4 | 18,204 | 1.20 h | 0.06 h | 5% |
| 2 | 10.0.4.12:3301 | 4 | 18,190 | 1.20 h | 0.07 h | 6% |
| 3 | 10.0.4.13:3301 | 4 | 6,410 | 0.90 h | 0.34 h | 38% |
| 4 | 10.0.4.14:3301 | 4 | 18,201 | 1.20 h | 0.05 h | 4% |
Ignore the absolute GC Time. A long-running executor accumulates GC time whether it is healthy or not. The ratio is the signal, because Task Time is the yardstick for how much work that executor actually did in the same window.
Executor 3 is not just slow β it did a third of the work of its peers while running hotter. That shape (one hot executor, everybody else fine) is the single most important pattern in this whole article, and it is not a heap-size problem. We come back to it in Signal 3.
Two caveats that have burned me:
GC Timeis cumulative since executor start, not a sliding window. An executor that thrashed for its first ten minutes shows an inflated total for the next three hours. If you have changed anything, kill the application and measure again from zero.GC Timesums all collectors, including fast young-gen collections. A high number with 40,000 tiny pauses is a different disease from a high number with six full GCs. You cannot tell those apart from the UI β that’s what the log is for.
How high is too high? Spark’s own guidance is qualitative: if full GCs fire multiple times before a task completes, you don’t have enough room. The ratio itself is a field convention, not a Spark rule. My working thresholds:
| GC / Task Time | Read |
|---|---|
| under 5% | Healthy. Stop looking here. |
| 5 β 10% | Watch it. Something is drifting. |
| 10 β 20% | Tune. You are leaving real money on the table. |
| over 20% | Incident. Expect instability, not just slowness. |
Here is what the same 60-second window looks like for a healthy heap versus a thrashing one. Watch the bars, not the caption: frequency on top, duration on the bottom.
One minute of executor GC, two heaps
Healthy β ~40 ms per pause, 22 pauses / 60 s. GC β 6% of task time.
Thrashing β ~4.5 s per pause, 4 pauses / 60 s. GC β 38% of task time.
Same 60 seconds. The healthy heap collects often and quickly, which is cheap. The thrashing heap collects rarely and catastrophically, which is where your SLA dies.
π Signal 2 β The GC log is the ground truth
The UI tells you that GC is expensive. Only the log tells you why, and the reason determines which lever you pull. Ten minutes of log reading beats a week of guesses.
Turn it on for executors. On Java 11 and newer, the unified logging flag writes to the executor’s stderr, which Spark already collects for you β so it shows up in the UI and in your cluster’s aggregated logs with no extra plumbing:
spark-submit \
--conf "spark.executor.extraJavaOptions=-Xlog:gc*:time,uptime,level,tags" \
--conf "spark.executor.defaultJavaOptions=-XX:+HeapDumpOnOutOfMemoryError" \
your_job.py
On Java 8 (Spark 3.x still runs plenty of it):
--conf "spark.executor.extraJavaOptions=-verbose:gc -XX:+PrintGCDetails -XX:+PrintGCDateStamps -XX:+PrintGCTimeStamps"
Two gotchas worth internalising:
- The
spark.executor.prefix is not optional. Set these onspark.driver.extraJavaOptionsby accident and you will faithfully log the driver’s GC while the executors β where the data lives β stay dark. - Never put
-Xmxhere. Spark sizes the executor heap fromspark.executor.memory. A hand-set-Xmxeither gets ignored or fights the framework. Size the heap with Spark config; use the Java options for behaviour, not capacity.
Now the anatomy of a modern line:
[2026-10-12T09:14:22.418+0000][7.842s][info][gc] GC(214) Pause Young (Normal) (G1 Evacuation Pause) 1841M->612M(4096M) 412.883ms
Read it left to right and you already have the whole diagnosis:
2026-10-12T09:14:22.418+0000β wall-clock time. Correlate it with the stage that was running.7.842sβ JVM uptime. A pause at 7 seconds is a startup artifact. The same pause at 7,842 seconds is a production problem.GC(214)β the collection’s sequence number. If this is in the low thousands after an hour of work, you are collecting constantly.Pause Young (Normal)β the kind.Pause Youngis cheap.Pause Fullis not, and its presence is the loudest thing in the file.(G1 Evacuation Pause)β the cause. This is the collector telling you why it decided to stop the world.1841M->612Mβ heap used before and after. This single pair is the diagnosis: how much did it actually reclaim?(4096M)β the committed heap size, so you can see how much headroom was left.412.883msβ how long the world stopped. Multiply by frequency to get your tax.
Read the log in one pass and classify. There are only three shapes that matter:
1841M -> 612M (4096M) 412 ms healthy young GC reclaimed 1229M β 67% of what it touched
1841M -> 1780M (4096M) 389 ms promotion pressure reclaimed 61M β 3% of what it touched
3902M -> 1104M (4096M) 4812 ms FULL GC reclaimed a lot, and cost 4.8 s of stop-the-world
That second line is the one people misread. It is fast (389 ms) and reclaims almost nothing β which looks contradictory until you remember what it means: the objects it swept were all still reachable, so they got promoted into the old generation instead of freed. You are not running out of memory. You are manufacturing long-lived objects, and the bill arrives later as a full GC.
Which brings us to the two rates you should actually compute:
$$R_{\text{alloc}} \approx \frac{\Delta\,\text{heap used}}{\Delta t} \qquad\qquad R_{\text{promotion}} = R_{\text{alloc}} \times P(\text{survives young gen})$$R_alloc is a property of your code β how many objects you create per second. P(survives young) is a property of your data flow β how long your objects stay alive. Almost every effective fix attacks one of those two, and almost every ineffective fix (adding heap) attacks neither. A bigger young generation lowers the frequency of collections without lowering the allocation rate; a bigger old generation just delays the full GC.
Here is the lifecycle the executor is actually looping through, and where each lever bites:
Young-gen lifecycle of a Spark task object
Behind that one-way arrow is a small, brutal piece of bookkeeping. A young collection never frees anything in place β it copies whatever is still reachable somewhere else, then wipes the rest:
Eden ββββββββββββββΊ every task object is born here
β
β minor GC: live objects are copied into the *other* survivor space
βΌ
Survivor S0 β· S1 ββΊ survivors get copied, swapped, copied again
β
β after N survivals they are promoted
βΌ
Old generation ββββΊ long-lived data starts piling up here
β
βΌ
full GC βββββββββββΊ stop the world, but only once this actually fills
Two things follow from that. Objects that die young are free β they are never copied, so the sweep costs only what survives. Objects that survive are copied, repeatedly, and every copy is CPU you paid for on the critical path. So object lifetime, not object count, is what turns a cheap young GC into an expensive one β which is exactly why a serialized cache and early filtering beat any heap-size change.
Every stage above Old generation is cheap. The whole game is stopping objects from reaching the bottom.
π§΅ Signal 3 β The skew trap
This is the section I would tattoo on the inside of every Spark engineer’s eyelids.
Look back at that Executors table. One executor at 38% GC, three at ~5%. The instinct is “that executor needs more heap”. Run the experiment and it fails, because the executor is not the unit of blame β the partition is.
Open Stages and compare the max task duration against the 75th percentile for the same stage:
| Stage 12 β 18,204 tasks | Duration | Shuffle read |
|---|---|---|
| min | 41 s | 118 MB |
| p25 | 48 s | 121 MB |
| median | 52 s | 126 MB |
| p75 | 55 s | 126 MB |
| max | 41 min | 61.2 GB |
That is data skew: one key has 500x the rows of its peers. The executor running that one giant partition is allocating objects at a rate its peers never see β so of course its GC time is enormous. Adding heap gives the giant partition more room to be giant, and the stage gets slower because the long tail is now longer.
Tuning GC before ruling out skew is the most common expensive mistake in Spark performance work. Here is the triage, in the order I actually run it:
GC Time / Task Time is high on the Executors tab
β
ββ Is it high on EVERY executor?
β ββ NO βββΊ go to Stages and check the task duration spread
β β ββ max >> p75 on the same stage βββΊ DATA SKEW
β β β fix: repartition or salt the hot key. Do not touch the JVM.
β β ββ durations even, one host slow βββΊ NODE problem
β β fix: noisy neighbour, local disk, network. Not GC.
β ββ YES βββΊ is Task Time also high on those executors?
β ββ NO βββΊ GC is not your bottleneck; look at I/O and shuffle spill
β ββ YES βββΊ real, cluster-wide GC pressure βββΊ go to the levers
β
ββ Does the GC log contain any "Pause Full"?
ββ YES βββΊ working set no longer fits, or too much is cached long-term
Notice that the first branch has nothing to do with the JVM. In my experience roughly half of “GC incidents” exit this flowchart at the top.
β Signal 4 β Three OutOfMemoryErrors, three different fixes
“Out of memory” is not one failure. It is three, and they need opposite remedies. Getting this wrong is how teams spend a sprint raising a number that was never the constraint.
| What you see | Where | What it actually is | What to change |
|---|---|---|---|
java.lang.OutOfMemoryError: Java heap space | executor stderr | The object graph genuinely does not fit β often one oversized partition, a collect(), or a toPandas() | More spark.executor.memory, smaller/more partitions, stop pulling data to the driver |
java.lang.OutOfMemoryError: GC overhead limit exceeded | executor stderr | The JVM spent ~98% of its time collecting and reclaimed ~2% of the heap. This is thrashing, by definition | Reduce churn: kill Python UDFs, serialized caching, more partitions. Heap alone will not fix this |
| Container killed… exit code 137 / 143, or “Container killed by YARN for exceeding memory limits” | cluster manager, executor lost | Off-heap: Python workers, netty shuffle buffers, page cache. The JVM heap is fine | Raise spark.executor.memoryOverhead or spark.memory.offHeap.size β never the heap |
The third row is the one that wastes the most time, because the failure never appears in a Java stack trace. A killed container means the process died, not the heap. If you respond by raising spark.executor.memory, you spend RAM on a heap that was never the problem, and the container gets killed again.
A useful rule: heap OOM and GC-thrash OOM are Java saying “my memory is wrong”; exit code 137 is the operating system saying “your other memory is wrong”.
π The diagnosis loop
Here is the whole procedure as one loop. Internalise this and you will never again start with spark.executor.memory.
The tuning loop I actually run
The discipline that makes it work is one lever per iteration. Change four configs at once and you learn nothing β you get a new number with no causal story, and the next person to touch the job has to rediscover everything you skipped.
π§ The levers, ranked by payoff
Ordered by how much they actually move the number, not by how easy they are to type.
1. Fix the data shape before the JVM
The most effective GC fix is usually not a GC fix. Fewer, smaller objects beat every tuning flag.
- More partitions. Target 128β200 MB per partition for the shuffle. Under
spark.sql.shuffle.partitions(default 200), a 400 GB aggregate gives you 2 GB partitions, and each of those is a heap bomb. Right-sized partitions give each task a working set the collector can actually handle. reduceByKeyovergroupByKey. They look equivalent and are not:groupByKeymaterialises every value for a key before reducing, which builds the exact long-lived object graph that causes promotion.reduceByKeycombines as it goes.- Never
collect()ortoPandas()on a large frame. Both pull every row into one JVM β the driver. If your GC graph shows the driver thrashing, this is almost certainly why. - Built-in functions over Python UDFs. A Python UDF round-trips every row through a serialized boundary and back. The JVM still allocates the row representation, and now pays for the marshalling too.
- Filter and project early. A column you dropped in line 4 was never allocated in line 40.
2. Make cached data cheap β or don’t cache it
If this is the only thing you take away, it is enough. Serialized caching is the single highest-leverage memory change available, and Spark’s own docs call it the first thing to try when GC is a problem.
from pyspark import StorageLevel
# Every cached object stays a live JVM object -> the collector must walk them all
df.persist(StorageLevel.MEMORY_AND_DISK)
# The same data as byte arrays: a handful of large objects instead of millions of small ones
df.persist(StorageLevel.MEMORY_AND_DISK_SER)
Serialized, a cached row is bytes in one array. Deserialized, it is a tree of objects with headers, pointers and boxed primitives β commonly 2β5x the raw data size, and every node of it is a live object the collector must traverse on every pass. If a DataFrame is cached and also shows up in every full GC, this is your answer.
And the harder question: do you need to cache it at all? Long-lived cached data is precisely what fills the old generation. A cache that gets reused twice is often cheaper to recompute than to keep alive.
3. Executor shape: many medium heaps beat few huge heaps
This is where the “obvious fix” does damage. A 32 GB executor heap does not collect more garbage faster β it collects the same garbage more slowly, because a stop-the-world pause scales with the live set. You have bought longer pauses and a bigger bill.
The shape I reach for first:
--executor-cores 4 # a small, well-understood unit of parallelism
--executor-memory 8g # 4-8 GB heaps: short pauses, sane promotion
--conf spark.executor.memoryOverhead=1g # room for Python + netty, off-heap
Roughly: 4β5 cores and 8β16 GB of heap per executor, with memory headroom proportional to how much Python work the job does. Below ~4 GB you start paying JVM startup and shuffle-fragmentation overhead. Above ~16 GB you are buying pause time you can’t afford.
Memory overhead deserves its own line. Spark defaults it to max(384 MB, 10% of executor memory), and that default is generous for pure DataFrame jobs and far too small for anything running pandas UDFs, PyArrow, or heavy shuffle. Under-provisioned overhead is the source of the exit-code-137 family above.
4. Memory layout
Once the shape is sane, these tune the split inside the heap:
| Config | Default | Meaning |
|---|---|---|
spark.memory.fraction | 0.6 | Share of (heap β 300 MB) available to execution + storage; the rest absorbs per-task overhead and metadata |
spark.memory.storageFraction | 0.5 | Within that pool, the slice storage may hold against eviction |
spark.memory.offHeap.enabled | false | Move execution/cache memory outside the GC’s reach |
spark.memory.offHeap.size | 0 | Size of that off-heap arena (required if enabled) |
Off-heap is underused and worth knowing: memory allocated off-heap is invisible to the garbage collector, so a cache-heavy job can park its hot data there and simply stop paying GC tax on it. The trade is that you now manage two budgets and must size the arena deliberately β off-heap memory counts against spark.executor.memoryOverhead.
spark.memory.fraction is also a real GC lever in the direction that surprises people. If the GC log shows the old generation filling with cached data, lowering this fraction reclaims heap for task execution and shortens pauses.
5. Serialization
--conf spark.serializer=org.apache.spark.serializer.KryoSerializer
Java serialization produces more bytes and more intermediate objects than Kryo. Smaller serialized payloads mean smaller shuffle buffers, smaller caches, and less to collect. Register your types explicitly for the best result:
conf.set("spark.kryo.registrationRequired", "true")
conf.set("spark.kryo.classesToRegister", "com.example.Row1,com.example.Row2")
registrationRequired=true turns unregistered classes into hard failures. That is a feature β it forces you to audit what is actually flowing through your pipeline β but it will break a job the first time you enable it. Do it in staging.
6. Choose a collector deliberately
Spark does not set a garbage collector; your JVM’s default applies. On Java 8 that is Parallel GC. On Java 9 and later β and Spark 4.x uses JDK 17 by default β that is G1.
Neither is universally right, and this is a measurement, not an opinion:
- Parallel GC (
-XX:+UseParallelGC) β stop-the-world, multi-threaded, no concurrent marking competing for CPU. For executor heaps in the 4β8 GB range doing throughput-oriented batch work, it is frequently the better choice, and it is the one teams never try. - G1 (
-XX:+UseG1GC) β concurrent marking, pause targets. Worth it for large heaps and latency-sensitive work.
If you stay on G1, three flags matter when Spark workloads hit it:
-XX:MaxGCPauseMillis=200 # G1's actual goal; without it G1 optimises for throughput
-XX:InitiatingHeapOccupancyPercent=35 # start concurrent marking earlier; late starts cause FULL GCs
-XX:G1HeapRegionSize=8m # keep big row batches from becoming "humongous" allocations
Humongous allocations are a G1-specific Spark hazard: any object larger than half a G1 region is allocated in contiguous regions and is expensive to collect. Wide row batches and large cached arrays hit this constantly. That is why Spark’s own docs call out raising the region size for large heaps.
7. Re-measure, and confirm the flags applied
Open Environment and search your JVM options. Half the “tuning did nothing” incidents I have debugged ended here β the flag was typo’d, applied to the driver, or silently dropped. Then compare the same two numbers you started with: GC/Task ratio, and the dominant stage’s duration.
π« The anti-pattern loop
I want to name this explicitly, because it is where almost everyone starts and it is self-reinforcing.
The loop that makes GC worse
Every step feels like progress and none of them is. Adding heap buys pause duration. Adding cores buys GC parallelism, but more concurrent workers on the same heap buys more allocation pressure β so you arrive back at the top of the loop with a larger invoice. The exit is always the same: reduce churn, or reduce retention.
π What it looks like on the bill
Numbers below are representative of a real batch pipeline after three changes: right-sized partitions, MEMORY_AND_DISK_SER, and switching executors from one 16-core/22 GB shape to 4-core/8 GB with Parallel GC. Your mileage will vary; the shape of the curve will not.
| Metric | Before | After |
|---|---|---|
| GC time / task time | 38% | 7% |
| Dominant stage duration | 41 min | 24 min |
| Executors (4 cores each) | 200 | 140 |
| End-to-end wall clock | 2 h 05 m | 1 h 26 m |
| Core-hours consumed | 1,664 | 801 |
The arithmetic that matters is the last row. Cluster cost is roughly cores multiplied by time, so:
before: 200 executors x 4 cores x 2.08 h = 1,664 core-hours
after: 140 executors x 4 cores x 1.43 h = 801 core-hours
---------------------------
~52% less compute
You did not delete work. You stopped paying full price for the fraction of each task the collector owned. That is what makes GC tuning a cost-optimisation exercise wearing a performance-engineering costume.
β The checklist
Run it in this order. Do not skip ahead.
- Executors tab β compute GC Time Γ· Task Time for every executor. Record it.
- Is it one executor or all of them? One is skew. All is pressure.
- Stages tab β compare max task duration to p75 on the dominant stage. Rule skew out before anything else.
- Turn on GC logging to executor stderr; confirm the flags in the Environment tab.
- Read one minute of log β count pauses, sum their duration, and check for
Pause Full. - Classify: churn (many minor GCs, low reclamation) or retention (promotion per cycle, old gen steadily filling).
- Pick exactly one lever from the ranked list and re-run.
- Compare the same two numbers you wrote down in step 1, and confirm the flags are live.