Menu

Apache Spark course · Lesson 8 of 9

Spark Memory, Executor Sizing and Cluster Tuning

How Spark's unified memory works, how to size executors with a worked example, and how to tune GC, diagnose spill, and use dynamic allocation and speculation.

  • Advanced
  • 20 min read
  • Updated Oct 2026
On this page
  1. Sample session
  2. Unified memory management
  3. What it is
  4. The borrowing rules
  5. Checking the numbers
  6. Pitfalls
  7. In interviews
  8. Executor sizing
  9. The trade-off
  10. A worked example
  11. Adjusting the result
  12. In interviews
  13. Cores and memory configuration
  14. The settings that matter
  15. Pitfalls
  16. In interviews
  17. GC tuning
  18. What it is
  19. How to approach it
  20. Pitfalls
  21. In interviews
  22. Diagnosing spill to disk
  23. What it is
  24. Seeing it
  25. Causes and fixes
  26. Pitfalls
  27. In interviews
  28. Dynamic allocation
  29. What it is
  30. How it works
  31. Pitfalls
  32. In interviews
  33. Speculative execution
  34. What it is
  35. How it is configured
  36. Pitfalls
  37. In interviews
  38. Practice questions
  39. Key takeaways

An executor is a JVM process with a fixed amount of memory and a fixed number of cores, and nearly every “my Spark job is slow or dies” problem eventually lands here: tasks spilling to disk, containers killed for exceeding their memory limit, long garbage-collection pauses, or a stage waiting on one slow machine. This lesson explains how Spark divides executor memory, how to size executors for a real cluster, and which settings to reach for when things go wrong.

Sample session

A local session with a 1 GB heap, so the memory numbers are easy to check by hand. spark.shuffle.spill.numElementsForceSpillThreshold is an internal setting that forces a spill after a given number of records; it is used here only so a spill happens on laptop-sized data. Do not set it in production.

import json, urllib.request
from pyspark.sql import SparkSession, functions as F

spark = (SparkSession.builder.master("local[2]").appName("memory-tuning")
         .config("spark.driver.memory", "1g")
         .config("spark.sql.shuffle.partitions", "2")
         .config("spark.shuffle.spill.numElementsForceSpillThreshold", "20000")
         .getOrCreate())
sc = spark.sparkContext

_opener = urllib.request.build_opener(urllib.request.ProxyHandler({}))
def api(path):
    port = sc.uiWebUrl.rsplit(":", 1)[1]
    url = f"http://localhost:{port}/api/v1/applications/{sc.applicationId}/{path}"
    return json.load(_opener.open(url))

Unified memory management

What it is

Since Spark 1.6, executor heap memory is managed by the unified memory manager: execution and storage share one pool and borrow from each other, instead of having fixed, separate slices. The heap (spark.executor.memory) is divided like this:

Region Size Used for
Reserved 300 MB, fixed Spark’s internal objects; also a safeguard for small heaps
Unified pool (M) (heap − 300 MB) × spark.memory.fraction (0.6) Execution and storage together
↳ Storage Up to all of M; spark.memory.storageFraction (0.5) of M is protected from eviction Cached blocks, broadcast variables, unrolling blocks
↳ Execution Up to all of M Shuffles, joins, sorts and aggregations (hash maps, sort buffers)
User memory The remaining 40% of (heap − 300 MB) Your objects, UDF data structures, Spark internal metadata, RDD dependency information

Outside the heap, each executor also gets memory overhead (spark.executor.memoryOverhead, by default 10% of executor memory with a 384 MB minimum) for JVM internals, thread stacks, native libraries, network buffers and, in PySpark, the Python worker processes unless spark.executor.pyspark.memory is set. Off-heap execution and storage memory (spark.memory.offHeap.enabled, spark.memory.offHeap.size) is extra again.

The borrowing rules

  • Execution can borrow free storage memory, and can evict cached blocks to reclaim memory, but only down to the protected storage region (storageFraction of M).
  • Storage can borrow free execution memory, but cannot evict execution: if execution needs memory back, storage must wait or drop blocks.
  • Within execution memory, each running task can use between 1/(2N) and 1/N of the execution pool, where N is the number of tasks currently running on the executor. A task that cannot get its minimum share waits; a task that runs out spills.

The asymmetry is deliberate: evicting a cached block costs a recomputation later, while forcing an aggregation to fail or spill heavily costs now.

Checking the numbers

The pool size Spark reports should equal (heap − 300 MB) × 0.6. In local mode the driver is also the executor, so its 1 GB heap is used:

