Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. Executor sizing for a cluster

PySpark · Memory, OOM & Tuning

Executor sizing for a cluster

Hardpyspark-57
scenarioexecutor-sizingyarnresources

Question

You have 10 nodes with 16 cores and 64 GB each. How would you size executors?

Solution

Start from what each node can really offer, then pick a number of cores per executor, then divide. Do it out loud with the reasoning, because the numbers matter less than the logic.

Step by step

Each node has 16 cores and 64 GB.

  • Leave resources for the operating system and cluster daemons: 1 core and about 1 GB per node. That leaves 15 cores and 63 GB per node.
  • Choose cores per executor. A common choice is about 5. More than that often hurts I/O throughput to HDFS or object storage, because too many tasks share one process. Fewer means many small executors.
  • Executors per node: 15 cores divided by 5 = 3.
  • Total executors: 3 per node times 10 nodes = 30. On YARN, leave one slot for the application master (and the driver in cluster mode), so request 29.
  • Memory per executor: 63 GB divided by 3 = 21 GB per container. Overhead is max(384 MB, 10 percent) of the heap, so the heap is about 19 GB.
--num-executors 29
--executor-cores 5
--executor-memory 19g

With these, you have 145 task slots. A stage with 1,450 tasks runs in about 10 waves.

Fat versus tiny executors

  • Tiny (1 core, 4 GB each): lots of JVMs, each with its own overhead. Broadcast data is copied many times. Little memory per task.
  • Fat (15 cores, 63 GB): many tasks share one JVM heap, so GC pauses get long, and a single executor loss costs a lot.
  • Medium (about 5 cores) is a compromise that has worked well in practice.

Caveats to mention

The 5-core figure is a rule of thumb from older HDFS client behaviour, not a law. Check the actual workload. Memory-heavy jobs may need fewer cores per executor to give each task more memory. Python-heavy jobs need more overhead. If the cluster is shared, dynamic allocation with a maximum can be better than a fixed number. And check what YARN allows per container, since the scheduler may cap container size.

🎯 Put this concept into practice

Solidify this answer with real hands-on interview drills in the browser studio.

Open related drill →
PreviousNext