skip to content

How do you size --executor-cores, --executor-memory and --num-executors for a YARN cluster?

level: seniorimportance: should knowfreq 58%

answer

  1. the container is bigger than the flag you set
  2. two YARN ceilings decide what is grantable
  3. neither one huge JVM nor many tiny ones
  4. leave room for the OS, the agent and the driver
  5. slots only help if partitions exist to fill them

basics

~20 s

Size an executor so its container fits a node: requested memory is executor memory plus spark.executor.memoryOverhead, and must stay under yarn.scheduler.maximum-allocation-mb. Keep cores per executor moderate, reserve capacity for the OS, NodeManager and driver, then set the count from available cluster capacity.

solid answer

~50 s

Work from the node inward. A YARN container asks for `spark.executor.memory` **plus** `spark.executor.memoryOverhead` (default `max(384 MiB, 0.10 × executor memory)`, plus off-heap if enabled), and that total must be within both `yarn.scheduler.maximum-allocation-mb` and what a NodeManager advertises in `yarn.nodemanager.resource.memory-mb` — otherwise the job is rejected at submit with an "above the max threshold of this cluster" error. Choose cores per executor next: `--executor-cores` is the number of task slots per executor, and the common guidance is roughly 4–5 rather than one giant executor, because a very large heap brings long GC pauses and many concurrent tasks in one JVM contend on I/O. One-core executors are the other extreme: they duplicate JVM overhead and cannot share broadcasts or cached blocks. Then divide node capacity by executor size, leaving a core and a gigabyte or so per node for the OS and NodeManager and one container for the driver. Total slots is `--num-executors × --executor-cores`; size it against the work, not the cluster.

code

text · 4 lines
text
Required executor memory (65536 MB), offHeap memory (0) MB, overhead (6553 MB), and
PySpark memory (0 MB) is above the max threshold (57344 MB) of this cluster!
Please check the values of 'yarn.scheduler.maximum-allocation-mb' and/or
'yarn.nodemanager.resource.memory-mb'.

go deeper

for a junior

Know that --executor-cores is how many tasks an executor runs at once and that --num-executors times --executor-cores gives your task slots. Do not expect to pick production numbers yet.

for a middle

Explain the container arithmetic: executor memory plus overhead, checked against the node's advertised capacity and the scheduler's maximum allocation. Be able to say why the job is rejected before any task runs.

for a senior

Show a worked sizing from real node capacity, including reservations for the OS, the NodeManager and the driver container, and connect symptoms — GC time, idle slots, container kills, ACCEPTED-but-waiting — back to the sizing decision that caused them.

for a principal

Own the defaults and the guardrails: queue capacity, maximum allocation, a sane default executor shape, and a policy on when jobs must use dynamic allocation instead of claiming fixed capacity on a shared cluster.