heap = sc._jvm.java.lang.Runtime.getRuntime().maxMemory()
unified = (heap - 300 * 1024 * 1024) * 0.6
executor = api("executors")[0]

print("JVM max heap bytes:     ", heap)
print("(heap - 300 MB) x 0.6:  ", int(unified))
print("Spark-reported pool:    ", executor["maxMemory"])
JVM max heap bytes:      1073741824
(heap - 300 MB) x 0.6:   455501414
Spark-reported pool:     455501414

The Executors tab of the Spark UI shows this as the “Storage Memory” column’s total (used versus available), because the whole unified pool is available to storage when execution is idle.

Pitfalls

  • Raising spark.memory.fraction to “give Spark more memory” takes it from user memory, where your UDF objects and Spark’s metadata live, and can cause heap out-of-memory errors. Leave it at 0.6 unless you have measured a reason.
  • Assuming cache is guaranteed memory. Cached blocks above the protected region are evicted when execution needs room.
  • Forgetting overhead. A container is killed for exceeding its memory limit, not just for heap exhaustion. In PySpark, Python workers live in overhead memory.

In interviews

“Explain Spark’s memory model.” Draw the regions: 300 MB reserved, 60% of the rest unified (execution plus storage, half of it protected for storage), 40% user memory, plus off-heap overhead outside the heap. Then explain the eviction rule (execution can evict storage, not the reverse) and what spill and out-of-memory mean in that picture.

Executor sizing

The trade-off

You choose how to cut each worker node into executors:

Shape Pros Cons
Tiny (1 core, small memory, many executors) Fine-grained scheduling Broadcast variables and cached data duplicated per executor; per-JVM overhead multiplied; no benefit from sharing memory between tasks
Fat (all cores and memory of a node in one executor) Few JVMs, shared caches and broadcasts Very large heaps mean long GC pauses; one executor failure loses a whole node’s work; with HDFS, many concurrent tasks per JVM can saturate client I/O
Balanced (a few cores each, several executors per node) Good parallelism per JVM, moderate heaps Requires a little arithmetic

A widely used rule of thumb is about 4–5 cores per executor, with heaps in the tens of gigabytes rather than hundreds. It is a starting point to test, not a rule in the Spark documentation.

A worked example

Cluster: 10 worker nodes, each with 16 cores and 128 GB of RAM, running on YARN, with the driver in cluster mode.

  1. Leave room for the operating system and node daemons. Reserve 1 core and about 8 GB per node. Available per node: 15 cores, 120 GB.
  2. Choose cores per executor. 5 cores → 15 / 5 = 3 executors per node.
  3. Split memory per executor. 120 GB / 3 = 40 GB per container. The container must hold heap plus overhead (10% by default, at least 384 MB): heap × 1.1 ≤ 40 GB, so heap ≈ 36 GB (spark.executor.memory=36g) and overhead ≈ 3.6 GB.
  4. Count executors. 10 nodes × 3 = 30. In cluster mode the driver also needs a container (or leave one for the YARN ApplicationMaster), so request 29 executors.
  5. Total parallelism. 29 × 5 = 145 task slots, so stages should have a few hundred partitions or more (2–3 × 145).
  6. Memory per task. Unified pool = (36 GB − 300 MB) × 0.6 ≈ 21.4 GB. With 5 tasks running, each can use between about 2.1 GB (1/2N) and 4.3 GB (1/N) of execution memory while nothing is cached. Size shuffle partitions so a partition’s working set fits comfortably in that.

The same calculation as code you can rerun with your own numbers:

nodes, cores_per_node, mem_per_node_gb = 10, 16, 128
reserve_cores, reserve_mem_gb = 1, 8
cores_per_executor, overhead_factor = 5, 0.10

usable_cores = cores_per_node - reserve_cores
usable_mem_gb = mem_per_node_gb - reserve_mem_gb
executors_per_node = usable_cores // cores_per_executor
container_gb = usable_mem_gb / executors_per_node
heap_gb = int(container_gb / (1 + overhead_factor))
executors = nodes * executors_per_node - 1             # one container for the driver
unified_gb = (heap_gb - 0.3) * 0.6

