Explore the future of AI-Native Data Management at Autonomous 26 | May 19 --> Save your spot
Acceldata recognized as an Exemplary Leader in 2026 ISG Buyers Guide™ for Data Quality and Data Observability. Read the Report→

Spark Executor Instances: How to Right-Size Them on Kubernetes

September 28, 2026
10 minutes

Key Takeaways

  • Executor sizing controls how much memory, CPU, and parallelism each Spark task gets, and a wrong number in either direction costs you something.
  • Over-provisioned executors quietly inflate cloud spend, while under-provisioned ones trigger garbage collection pressure, disk spill, or an outright kill.
  • Instances, cores per executor, and memory per executor are one sizing decision spread across three settings, not three knobs tuned in isolation.
  • A configuration that was correct in January can be wrong by June, since data volume and job logic keep changing while the config file does not.

A job that finished in 20 minutes last month is now stalling past an hour, and nobody touched the code. The data grew, the executors did not, and now half of them are spilling to disk while the other half sit idle.

Somewhere else, a cluster bill climbed last month without a single new pipeline to explain it. Both problems trace back to the same decision: Spark executor instances sized for a workload that no longer exists. Right-sizing them is a repeatable calculation, something you can run, check, and rerun as the workload moves.

What Does spark.executor.instances Control?

spark.executor.instances sets the number of executor processes Spark requests for a job, and on its own that number tells you almost nothing.

An executor's actual capacity depends on how many cores and how much memory each one gets, so instances, executor cores, and executor memory form one sizing decision spread across three configuration keys rather than three settings tuned in isolation.

Get the split wrong, and adding more instances without adjusting the other two can slow a job down instead of speeding it up, since more small executors mean more scheduling overhead and less memory for each task. Getting that split wrong is also exactly how sizing mistakes happen in the first place.

Why Executor Sizing Goes Wrong So Often

Executor sizing goes wrong in two directions at once, and most teams only watch for one of them.

Over-provisioning

A cluster running executors sized well above what any task actually uses does not throw an error; it just runs at a fraction of its potential utilization while the invoice keeps arriving in full. Nobody gets paged for a job that finishes on time using twice the compute it needed, so the waste compounds for months before it shows up in a cost review.

Under-provisioning

Under-provisioning shows up in specific, diagnosable ways. Executors starved of memory spend a growing share of their time in garbage collection instead of running task code, which slows every job on the cluster without crashing any of them outright.

Push memory pressure further, and Spark starts spilling shuffle data to disk, trading RAM it does not have for I/O it did not budget for. Push it further still, and the executor gets killed outright, forcing Spark to recompute lost partitions and cascade into retries across the whole stage.

Both directions trace back to the same starting point: a number picked without a method behind it.

How Do You Calculate a Starting Point for Executor Instances?

Calculate a starting point by dividing total usable cluster cores by the cores you plan to assign to each executor, after reserving headroom for node-level overhead. That calculation breaks into three steps:

  • Reserve roughly 1 core and 1 to 2 GB of memory per node for the OS, the node agent, and Kubernetes system pods before counting anything else as usable.
  • Divide the remaining usable cores by your chosen cores per executor, typically 4 to 5, to get a starting instance count.
  • Treat that number as a starting point to observe and refine, using the signals below to guide the adjustment.

Going meaningfully higher than 5 cores per executor concentrates too many tasks inside a single JVM, and each additional task adds pressure to the same heap, which pushes garbage collection pauses longer and more frequent.

What Signals Tell You Your Executors are the Wrong Size?

The clearest signals are exit code 137, rising garbage collection time, disk spill during shuffle, and long idle stretches between tasks. Each one points to a different fix:

Signal What it means
Exit code 137 The kernel's out-of-memory killer stepped in, and executors are undersized for the data they are handling.
Garbage collection time above 10-15% of task time Executors are spending real work capacity managing memory instead of running task code.
Disk spill during shuffle Spark ran out of memory to hold shuffle data and fell back to writing it to disk.
Long idle stretches between tasks Too many small executors are sitting mostly unused while a handful of tasks do the real work.

‍

Exit code 137 is one of the clearest markers that OOMKilled Spark executors are undersized, and this full set is exactly the kind of signals Spark on Kubernetes teams miss when they rely on pod restart counts alone instead of watching memory and shuffle behavior directly.

See how Acceldata's data and AI observability platform brings infrastructure, performance, and anomaly signals into one system instead of forcing a switch between separate dashboards for each one.

Watching for these signals gets harder once Kubernetes enters the picture, since the platform changes how each one shows up and where it has to be caught.

How Does Executor Tuning Change When Spark Runs on Kubernetes?

Executor tuning changes on Kubernetes because pod resource requests and limits replace the container-level allocations YARN used to manage, and there is no persistent cluster sitting idle to absorb a bad estimate.

Pod requests and limits

Pod requests and limits do the job container memory settings used to do, and setting them too tight gets an executor evicted by the kubelet before Spark ever sees a Java-level out-of-memory error.

Memory overhead for non-JVM workloads

spark.kubernetes.memoryOverheadFactor behaves differently depending on the workload: it defaults to 10% of executor memory for JVM-based jobs but jumps to 40% for non-JVM workloads such as Python or R processes, since those run outside the JVM heap and need more headroom reserved outside it.