## Two limits and one arithmetic Every executor is a YARN container, and YARN grants containers only within limits the cluster operator set. Two matter: - `yarn.nodemanager.resource.memory-mb` and `yarn.nodemanager.resource.cpu-vcores` — what one node offers Spark, which is deliberately less than the physical machine so the OS, the DataNode and the NodeManager keep something. - `yarn.scheduler.maximum-allocation-mb` and `yarn.scheduler.maximum-allocation-vcores` — the largest single container the scheduler will ever grant. Spark's request per executor is not `--executor-memory`. It is: ``` container = spark.executor.memory + spark.executor.memoryOverhead # default max(384 MiB, 0.10 x memory) + off-heap, if spark.memory.offHeap.enabled + PySpark memory, if configured ``` Ask for `--executor-memory 20g` and you are really asking for about 22 GB. If that exceeds the maximum allocation, Spark refuses before a single task runs, naming both YARN keys in the error. This arithmetic is the single most common cause of "my job will not even start". ## Cores per executor `--executor-cores` sets `spark.executor.cores`, which is how many tasks that executor runs concurrently — one task per slot, each task processing one partition. Two failure modes bracket the sensible range: - **Too fat.** One executor with all 16 cores of a node and a 100 GB heap runs sixteen tasks that share one JVM's garbage collector; full GCs stall all sixteen at once, and sixteen threads hammering the same disks and the same HDFS client saturate I/O. Widely used guidance caps this around four or five cores per executor. - **Too thin.** One core per executor means every task pays a full JVM's overhead, broadcast variables are replicated once per executor rather than shared by several tasks, and cached partitions cannot be reused across tasks in the same JVM. You also multiply the number of shuffle connections. Cores are also where the driver and ApplicationMaster hide: in cluster mode the driver occupies a container of its own, sized by `--driver-memory` plus `spark.driver.memoryOverhead`, and it must fit the same limits. ## A worked example Take nodes with 16 cores and 64 GB, where the operator advertises 15 vcores and 56 GB to YARN (the rest is OS, DataNode, NodeManager). - Pick 5 cores per executor → 3 executors per node (15 / 5). - Memory budget per executor: 56 GB / 3 ≈ 18.6 GB *including* overhead. - With the default 10% factor, `--executor-memory 16g` requests 16 + 1.6 ≈ 17.6 GB. That fits, with headroom. - Across 10 such nodes: 30 executors, 150 task slots. Reserve one executor slot's worth of room for the driver container, so `--num-executors 29` if you intend to fill the cluster. None of these numbers are laws; they are the shape of the calculation. Redo it against your own `yarn.nodemanager.resource.*` values. ## Sizing the count against the work, not the cluster `--num-executors × --executor-cores` gives concurrent task slots. Useful parallelism is bounded by the number of partitions a stage actually has: 150 slots and 40 partitions leaves 110 slots idle, while 150 slots and 20,000 tiny partitions wastes time on scheduling. Aim for each slot to get several tasks over the stage so stragglers can be absorbed, and remember that a shuffle stage's partition count comes from `spark.sql.shuffle.partitions` rather than from your executor count. Taking the whole cluster because it is there is anti-social on a shared platform; queue capacity in YARN exists precisely to stop it, and a job that requests more than its queue allows will simply wait. ## When to stop hand-sizing If the job's shape varies run to run — a backfill, a skewed daily volume — dynamic allocation replaces `--num-executors` with `spark.dynamicAllocation.minExecutors` / `maxExecutors` and lets Spark scale between them. Setting `spark.executor.instances` alongside dynamic allocation is a configuration conflict; express the ceiling as `maxExecutors` and the warm start as `initialExecutors` instead. ## Signals you sized wrong - Job never starts, error names `yarn.scheduler.maximum-allocation-mb` → container request too big. - Executors are killed by YARN for exceeding memory limits → overhead too small for off-heap, Python, or netty buffers. - Long GC time per task in the Spark UI's executors tab → heap too big for one JVM, or too many concurrent tasks in it. - Many executors sitting at zero active tasks → more slots than partitions. - Queue shows the application `ACCEPTED` for minutes → the request exceeds available queue capacity, not node capacity.

  • Why is a request for --executor-memory 20g rejected on a cluster whose maximum allocation is 21 GB?
    Because the container request adds overhead. The default is `max(384 MiB, 0.10 × executor memory)`, so 20 GB becomes roughly 22 GB, over the 21 GB ceiling, and Spark fails at submit naming `yarn.scheduler.maximum-allocation-mb` and `yarn.nodemanager.resource.memory-mb`. Either lower `--executor-memory` or raise the YARN limit if the node genuinely has the RAM.
  • How does --executor-cores relate to the number of tasks running at once?
    Each core is a task slot, so an executor runs up to `spark.executor.cores` tasks concurrently, each on one partition. Application-wide concurrency is `num-executors × executor-cores`. Raising cores adds concurrency inside a shared JVM — same heap, same garbage collector, same local disks — which is why more is not linearly better.
  • What changes about sizing when the job runs on Kubernetes rather than YARN?
    The arithmetic is the same but the ceilings move: the pod request is `spark.executor.memory` plus overhead, and it must fit a node's allocatable resources and any namespace ResourceQuota or LimitRange. Executor counts come from `spark.executor.instances` or dynamic allocation, and cores map to pod CPU requests, so fractional CPU settings need care.

saying these in an interview costs you the question

  • Thinks the YARN container equals --executor-memory exactly
  • Gives one executor every core on the node
  • Sets one core per executor to maximise parallelism
  • Ignores the driver container when dividing node capacity
  • Sizes executors from cluster size rather than from the work

context