Garbage collection (GC) is the JVM freeing memory from objects that are no longer used. When it takes a large share of task time, Spark is spending effort cleaning up instead of computing. You see it as slow tasks and sometimes executors that look frozen.
How to recognise it
In the stage page of the Spark UI, the task table has a "GC Time" column. If GC time is more than about 10 percent of task duration, there is a problem. In the executor logs you may see long "Full GC" pauses, and heartbeats may time out, which makes the driver think the executor is lost.
Usual causes
- Too much data per executor, so the heap is nearly full and the collector runs constantly.
- Millions of small objects. RDDs of Python or Java objects, or
Rowobjects, are much heavier than Tungsten's binary rows. - Caching deserialized objects (
MEMORY_ONLYon RDDs), which fills the old generation. - Huge collections built in user code, such as large lists in a UDF.
What to do
- Use DataFrames and built-in functions so data stays in the compact binary format.
- Cache DataFrames in serialized form or with
MEMORY_AND_DISK, and only cache what you reuse. - Use the G1 garbage collector, which is the usual choice on modern JVMs, and handles large heaps better than the older collectors.
- Reduce partition size so each task holds less at once.
- Try smaller executors. A 64 GB heap with 15 cores can have long pauses, while several 16 GB heaps often pause less, at the cost of copying broadcast data more times.
- Look for leaks of references, such as a growing Python list in a driver loop.
--conf spark.executor.extraJavaOptions=-XX:+UseG1GC
What to say in an interview
Name the symptom (GC time column), the likely cause (too many objects or too much per executor), and the main fix (DataFrames, smaller partitions, tuning executor size). Do not start with GC flags. Changing data shape helps more than flags.