skip to content

Cluster Execution and Tuning

How a Spark application maps onto executors, cores and memory once it leaves your laptop, and which knobs actually help when it does not fit. Tuning questions are how interviewers separate people who ran Spark in production from people who read about it.

on this pageshow

explore

questions

18

In spark-submit, what do the --master, --class and application-jar arguments specify?

level: juniorimportance: must knowfreq 55%

answer

  1. one launcher for every cluster manager
  2. where to run, what to start, what to ship
  3. local[*], yarn, k8s://, spark://host:7077
  4. fully-qualified class holding main()
  5. code beats flags beats defaults file

basics

~20 s

In spark-submit, --master names the cluster manager and its address (local[*], yarn, k8s://https://host:6443, spark://host:7077), --class is the fully-qualified class holding main() inside the jar, and the trailing application jar is the code Spark ships to the cluster.

solid answer

~40 s

`spark-submit` is the single launcher for every Spark deployment target, and three arguments carry the essentials. `--master` selects the cluster manager and where it lives: `local[*]` for a single JVM, `yarn`, `k8s://https://api-server:6443`, or `spark://host:7077` for a Spark standalone master. `--class` is the fully-qualified name of the class whose `main` method starts the application — Scala/Java only; for PySpark you pass the `.py` file instead and ship dependencies with `--py-files` or `--archives`. The last positional argument is the application jar or script, and anything after it is handed to your `main` as program arguments. Around those sit `--deploy-mode`, `--executor-memory`, `--executor-cores`, `--num-executors`, `--jars`, `--files` and arbitrary `--conf key=value` pairs. Precedence is worth memorising: properties set on a `SparkConf` in code beat `spark-submit` flags, which beat `conf/spark-defaults.conf`.

code

bash · 9 lines
bash
spark-submit \
  --class com.example.DailyRollup \
  --master yarn \
  --deploy-mode cluster \
  --executor-memory 8g \
  --executor-cores 4 \
  --num-executors 20 \
  --conf spark.sql.shuffle.partitions=400 \
  rollup-1.4.0.jar 2026-08-01

go deeper

for a junior

Be ready to write a working spark-submit line from memory: master, deploy mode, class, sizing flags, then the artifact. Know that arguments after the jar belong to your application, not to Spark.

for a middle

Explain what submission actually stages to the cluster and why --jars, --files and --py-files exist. Know the precedence order between code, flags and spark-defaults.conf, and the client-mode driver-memory exception.

for a senior

An interviewer expects you to standardise submission for a team: a properties file for cluster-wide defaults, artifacts fetched from a repository rather than a laptop, and dependency resolution that does not hit Maven Central from every job.

for a principal

Own the argument surface as an interface. Decide what the platform fixes centrally versus what teams may override, and treat --packages resolution at submit time as a build-time concern to move into the artifact instead.

## One launcher, every target `spark-submit` (in `$SPARK_HOME/bin`) is the only supported way to start a Spark application, and it is deliberately the same command whether the job runs in a single local JVM, on a Hadoop YARN cluster, on Kubernetes, or on a Spark standalone cluster. It builds the driver's classpath, resolves dependencies, uploads what needs uploading, and asks the chosen cluster manager to start the application. Learning its argument shape is the everyday skill; everything else about Spark deployment is configuration layered on top. The canonical invocation looks like this: ```bash spark-submit \ --class com.example.DailyRollup \ --master yarn \ --deploy-mode cluster \ --executor-memory 8g \ --executor-cores 4 \ --num-executors 20 \ --conf spark.sql.shuffle.partitions=400 \ s3://artifacts/rollup-1.4.0.jar 2026-08-01 ``` ## --master: which cluster manager, and where `--master` is a URL-shaped string naming the resource manager that will hand Spark its containers: - `local`, `local[4]`, `local[*]` — no cluster at all; driver and "executors" are threads in one JVM. This is what unit tests use. - `yarn` — the address comes from the Hadoop configuration (`HADOOP_CONF_DIR`/`YARN_CONF_DIR`), which is why you do not write a host here. - `k8s://https://api-server:6443` — the Kubernetes API server endpoint; Spark creates a driver pod and executor pods. - `spark://host:7077` — a Spark standalone master (7077 is its default port). If you omit `--master`, Spark falls back to `spark.master` from the properties file, and many gateway installations set that for you. Getting `--master` wrong is the classic first-day error: the job silently runs `local[*]` on the edge node and takes forever on data it can barely read. ## --class and the application artifact For a JVM application you package a jar and point `--class` at the fully-qualified class containing `def main(args: Array[String])` (Scala) or `public static void main(String[])` (Java). Spark does not scan the jar for a main class unless the jar's manifest declares one, so the flag is normally required. The jar itself is the final positional argument, and it may be a local path, an HDFS path, or an object-store URI — Spark distributes it to the cluster as part of submission. PySpark inverts this: there is no `--class`. You pass `app.py` as the positional argument, add pure-Python dependencies with `--py-files deps.zip`, and ship a whole packed interpreter with `--archives venv.tar.gz#environment` when the cluster's Python does not match yours. Any JVM connectors a Python job needs (a Kafka or JDBC source, say) still arrive through `--jars` or `--packages` Maven coordinates. Everything after the artifact is application arguments, not Spark arguments. `spark-submit ... app.jar --date 2026-08-01` passes `--date 2026-08-01` to your `main`; putting it before the jar makes `spark-submit` reject it. ## The supporting flags A handful of flags cover most real submissions: - `--deploy-mode client|cluster` — where the driver process itself runs. - `--driver-memory`, `--executor-memory`, `--executor-cores`, `--num-executors` — sizing. - `--jars a.jar,b.jar` — extra jars added to the classpath of **both** driver and executors and distributed for you. - `--packages group:artifact:version` — Maven coordinates resolved at submit time. - `--files config.yaml` — files copied into each executor's working directory. - `--conf spark.some.key=value` — any Spark property; the escape hatch for everything without a dedicated flag. - `--properties-file` — override the default `conf/spark-defaults.conf`. ## Configuration precedence Three layers can set the same property, and they resolve in this order, strongest first: values set programmatically on a `SparkConf`/`SparkSession.builder`, then `spark-submit` flags and `--conf`, then the defaults file. The important exception is anything that configures the driver JVM itself — `spark.driver.memory`, `spark.driver.extraJavaOptions` — in client mode: that JVM is already running by the time your code executes, so the value must come from the command line or the properties file to have any effect. ## What submission actually does `spark-submit` runs a small JVM (`SparkSubmit`) that parses the arguments, resolves `--packages` through Ivy, and then either starts your driver in that same process (client mode) or asks the cluster manager to start it elsewhere (cluster mode). In both cases the application jar, `--jars`, `--files` and `--py-files` are staged so executors can fetch them. Understanding that staging step explains most "class not found on the executor" failures: a path that exists on your laptop means nothing to a container three racks away unless you told `spark-submit` to distribute it.

  • How do you submit a PySpark application, given there is no jar and no main class?
    Pass the `.py` file as the final positional argument instead of a jar and drop `--class` entirely. Pure-Python dependencies go in `--py-files deps.zip`; a packed virtual environment goes in `--archives venv.tar.gz#env` with `spark.pyspark.python` pointed inside it; JVM connectors still come from `--jars` or `--packages`. Program arguments still follow the script path.
  • If the same property is set in spark-defaults.conf, in a --conf flag, and on a SparkConf in code, which wins?
    Code wins, then the `spark-submit` flag, then the defaults file. The exception is driver-JVM settings such as `spark.driver.memory` in client mode: that JVM has already started before your code runs, so setting it in a `SparkConf` is silently ignored and you must use `--driver-memory`.
  • What does --jars do that simply adding a jar to the driver's classpath does not?
    `--jars` both adds the jar to the driver and executor classpaths and distributes the file to every executor's working directory. Putting a jar only on the submitting machine's classpath leaves executors without it, which surfaces as `ClassNotFoundException` inside tasks while the driver plans the job happily.

saying these in an interview costs you the question

  • Thinks spark-submit only works with Hadoop YARN
  • Passes --class when submitting a PySpark script
  • Confuses --master with --deploy-mode
  • Assumes a local jar path is visible to executors automatically
  • Believes spark.driver.memory set in code always takes effect

context

open as a page

In Spark, what is the difference between a job, a stage and a task?

level: juniorimportance: must knowfreq 82%

basics

~10 s

In Spark, an action submits a job; the driver cuts that job into stages at shuffle boundaries; each stage runs one task per partition. Tasks are the smallest unit executors actually execute.

open as a page

In Spark, how do the MEMORY_ONLY and MEMORY_AND_DISK persist levels differ when a partition will not fit?

level: juniorimportance: must knowfreq 70%

basics

~20 s

With MEMORY_ONLY, a partition that does not fit in storage memory is simply not cached and is recomputed from lineage the next time it is needed. With MEMORY_AND_DISK, that partition is written to the executor's local disk and read back instead of recomputed.

open as a page

In spark-submit, what is the difference between --deploy-mode client and --deploy-mode cluster?

level: middleimportance: must knowfreq 78%

basics

~20 s

--deploy-mode decides where the Spark driver runs. In client mode the driver is the spark-submit process on the submitting machine; in cluster mode the driver runs inside the cluster, in a YARN ApplicationMaster container or a Kubernetes driver pod.

open as a page

In Spark, what determines where the DAG scheduler places a stage boundary?

level: middleimportance: must knowfreq 68%

basics

~10 s

Spark cuts a stage at every wide dependency: any point where a row's destination depends on its key, so data must be redistributed. Narrow operations are pipelined into the current stage instead.

open as a page

In a Spark executor, what does spark.memory.fraction control and how do execution and storage share it?

level: middleimportance: must knowfreq 62%

basics

~20 s

spark.memory.fraction (default 0.6) sizes the unified pool for execution and storage out of the executor heap left after a fixed 300 MB reservation. Execution and storage borrow from each other on demand; only storage's protected floor is safe from eviction.

open as a page

In Spark, how many tasks can one executor run at once, and what sets that limit?

level: middleimportance: should knowfreq 60%

basics

~10 s

An executor runs spark.executor.cores divided by spark.task.cpus tasks concurrently — one task per slot, each in its own thread of the same JVM. Total application parallelism is that number times the executor count.

open as a page

In the Spark UI, one task reports Spill (Memory) of 48 GB on an executor with an 8 GB heap — what is that measuring?

level: middleimportance: should knowfreq 45%

basics

~20 s

Spill (Memory) is the cumulative in-memory, deserialized size of every buffer the task flushed to disk, summed over all spill events. It is not a peak, so a task that spills repeatedly can report far more than the executor heap ever held.

open as a page

In cluster mode on YARN, where do you find the Spark driver's stack trace after a failed job?

level: seniorimportance: should knowfreq 48%

basics

~20 s

In YARN cluster mode the driver runs inside the ApplicationMaster container, so its stdout and stderr are container logs. Fetch them with yarn logs -applicationId <appId>, or open the ApplicationMaster log link in the ResourceManager UI.

open as a page

Why does spark.dynamicAllocation.enabled require an external shuffle service or shuffle tracking?

level: seniorimportance: should knowfreq 52%

basics

~20 s

Dynamic allocation removes idle executors, but an executor's local disk holds the shuffle files later stages must fetch. Killing it destroys that map output, so Spark needs an external shuffle service to serve the blocks after exit, or shuffle tracking to refuse to remove executors still holding them.

open as a page

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

level: seniorimportance: should knowfreq 58%

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.

open as a page

In Spark, which job patterns push work onto the driver and stall the application?

level: seniorimportance: should knowfreq 52%

basics

~20 s

Anything that pulls data or decisions back to one JVM: collect and toPandas, huge broadcasts, listing millions of input files, and stages with tens of thousands of tiny tasks. Executors idle while the single driver works.

open as a page

In Spark, how does a task retry differ from a stage retry after FetchFailedException?

level: seniorimportance: should knowfreq 44%

basics

~20 s

A task retry re-runs one failed task on another executor, up to spark.task.maxFailures (default 4). A FetchFailedException means map output is gone, so Spark fails the stage and resubmits the parent stage to regenerate it.

open as a page

Why does YARN kill a Spark executor for exceeding memory limits when its JVM heap never filled up?

level: seniorimportance: should knowfreq 58%

basics

~20 s

The container limit covers the whole process, not just the heap: spark.executor.memory plus spark.executor.memoryOverhead. Off-heap consumers — shuffle network buffers, JVM metaspace and thread stacks, native libraries and Python workers — burst past the overhead allowance while the heap sits half empty.

open as a page

For a new shared Spark platform, how would you choose between YARN and Kubernetes as the cluster manager?

level: principalimportance: should knowfreq 32%

basics

~20 s

Choose on the surrounding estate, not on Spark. YARN wins where HDFS, Kerberos and queue-based multi-tenancy already exist; Kubernetes wins where dependencies ship as per-job container images and isolation comes from namespaces, at the cost of replacing the external shuffle service.

open as a page

When several downstream Spark jobs read one expensive DataFrame, how do you choose between persist() and writing it out?

level: principalimportance: should knowfreq 38%

basics

~20 s

Scope decides it. persist() reuses data only within one Spark application and dies with its executors; writing the result to Parquet makes it readable by every later job, survives failures, and keeps executor memory for execution. Cache inside an application, materialize across them.

open as a page

In Spark, what does spark.speculation do, and when does it make a job worse?

level: seniorimportance: nice to knowfreq 36%

basics

~20 s

It relaunches a duplicate of any task running far longer than its stage's median, and keeps whichever attempt finishes first. It is off by default and wastes slots when the slowness is skew rather than a bad node.

open as a page

What does setting spark.memory.offHeap.enabled to true change about a Spark executor's memory?

level: seniorimportance: nice to knowfreq 28%

basics

~20 s

It gives Spark a second memory pool of spark.memory.offHeap.size bytes, allocated outside the JVM heap and used by Tungsten for execution and by the OFF_HEAP storage level. That memory is not garbage collected, so large working sets stop driving GC pauses.

open as a page