Skipping that adjustment for a PySpark job is one of the most common causes of a pod that gets killed despite memory numbers that look correct on paper.

No idle cluster to fall back on

A YARN cluster with idle capacity would run an oversized job a little less efficiently, while a Kubernetes cluster provisions pods on demand, so an oversized request either fails to schedule or forces the autoscaler to add nodes it would not otherwise need, turning a sizing error into an immediate cost or scheduling problem instead of a hidden one.

This shift matters at scale. Kubernetes now runs in production for 82% of organizations using containers, up from 66% in 2023, according to the 2025 CNCF Annual Cloud Native Survey, and Spark workloads are increasingly part of that shift.

Getting comfortable with how to monitor Spark on Kubernetes is no longer optional for a platform team, since the signals that used to come from a YARN resource manager now have to be gathered differently. Kubernetes is not the only place the sizing math shifts underneath you. Spot capacity changes it again.

How Do Spot Instances Change the Sizing Decision?

Spot instances change the sizing decision because capacity can disappear mid-job with only a short warning, so sizing has to account for graceful executor loss, not just steady-state throughput.

Building executors for spot capacity comes down to a few practices:

  • Expect roughly two minutes' notice before a spot reclaim terminates an executor, enough time to stop accepting new tasks and checkpoint in-flight work.
  • Keep executors small enough that losing several at once does not stall the whole job waiting for their tasks to be recomputed elsewhere.
  • Pair spot instances with a smaller on-demand baseline so critical stages always have somewhere stable to land.
  • Avoid running large executors on spot capacity, since that concentrates more work behind every interruption.

The tradeoffs involved in running Spark on spot instances go deeper than sizing alone, but sizing is where most teams get the strategy wrong first.

Making Executor Right-Sizing a Continuous Practice with Acceldata

A configuration that is correct in January can be wrong by June, since data volume grows, job logic changes, and new pipelines get added to a shared cluster without anyone revisiting the executor math that used to work. Right-sizing Spark executor instances works best as an ongoing habit built on the same signals covered above, checked again every time the workload shifts rather than left to drift until something breaks.

Keeping that habit alive without turning it into a manual audit every quarter comes down to a few practices:

  • Recheck instance count, cores, and memory together whenever data volume or job logic changes materially, before a job fails outright rather than after.
  • Watch garbage collection time, disk spill, and OOMKilled exit codes continuously, since each one points to a different correction.
  • Treat a rising cluster bill with no new workload as a signal worth investigating rather than a cost anomaly to write off.

The signals worth watching rarely show up alone, which is also why four common Spark issues tend to compound into each other once executor sizing drifts.

Acceldata's xLake is built for that second habit: it correlates Spark application signals with Kubernetes pod events automatically, connecting OOMKills, evictions, and scheduling failures directly to the Spark job they affected, in one control plane, instead of the manual, tool-switching investigation those signals usually require.

Book a demo and see how xLake catches a mis-sized Spark executor before it becomes a production issue.

FAQs: Spark executor instances and right-sizing

What's a reasonable spark.executor.instances value to start with on a brand-new job?

For a brand-new job with no history to check against, start from the core-based calculation in this guide using 5 cores per executor, then run the first attempt at roughly 75% of the cluster's usable capacity rather than 100%, so there is headroom to observe behavior without immediately competing for every available core. Treat the first real run as a diagnostic pass, and adjust from the signals it produces.

Does adding more executor instances always make a Spark job faster?

No. Past a certain point, more instances mean more small executors coordinating over the network, and that coordination overhead can outweigh the parallelism gained, particularly for jobs with heavy shuffle. A job bottlenecked on shuffle or driver coordination often gets slower as instances increase, since the added overhead compounds until the underlying bottleneck is addressed directly.

How many cores should each Spark executor have?

Most teams do well starting at 4 to 5 cores per executor. That range balances enough parallelism per executor against the garbage collection cost of a single JVM managing too many concurrent tasks against one heap. Going much higher rarely improves throughput and tends to show up first as longer, more frequent garbage collection pauses, a more reliable signal to watch than any fixed cap on the number itself.

Can dynamic allocation replace manual executor right-sizing?

Dynamic allocation changes what needs sizing; it does not remove the need to size anything. With it enabled, Spark adds and removes executors within a floor and ceiling you still set, along with the same cores-per-executor and memory-per-executor decisions covered in this guide. It solves the problem of a fixed instance count sitting idle or maxed out for the entire life of a job, but the underlying per-executor math stays the same.

Does executor sizing work differently for streaming jobs than for batch jobs?

Yes, mainly around how steady the workload is. A batch job's resource needs are usually visible up front from the size of the input data, while a streaming job has to be sized for its peak load, since an executor fleet sized for average throughput falls behind during a burst and rarely catches back up on its own. Streaming jobs also benefit more from headroom in cores and memory than batch jobs do, since a backlog compounds instead of simply finishing later.

About Author

Shreya Bose

Shreya Bose has been writing for technology companies across software testing, developer tools, fintech, and recruitment technology for over six years. She has spent nearly three years at BrowserStack building the BrowserStack Guide, and has written for TestRail, 100ms (where her long-form strategy drove 11X organic traffic growth), and others.

LinkedIn: linkedin.com/in/shreya-bose-writes

Similar posts