Skip to content
ranjit.data
NEW deep dive October 12, 2026 19 min read

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.

#SignalWhere it livesWhat it proves
1GC Time / Task Time ratioExecutors tabThat GC is expensive, and on which executors
2GC log pause patternexecutor stderrWhy it’s expensive: churn vs retention vs full GC
3Task-duration spreadStages tabWhether it’s really GC, or data skew wearing a GC costume
4The OOM messageexecutor log / exit codeWhich 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:

ExecutorAddressCoresComplete tasksTask TimeGC TimeGC / Task
110.0.4.11:3301418,2041.20 h0.06 h5%
210.0.4.12:3301418,1901.20 h0.07 h6%
310.0.4.13:330146,4100.90 h0.34 h38%
410.0.4.14:3301418,2011.20 h0.05 h4%

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 Time is 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 Time sums 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 TimeRead
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

HEALTHYTHRASHING0 s60 s
Healthy heap
Healthy β€” ~40 ms per pause, 22 pauses / 60 s. GC β‰ˆ 6% of task time.
GC thrashing
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:

  1. The spark.executor. prefix is not optional. Set these on spark.driver.extraJavaOptions by accident and you will faithfully log the driver’s GC while the executors β€” where the data lives β€” stay dark.
  2. Never put -Xmx here. Spark sizes the executor heap from spark.executor.memory. A hand-set -Xmx either 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 Young is cheap. Pause Full is 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

Eden fillsMinor GCPromotionOld genFull GC

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 tasksDurationShuffle read
min41 s118 MB
p2548 s121 MB
median52 s126 MB
p7555 s126 MB
max41 min61.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 seeWhereWhat it actually isWhat to change
java.lang.OutOfMemoryError: Java heap spaceexecutor stderrThe 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 exceededexecutor stderrThe JVM spent ~98% of its time collecting and reclaimed ~2% of the heap. This is thrashing, by definitionReduce 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 lostOff-heap: Python workers, netty shuffle buffers, page cache. The JVM heap is fineRaise 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

Spark UIGC logsHeap mathOne leverRe-measure

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.
  • reduceByKey over groupByKey. They look equivalent and are not: groupByKey materialises every value for a key before reducing, which builds the exact long-lived object graph that causes promotion. reduceByKey combines as it goes.
  • Never collect() or toPandas() 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:

ConfigDefaultMeaning
spark.memory.fraction0.6Share of (heap βˆ’ 300 MB) available to execution + storage; the rest absorbs per-task overhead and metadata
spark.memory.storageFraction0.5Within that pool, the slice storage may hold against eviction
spark.memory.offHeap.enabledfalseMove execution/cache memory outside the GC’s reach
spark.memory.offHeap.size0Size 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

GC time highAdd heapLonger pausesAdd coresBigger bill

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.

MetricBeforeAfter
GC time / task time38%7%
Dominant stage duration41 min24 min
Executors (4 cores each)200140
End-to-end wall clock2 h 05 m1 h 26 m
Core-hours consumed1,664801

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.

  1. Executors tab β€” compute GC Time Γ· Task Time for every executor. Record it.
  2. Is it one executor or all of them? One is skew. All is pressure.
  3. Stages tab β€” compare max task duration to p75 on the dominant stage. Rule skew out before anything else.
  4. Turn on GC logging to executor stderr; confirm the flags in the Environment tab.
  5. Read one minute of log β€” count pauses, sum their duration, and check for Pause Full.
  6. Classify: churn (many minor GCs, low reclamation) or retention (promotion per cycle, old gen steadily filling).
  7. Pick exactly one lever from the ranked list and re-run.
  8. Compare the same two numbers you wrote down in step 1, and confirm the flags are live.

The one-sentence version

GC time is a tax on object churn and object lifetime. Adding heap raises the tax ceiling without lowering the rate β€” so fix the partition sizes, the cache representation, and the executor shape first, and reach for the JVM flags only after the logs have told you which of the three GC failure modes you’re actually in.

Activity Feed

GC Pressure in Spark: How to Find It, Read It, and Actually Fix It

Oct 12, 2026 · New Article

Welcome to the Field Notes

Oct 11, 2026 · New Article

The Real Math Behind PySpark Resource Sizing

Oct 11, 2026 · New Article

The Ultimate SQL Cheatsheet

Oct 11, 2026 · Cheatsheet Updated

Apache Airflow

Oct 11, 2026 · Cheatsheet Updated

View All Activity →