print(f"executors per node: {executors_per_node}, container: {container_gb:.1f} GB")
print(f"--executor-memory {heap_gb}g, overhead ~{heap_gb * overhead_factor:.1f} GB")
print(f"--num-executors {executors}, --executor-cores {cores_per_executor}")
print(f"task slots: {executors * cores_per_executor}")
print(f"unified pool per executor: {unified_gb:.1f} GB, per task {unified_gb / (2 * cores_per_executor):.1f}-{unified_gb / cores_per_executor:.1f} GB")
executors per node: 3, container: 40.0 GB
--executor-memory 36g, overhead ~3.6 GB
--num-executors 29, --executor-cores 5
task slots: 145
unified pool per executor: 21.4 GB, per task 2.1-4.3 GB

Adjusting the result

  • PySpark with heavy Python UDFs or pandas UDFs: Python workers use overhead memory, so raise spark.executor.memoryOverhead (or set spark.executor.pyspark.memory) and shrink the heap accordingly. On Kubernetes, the documentation gives a default overhead factor of 0.40 for non-JVM jobs instead of 0.10.
  • Memory-hungry stages (big joins, wide windows): fewer cores per executor, so each task gets a larger share of execution memory.
  • Cloud instances with local NVMe: spills and shuffle files are cheaper, so slightly smaller memory per task is acceptable.
  • Managed platforms (Databricks, EMR Serverless, Dataproc) often fix executor shapes per instance type; you then choose instance types instead.

In interviews

Interviewers love “you have N nodes with C cores and M GB; how do you size executors?” Walk through the steps aloud: reserve for the OS, pick about 5 cores, divide memory, subtract overhead, subtract the driver, then check task slots and memory per task against the data. State that you would validate with the Spark UI (GC time, spill, executor losses) and adjust.

Cores and memory configuration

The settings that matter

Setting spark-submit flag Default Notes
spark.executor.memory --executor-memory 1g JVM heap per executor
spark.executor.cores --executor-cores 1 on YARN and Kubernetes; all free cores of a worker on standalone Concurrent tasks per executor (with spark.task.cpus = 1)
spark.executor.instances --num-executors (cluster manager default) Fixed executor count when dynamic allocation is off
spark.executor.memoryOverhead executor memory × spark.executor.memoryOverheadFactor Non-heap memory per executor
spark.executor.memoryOverheadFactor 0.10 (0.40 for non-JVM jobs on Kubernetes per the docs)
spark.executor.minMemoryOverhead 384m Floor for the overhead
spark.executor.pyspark.memory not set Caps Python worker memory; added to the container size
spark.driver.memory --driver-memory 1g Must be set before the driver JVM starts
spark.driver.cores --driver-cores 1 Cluster mode only
spark.driver.maxResultSize 1g Cap on serialized results of one action returned to the driver
spark.task.cpus 1 Cores per task; raise for multi-threaded tasks
spark.memory.fraction 0.6 Unified pool share
spark.memory.storageFraction 0.5 Protected storage share of the pool
spark.memory.offHeap.enabled / .size false / 0 Off-heap execution and storage
spark.kubernetes.executor.request.cores not set Kubernetes CPU request; can be fractional (for example 0.5) and differ from spark.executor.cores

A cluster-mode submission using the worked example:

spark-submit --master yarn --deploy-mode cluster \
  --num-executors 29 --executor-cores 5 --executor-memory 36g \
  --conf spark.executor.memoryOverhead=4g \
  --driver-memory 8g --driver-cores 2 \
  --conf spark.sql.shuffle.partitions=600 \
  jobs/daily_sales.py

Pitfalls

  • Setting spark.driver.memory in code after the JVM has started has no effect; in client mode set it with --driver-memory or spark-defaults.conf.
  • Putting heap size in extraJavaOptions. Spark rejects -Xmx there; use spark.executor.memory.
  • A huge spark.driver.maxResultSize hides a design problem: collect() of large data. Write results to storage instead.
  • Over-requesting cores. If spark.executor.cores exceeds what a node can offer, executors never start and the application waits for resources.

In interviews

Know the names and the difference between heap and container size: --executor-memory 36g with 10% overhead asks the cluster manager for about 40 GB per executor. A frequent question is “why was my container killed by YARN for exceeding memory limits?” The answer is overhead, often Python workers or native buffers, so raise spark.executor.memoryOverhead rather than the heap.

GC tuning

What it is

The JVM’s garbage collector reclaims objects your tasks no longer use. When the heap fills with long-lived data (cached blocks, big hash maps, wide rows) or churns through many short-lived objects, GC runs often or pauses for long periods, and tasks stall. Spark 4.0 and later use JDK 17 by default, which makes G1GC the default collector.

How to approach it

  1. Measure first. The Executors and stage pages of the Spark UI show GC time per executor and per task. As a rough signal, GC time above about 10% of task time is worth investigating.
  2. Collect GC logs for a closer look by adding JVM flags through spark.executor.extraJavaOptions (or spark.executor.defaultJavaOptions). On Java 17, use unified logging (-Xlog:gc*); the older -verbose:gc -XX:+PrintGCDetails style flags in some guides come from Java 8. Logs appear in the executors’ stdout or stderr.
  3. Reduce garbage before tuning the collector.
    • Prefer DataFrames: Tungsten’s binary rows create far fewer objects than RDDs of Java or Python objects.
    • Cache less, or cache serialized or on disk; unpersist datasets you no longer need.
    • Avoid huge per-task data structures in UDFs and mapPartitions (for example building a list of a whole partition).
    • Use smaller partitions so each task holds less at once.
    • Consider off-heap memory for very large executors.
  4. Then tune the collector. With G1 and large heaps, the Spark tuning guide suggests increasing the region size with -XX:G1HeapRegionSize. Very large heaps (well beyond ~64 GB) also argue for smaller executors rather than more GC flags.
spark-submit ... \
  --conf "spark.executor.extraJavaOptions=-Xlog:gc*:stdout:time,uptime -XX:G1HeapRegionSize=16m" \
  jobs/daily_sales.py

Pitfalls

  • Copying old CMS or ParallelGC tuning recipes onto Java 17: CMS was removed from the JDK, and those flags fail at JVM start.
  • Treating GC tuning as the first fix. Most GC problems are data-volume problems: too much cached, partitions too big, too many objects.
  • Ignoring the driver: a driver collecting large results or holding large broadcast values also suffers GC pauses, which stall scheduling for the whole application.

In interviews

“A stage shows high GC time; what do you do?” Confirm in the Executors tab and task metrics, check what is cached and how big partitions are, move from RDDs or heavy UDFs to DataFrame functions, reduce per-task memory (more partitions, fewer cores per executor), and only then tune G1 (region size, logging). Mention that Spark 4 runs on Java 17 with G1 as the default collector.

Diagnosing spill to disk

What it is

Spill happens when an operator (sort, hash aggregation, sort-merge join, window) runs out of execution memory: it writes its in-memory data to local disk and continues, merging the spilled files later. Spill keeps the job alive, but costs serialization, disk writes and reads, and often GC.

The stage page shows two metrics:

  • Spill (Memory): the size of the spilled data as it was in memory (deserialized).
  • Spill (Disk): the size written to disk (serialized and compressed), usually much smaller.

Seeing it

A sort of one million random numbers into two partitions, with the forced-spill threshold set at session start:

sc.setJobDescription("sort with spill")
(spark.range(0, 1_000_000, numPartitions=2)
      .withColumn("k", F.rand(1))
      .orderBy("k")
      .write.format("noop").mode("overwrite").save())
sc.setJobDescription(None)

for s in sorted(api("stages"), key=lambda s: s["stageId"]):
    if s.get("description") == "sort with spill":
        print(f'stage {s["stageId"]}: tasks {s["numTasks"]}, '
              f'spill memory {s["memoryBytesSpilled"]:,} B, spill disk {s["diskBytesSpilled"]:,} B')
stage 0: tasks 2, spill memory 0 B, spill disk 0 B
stage 1: tasks 2, spill memory 0 B, spill disk 0 B
stage 3: tasks 2, spill memory 411,041,008 B, spill disk 20,985,559 B

The first stages compute the data and sample it to find range boundaries for the sort (stage 2 was skipped because its output already existed). The final stage sorted and spilled: about 411 MB in memory became about 21 MB on disk, which is why “Spill (Memory)” always looks alarming next to “Spill (Disk)”.

Causes and fixes

Cause How it shows Fix
Partitions too large Spill on most tasks of a stage More shuffle partitions (or a higher initialPartitionNum with AQE); smaller spark.sql.files.maxPartitionBytes for scans
Skew Spill only on the few slowest tasks Skew fixes: AQE skew join, broadcast, salting
Too many concurrent tasks per executor Spill and GC across the board Fewer cores per executor, or more memory per executor
Cached data crowding execution Storage memory near full in the Executors tab Unpersist, cache less, or cache on disk
Wide rows or exploded arrays Large shuffle sizes relative to input Select fewer columns before the shuffle; explode later

Pitfalls

  • Treating any spill as a failure. Small spill on a few tasks is normal; chase spill when it dominates task time or comes with GC and executor loss.
  • Fixing spill by raising executor memory when the real cause is one skewed key: the hot partition still does not fit.
  • Forgetting local disk capacity: heavy spill and shuffle can fill nodes’ local disks and kill executors.

In interviews

“How do you diagnose spill?” Point to the stage page metrics (Spill Memory and Disk, per task in the summary percentiles), decide whether it is uniform (partition size or memory per task) or concentrated (skew), and pick the matching fix. Explaining why the memory figure is much larger than the disk figure is a good detail.

Dynamic allocation

What it is

With dynamic allocation, Spark adds executors when tasks are waiting and removes executors that have been idle, so an application uses only the resources it needs. It suits shared clusters and jobs with uneven stages. It is off by default (spark.dynamicAllocation.enabled=false).

How it works

  • Scaling up: when tasks have been pending for spark.dynamicAllocation.schedulerBacklogTimeout (1s), Spark requests more executors, then keeps requesting every sustainedSchedulerBacklogTimeout while the backlog lasts, increasing the request exponentially (1, 2, 4, 8 …) up to maxExecutors.
  • Scaling down: an executor idle for spark.dynamicAllocation.executorIdleTimeout (60s) is removed. Executors holding cached data are kept for spark.dynamicAllocation.cachedExecutorIdleTimeout, which defaults to infinity.
  • Shuffle files: removing an executor would delete shuffle output that later stages need. Either run an external shuffle service (spark.shuffle.service.enabled=true, typical on YARN) so files outlive executors, or rely on shuffle tracking (spark.dynamicAllocation.shuffleTracking.enabled, true by default), which keeps executors that hold live shuffle data. Since Spark 3.4 the driver tracks shuffle data when dynamic allocation runs without a shuffle service, which is the usual setup on Kubernetes.
Setting Default
spark.dynamicAllocation.enabled false
spark.dynamicAllocation.minExecutors 0
spark.dynamicAllocation.maxExecutors infinity
spark.dynamicAllocation.initialExecutors minExecutors
spark.dynamicAllocation.executorIdleTimeout 60s
spark.dynamicAllocation.cachedExecutorIdleTimeout infinity
spark.dynamicAllocation.schedulerBacklogTimeout 1s
spark.dynamicAllocation.shuffleTracking.enabled true
spark-submit ... \
  --conf spark.dynamicAllocation.enabled=true \
  --conf spark.dynamicAllocation.minExecutors=2 \
  --conf spark.dynamicAllocation.maxExecutors=50 \
  --conf spark.dynamicAllocation.executorIdleTimeout=120s \
  --executor-cores 5 --executor-memory 36g \
  jobs/daily_sales.py

Pitfalls

  • No maxExecutors: one large stage can grab the whole shared cluster.
  • Caching with dynamic allocation: executors with cached blocks never time out by default, so the application holds resources. Unpersist when done, or set cachedExecutorIdleTimeout.
  • Shuffle tracking keeps executors alive as long as their shuffle data might be needed, so scale-down can be slower than the idle timeout suggests.
  • Start-up latency: requesting executors takes seconds to minutes (pod or container scheduling), so very short jobs may be faster with a fixed allocation.

In interviews

“How does dynamic allocation work and what problem does the external shuffle service solve?” Explain backlog-driven scale-up and idle-timeout scale-down, then the shuffle-file problem and its two solutions (external shuffle service, shuffle tracking). Mention maxExecutors as a guard on shared clusters.

Speculative execution

What it is

Speculative execution launches a second copy of a task that is running much slower than its peers, on another executor, and takes whichever finishes first. It targets stragglers caused by bad machines (a failing disk, a noisy neighbour), not slow tasks caused by data.

How it is configured

Setting Default Meaning
spark.speculation false Turn speculation on
spark.speculation.interval 100ms How often Spark checks for slow tasks
spark.speculation.quantile 0.9 Fraction of tasks in the stage that must finish before speculation is considered
spark.speculation.multiplier 3 A task is a candidate when it runs this many times longer than the median
spark.speculation.minTaskRuntime 100ms Minimum runtime before a task can be speculated

Spark 4.0 made speculation less aggressive: the defaults changed to multiplier=3 and quantile=0.9. To restore the older behaviour, set spark.speculation.multiplier=1.5 and spark.speculation.quantile=0.75.

Pitfalls

  • Speculation does not fix skew. A task that is slow because it has ten times more data will be equally slow on another executor; the copy just wastes resources.
  • Non-idempotent side effects. Two copies of a task may both write to an external system (a database, an API). Output committers protect file writes, but custom writes in foreachPartition must be idempotent.
  • Accumulators in transformations may count both copies; see the accumulators section.
  • Busy clusters: duplicate tasks compete for slots. Speculation helps most on large clusters with occasional bad nodes.

In interviews

“What is speculative execution and when would you turn it on?” Define it, give the trigger (quantile of tasks done, multiplier over the median), say it addresses hardware or node stragglers rather than data skew, and warn about non-idempotent writes. Knowing that the Spark 4.0 defaults became less aggressive shows up-to-date knowledge.

Practice questions

An executor has spark.executor.memory=10g. How much memory is available for execution and storage, and how much does the container request?

Unified pool = (10 GB − 300 MB) × 0.6 ≈ 5.8 GB, shared by execution and storage, of which about 2.9 GB is protected for storage. User memory is about 3.9 GB. The container request is heap plus overhead: by default max(384 MB, 10% of 10 GB) = 1 GB, so about 11 GB (plus spark.executor.pyspark.memory and off-heap memory if configured).

Size executors for 20 nodes with 32 cores and 256 GB each.

Reserve 1 core and about 8–16 GB per node: 31 cores and about 240 GB usable. At 5 cores per executor that is 6 executors per node (30 cores, one core spare). 240 / 6 = 40 GB per container; heap ≈ 40 / 1.1 ≈ 36 GB, overhead ≈ 3.6 GB. Total 20 × 6 = 120 executors, minus one for the driver or ApplicationMaster: 119 executors with 5 cores and 36g, giving 595 task slots. Then validate with GC time, spill and executor loss in the UI.

YARN kills containers with “exceeding physical memory limits” in a PySpark job, but the heap is far from full. What is going on?

The container limit covers heap plus overhead. Python workers, Arrow buffers and native memory live outside the heap, in overhead. Heavy pandas UDFs or large Python objects push past the limit. Increase spark.executor.memoryOverhead (or set spark.executor.pyspark.memory), possibly reducing the heap to keep the container size, and reduce per-task Python memory (smaller batches, fewer concurrent tasks per executor).

A stage shows Spill (Memory) 400 GB and Spill (Disk) 20 GB. Is that 400 GB written to disk?

No. Spill (Memory) is the deserialized in-memory size of the data at the moment it was spilled; Spill (Disk) is what was actually written after serialization and compression. The disk figure measures I/O; both indicate that execution memory was insufficient for the partitions processed.

Dynamic allocation is on, but the application never releases executors after a big stage. Why?

Likely causes: executors hold cached blocks (cachedExecutorIdleTimeout is infinite by default), shuffle tracking keeps executors whose shuffle output may still be read, or tasks are still pending. Unpersist cached data, set cachedExecutorIdleTimeout, use an external shuffle service where available, and check minExecutors.

Would speculative execution help a stage whose slowest task processes a hot key?

No. The task is slow because of its data, so a copy on another executor takes just as long and wastes a slot. Fix the skew (AQE skew join, broadcast, salting). Speculation helps with stragglers caused by slow or failing nodes.

Key takeaways

  • Executor heap = 300 MB reserved + 60% unified (execution and storage, half protected for storage) + 40% user memory; overhead sits outside the heap and holds Python workers.
  • Size executors by reserving node resources, choosing about 5 cores, splitting memory, subtracting overhead and the driver, then checking task slots and memory per task.
  • Container size is heap plus overhead; out-of-memory kills from the cluster manager usually need more overhead, not more heap.
  • Spark 4 runs on Java 17 with G1GC by default; measure GC time and reduce garbage before tuning flags.
  • Spill comes from partitions too big for execution memory: uniform spill means resize partitions or executors, concentrated spill means skew.
  • Dynamic allocation scales executors by backlog and idle time and needs a shuffle service or shuffle tracking; speculation re-runs node stragglers, with less aggressive defaults since Spark 4.0.

By DataDank Editorial · Last reviewed Oct 2026 · The memory-region and spill examples run on PySpark 4.2.0 in local mode (local[2], 1 GB driver heap); the spill example forces spills with an internal threshold so it happens on small data. Executor sizing, GC flags, dynamic allocation and speculation apply to real clusters and are not executed here; configuration names and defaults were checked against the Spark 4.2 runtime and documentation.

Progress is saved in this browser only. No account needed.

Search
Filter by type