All posts

Top Apache Spark interview questions, with answers

Two hundred questions from first principles to production scenarios, each answered in depth with the mechanism, an example and the trap. Grouped by topic and collapsed, so you can drill one section at a time.

206 min read Spark

TL;DR

  • 200 questions in 11 sections, each answered properly: the direct answer, the mechanism underneath, a worked example, and the trap. Click any question to open it.
  • The last 15 are scenario questions: a symptom, and how to reason to the cause. Practise those out loud, because they are what separate people who have run Spark in production.
  • Where a default surprises people the number is given, because a confidently wrong default is worse in an interview than an honest “I would look it up”.
  • The answers say where the textbook version is wrong: AQE skew handling needs a partition over 256 MB before it fires, input partitions do not come from file count, accumulators double-count outside actions, and coalesce(1) can make an entire pipeline single-threaded.
  • No speedup ratios anywhere. Where performance is the point, the answer names the mechanism and leaves the measuring to you.

Interview lists tend to collect questions without checking the answers, so they repeat things that were true in Spark 1.6 and quote defaults that changed three releases ago. They also tend to answer in two lines, which is enough to recognise a term and not enough to defend it when someone asks “why”.

This one is built the other way round. The answers came first, from source and documentation at v4.2.0, and each one is written to survive a follow-up question: what the mechanism is, what the numbers actually are, and where the common version of the answer is wrong.

Use it two ways. Read a section end to end to find the gaps in your own model, or open a single question when you want the defensible version of something you already half know. The sections run roughly from basics to advanced, and the scenario section at the end is the one worth rehearsing aloud.

A note on what is deliberately absent: there are no benchmark numbers or speed multipliers. Where performance is the point, the answer explains the mechanism that makes one path cheaper and leaves the measuring to you, on your data and your cluster.

How this is organised

Section Questions Focus
Spark fundamentals 1-20 Vocabulary and the execution model
RDDs 21-40 Lineage, partitioning, shared variables
DataFrames, Datasets and Spark SQL 41-62 Catalyst, plans, formats
Architecture and execution 63-84 Driver, executors, schedulers
Partitioning and shuffle 85-104 The main source of slowness
Memory, caching and persistence 105-122 The unified memory model
Joins, Catalyst and adaptive execution 123-142 Strategy selection and AQE
Performance tuning and troubleshooting 143-160 Diagnosis before knobs
Structured Streaming 161-175 Watermarks, state, guarantees
Deployment, cluster managers and operations 176-185 Running it for real
Scenario-based questions 186-200 Symptom to cause

The diagrams render when you open a question, so nothing is drawn until you ask for it.

Spark fundamentals

Start here. These come up in the first ten minutes of almost every interview, and the answers below give you the version you can defend rather than the version everyone recites.

1. What is Apache Spark?

Apache Spark is a distributed engine for large-scale data processing. You describe a computation over a dataset, and Spark splits both the data and the work across a cluster of machines, coordinating the whole thing from a single process called the driver.

Three properties define it:

  • Distributed execution. A dataset is divided into partitions, and each partition is processed independently, in parallel, on whichever machine has capacity.
  • Lazy evaluation. Operations build a plan rather than running immediately, which lets the engine optimise the whole computation before any of it starts.
  • Fault tolerance through lineage. Spark records how each partition was derived, so a lost one is rebuilt rather than the job being restarted.

It ships APIs for Scala, Java, Python and R, plus libraries for SQL, structured streaming, machine learning and graph processing, all on the same execution engine.

The framing that lands well in an interview: Spark is a compute engine, not a storage system and not a table format. It reads from HDFS, S3, JDBC or Kafka, and writes back out. Confusing Spark with Hudi, Iceberg or Delta is a common mistake, and being crisp about the boundary is a good signal.

2. Why use Spark instead of a single machine?

Because one of two things has become true: the data no longer fits the memory, CPU or disk of one machine, or it fits but takes unacceptably long to process serially.

Spark’s answer to both is partitioning. A 1 TB dataset might become 8,000 partitions of roughly 128 MB, and if the cluster has 200 cores then 200 of them are processed at a time.

1 TB input
     |
     v  split into partitions
+------+------+------+     +--------+
|  P1  |  P2  |  P3  | ... |  P8000 |
+------+------+------+     +--------+
     |                          |
     v   processed in parallel  v
        across 200 cores
     |
     v
   result

The part candidates usually miss is the cost. Distribution is not free: there is serialization, network transfer during shuffles, scheduling overhead per task, and the coordination cost of the driver. For a dataset that fits comfortably on one machine, a single-process tool such as pandas, DuckDB or Polars will routinely beat a Spark cluster, because Spark’s overhead is not amortised.

A good answer names the crossover point rather than claiming Spark is always faster. Spark earns its keep when the work genuinely exceeds one machine, or when you need fault tolerance across a long-running job.

3. What is the difference between Spark and Hadoop MapReduce?

The structural difference is where intermediate results live and how much of the computation the engine can see at once.

MapReduce runs one map phase and one reduce phase, and writes the intermediate result to disk, usually replicated, between them. A multi-step algorithm becomes a chain of MapReduce jobs, each re-reading the previous one’s output from disk.

Spark keeps intermediate data in memory where it can, and expresses an entire chain of operations as one plan before running any of it. That gives it two advantages MapReduce structurally cannot have:

  MapReduce Spark
Intermediate data Written to disk between stages Kept in memory where possible
Unit of work One map plus one reduce An arbitrary DAG of stages
Optimisation scope Per job The whole plan, via Catalyst
Iterative workloads Re-reads input each pass Can cache and reuse
APIs Map and reduce SQL, DataFrames, streaming, ML, graph

The iterative case is where the gap is widest. Ten passes of an algorithm in MapReduce means ten round trips through disk; in Spark the dataset can be cached once and reused.

The honest framing for an interview: the advantage is in-memory intermediate data plus whole-plan optimisation, not a fixed speed multiplier. The often-quoted “100x” comes from a specific iterative benchmark and does not describe a typical batch job. Quote a number only if you measured it on your own data, and otherwise describe the mechanism.

4. What is lazy evaluation?

Transformations do not compute anything when you call them. Each one records what to do and returns a new dataset describing the result. Nothing executes until an action asks for a value.

The RDD programming guide states it directly:

All transformations in Spark are lazy, in that they do not compute their results right away. Instead, they just remember the transformations applied to some base dataset.

Why it matters is the important half. Because Spark holds the whole chain before running it, Catalyst can rewrite it: push a filter down into the scan so fewer rows are read, drop columns nobody selected, fuse adjacent operations into one generated loop, and choose a join strategy knowing what comes next. None of that is possible for an engine that executes each call as it arrives.

flowchart TB
  T["transformations called:<br/>read, filter, withColumn, groupBy"]
  T --> P[("plan recorded<br/><b>0 jobs, 0 tasks</b>")]
  P -.->|"nothing has run"| W["cluster is idle"]
  A["an action is called:<br/>count()"] --> R["<b>1 job</b>: Catalyst optimises<br/>the whole plan, then runs it"]
  P --> R

There are two consequences to be ready for:

  • A dataset is recomputed every time an action needs it, unless you persist it. The guide is explicit: “each transformed RDD may be recomputed each time you run an action on it.”
  • Errors surface at the action, not at the transformation. A typo in a column name in a filter may only raise when you call count(), which is why analysis errors seem to appear far from their cause.
5. What is the difference between a transformation and an action?

A transformation creates a new dataset from an existing one and runs nothing. An action returns a value to the driver or writes to storage, and that is what submits a job.

  Transformation Action
Examples map, filter, select, groupBy, join, repartition count, collect, show, first, take, write, foreach
Returns A new dataset A value to the driver, or a write
Triggers execution No Yes, one job

Transformations divide further, and the distinction drives everything about performance:

  • Narrow transformations, where each input partition feeds at most one output partition, so the work pipelines inside a single task: map, filter, mapPartitions.
  • Wide transformations, where an output partition draws from many input partitions, so data must move across the network: groupByKey, join, distinct, repartition.

The count of actions tells you how many jobs you will run, and the count of wide transformations tells you how many stages each job will have.

One trap worth naming: show() and take(n) are actions, so calling them during development runs the whole chain up to that point. A notebook with ten show() calls has run ten jobs, and if nothing is cached, much of the work is repeated each time.

6. How can you prove that transformations are lazy?

Count the jobs. Chaining transformations submits none; the action submits one.

sc = spark.sparkContext
rdd = sc.parallelize(range(4), 2)

doubled = rdd.map(lambda x: x * 2).filter(lambda x: x > 2)
print(len(sc.statusTracker().getJobIdsForGroup() or []))   # 0

doubled.count()
print(len(sc.statusTracker().getJobIdsForGroup() or []))   # 1

Two transformations, zero jobs. One action, one job.

A second demonstration is more visual: put a side effect in a transformation and watch when it happens.

acc = sc.accumulator(0)
counted = sc.parallelize(range(100), 4).map(lambda x: (acc.add(1), x)[1])

print(acc.value)      # 0, the map has not run
counted.count()
print(acc.value)      # 100, it ran when the action asked

The same trick exposes the other half of laziness. Call count() a second time and the accumulator reads 200, for the same 100 rows, because the uncached dataset was recomputed from scratch. Add .cache() and it stays at 100.

In an interview, the accumulator version is the stronger answer, because it demonstrates laziness and recomputation in one go.

7. What is a partition?

A partition is the unit of parallelism: a slice of the dataset that one task processes on one core. A stage runs exactly one task per partition, so the partition count is the maximum parallelism available to that stage.

flowchart TB
  subgraph S["one stage"]
    P1["partition 0"] --> T1["task 0"]
    P2["partition 1"] --> T2["task 1"]
    P3["partition 2"] --> T3["task 2"]
    P4["partition 3"] --> T4["task 3"]
  end
  T1 --> E["executors run<br/>as many as they have slots"]
  T2 --> E
  T3 --> E
  T4 --> E

Three numbers follow from this, and almost every tuning question is one of them:

  • How many partitions. Too few and cores sit idle; too many and per-task scheduling overhead dominates. A useful target is partitions in the low hundreds of megabytes, with several tasks per core so stragglers average out.
  • How evenly the data is spread. One partition holding most of the rows is skew, and because a task cannot be split, that one task sets the stage duration.
  • Where they come from. On read, from file splitting; after a shuffle, from spark.sql.shuffle.partitions or adaptive coalescing.

Check the number rather than assume it:

print(df.rdd.getNumPartitions())
8. What is `SparkSession` and how does it relate to `SparkContext`?

SparkSession is the entry point for the DataFrame and SQL APIs, and it is what modern code creates. SparkContext is the older, lower-level entry point for RDDs and for cluster coordination.

from pyspark.sql import SparkSession

spark = (SparkSession.builder
    .appName("trips-ingest")
    .config("spark.sql.shuffle.partitions", "400")
    .getOrCreate())

sc = spark.sparkContext          # the SparkContext inside it
  SparkSession SparkContext
Introduced 2.0 1.0
Gives you DataFrames, SQL, catalog, streaming RDDs, broadcasts, accumulators
Relationship Wraps one SparkContext One per application

A SparkSession wraps exactly one SparkContext, and one application has one SparkContext. getOrCreate() returns the existing one rather than making a second, which is why it is the idiomatic call.

Two details interviewers probe:

  • Creating the session is what starts negotiating for executors. Before it, nothing has been requested from the cluster manager.
  • newSession() gives an isolated SQL session, with its own temporary views and SQL configs, sharing the same SparkContext and therefore the same executors. That is how a multi-tenant server such as the Thrift server isolates users.
9. What languages does Spark support, and does the choice affect performance?

Scala, Java, Python and R, with a community-maintained C# binding outside the project. For DataFrame and SQL work the choice barely affects performance, because operations in any language compile to the same logical plan and execute in the JVM.

It matters in two places:

  • RDD code, where the logic itself is a closure in your language. A Python RDD map ships every row to a Python worker process and back.
  • Python UDFs, which sit inside a plan that is otherwise all JVM. Rows must be serialized out to Python and results serialized back, and because Catalyst cannot see inside the function, it cannot push filters through it or fuse it into generated code.
# stays in the JVM, fully optimisable
df.filter(F.col("fare_amount") > 50)

# leaves the JVM per batch, breaks codegen across this point
@F.pandas_udf("double")
def adjust(v: pd.Series) -> pd.Series:
    return v * 1.1

The ordering to remember, best first: built-in functions and SQL, then pandas (vectorised) UDFs which amortise the crossing over an Arrow batch, then plain Python UDFs which pay it per row.

Scala and Java additionally get the typed Dataset API. Python does not, since it has no compile-time types to check.

10. What is a job, a stage and a task?

Three nested units of work, each created by a different component.

Term Definition Created by
Job All the work for one action The action call
Stage A group of tasks with no shuffle between them DAGScheduler, cut at shuffle boundaries
Task One partition of one stage, on one core TaskSchedulerImpl
flowchart TB
  A["action: write()"] --> J["1 job"]
  J --> S0["stage 0<br/>read, filter, partial aggregate"]
  S0 -->|"shuffle"| S1["stage 1<br/>final aggregate, write"]
  S0 --> T0["8 tasks<br/>one per input partition"]
  S1 --> T1["200 tasks<br/>one per shuffle partition"]

Read it as a hierarchy: a job has one or more stages, and a stage has one task per partition. The arithmetic is worth internalising, because it is how you predict what the UI will show:

  • Jobs equals the number of actions you called.
  • Stages per job equals the number of shuffles, plus one.
  • Tasks per stage equals that stage’s partition count, unrelated to how many executors you have.

A nuance that catches people: a single write can produce more than one job, because some writers run extra jobs for committing or for computing statistics. “One action, one job” is the right model, with that exception noted.

11. What is a DAG in Spark?

A directed acyclic graph of dependencies between datasets. Directed because every edge points from an input to an output; acyclic because nothing can depend on itself, which is what makes the graph safe to walk.

Spark builds it as you chain transformations, and DAGScheduler uses it for two jobs:

  1. Cutting stages. It walks the graph backwards from the action and starts a new stage at every shuffle boundary, because everything narrow can pipeline into the stage it is already in.
  2. Recovery. Because each node records how it was derived, a lost partition can be rebuilt by replaying just the path that produced it.
flowchart TB
  A["read trips"] --> B["filter"]
  C["read cities"] --> D["filter"]
  B --> J["join"]
  D --> J
  J --> K["groupBy city"]
  J --> L["groupBy rider"]
  K --> M["write by city"]
  L --> N["write by rider"]

The DAG is the reason Spark is not restricted to a map phase and a reduce phase: any shape of dependency graph is expressible, so a job can branch, rejoin and have many stages.

You can print it for an RDD with toDebugString(), and for a DataFrame the physical plan from explain("formatted") is the equivalent. In the UI, the DAG visualisation on the job page shows the same thing graphically, with each box a stage.

12. What is the driver?

The driver is the process that holds the SparkSession and makes every decision. In order, it:

  1. Creates the SparkSession, and with it the SparkContext, SparkEnv and the RPC endpoints executors will call back on.
  2. Requests executors from the cluster manager through SchedulerBackend.
  3. Turns your code into a logical plan, has Catalyst optimise it, and selects a physical plan.
  4. Cuts the plan into stages in DAGScheduler, at every ShuffleDependency.
  5. Places individual tasks on executors via TaskSchedulerImpl and TaskSetManager, preferring the best available locality.
  6. Tracks every attempt, retrying failures up to spark.task.maxFailures, default 4.
  7. Maintains BlockManagerMaster, the index of every cached and shuffle block in the cluster.
  8. Receives results, bounded in total by spark.driver.maxResultSize, default 1g.
  9. Serves the UI on port 4040 and writes the event log the history server replays.

What it does not do is process your data. It is a coordinator.

Two consequences follow, and both are common interview questions:

  • spark.driver.memory defaults to 1g, which is a coordinator’s budget. Anything that pulls data back, collect() or toPandas(), needs more, or needs not to do that.
  • The driver is the one process with no redundancy. The plan, the scheduler state and the block index exist nowhere else, so if the driver dies the application dies, whereas losing an executor costs only recomputable work.
13. What is an executor?

An executor is a JVM launched for one specific application that runs tasks and stores blocks. The Spark source describes it in a line: a “Spark executor, backed by a threadpool to run tasks”.

Its lifecycle:

  1. Registers. CoarseGrainedExecutorBackend sends RegisterExecutor and waits for RegisteredExecutor before accepting work. Until that round trip finishes, a granted container is doing nothing.
  2. Accepts LaunchTask messages and runs each in a pool thread. Concurrency is spark.executor.cores, which defaults to 1 on YARN and to every core on the machine in standalone mode.
  3. Reads input, runs the generated code, and for a shuffle map stage writes shuffle output to local disk.
  4. Reports back with StatusUpdate per task and a heartbeat every spark.executor.heartbeatInterval, default 10s. Miss enough and the driver declares it lost.
  5. Stores and serves blocks through its BlockManager, which is what makes a cached DataFrame reusable and lets the next stage fetch shuffle data from its peers.

The contract is deliberately narrow: an executor never makes a scheduling decision, never sees the query plan, and never coordinates with another executor about what to run. It receives tasks, runs them, and serves bytes.

That narrowness is what makes executors disposable. Lose one and its tasks are retried elsewhere and its cached partitions are recomputed from lineage. Executors are also per-application: two applications never share one, which is why memory is reserved per application rather than pooled.

14. What is a cluster manager, and which ones does Spark support?

A cluster manager is a separate service whose only job is to decide who gets containers. At Spark 4.x the supported options are Standalone, YARN and Kubernetes.

Cluster manager What grants the container What Spark runs there
Standalone Master, tracking registered Worker daemons The worker launches an executor JVM per grant
YARN The ResourceManager, asked by Spark’s ApplicationMaster through YarnAllocator An executor per allocated container
Kubernetes The API server. The driver’s own ExecutorPodsAllocator creates the pods An executor per pod

Mesos support was removed in the 4.x line. At v4.2.0 the resource-managers/ directory in the source tree contains only kubernetes and yarn, so the three rows above are the whole list.

The division worth internalising is that the cluster manager never schedules your tasks. It hands out containers; executors register back with the driver; and from that point the driver drives everything. That single fact explains a lot of otherwise puzzling behaviour:

  • Why an under-provisioned driver throttles a large cluster.
  • Why YARN’s own scheduler settings change what you are given but never how your tasks are distributed.
  • Why Kubernetes needs no Spark daemon at all: the driver talks to the API server directly.

Kubernetes is the clearest illustration of the split, because the component granting resources and the component scheduling work are visibly different things.

15. Is Spark a master-slave architecture?

In the general sense yes: one process coordinates and many do the work. But Spark’s own vocabulary is worth getting exactly right, because “master” is overloaded and people conflate three different things with it.

Term What it actually is What it is not
Driver The coordinator: holds the plan, schedules every task Not called the master anywhere in Spark
--master A URL naming the cluster manager (yarn, k8s://…, local[*]) Not a machine, and not the driver
Standalone Master A daemon that tracks workers and grants resources Standalone mode only; YARN and Kubernetes have their own
Worker A machine, or the standalone daemon on it Not an executor; a worker hosts executors
Executor A JVM doing the work Not a machine

Spark has also moved away from the older terminology in its own scripts. At v4.2.0 there is sbin/start-workers.sh and conf/workers.template, while start-slaves.sh and slaves.template are gone, and the standalone documentation speaks of “a master and workers”.

Where the analogy breaks is the interesting part, and it is a strong thing to volunteer. In a classic master-worker system the master assigns work. In Spark the cluster manager does not assign your tasks: it grants containers, and the driver schedules into them. The coordinator role is split, resources from one component and work from another.

16. What is `spark-submit`?

The script that launches an application against any supported cluster manager through one uniform interface, so the same jar runs on YARN, Kubernetes or standalone with only --master changing.

spark-submit \
  --class org.apache.spark.examples.SparkPi \
  --master yarn \
  --deploy-mode cluster \
  --executor-memory 4g \
  --executor-cores 4 \
  --conf spark.sql.shuffle.partitions=400 \
  examples/jars/spark-examples.jar 1000

The options fall into three groups: what to run (--class, the jar), where to run it (--master, --deploy-mode), and how much to ask for (--executor-memory, --executor-cores, --num-executors).

The mistake to know: every Spark-facing option must come before the application jar. Anything after the jar is passed to your application’s main as an argument instead, so a misplaced --conf is silently ignored by Spark and handed to your code.

Configuration precedence, highest first:

  1. Values set on SparkConf in your code.
  2. Flags passed to spark-submit.
  3. spark-defaults.conf.

--properties-file names an alternative defaults file, and --load-spark-defaults makes Spark read spark-defaults.conf as well as that file. --verbose prints where each value came from, which settles most “why is my config not applying” questions.

17. What is the difference between client and cluster deploy mode?

Where the driver runs. That is the whole difference, and everything else follows from it.

flowchart TB
  subgraph C["client mode"]
    C1["submitting machine<br/><b>driver here</b>"] --> C2["cluster manager"]
    C2 --> C3["executors in the cluster"]
    C1 <-->|"all scheduling traffic"| C3
  end
  subgraph K["cluster mode"]
    K1["submitting machine<br/>exits after submit"] --> K2["cluster manager"]
    K2 --> K3["<b>driver in the cluster</b>"]
    K3 <--> K4["executors in the cluster"]
  end
  Client Cluster
Driver runs In the submitting process In a container in the cluster
Console output Attached to your terminal In the driver’s log
Submitting machine Must stay up for the whole job Can disconnect
Network Driver-executor traffic crosses to your machine Stays inside the cluster
Suited to Shells, notebooks, debugging Production jobs

Two practical consequences:

  • Client mode from a laptop over a VPN is slow and fragile, because every scheduling message and every collected result crosses that link. Cluster mode keeps it local.
  • In client mode the driver JVM already exists by the time your code runs. So spark.driver.extraJavaOptions set in SparkConf arrives too late to affect the JVM. Use --driver-java-options or a properties file instead. This catches people trying to add a logging config or a heap dump flag.
18. What is the default parallelism, and what sets it?

There are two separate settings, and confusing them is a classic slip.

spark.default.parallelism governs RDD operations. Its documented default depends on the operation: for a distributed shuffle such as reduceByKey it is the largest number of partitions in the parent RDD, and for an operation with no parent, such as parallelize, it follows the cluster manager, which on local mode is the number of cores.

spark.sql.shuffle.partitions governs DataFrame and SQL shuffles, and it defaults to 200 regardless of your data or your cluster.

print(spark.conf.get("spark.sql.shuffle.partitions"))   # 200
  spark.default.parallelism spark.sql.shuffle.partitions
Applies to RDD operations DataFrame and SQL shuffles
Default Depends on parent and cluster manager 200

So if you are tuning a DataFrame job and you change spark.default.parallelism, nothing happens, which is a frequent source of confusion. The other thing to know is that input read parallelism is governed by neither: it comes from file splitting via spark.sql.files.maxPartitionBytes.

19. What is `spark.sql.shuffle.partitions` and why does the default cause trouble?

It is the number of partitions produced by a DataFrame or SQL shuffle, and it defaults to 200.

The trouble is that 200 is a constant, while the right value depends on your data volume and your core count. Both directions hurt:

  • Too few for the data. 200 partitions over 1 TB means 5 GB per task, which will spill to disk and may fail outright.
  • Too many for the data. 200 partitions over 100 MB means 500 KB per task, so scheduling and task startup dominate and 200 near-empty files get written.
  • Mismatched to cores. With 32 cores, 200 tasks is a sensible several-waves-per-core; with 2,000 cores, most sit idle.

A reasonable target is partitions in the low hundreds of megabytes, and a total that is a small multiple of your core count so stragglers average out.

Adaptive query execution changes the calculus. With spark.sql.adaptive.enabled and spark.sql.adaptive.coalescePartitions.enabled both true by default, Spark writes 200 shuffle partitions and then merges the small ones at runtime toward spark.sql.adaptive.advisoryPartitionSizeInBytes, default 64MB.

spark.conf.set("spark.sql.adaptive.advisoryPartitionSizeInBytes", "128MB")

So in modern Spark the better habit is to set the advisory size and leave the partition count high, letting coalescing do the arithmetic, rather than guessing a single number for every stage in the query.

20. Which Spark version introduced Log4j 2, and why does it matter?

Spark 3.3.0. The configuration file became log4j2.properties and the JVM flag became -Dlog4j.configurationFile=.

It matters because the syntax changed completely, and most logging guides on the internet predate it. Log4j 1 used a flat log4j. prefix:

# Log4j 1, ignored by Spark 3.3.0 and later
log4j.rootLogger=INFO, console
log4j.logger.org.apache.spark=WARN

Log4j 2 uses named components:

# Log4j 2
rootLogger.level = info
rootLogger.appenderRef.stdout.ref = console
appender.console.type = Console
appender.console.name = console
appender.console.layout.type = PatternLayout
appender.console.layout.pattern = %d{HH:mm:ss} %-5level %logger{1} - %msg%n
logger.spark.name = org.apache.spark
logger.spark.level = warn

At v4.2.0 the only template Spark ships is conf/log4j2.properties.template; the Log4j 1 template is gone. So a config written in the old syntax is not an error you will see, it is a file that is silently ignored, which is the worst kind of failure.

Ship it to both sides, since driver and executor answer different questions:

spark-submit \
  --files /etc/spark/conf/log4j2.properties \
  --conf "spark.driver.extraJavaOptions=-Dlog4j.configurationFile=log4j2.properties" \
  --conf "spark.executor.extraJavaOptions=-Dlog4j.configurationFile=log4j2.properties" \
  app.jar

RDDs

The original API. Interviewers still ask, because the concepts underneath the DataFrame API are these, and because RDD questions expose whether you understand lineage and partitioning or just the syntax.

21. What is an RDD?

A Resilient Distributed Dataset: an immutable, partitioned collection of records that can be operated on in parallel, carrying enough information to rebuild any lost partition.

Unpack the name, because each word is a design decision:

  • Resilient. Spark records how each partition was derived, not just the data. Lose one and it is recomputed from its parents rather than the job failing.
  • Distributed. The collection is split into partitions spread across executors.
  • Dataset. A collection of records, with no schema and no interpretation of their contents.

It is also immutable. Every transformation produces a new RDD rather than modifying the old one, which is what makes lineage a reliable recovery record: a parent cannot have changed underneath you.

rdd = spark.sparkContext.parallelize([("sf", 24.5), ("nyc", 61.75)], 2)
totals = rdd.reduceByKey(lambda a, b: a + b)

The trade against DataFrames is visibility. Because an RDD holds opaque objects, Spark cannot see inside your closures, so there is no schema, no column pruning, no predicate pushdown and no Catalyst. You get control and lose optimisation, which is why the DataFrame API is the default for structured data.

22. What are the five properties of an RDD?

From RDD.scala, every RDD is characterised by five things, three required and two optional:

  1. A list of partitions. getPartitions, which defines the parallelism.
  2. A function for computing each partition. compute, which is the actual work.
  3. A list of dependencies on other RDDs. getDependencies, which is the lineage.
  4. Optionally, a partitioner for key-value RDDs, such as HashPartitioner.
  5. Optionally, preferred locations for each partition, such as the HDFS block locations.

The first three make the RDD work; the last two make it fast, and they are the ones worth talking about:

  • The partitioner is what lets Spark skip a shuffle. If two RDDs share a partitioner on the join key, matching keys are already co-located, so the join is narrow.
  • Preferred locations are what data locality is built on. TaskSchedulerImpl tries to place a task where its input already is, which is why locality levels exist.
rdd = spark.sparkContext.parallelize(range(10), 4)
print(rdd.getNumPartitions())                    # 1
print(rdd.partitioner)                           # 4, None for a non-pair RDD

Naming all five, and then saying which two are optional and why they matter, is a much stronger answer than listing them.

23. How does Spark achieve fault tolerance for RDDs?

Through lineage: the chain of transformations from a source to the dataset in hand. Because Spark recorded how to build every partition, not just the partition itself, a lost one is rebuilt on a surviving executor.

lineage = (spark.sparkContext.parallelize(range(10), 2)
    .map(lambda x: x + 1)
    .filter(lambda x: x % 2 == 0))
print(lineage.toDebugString().decode())

Two things make this cheaper than it sounds:

  • Recovery is per partition, not per job. Losing one executor costs the partitions it held, which are recomputed while everything else stands.
  • Only the path matters. Spark replays the branch that produced the lost partition, not the whole graph.
flowchart TB
  S["source"] --> A["map"] --> B["filter"] --> C["partition 3 lost"]
  C -.->|"recompute just this path"| D["partition 3 rebuilt<br/>on another executor"]

The limits are worth volunteering:

  • Shuffle output is the exception. If the shuffle files a completed stage wrote die with their executor, DAGScheduler resubmits that entire stage, because the next stage has nothing to fetch. That is a stage rerun rather than a task retry, and it is why a job can appear to go backwards.
  • Long lineages get expensive to replay, which is what checkpointing exists to fix: it writes the data to reliable storage and truncates the chain.
  • A non-deterministic transformation breaks the guarantee, because a recomputed partition may not match the lost one.
24. What is the difference between a narrow and a wide dependency?

In a narrow dependency each parent partition feeds at most one child partition, so the work can pipeline inside a single task with no data movement. In a wide dependency a child partition draws from many parents, which requires a shuffle.

flowchart TB
  subgraph N["narrow: map, filter"]
    P0["P0"] --> C0["C0"]
    P1["P1"] --> C1["C1"]
    P2["P2"] --> C2["C2"]
  end
  subgraph W["wide: groupByKey, join"]
    Q0["P0"] --> D0["C0"]
    Q0 --> D1["C1"]
    Q1["P1"] --> D0
    Q1 --> D1
    Q2["P2"] --> D0
    Q2 --> D1
  end
  Narrow Wide
Classes OneToOneDependency, RangeDependency ShuffleDependency
Examples map, filter, flatMap, mapPartitions, union groupByKey, reduceByKey, join, distinct, repartition
Data movement None Across the network
Stage Pipelines into the current stage Creates a boundary
On failure Recompute one parent partition May need to recompute many, or rerun the stage

The recovery difference is the part people miss. A lost narrow partition needs one parent partition recomputed. A lost partition after a wide dependency may need data from every parent partition, which is far more expensive, and is another reason shuffles are costly beyond their immediate network cost.

coalesce is the interesting edge case: it is narrow, because several parent partitions feed one child without a shuffle, which is exactly why it can only reduce the count and may leave partitions uneven.

25. Where exactly does Spark cut a stage?

At every ShuffleDependency in the lineage, and nowhere else.

DAGScheduler walks the graph backwards from the action. Each time it crosses a ShuffleDependency it ends the current stage and starts a new one; everything narrow is absorbed into the stage it is already in.

flowchart TB
  subgraph S0["stage 0"]
    R["read"] --> F["filter"] --> W["withColumn"] --> PA["partial aggregate"]
  end
  PA -->|"ShuffleDependency"| SH[("shuffle files<br/>on local disk")]
  subgraph S1["stage 1"]
    SH --> FA["final aggregate"] --> O["write"]
  end

The practical consequences:

  • A long chain of narrow operations costs no extra stages. Twenty chained filter and withColumn calls still run in one stage, fused by whole-stage codegen into a single loop.
  • The stage count is the shuffle count plus one, per job. Counting Exchange nodes in a physical plan gives you the same number.
  • Every stage but the last is a ShuffleMapStage, whose purpose is to produce shuffle output. The last is a ResultStage, which computes what the action asked for.

Stages also enforce a barrier: stage 1 cannot start until stage 0 has finished writing all its output, because a reduce task needs its slice from every map task. That is why one straggler in stage 0 delays the entire next stage, and why skew is so expensive.

26. What is the difference between `reduceByKey` and `groupByKey`?

reduceByKey combines values on the map side before the shuffle; groupByKey shuffles every value and combines afterwards.

# prefer this: combines locally first
pairs.reduceByKey(lambda a, b: a + b)

# shuffles every value across the network
pairs.groupByKey().mapValues(sum)
flowchart TB
  subgraph R["reduceByKey"]
    A1["('sf',1) ('sf',1) ('sf',1)"] --> A2["local combine<br/>('sf',3)"] --> A3["shuffle<br/>1 record"]
  end
  subgraph G["groupByKey"]
    B1["('sf',1) ('sf',1) ('sf',1)"] --> B3["shuffle<br/>3 records"] --> B4["combine after"]
  end

With a hot key, the difference is the number of records crossing the network: one per key per partition for reduceByKey, versus every record for groupByKey. groupByKey also has a correctness-adjacent risk: all values for one key must fit in memory on the reduce side, so a hot key can cause an OutOfMemoryError that reduceByKey would not.

The subtlety that makes a strong answer: a map-side combine disqualifies Spark’s two faster shuffle writers, because they cannot combine. So reduceByKey shuffles less data through a more expensive writer. It is still almost always the right choice, because network volume dominates, but the trade-off is real and worth naming.

Use groupByKey only when you genuinely need every value, for example to collect them into a list, and even then consider aggregateByKey which combines into a buffer as it goes.

27. What is the difference between `map` and `mapPartitions`?

map applies your function once per element. mapPartitions applies it once per partition, handing you an iterator over that partition’s elements and expecting an iterator back.

# once per row: a connection per row, which is fatal
rdd.map(lambda r: lookup(open_conn(), r))

# once per partition: one connection amortised over thousands of rows
def enrich(rows):
    conn = open_conn()
    for r in rows:
        yield lookup(conn, r)
    conn.close()

rdd.mapPartitions(enrich)

Use mapPartitions when there is per-partition setup worth amortising: a database or HTTP connection, loading a model, or compiling a regex.

Two cautions:

  • Do not materialise the partition. list(rows) inside the function loads the whole partition into memory and defeats the streaming design. Use a generator, as above.
  • Return an iterator, not a single value. Returning a list works but holds everything in memory.

mapPartitionsWithIndex additionally gives you the partition number, which is the standard trick for inspecting how data is distributed:

print(rdd.mapPartitionsWithIndex(
    lambda i, rows: [(i, sum(1 for _ in rows))]).collect())

That prints rows per partition, which is how you confirm skew on an RDD.

28. What is the difference between `map` and `flatMap`?

map returns exactly one output per input. flatMap returns an iterable per input and flattens the results, so it can return zero, one or many.

lines = sc.parallelize(["a b c", "d e"])

lines.map(lambda s: s.split()).collect()
# [['a','b','c'], ['d','e']]        two elements, each a list

lines.flatMap(lambda s: s.split()).collect()
# ['a','b','c','d','e']             five elements

The “zero” case is the one worth knowing, because it makes flatMap a combined map and filter:

# parse, and silently drop anything unparseable, in one pass
rdd.flatMap(lambda s: [parse(s)] if valid(s) else [])

Both are narrow transformations, so neither causes a shuffle, and both pipeline into the surrounding stage. The DataFrame equivalent of flatMap for a column is explode, which turns an array column into one row per element.

29. What is the difference between `repartition` and `coalesce`?

repartition does a full shuffle to reach any target count and balances the result. coalesce merges existing partitions without a full shuffle, so it can only reduce the count and may leave partitions uneven.

  repartition(n) coalesce(n)
Shuffle Full None by default
Direction Up or down Down only
Balance Even Possibly uneven
Dependency Wide, so a new stage Narrow
df.coalesce(10).write.parquet(path)      # cheap, before a write
df.repartition(400).groupBy("k").count() # balanced, before heavy work
df.repartition("city_id")                # hash-partitioned by a column

The trap is what coalesce does upstream, and it is a favourite interview question. Because coalesce is narrow, it does not create a stage boundary, so the reduced parallelism propagates backwards into the stage feeding it. coalesce(1) before a heavy transformation can make that entire transformation run in a single task.

flowchart TB
  subgraph C["coalesce(2): merge, no shuffle"]
    A1["P0 10 MB"] --> M1["P0 + P1<br/>30 MB"]
    A2["P1 20 MB"] --> M1
    A3["P2 5 MB"] --> M2["P2 + P3<br/>45 MB"]
    A4["P3 40 MB"] --> M2
  end
  subgraph R["repartition(2): shuffle, balanced"]
    B1["P0 10 MB"] --> SH[("shuffle")]
    B2["P1 20 MB"] --> SH
    B3["P2 5 MB"] --> SH
    B4["P3 40 MB"] --> SH
    SH --> N1["37 MB"]
    SH --> N2["38 MB"]
  end

So if you need one output file, repartition(1) pays for a shuffle but keeps the upstream stage parallel. coalesce(1) is only safe when the work before it is trivial.

30. What is the difference between `cache` and `persist`?

cache() is exactly persist() with the default storage level. persist(level) lets you choose the level.

df.cache()                                    # default level
df.persist(StorageLevel.MEMORY_AND_DISK_SER)  # explicit

The default differs by API, which is the part that surprises people:

API Default level Behaviour when it does not fit
rdd.cache() Memory only Partitions are dropped and recomputed
df.cache() Disk and memory, deserialized Partitions spill to disk

So an RDD cache that exceeds memory silently recomputes, which looks like the cache doing nothing at all. A DataFrame cache spills instead, which is slower than memory but still avoids recomputation. If an RDD cache seems to have no effect, that asymmetry is usually why.

Both are lazy: they mark the dataset and populate it on the next action. The exception is SQL’s CACHE TABLE, which is eager by default.

Release it explicitly when you are done, because cached blocks hold storage memory for the life of the application and squeeze execution memory from every later stage:

df.unpersist()
spark.catalog.clearCache()       # everything
31. What are the RDD storage levels?

The levels combine three choices: memory or disk, serialized or not, and replicated or not.

Level Memory Disk Serialized Replicas
MEMORY_ONLY Yes No No 1
MEMORY_ONLY_SER Yes No Yes 1
MEMORY_AND_DISK Yes Spills No 1
MEMORY_AND_DISK_SER Yes Spills Yes 1
DISK_ONLY No Yes Yes 1
*_2 variants as above as above as above 2
OFF_HEAP Off-heap No Yes 1

How to choose:

  • Serialized forms are far more compact, so more fits in memory, at the cost of CPU to deserialize on every read. Worth it when memory is the constraint.
  • _2 replicas put a copy on a second node, so losing an executor does not force recomputation. Worth it when recomputation is very expensive; otherwise it doubles the memory cost.
  • DISK_ONLY is right when recomputing is more expensive than reading back, which is the case after a long or costly lineage.
  • OFF_HEAP avoids garbage collection pressure but needs spark.memory.offHeap.enabled, which defaults to false, and a size.

Read the storage tab of the UI to check what actually happened. A dataset showing well below 100% in memory is being partly recomputed or read from disk on every use, which usually means the cache is not paying for itself.

32. What are broadcast variables?

A broadcast variable is a read-only value sent to each executor once and cached there, rather than serialized into every task.

The guide describes them as keeping “a read-only variable cached on each machine rather than shipping a copy of it with tasks”.

sc = spark.sparkContext
cities = sc.broadcast({"sf": "San Francisco", "nyc": "New York"})

resolved = sc.parallelize(["sf", "nyc", "sf"], 2).map(lambda c: cities.value[c])
print(resolved.collect())   # ['San Francisco', 'New York', 'San Francisco']

Without the broadcast, that dictionary is captured by the closure and serialized with every task. With 10,000 tasks and a 10 MB lookup table, that is 100 GB of serialization instead of one copy per executor.

flowchart TB
  subgraph W["captured variable"]
    D1["driver"] -->|"copy per task<br/>x 10,000"| T1["tasks"]
  end
  subgraph B["broadcast"]
    D2["driver"] -->|"one copy per executor<br/>x 20"| E2["executors"] --> T2["all tasks read locally"]
  end

Practical notes:

  • Read-only. Mutating .value on an executor changes only that copy and nothing comes back.
  • Call .unpersist() when finished, or .destroy() to release it permanently, if it is large and no longer needed.
  • This is the same machinery behind a broadcast hash join, which is why spark.sql.autoBroadcastJoinThreshold exists and why broadcasting something large pressures the driver, which must assemble it first.
33. What are accumulators, and what is the catch?

Accumulators are write-only-from-tasks counters that the driver can read. Tasks add to them; only the driver should read the value.

errors = sc.accumulator(0)

def parse(line):
    try:
        return json.loads(line)
    except ValueError:
        errors.add(1)
        return None

parsed = rdd.map(parse).filter(lambda x: x is not None)
parsed.count()
print(errors.value)

The catch is the guarantee, and it is the whole question. The guide is explicit:

For accumulator updates performed inside actions only, Spark guarantees that each task’s update to the accumulator will only be applied once, i.e. restarted tasks will not update the value. In transformations, users should be aware of that each task’s update may be applied more than once if tasks or job stages are re-executed.

This is not theoretical, and it follows directly from lazy evaluation:

acc = sc.accumulator(0)
counted = sc.parallelize(range(100), 4).map(lambda x: (acc.add(1), x)[1])

counted.count(); print(acc.value)   # 100
counted.count(); print(acc.value)   # 200, for the same 100 rows

The uncached dataset is recomputed for the second action, so the accumulator increments again. Add .cache() and it stays at 100. A task retry has the same effect.

So use accumulators for diagnostics you can tolerate being approximate, and never as the source of a number the output depends on. If you need an exact count, compute it as data with an aggregation.

34. What is a partitioner, and when does it help?

A partitioner is the function mapping a key to a partition number. Spark ships two:

  • HashPartitioner, which uses hash(key) % numPartitions. The default for reduceByKey and friends.
  • RangePartitioner, which samples the data to find range boundaries so that partitions are roughly equal and ordered. Used by sortByKey and by a global orderBy.

It helps because a known partitioner lets Spark skip a shuffle. If two RDDs share the same partitioner on the join key, matching keys are already co-located, so the join becomes a narrow dependency.

part = 100
a = rdd_a.partitionBy(part).cache()      # note the cache
b = rdd_b.partitionBy(part).cache()
joined = a.join(b)                       # no shuffle

The cache is not optional. Without persisting, the partitioning is recomputed each time the RDD is used, so you pay the shuffle you were trying to avoid.

You can inspect and set it:

print(pairs.partitioner)                      # None until you set one
tagged = pairs.partitionBy(HashPartitioner(100))

The DataFrame equivalents are repartition(col) for the same effect within a job, and bucketing for the effect to persist across jobs, recorded in the metastore.

35. What does `toDebugString` tell you?

It prints the lineage of an RDD as an indented tree: each level is an RDD, with its type, its id and its partition count.

lin = (sc.parallelize(range(100), 4)
    .map(lambda x: (x % 5, x))
    .reduceByKey(lambda a, b: a + b))
print(lin.toDebugString().decode())
(4) PythonRDD[3] at RDD at PythonRDD.scala:53 []
 |  MapPartitionsRDD[2] at mapPartitions at ... []
 |  ShuffledRDD[1] at partitionBy at ... []
 +-(4) PairwiseRDD[0] at reduceByKey at ... []
    |  ParallelCollectionRDD ... []

How to read it:

  • The number in parentheses is the partition count at that level.
  • Indentation changes, marked by +-, are shuffle boundaries. Count them and you have the stage count minus one.
  • Read bottom up: the source is at the bottom, the RDD you called it on is at the top.

It is the fastest way to answer “how many stages will this take” for RDD code, and to spot an accidental shuffle. For DataFrames the equivalent is df.explain("formatted"), where each Exchange is the same signal.

36. When should you still use RDDs?

When you need something the structured APIs do not expose. The honest list is short:

  • Custom partitioning logic beyond hash or range, for example routing by a business rule.
  • Per-partition resource handling via mapPartitions, although the DataFrame API has mapInPandas for much of this now.
  • Genuinely unstructured data with no schema to impose: raw text, binary records, custom parsing where a row shape only emerges after processing.
  • Low-level control over lineage and checkpointing, in iterative algorithms.

For everything with a schema, the DataFrame API wins, and the reason is structural rather than a matter of taste: Catalyst can see a DataFrame plan and cannot see inside an RDD closure. That means no column pruning, no predicate pushdown, no join strategy selection and no whole-stage codegen for RDD code. Tungsten’s compact binary rows also avoid materialising JVM objects, which RDD code cannot avoid.

A good answer also notes that RDDs did not go away. DataFrames are built on them: df.rdd gives you the underlying RDD, and the scheduler works entirely in RDD terms. So understanding RDDs is still how you understand what a DataFrame job actually does.

37. What is the difference between `collect` and `take`?

collect() brings every row to the driver. take(n) returns the first n rows and tries to read as little as possible to get them.

rows = df.collect()     # everything, bounded only by driver memory
head = df.take(5)       # 5 rows, usually from one partition
df.show(5)              # take plus formatting

take(n) starts by scanning one partition. If that yields fewer than n rows it scans more, increasing the number it tries each round. So on a large dataset it usually touches a tiny fraction of the data.

collect() is bounded by spark.driver.maxResultSize, default 1g, and by the driver heap, spark.driver.memory, default 1g. Exceeding the first raises an error naming the config; exceeding the second is an OutOfMemoryError.

This is the classic way to kill a driver, and the answer an interviewer wants is what to do instead:

  • To inspect, use show(), take() or limit().
  • To count, use count(), which aggregates rather than transferring rows.
  • To export, use write so executors write in parallel and only a status returns.
  • To hand data to Python, toPandas() has the same problem as collect(); prefer writing out or using mapInPandas to keep the work distributed.
38. What is the difference between `reduce` and `fold`?

Both aggregate an RDD to a single value using a binary function. The difference is the zero value.

sc.parallelize([1, 2, 3, 4]).reduce(lambda a, b: a + b)      # 10
sc.parallelize([1, 2, 3, 4]).fold(0, lambda a, b: a + b)     # 10
sc.parallelize([]).reduce(lambda a, b: a + b)                # raises
sc.parallelize([]).fold(0, lambda a, b: a + b)               # 0

So fold is defined on an empty RDD and reduce is not.

The trap is that the zero value must be a true identity, because Spark applies it once per partition as well as when combining partitions. With four partitions and a zero of 1 for addition, you would add 1 five times, not once.

aggregate generalises both: it takes a zero, a function to fold a partition into an accumulator, and a function to combine accumulators. That is what you need when the accumulator type differs from the element type, such as computing a mean:

total, n = sc.parallelize([1, 2, 3, 4]).aggregate(
    (0, 0),
    lambda acc, x: (acc[0] + x, acc[1] + 1),      # fold within a partition
    lambda a, b: (a[0] + b[0], a[1] + b[1]))      # combine partitions
print(total / n)

Note also that all three are actions, returning to the driver, whereas reduceByKey is a transformation returning a distributed RDD.

39. Why are RDD operations on key-value pairs special?

Because only pair RDDs can be partitioned by key, and key-partitioning is what enables the optimisations that matter.

Being a pair RDD unlocks:

  • partitionBy, and therefore shuffle-free joins between co-partitioned RDDs.
  • Map-side combines: reduceByKey, aggregateByKey, combineByKey.
  • Key-wise operations: groupByKey, join, cogroup, lookup, sortByKey.
pairs = rdd.map(lambda r: (r.city_id, r.fare_amount))   # now a pair RDD
pairs.reduceByKey(lambda a, b: a + b)

In Scala these arrive by implicit conversion on RDD[(K, V)] through PairRDDFunctions, which is why they are not on the RDD class itself. In Python any RDD whose elements are 2-tuples qualifies.

The distinction matters in an interview because it explains why the API is shaped this way: partitioning is defined on keys, so anything that benefits from co-location has to be expressed as a pair. The DataFrame equivalent is that these optimisations attach to columns you group or join on, which is why repartition(col) and bucketing exist.

40. What happens if a task fails?

It is retried, and if retries are exhausted the stage and then the job fail.

The sequence:

  1. The task fails and reports back to the driver.
  2. TaskSetManager reschedules it, preferring a different executor if the failure looked executor-specific.
  3. Retries continue up to spark.task.maxFailures, default 4, counted per task across the stage.
  4. On the fourth failure the stage fails, and the job fails with it.

Two related mechanisms are worth naming:

  • Blacklisting, now called excludeOnFailure. If one executor or node keeps failing tasks, Spark can stop scheduling onto it, which prevents a single bad machine from burning through the retry budget of every task.
  • Speculation. spark.speculation, default false, duplicates tasks that are running much slower than their peers and takes whichever finishes. It helps with a slow node and hurts with skew, because the duplicate is equally slow.

The distinct case is losing shuffle output rather than a task:

flowchart TB
  T1["task attempt 1 fails"] --> T2["attempt 2, preferably<br/>on another executor"]
  T2 --> T3["attempt 3"]
  T3 --> T4["attempt 4"]
  T4 --> F["spark.task.maxFailures reached:<br/><b>stage fails, job fails</b>"]
  T2 -.->|"same executor keeps failing"| X["excludeOnFailure:<br/>stop scheduling there"]
  T2 -.->|"slow, not failing"| SP["speculation may<br/>duplicate the task"]

That is a stage rerun, and it is why a job can appear to go backwards in the UI. An external shuffle service avoids it by serving shuffle blocks independently of the executor that wrote them.

DataFrames, Datasets and Spark SQL

The API you should reach for by default. These questions separate people who use it from people who understand what it does underneath.

41. What is a DataFrame?

A distributed collection of rows organised into named columns, with a schema Spark knows at plan time. That last part is the whole point.

df = spark.read.parquet("s3a://lakehouse-prod/warehouse/trips")
df.printSchema()

Because Spark knows the schema and the operations, a DataFrame plan goes through Catalyst, which can:

  • Prune columns nobody selected, so they are never read from disk.
  • Push predicates into the scan, so rows are filtered at the source.
  • Choose a join strategy based on estimated sizes.
  • Fuse operators into a single generated loop via whole-stage codegen.
  • Store rows in Tungsten’s binary format, avoiding JVM objects entirely.

An RDD of the same data gets none of that, because its records are opaque objects and its closures are opaque functions.

Two clarifications interviewers like:

  • A DataFrame in Scala is Dataset[Row]. It is not a separate type; Row is the untyped record.
  • It is still immutable and lazy. withColumn returns a new DataFrame, and nothing runs until an action.
42. What is a Dataset, and how does it differ from a DataFrame?

A Dataset is a typed distributed collection that still runs through Spark SQL’s optimised engine. It gives you compile-time type safety and your own case classes instead of generic rows.

case class Trip(tripId: String, fareAmount: Double, cityId: String)

val trips: Dataset[Trip] = spark.read.parquet(path).as[Trip]
trips.filter(_.fareAmount > 50).map(_.cityId)   // typed, checked at compile time

In Scala, DataFrame is simply Dataset[Row], so they are the same machinery with different static types.

  DataFrame Dataset
Element type Row, untyped Your class, typed
Errors caught At analysis, runtime At compile time
Available in All four languages Scala and Java only
Encoder Built in Generated for your type

Python has no typed Dataset API, and the reason is worth stating: the API exists to give compile-time checking, and Python has no compile step to check in. In PySpark, DataFrame is what you get.

The cost of typed Datasets is real too: a lambda like _.fareAmount > 50 is opaque to Catalyst in the same way an RDD closure is, so it cannot be pushed down. Typed filter and map trade optimisation for type safety, which is why even in Scala, column expressions are often preferred for hot paths.

43. Compare RDD, DataFrame and Dataset.
  RDD DataFrame Dataset
Level Low Structured Typed structured
Schema None Yes Yes, from your type
Catalyst optimisation No Yes Yes, for column expressions
Whole-stage codegen No Yes Yes
Tungsten binary format No Yes Yes
Type safety Compile time (Scala) Runtime Compile time
Python support Yes Yes No typed API
Serialization Java or Kryo Tungsten encoders Tungsten encoders

How to choose, in one line each:

  • DataFrame or SQL for almost everything. It is the fastest path because Catalyst can see it all.
  • Dataset in Scala or Java when compile-time safety on a domain type is worth losing pushdown through typed lambdas.
  • RDD only for custom partitioning, per-partition resource handling, or data with no schema.

The evolution is worth a sentence: RDDs came first, DataFrames added a schema and the optimiser, and Datasets added types back on top of the optimiser. They are layered, not alternatives, and df.rdd drops down a level whenever you need it.

The strong closing point: hand-written RDD code loses to DataFrame code on structured data not because the author wrote it badly, but because the engine is allowed to rewrite the DataFrame version and is not allowed to touch the RDD one.

44. What is Catalyst?

Catalyst is Spark SQL’s query optimiser: the component that turns a query you wrote into a physical plan the engine runs.

It is a rule-based tree transformation framework. Plans are trees, rules are functions from tree to tree, and the optimiser applies batches of rules until the tree stops changing. That design is why Spark SQL is extensible: a new optimisation is a new rule, and you can add your own through session extensions.

flowchart TB
  A["SQL or DataFrame"] --> B["unresolved<br/>logical plan"]
  B -->|"catalog"| C["analysed<br/>logical plan"]
  C -->|"rules"| D["optimised<br/>logical plan"]
  D -->|"strategies + cost"| E["physical plan"]
  E -->|"codegen"| F["JVM bytecode"]

What Catalyst gives you in practice, in rough order of value:

  • Column pruning, so unread columns are never touched.
  • Predicate pushdown, so filtering happens at the source.
  • Join strategy selection, so a small side gets broadcast.
  • Constant folding and expression simplification.
  • Whole-stage code generation, producing tight loops rather than operator-by-operator iteration.

The idea underneath all of it is lazy evaluation. Catalyst can only optimise a plan it can see in full, which is exactly what laziness provides.

45. What are Catalyst's phases?

Five, and naming them in order with what each does is a strong answer.

  1. Parse. SQL text or DataFrame calls become an unresolved logical plan: a tree with relation and column names that have not been checked against anything.
  2. Analyse. The analyzer resolves table names, column names, data types and function calls against the catalog. This is where AnalysisException comes from, and why a typo in a column name surfaces here rather than at parse time.
  3. Logical optimisation. Rule batches rewrite the plan: predicate pushdown, column pruning, constant folding, boolean simplification, filter and projection reordering, subquery flattening. The output is still logical, meaning it says what, not how.
  4. Physical planning. Strategies turn each logical operator into one or more physical ones, for example a logical join into a broadcast hash join, shuffle hash join or sort-merge join. Costs choose between candidates.
  5. Code generation. Whole-stage codegen fuses the operators in a stage into one generated Java method, compiled with Janino, so intermediate rows never become objects.

You can see each stage:

df.explain("extended")   # parsed, analysed, optimised, physical, in that order

The distinction to be crisp about is logical versus physical. Logical optimisation is algebraic and knows nothing about the cluster; physical planning is where partitioning, shuffles and join implementations are decided.

46. What is predicate pushdown?

Moving a filter as close to the data as possible, so that fewer rows are ever read or carried up the plan.

There are two levels, and distinguishing them is a good signal:

  • Within the plan. Catalyst moves a Filter below a Project or through a join, so downstream operators see fewer rows.
  • Into the data source. For Parquet or ORC, the predicate is handed to the reader, which skips entire row groups using footer statistics. For JDBC, it becomes a WHERE clause in the SQL sent to the database.

You can confirm it in the plan:

(spark.read.parquet(path)
   .filter(F.col("fare_amount") > 50)
   .select("city_id")
   .explain("formatted"))
FileScan parquet [fare_amount#19,city_id#21]
  DataFilters: [isnotnull(fare_amount#19), (fare_amount#19 > 50.0)]
  PushedFilters: [IsNotNull(fare_amount), GreaterThan(fare_amount,50.0)]
  ReadSchema: struct<fare_amount:double,city_id:string>

PushedFilters is the line that matters: those predicates reached the file format. spark.sql.parquet.filterPushdown defaults to true.

What blocks it is the practically useful half:

  • Wrapping the column in a function. to_date(ts) = '2026-09-17' cannot push down; ts >= '...' AND ts < '...' can.
  • A UDF in the predicate, because Catalyst cannot see inside it.
  • Filtering after a non-deterministic operation, which would change semantics.
47. What is column pruning?

Reading only the columns the query actually needs, and dropping the rest as early as possible in the plan.

On a columnar format such as Parquet or ORC, unread columns are never touched on disk at all, so a table with 200 columns queried for 3 reads roughly 1.5% of the bytes. On a row format such as JSON or CSV the whole row must still be read and parsed, so pruning helps less.

Confirm it by reading ReadSchema in the plan:

ReadSchema: struct<fare_amount:double,city_id:string>

Two columns, not the whole table.

What defeats it:

  • select("*"), explicitly asking for everything.
  • df.rdd or a typed Dataset lambda, which forces Spark to materialise the whole row because it cannot see which fields you touch.
  • Caching before selecting. df.cache() caches the columns the plan needed at that point, so caching a wide DataFrame and then selecting three columns stores all of them.

The practical habit: select the columns you need immediately after reading, not at the end of the pipeline. It costs nothing and it is the difference between reading 1.5% and 100% of a wide table.

48. What is whole-stage code generation?

An optimisation where Catalyst fuses all the operators in a stage into a single generated Java method, rather than running each operator as a separate iterator that passes rows to the next.

Without it, a scan-filter-project chain is three objects calling next() on each other, with a virtual call and an intermediate row per step. With it, Spark emits roughly the loop you would write by hand:

// conceptually, for filter + project over a scan
while (scan.hasNext()) {
  InternalRow row = scan.next();
  double fare = row.getDouble(0);
  if (!row.isNullAt(0) && fare > 50.0) {
    out.write(row.getUTF8String(1));
  }
}

You can see it in the plan. Operators marked with * and a number share one generated stage:

*(1) HashAggregate(keys=[city_id#21], ...)
+- Exchange hashpartitioning(city_id#21, 200)
   +- *(2) HashAggregate(keys=[city_id#21], ...)
      +- *(2) Project [city_id#21]
         +- *(2) Filter (fare_amount#19 > 50.0)
            +- *(2) FileScan parquet

Everything marked *(2) was fused into one method.

What breaks the fusion is the useful part to know: a Python UDF, because rows must leave the JVM; some complex expressions; and plans that exceed the generated-method size limit, at which point Spark falls back to the interpreted path. Losing codegen is a large part of why Python UDFs cost what they do.

49. How do you read a query plan?

Start with the mode that matches the question:

df.explain("formatted")   # the readable one: tree plus per-node details
df.explain("extended")    # parsed, analysed, optimised, physical
df.explain("cost")        # with statistics, where available
df.explain(True)          # same as extended

Then read the physical plan bottom up, because that is execution order: scans at the bottom, then filters and projections, then exchanges, then the final operator.

What to look for, in the order that usually matters:

In the plan Means
Exchange A shuffle, therefore a stage boundary. Count them
PushedFilters Predicates that reached the data source
ReadSchema Which columns are actually read
BroadcastHashJoin The small side is being broadcast
SortMergeJoin Both sides shuffled and sorted
*(n) Whole-stage codegen; same n means fused
CartesianProduct Almost always a mistake
AdaptiveSparkPlan isFinalPlan=false AQE will re-plan at runtime

The last row matters for reading plans in modern Spark. Before execution the plan is wrapped in AdaptiveSparkPlan with isFinalPlan=false, meaning what you see is the planned shape. To see what actually ran, look at the SQL tab of the UI after the query finishes, where the final plan shows the strategies AQE settled on.

50. What is the difference between `select` and `withColumn`?

select projects a set of columns, replacing the projection entirely. withColumn adds or replaces one column, keeping everything else.

df.select("trip_id", (F.col("fare_amount") * 1.1).alias("fare_with_fee"))
df.withColumn("fare_with_fee", F.col("fare_amount") * 1.1)

For one or two columns the difference is style. For many it is performance, and this is a real production issue:

# 50 chained calls: a 50-level deep plan, slow to analyse and optimise
for i in range(50):
    df = df.withColumn(f"c{i}", F.col("x") + i)

# one projection: a flat plan
df = df.select("*", *[(F.col("x") + i).alias(f"c{i}") for i in range(50)])

Each withColumn adds a Project node. Catalyst will usually collapse them, but the analysis and optimisation passes still walk the deep tree first, and with enough levels you get very slow planning or a StackOverflowError from the recursive walk.

The rule of thumb: a handful of withColumn calls is fine and more readable; dozens should become one select. If planning time is a noticeable fraction of your job, count the withColumn calls first.

51. What is the difference between a temporary view and a global temporary view?

Scope and lifetime.

df.createOrReplaceTempView("trips")              # this SparkSession only
df.createGlobalTempView("trips_global")          # all sessions in the app
spark.sql("SELECT * FROM global_temp.trips_global")
  Temporary view Global temporary view
Visible to The creating SparkSession Every session in the application
Database The current one global_temp, always qualified
Lives until The session ends The application ends

Neither persists metadata: both disappear when the process does. For something that survives, you need a catalog table via saveAsTable, which writes to the metastore.

The reason global temp views exist is multi-session applications. newSession() gives an isolated SQL session sharing the same SparkContext, which is how the Thrift server isolates users. A global temp view is how you deliberately share something across those sessions.

Worth noting: a view stores the query, not the result. Reading a view re-executes its plan. If you want the result stored, cache it or write a table.

52. What is a UDF, and why avoid it?

A user-defined function: your own logic registered as a callable expression.

@F.udf("double")
def add_fee(fare):
    return fare * 1.1 if fare else None

df.withColumn("total", add_fee("fare_amount"))

Avoid it when a built-in exists, for three compounding reasons:

  1. Catalyst cannot see inside it. It is a black box, so no predicate can be pushed through it, no constant folded, no expression simplified.
  2. It breaks whole-stage codegen across that point, so the tight generated loop becomes iterator-by-iterator.
  3. In PySpark it leaves the JVM. Rows are serialized to a Python worker and results serialized back, per row for a plain UDF.

Almost everything people write a UDF for exists already:

# instead of a UDF
F.when(F.col("fare") > 100, "high").otherwise("mid")
F.regexp_extract("path", r"city_id=(\w+)", 1)
F.coalesce("a", "b", F.lit(0))
F.date_format("started_at", "yyyy-MM-dd")

The ordering to remember: built-in functions, then SQL expressions, then a pandas (vectorised) UDF, then a Scala UDF if you are on the JVM, and a plain Python UDF last. If you must write one, keep it narrow and apply it after filtering, so it sees as few rows as possible.

53. What is a pandas UDF and when does it help?

A vectorised UDF that receives an Arrow batch as a pandas Series and returns one, rather than being called per row.

import pandas as pd
from pyspark.sql.functions import pandas_udf

@pandas_udf("double")
def adjust(fare: pd.Series) -> pd.Series:
    return fare * 1.1

The cost of a Python UDF is dominated by crossing the JVM-Python boundary and by per-row interpreter overhead. A pandas UDF amortises the crossing over a whole batch and lets the work happen in vectorised NumPy or pandas code, so the per-row cost falls sharply.

It still breaks whole-stage codegen, because the data still leaves the JVM. So the ordering is: built-ins beat pandas UDFs, and pandas UDFs beat plain Python UDFs.

Variants worth knowing:

  • Series to Series, the scalar case above.
  • Iterator of Series, when you need expensive per-partition setup such as loading a model once.
  • Grouped map, via df.groupBy(...).applyInPandas(fn, schema), which hands you a whole group as a DataFrame. Watch memory: one group must fit.
  • mapInPandas, which streams batches of an entire DataFrame without grouping.

Arrow is what makes this work, and it is also what makes toPandas() faster when enabled.

54. How do you handle nulls correctly?

Use null-aware operators, because in SQL semantics any comparison with null is null, not false.

df.filter(F.col("city_id").isNull())
df.filter(F.col("city_id").isNotNull())
df.filter(F.col("a").eqNullSafe(F.col("b")))     # null == null is true here
df.na.fill({"fare_amount": 0.0, "city_id": "unknown"})
df.na.drop(subset=["trip_id"])
df.withColumn("city", F.coalesce("city_id", F.lit("unknown")))

The traps that produce silently wrong results:

  • col == value never matches nulls. A filter meant to exclude one city also drops every row whose city is null.
  • An inner join drops null keys entirely, because null never equals null.
  • Aggregations skip nulls. count("col") counts non-null values while count("*") counts rows, and avg ignores nulls rather than treating them as zero.
  • A large number of null keys is skew, because they all hash to one partition.

And one version-specific change to know: from Spark 4.0 spark.sql.ansi.enabled is on by default, so invalid casts and overflows now raise instead of quietly producing null. That converts a class of silent null-generation into a visible error, which is usually an improvement even when it breaks a job on upgrade.

55. What changed with ANSI mode in Spark 4?

It became the default. The migration guide is explicit:

Since Spark 4.0, spark.sql.ansi.enabled is on by default.

What changes with it on:

Operation Legacy behaviour ANSI behaviour
CAST('abc' AS INT) null Raises an error
Integer overflow Wraps silently Raises an error
Division by zero null Raises an error
Reserved keywords Permitted as identifiers Enforced

To restore the old behaviour:

spark.conf.set("spark.sql.ansi.enabled", "false")

There is also a SPARK_ANSI_SQL_MODE environment variable, which the source reads when computing the default.

The right framing for an interview is that this converts silent data corruption into loud failure. Under the old behaviour a malformed value became null, flowed through the pipeline, and quietly changed a downstream aggregate. Under ANSI the job fails and you find out.

So when a job breaks on upgrade, treat it as a finding rather than a regression: something in the data was already not what the schema claimed. Fix the cast explicitly with try_cast where null genuinely is the right answer, rather than disabling ANSI globally.

56. What is the difference between `union` and `unionByName`?

union matches columns by position. unionByName matches them by name.

a = spark.createDataFrame([(1, "sf")], ["id", "city"])
b = spark.createDataFrame([("nyc", 2)], ["city", "id"])

a.union(b).show()          # id=nyc, city=2   <- silently wrong
a.unionByName(b).show()    # correct

This is one of the few Spark behaviours that produces silently wrong data rather than an error, because the types happened to be compatible. If the types had differed you would have got an analysis error and noticed.

unionByName also handles schema differences:

a.unionByName(b, allowMissingColumns=True)   # missing columns become null

Without that flag, a name present on one side and not the other is an error, which is usually what you want.

The recommendation is to default to unionByName, and to reach for union only when you have deliberately aligned the schemas and want positional semantics. In a code review, a bare union on two DataFrames built by different code paths is worth questioning.

Note also that both are narrow transformations: they do not shuffle, they just concatenate partitions. Neither deduplicates, unlike SQL’s UNION, which is union followed by distinct.

57. How do you write partitioned output, and what is the trap?

partitionBy writes a directory per distinct value, which lets later readers prune whole directories.

(df.write
   .partitionBy("city_id")
   .mode("overwrite")
   .parquet("s3a://lakehouse-prod/warehouse/trips"))
trips/city_id=sf/part-00000-....parquet
trips/city_id=nyc/part-00000-....parquet

The trap is file count, and it is arithmetic: every task that holds rows for a partition value writes a file into it. With 200 shuffle partitions and 250 cities you can get up to 50,000 files.

The fix is to repartition by the same columns first, so each value is held by one task:

(df.repartition("city_id")
   .write.partitionBy("city_id")
   .mode("overwrite")
   .parquet(path))

Other things to get right:

  • Cardinality. Partition on a low-cardinality column. Partitioning by a timestamp or a user id creates a directory per value, which makes listing slower than scanning.
  • Dynamic overwrite. spark.sql.sources.partitionOverwriteMode set to dynamic makes overwrite replace only the partitions present in the data, instead of the whole table.
  • File size. spark.sql.files.maxRecordsPerFile, default 0 meaning unlimited, caps rows per file when one partition is huge.
58. What is the difference between `partitionBy` and `bucketBy`?

Both pre-organise data on write, but they solve different problems.

  partitionBy bucketBy
Layout A directory per value A fixed number of files per partition
Good for Pruning on read Avoiding a shuffle on join or aggregate
Cardinality Low Any, hashed into n buckets
Requires a metastore No Yes, so saveAsTable only
Read benefit Skips directories Skips the shuffle
(df.write
   .bucketBy(64, "rider_id")
   .sortBy("rider_id")
   .mode("overwrite")
   .saveAsTable("lakehouse.trips_bucketed"))

Bucketing works because the bucket count and key are recorded in the metastore, so when two tables are bucketed the same way on the join key, Spark knows matching rows are already co-located and skips the exchange. That is the only way to avoid a shuffle on a large-to-large join.

The conditions are strict, and naming them is a good answer: both sides bucketed, on the same column, with the same number of buckets. Mismatched bucket counts means Spark shuffles anyway.

The usual pattern is to combine them: partitionBy a date for pruning, and bucketBy a join key within each partition.

59. What are the save modes?

Four, set with .mode():

Mode If the target exists
error or errorifexists Raise. The default
append Add new files alongside
overwrite Replace
ignore Do nothing, silently
df.write.mode("append").parquet(path)

The two that need care:

overwrite on a partitioned table replaces the whole table by default, even if your DataFrame only contains one day. Setting

spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")

makes it replace only the partitions present in the data, which is almost always what an incremental job wants.

append has no transactional guarantee on plain Parquet. Two concurrent appends can interleave, and a failure mid-write leaves partial files with no way to roll back, because a directory of Parquet files has no commit protocol. That is the structural reason lakehouse formats exist: Hudi, Iceberg and Delta add atomic commits over the same files.

ignore is worth flagging in review, because it silently does nothing, which hides a logic error rather than surfacing it.

60. Why is Parquet usually the right format?

Three properties that compound:

  1. Columnar. Columns are stored separately, so reading 3 of 200 columns touches roughly 1.5% of the bytes. Row formats must read and parse the whole row.
  2. Statistics per row group. Each row group carries min, max and null counts per column, so a pushed-down predicate can skip whole groups without decoding them.
  3. Compression that works. A column holds like-typed, often similar values, so dictionary and run-length encoding are very effective, which reduces both storage and read time.

It also carries a schema, including nested structures, so there is no inference step and no ambiguity about types.

spark.read.parquet(path)     # schema comes from the files

spark.sql.parquet.filterPushdown defaults to true, so the statistics are used by default.

When it is not the right choice:

  • Write-once-read-never logs, where a cheaper row format is fine.
  • Row-level random access, which no scan format does well.
  • Very small datasets, where the per-file overhead dominates.
  • Interchange with tools that cannot read it, where CSV wins on ubiquity alone.

ORC is broadly comparable and stronger in some Hive-centric stacks; the columnar argument is the same. Avro is row-oriented and better for streaming records and schema evolution on the write path, which is why Hudi uses it for log blocks.

61. How do you read from JDBC in parallel?

Give Spark a column to split on, plus bounds and a partition count. Without them the entire table is read by a single task, whatever your cluster size.

trips = (spark.read.format("jdbc")
    .option("url", "jdbc:postgresql://db.internal:5432/rides")
    .option("dbtable", "public.trips")
    .option("user", "etl")
    .option("partitionColumn", "trip_seq")
    .option("lowerBound", "1")
    .option("upperBound", "40000000")
    .option("numPartitions", "16")
    .option("fetchsize", "10000")
    .load())

Spark turns this into 16 queries with disjoint ranges on trip_seq.

The rules:

  • partitionColumn must be numeric, date or timestamp. A string primary key will not work, which catches people with UUID keys.
  • All four options go together. Setting numPartitions alone does nothing.
  • The bounds are for splitting, not filtering. Rows outside them are still read, in the first and last partitions, so badly chosen bounds create skew rather than dropping data.
  • fetchsize controls JDBC round trips and often matters more than partitioning for a moderate table.

And the caution to volunteer: 16 parallel connections is 16 concurrent queries against a production database. Parallelism here is a load decision on the source system, not just a Spark setting.

62. What is the difference between `cache` on a DataFrame and a temporary view?

They do unrelated things. Caching stores a computed result; a view stores a query.

df.createOrReplaceTempView("trips")   # names a plan
df.cache()                            # stores the data
spark.sql("CACHE TABLE trips")        # names a plan AND stores the data

So reading a temporary view re-executes its plan every time. If the plan is expensive and you read it repeatedly, you want caching as well, which is what CACHE TABLE does.

One behavioural difference worth knowing: CACHE TABLE is eager by default, populating immediately, while df.cache() is lazy and populates on the next action. CACHE LAZY TABLE makes the SQL form lazy.

Two related notes:

  • Caching is by plan, not by name. Two DataFrames with identical plans share a cache entry, which is why caching then re-deriving the same query can hit the cache unexpectedly.
  • A cached DataFrame is pinned until released. unpersist() or spark.catalog.clearCache(); otherwise it holds storage memory for the life of the application and squeezes execution memory from every later stage.

Architecture and execution

How a line of code becomes work on a cluster. Depth here is the strongest signal an interviewer gets.

63. Walk through what happens when you submit a Spark application.

Ten steps, and being able to narrate them in order is the answer.

flowchart TB
  A["<b>Startup</b> (1-4)<br/>submit, driver starts,<br/>executors requested and registered"]
  B["<b>Plan</b> (5-7)<br/>code builds a plan, an action submits a job,<br/>Catalyst optimises, DAGScheduler cuts stages"]
  C["<b>Execute</b> (8-9)<br/>TaskSchedulerImpl places tasks,<br/>executors run them and write shuffle output"]
  D{"more stages?"}
  E["<b>Finish</b> (10)<br/>results to the driver,<br/>or written by the executors"]
  A --> B --> C --> D
  D -->|"yes"| C
  D -->|"no"| E

Two details that make the answer stronger:

  • Step 5 is where laziness lives. Everything between the session and the first action is plan construction, so a job that fails at step 6 may have a bug written many lines earlier.
  • Step 10 forks. collect() pulls every row back to the driver; write has each executor write its own partition and return only a status. Same job shape, completely different driver cost.

Steps 8 and 9 repeat per stage, and that loop is where adaptive query execution intervenes, because between two passes it has real statistics from the stage that just finished.

64. What does the driver actually do, step by step?

Nine responsibilities, in roughly the order they happen:

  1. Creates the SparkSession, bringing up SparkContext, SparkEnv, the block manager and the RPC endpoints executors call back on.
  2. Asks for executors. SchedulerBackend requests resources from whichever cluster manager --master named.
  3. Turns code into a plan. Unresolved, then analysed against the catalog, then optimised by Catalyst, then a physical plan chosen by cost.
  4. Cuts the plan into stages. DAGScheduler walks the RDD lineage and breaks it at every ShuffleDependency, producing a TaskSet per stage.
  5. Places every task. TaskSchedulerImpl and TaskSetManager decide which executor runs which task, preferring the best available locality.
  6. Tracks every attempt. Retries up to spark.task.maxFailures, default 4; speculates slow tasks if enabled; resubmits stages whose shuffle output was lost.
  7. Keeps the block index. BlockManagerMaster knows where every cached partition and shuffle block lives, so executors can fetch from each other.
  8. Receives results, bounded in total by spark.driver.maxResultSize, default 1g.
  9. Serves the UI on port 4040 and writes the event log the history server replays.

The two facts that follow: spark.driver.memory defaults to 1g, which is a coordinator’s budget and not a data-processing one; and the driver is the one process with no redundancy, because items 3 to 7 exist nowhere else.

65. What does an executor do, step by step?

Much less than the driver, deliberately.

  1. Registers. CoarseGrainedExecutorBackend sends RegisterExecutor and waits for RegisteredExecutor before accepting work. Until that round trip completes, a granted container is idle.
  2. Accepts LaunchTask and runs each task in a thread from its pool. Concurrency is spark.executor.cores, which defaults to 1 on YARN and to all cores in standalone mode.
  3. Reads input, runs the generated code, and on a shuffle map stage writes partitioned output to local disk.
  4. Reports back with StatusUpdate per task and a heartbeat every spark.executor.heartbeatInterval, default 10s. Miss enough and the driver declares it lost.
  5. Stores and serves blocks through its BlockManager, which makes a cached DataFrame reusable and lets the next stage fetch shuffle data from peers.

The contract is narrow on purpose: an executor never makes a scheduling decision, never sees the query plan, and never coordinates with another executor about what to run. It receives tasks, runs them, serves bytes.

That is exactly what makes executors disposable. Lose one and its tasks are retried elsewhere with its cached partitions recomputed from lineage. It is also why executors are per-application rather than shared: memory and slots are reserved for one application’s work.

66. What does the cluster manager do?

One thing: decide who gets containers. It is a resource manager, not a work scheduler.

Cluster manager What grants the container What Spark runs there
Standalone Master, tracking registered Worker daemons The worker launches an executor JVM per grant
YARN The ResourceManager, asked by Spark’s ApplicationMaster through YarnAllocator An executor per allocated container
Kubernetes The API server. The driver’s own ExecutorPodsAllocator creates pods An executor per pod

Kubernetes is the clearest illustration of the split, because there is no Spark cluster-manager daemon at all: the driver talks to the API server directly and asks for pods. The component granting resources and the component scheduling work are simply different things, however they are packaged.

Two defaults follow from this being a separate concern: spark.dynamicAllocation.enabled is false, so absent a decision your application holds its executors from first request to exit, idle or not. And Mesos was removed in the 4.x line, so the three rows above are the whole list.

The practical consequence people find surprising: tuning YARN’s scheduler changes what you are given, never how your tasks are distributed across what you were given. That is the driver’s job.

67. What is `DAGScheduler` responsible for?

Stage-level scheduling. Given a job, it:

  1. Walks the RDD lineage backwards from the action.
  2. Cuts a new stage at every ShuffleDependency, absorbing all narrow dependencies into the current stage.
  3. Determines the order stages must run in, and which can run in parallel.
  4. Submits each ready stage to TaskScheduler as a TaskSet.
  5. Handles FetchFailedException by resubmitting the upstream stage, because lost shuffle output cannot be re-fetched.
  6. Computes preferred locations for tasks from the RDD’s getPreferredLocations.

The distinction from TaskScheduler is the point of the question: DAGScheduler decides what stages exist and in what order; TaskSchedulerImpl decides which executor runs which individual task. One is about the shape of the job, the other about placement.

flowchart TB
  J["job"] --> D["DAGScheduler<br/>cuts stages at shuffles"]
  D --> TS["TaskSet per stage"]
  TS --> T["TaskSchedulerImpl<br/>places each task"]
  T --> B["SchedulerBackend"] --> E["executors"]

The stage-resubmission behaviour is the one to volunteer, because it explains a confusing UI: a job that appears to go backwards, re-running a stage that had already completed, is DAGScheduler rebuilding shuffle output that died with an executor.

68. What is `TaskSchedulerImpl` responsible for?

Task-level scheduling: taking a TaskSet and deciding which executor runs each task in it.

Its concerns:

  • Placement by locality. It offers each task the best available locality level, waiting briefly (spark.locality.wait) before settling for a worse one.
  • Retries. Failed tasks are rescheduled up to spark.task.maxFailures, default 4, preferring a different executor when the failure looked executor-specific.
  • Speculation. If enabled, duplicating tasks that lag their peers.
  • Excluding bad executors. If one executor or node keeps failing tasks, it can stop scheduling there so a single bad machine does not burn every task’s retry budget.
  • Fair or FIFO ordering between concurrent TaskSets, per spark.scheduler.mode, default FIFO.

It works with a TaskSetManager per stage, which holds the per-task state, and talks to the cluster through SchedulerBackend.

The reason to know the split: when tasks are not being placed where you expect, the question is a TaskSchedulerImpl and locality question. When stages are not in the shape you expect, it is a DAGScheduler question. Knowing which component owns which symptom is most of debugging.

69. What is `SchedulerBackend`?

The single component that talks to the cluster manager, with an implementation per manager. It is the boundary between Spark’s scheduling and the outside world.

Its job:

  • Request and release executors. Including the dynamic-allocation calls to scale up and down.
  • Receive executor registration and tell the scheduler how many slots exist.
  • Launch tasks onto executors, and relay their status updates back.

Implementations include CoarseGrainedSchedulerBackend as the common base, with YARN, Kubernetes and standalone specialisations, plus LocalSchedulerBackend for local[*].

“Coarse-grained” is worth explaining because the name is opaque: the executor is acquired once and held for the application’s lifetime, running many tasks. The alternative, fine-grained acquisition per task, existed under Mesos and is gone.

This is the layer that makes Spark portable. The whole scheduler above it is identical on every cluster manager, and only SchedulerBackend changes. That is why the same jar runs on YARN and Kubernetes with nothing but --master different.

70. What is a `ShuffleMapStage` versus a `ResultStage`?

Every stage but the last in a job is a ShuffleMapStage; the last is a ResultStage.

  ShuffleMapStage ResultStage
Output Shuffle files on local disk The action’s result
Purpose Feed the next stage Answer the action
Consumed by The next stage’s fetch The driver, or a write
Rerunnable Yes, and resubmitted if output is lost Yes

The reason the distinction exists in the scheduler is recovery. ShuffleMapStage output is registered with MapOutputTracker on the driver, so Spark knows where every map output block lives. If an executor dies and takes blocks with it, DAGScheduler knows which stage to resubmit and which partitions of it are missing.

A ResultStage has nothing registered, because its output goes to the driver or to storage.

This is also why the UI shows the final stage differently, and why a ShuffleMapStage can be skipped: if its output is still registered from a previous job, Spark reuses it rather than recomputing. Those grey “skipped stages” in the UI are a feature, not a problem, and they are what makes reusing a shuffle across two actions cheap.

71. What are the task locality levels?

Five, best first:

Level Meaning
PROCESS_LOCAL The data is in this executor’s memory already
NODE_LOCAL On the same machine, another executor or local disk or HDFS
NO_PREF No preference, for example a JDBC source
RACK_LOCAL On the same rack
ANY Anywhere

TaskSchedulerImpl offers each task the best level available, and if nothing at that level is free it waits spark.locality.wait, default 3s, before settling for the next level down.

flowchart LR
  A["task ready"] --> B{"PROCESS_LOCAL<br/>slot free?"}
  B -->|"yes"| R["run"]
  B -->|"no"| C["wait spark.locality.wait"]
  C --> D{"NODE_LOCAL free?"}
  D -->|"yes"| R
  D -->|"no"| E["degrade further, to ANY"]

The modern caveat is the useful part. Locality was designed for HDFS, where compute ran on the machines holding the blocks. Reading from S3 or another object store means no partition has a meaningful preferred location, so most tasks are NO_PREF or ANY and the waiting buys nothing. On such clusters, lowering spark.locality.wait can measurably reduce scheduling delay.

Where locality still matters is cached data: a PROCESS_LOCAL hit on a cached partition avoids a network fetch entirely.

72. Where does the input partition count come from?

Not from the number of files, which is the assumption that costs people the most time.

FilePartition.maxSplitBytes decides the split size first:

Math.min(defaultMaxSplitBytes, Math.max(openCostInBytes, bytesPerCore))
  • defaultMaxSplitBytes is spark.sql.files.maxPartitionBytes, default 128MB.
  • openCostInBytes is spark.sql.files.openCostInBytes, default 4MB, charged per file so many small files are packed together rather than each getting a task.
  • bytesPerCore is the total size divided by spark.sql.files.minPartitionNum, which falls back to the default parallelism.

The consequence is measurable. Write the same data as 1, 8 and 16 Parquet files, then read each back:

Files written Partitions on read
1 1
8 2
16 2

File count did not determine partition count in any of them, because the files were small enough to pack.

Two directions follow: a directory of small files collapses into far fewer partitions than it has files, and a single file larger than the split size is divided into several. So read the number rather than predicting it:

print(df.rdd.getNumPartitions())
73. How many tasks will a stage have?

One per partition of that stage. Nothing else affects it.

Where that partition count comes from depends on the stage:

  • The first stage takes it from input splitting, so spark.sql.files.maxPartitionBytes and file sizes.
  • A stage after a shuffle takes it from spark.sql.shuffle.partitions, default 200, or whatever adaptive coalescing reduced it to at runtime.
  • After an explicit repartition(n), it is n.

It is not related to the number of executors, and that confusion is the reason the question gets asked. With 200 tasks and 10 executors of 4 cores each, you have 40 slots, so the stage runs in 5 waves. Tasks queue for slots; they are not divided among executors up front.

The arithmetic worth having ready:

slots = num_executors x spark.executor.cores
waves = ceil(tasks / slots)

A stage with fewer tasks than slots leaves cores idle. A stage with thousands of very short tasks spends its time in scheduling. The target is several waves of tasks each doing a meaningful amount of work, which is what “partitions in the low hundreds of megabytes” is really a proxy for.

74. What happens when an executor is lost?

The work is redone, and the reason it can be redone is lineage.

Recovery is per partition, not per job: losing one executor costs the partitions it held, which are rebuilt on whichever executors remain while everything else stands. Its running tasks are retried elsewhere, counting against spark.task.maxFailures.

Cached blocks are recomputed from lineage, unless the storage level replicated them (MEMORY_ONLY_2 and friends), in which case the replica is promoted.

Shuffle output is the special case worth volunteering:

flowchart LR
  A["stage 0 completed"] --> B["executor dies with<br/>its shuffle files"]
  B --> C["stage 1 fetch fails<br/>FetchFailedException"]
  C --> D["DAGScheduler resubmits<br/><b>stage 0</b>"]

That is a stage rerun rather than a task retry, and it is why a job can appear to go backwards in the UI. An external shuffle service avoids it by serving shuffle blocks from a node-level daemon that outlives the executor, which is also why dynamic allocation wants one.

The contrast with the driver is the closing point: executors hold only recomputable state, so they are disposable; the driver holds the plan, the scheduler and the block index, so it is not.

75. Why does losing the driver kill the application, but losing an executor not?

Because of what each holds that exists nowhere else.

  Driver Executor
Holds Plan, scheduler state, block index, SparkContext Task threads, cached blocks, shuffle files
Recoverable from elsewhere No Yes, from lineage or a peer
On failure Application ends Tasks retried, partitions recomputed

An executor’s contents are either derivable (cached partitions, from lineage) or reproducible (shuffle output, by rerunning the stage). So Spark can always reconstruct them.

The driver’s contents are not derivable from anything. It is the only process that knows what the job is, which stages remain, which tasks succeeded, and where every block lives. There is no second copy, so there is nothing to recover from.

What you can do about it:

  • Cluster deploy mode puts the driver inside the cluster, where the manager can restart it. On standalone, --supervise restarts a failed driver.
  • Structured Streaming checkpoints let a restarted driver resume from recorded offsets and state, which is the closest thing to driver fault tolerance in practice.
  • Keep the driver unloaded. Most driver failures are self-inflicted: collect() on a large result, a broadcast that is not small, or hundreds of thousands of tasks to track.
76. What is the external shuffle service, and why use it?

A daemon running on each node that serves shuffle blocks on behalf of the executors that wrote them. spark.shuffle.service.enabled defaults to false.

Without it, shuffle output lives inside the executor process, so the executor must stay alive for as long as anything might fetch from it. With it, the blocks are on the node’s disk and served by the daemon, so the executor can exit freely.

flowchart TB
  subgraph N["without the service"]
    E1["executor<br/>holds shuffle files"] -->|"must stay alive"| F1["reduce tasks fetch"]
  end
  subgraph W["with the service"]
    E2["executor<br/>can exit"] --> D2["node shuffle service<br/>serves the files"] --> F2["reduce tasks fetch"]
  end

Why it matters:

  • Dynamic allocation needs it, or an equivalent. Otherwise Spark cannot release an idle executor that still holds shuffle data someone might want. The alternative is spark.dynamicAllocation.shuffleTracking.enabled, default true, which tracks which executors hold live shuffle blocks and keeps only those.
  • It reduces FetchFailedException, because an executor dying no longer destroys its shuffle output.
  • It offloads serving work from the executor JVM, which is otherwise doing your computation.

The cost is an extra daemon to deploy and operate per node, which on YARN means an auxiliary service in the NodeManager.

77. What is dynamic allocation?

Adding and removing executors according to pending work, rather than holding a fixed set.

spark.dynamicAllocation.enabled defaults to false, so absent a decision your application holds its executors from first request to exit, idle or not.

How it behaves:

  • Scales up when tasks are pending, adding executors exponentially up to maxExecutors.
  • Scales down when an executor has been idle for executorIdleTimeout.
  • Respects minExecutors as a floor, and initialExecutors at start.
--conf spark.dynamicAllocation.enabled=true --conf spark.dynamicAllocation.minExecutors=2 --conf spark.dynamicAllocation.maxExecutors=100

The shuffle problem is the part to know. An idle executor may still hold shuffle blocks a later stage needs, so releasing it would cause fetch failures. Two solutions: an external shuffle service, or spark.dynamicAllocation.shuffleTracking.enabled, which defaults to true and keeps executors that hold live shuffle data.

When it helps: interactive and bursty workloads, notebooks, multi-tenant clusters, and jobs whose parallelism varies a lot between stages.

When it hurts: short jobs, where the ramp-up costs more than it saves, and jobs with strict latency needs, where waiting for executors is worse than holding them. Note also that setting --num-executors alongside it is contradictory: the static value only acts as an initial size.

78. What is speculative execution?

Re-running tasks that are running much slower than their peers, and taking whichever attempt finishes first. spark.speculation defaults to false.

Spark watches the distribution of task durations within a stage. Once a configurable fraction of tasks has completed, any task running more than spark.speculation.multiplier times the median is duplicated on another executor.

When it helps: a genuinely slow machine. A failing disk, a noisy neighbour, a node with thermal throttling. The duplicate runs elsewhere and finishes normally.

When it hurts, and this is the answer interviewers want: data skew. If the task is slow because its partition holds ten times the rows, the duplicate has exactly the same work to do and is equally slow. You have now consumed two slots to get one result no sooner. Worse, it can mask the skew, because the symptom looks like flaky infrastructure rather than an uneven distribution.

So the diagnostic order is: check whether slow tasks also have larger shuffle read. If they do, it is skew and speculation is the wrong tool. If durations vary while input sizes are equal, it is the environment and speculation helps.

One more caution: with a non-idempotent sink, two attempts writing the same output can duplicate data. Speculation assumes writes are safe to repeat.

79. What is the default scheduler mode within an application?

spark.scheduler.mode defaults to FIFO.

That governs how concurrent jobs within one application share executor slots, which only matters when you submit jobs from multiple threads or serve multiple users from one SparkSession.

Mode Behaviour
FIFO Jobs run in submission order; the first takes all the slots it can use
FAIR Slots are shared round-robin between active jobs

Under FIFO, a large job submitted first can starve a small interactive query submitted second, because it holds every slot until it is done. That is precisely the wrong behaviour for a notebook server or a Thrift server, where a quick query should not wait behind an hour-long batch.

--conf spark.scheduler.mode=FAIR --conf spark.scheduler.allocation.file=/etc/spark/fairscheduler.xml

FAIR also supports pools with weights and minimum shares, so you can give an interactive pool priority over a batch pool. A job joins a pool by setting a thread-local property:

sc.setLocalProperty("spark.scheduler.pool", "interactive")

Keep the levels distinct: this is scheduling within an application. Sharing between applications is the cluster manager’s job, with YARN queues or Kubernetes quotas.

80. How do you set JVM options for the driver and executors?

Two configs, one per side:

spark-submit \
  --conf "spark.driver.extraJavaOptions=-Dlog4j.configurationFile=log4j2.properties" \
  --conf "spark.executor.extraJavaOptions=-Dlog4j.configurationFile=log4j2.properties -verbose:class" \
  app.jar

They are the delivery mechanism for almost every JVM-level change: logging configuration, -verbose:class for class-loading problems, -Xss for stack size, heap dump flags, proxy properties, and GC settings.

The client-mode trap is the part that catches people. In client mode the driver JVM is already running by the time your SparkConf executes, so a driver JVM option set in code arrives too late to have any effect. Use --driver-java-options, or a properties file, or set it before the JVM starts.

Two more notes:

  • Do not set -Xmx here. Heap comes from spark.driver.memory and spark.executor.memory; setting it in extra options conflicts with them.
  • Avoid secrets. The Environment tab of the UI displays configs, so a password in extraJavaOptions is visible to anyone who can reach the UI.

To confirm a flag took effect, the Environment tab is authoritative, and --verbose on submit prints where each value came from.

81. What is `spark.driver.maxResultSize` and why does it exist?

A cap on the total serialized size of results returned to the driver, default 1g. Exceeding it aborts the job with an error naming the config.

It exists because without it, collect() on a large dataset would simply exhaust the driver heap and produce an OutOfMemoryError, which is harder to diagnose and can take the application down mid-flight. The cap converts that into a clear, early, named failure.

spark.conf.get("spark.driver.maxResultSize")     # 1g

Note it counts the total across all tasks in an action, not per task, and it is checked as results arrive, so a job can fail partway through collecting.

The right response is almost never to raise it. If you are hitting a 1 GB limit, the question is why a gigabyte of rows needs to be in the driver:

  • Writing out? Use write, so executors write in parallel and only statuses return.
  • Inspecting? show(), take() or limit().
  • Counting or aggregating? Do it as an aggregation, which returns one row.
  • Handing to Python? toPandas() has the same problem; consider mapInPandas to keep the work distributed.

Setting it to 0 disables the check, which trades a clear error for an OutOfMemoryError. Raising it is defensible only when you genuinely need a large result on the driver and have sized spark.driver.memory accordingly.

82. What port does the Spark UI use, and what happens after the application ends?

spark.ui.port defaults to 4040. If that port is taken, Spark increments and tries 4041, 4042 and so on, which is why a second concurrent application appears on a different port.

The UI dies with the application. It is served by the driver, so when the driver exits there is nothing left to serve it. That is a problem, because most investigation happens after a job has failed or finished.

The solution is two-part:

--conf spark.eventLog.enabled=true --conf spark.eventLog.dir=hdfs:///spark-events

The driver writes an event log, and a separately running history server reads that directory and replays the logs, giving you the same UI for completed applications.

Without event logging, a job that failed overnight leaves you with only the driver’s stdout, which is why enabling it is close to mandatory in production.

Two practical notes: event logs grow, so they need retention and ideally compression (spark.eventLog.compress); and on YARN the ResourceManager UI links to the history server, so the “tracking URL” for a finished application lands there rather than at a dead port 4040.

83. How do broadcast variables and accumulators relate to the driver and executors?

They are the two sanctioned ways to move data between driver and executors, one in each direction.

A task runs in another JVM, so an ordinary variable captured by a closure is serialized, shipped and then discarded. Updates to it never come back, which is why both mechanisms exist.

flowchart LR
  D["driver"] -->|"broadcast: one copy per executor"| E["executors"]
  E -->|"accumulator: merged updates"| D

Broadcast, outward. A read-only value cached once per executor rather than serialized per task. The saving is per-task-count: with 10,000 tasks and a 10 MB table, one copy per executor instead of 10,000 copies.

Accumulators, inward. Tasks add, the driver reads. The guarantee is exactly-once only for updates inside actions; in a transformation, a recompute or retry applies the update again.

What neither does: give you shared mutable state. There is no way for two tasks to coordinate through a variable, which is deliberate. If tasks need to agree on something, that is what a shuffle is for.

The common mistake is closing over a large object without broadcasting it, which silently multiplies serialization cost by the task count and often shows up as inexplicably slow task launch.

84. What is Spark Connect?

A client-server architecture where the client builds an unresolved logical plan and sends it to a remote Spark server over gRPC, instead of the client being the driver.

Traditionally your application is the driver: it holds the SparkContext, it must be a JVM (or a JVM plus a Python process), and it must be network-close to the executors. Spark Connect splits that.

flowchart LR
  C["thin client<br/>Python, Scala, Go, JVM"] -->|"unresolved plan<br/>over gRPC"| S["Spark Connect server<br/>holds the real driver"]
  S --> E["executors"]

What it buys:

  • Thin clients. No local JVM, no cluster-side dependencies, so a notebook, an IDE or a web service can drive a cluster.
  • Isolation. User code crashing or leaking memory no longer takes down the driver, because it is not in the driver.
  • Upgradeability. Client and server versions are decoupled, so the cluster can move independently.
  • Multi-language. The protocol, not the JVM, is the interface, which is how non-JVM clients become possible.

The trade is that anything depending on being inside the driver does not work: SparkContext internals, RDD APIs in general, and code that assumes JVM colocation. So it is a DataFrame and SQL interface, which covers most modern usage.

Partitioning and shuffle

The single biggest source of Spark performance problems, and therefore of interview questions. Expect to be pushed for detail here.

85. What is a shuffle?

A redistribution of data across partitions so that rows which must be processed together end up in the same partition.

It has two halves:

  1. The map side writes each record into a bucket for its target partition, sorts or spills as needed, and writes partitioned output to local disk, along with an index of where each partition’s bytes start.
  2. The reduce side fetches its own slice from every map output, over the network, then merges them.
flowchart TB
  subgraph M["map side, stage 0"]
    M1["task 0"] --> F1[("file + index")]
    M2["task 1"] --> F2[("file + index")]
    M3["task 2"] --> F3[("file + index")]
  end
  subgraph R["reduce side, stage 1"]
    F1 --> R1["task 0 fetches<br/>its slice from all 3"]
    F2 --> R1
    F3 --> R1
    F1 --> R2["task 1"]
    F2 --> R2
    F3 --> R2
  end

Three consequences worth stating:

  • It always touches disk, even when nothing spills, because map output is written to local disk before being served. Shuffle is not an in-memory operation.
  • It is all-to-all. With m map tasks and r reduce tasks there are m x r block fetches, which is why very high parallelism on both sides is its own cost.
  • It creates a barrier. Stage 1 cannot start until stage 0 has finished writing every partition, so one straggler delays everything.
86. What triggers a shuffle?

Any wide dependency. In practice:

Operation Why
groupBy, groupByKey Rows with the same key must meet
join without a broadcast Both sides must be co-partitioned on the key
distinct, dropDuplicates Duplicates must meet to be removed
repartition By definition
orderBy, sort Global ordering needs range partitioning
reduceByKey, aggregateByKey Combines locally first, but still shuffles
Window functions Unless already partitioned by the same key

What does not shuffle: map, filter, select, withColumn, union, coalesce, and a join where one side is broadcast.

The reliable way to count them is the plan, because every shuffle is an Exchange:

df.explain("formatted")
Exchange hashpartitioning(city_id#21, 200), ENSURE_REQUIREMENTS

The description tells you the key and the target partition count, so you can see not just that it shuffles but on what.

The optimisation this enables: if two consecutive operations shuffle on the same key, Spark needs only one exchange. So aggregating and joining on the same column is much cheaper than on different columns, and that is often something you can arrange.

87. Why is a shuffle expensive?

Five costs, and naming several of them is the difference between a shallow and a deep answer.

  1. Serialization and deserialization. Every record crossing the boundary must be encoded and decoded.
  2. Disk I/O. Map output is always written to local disk, and spills add more.
  3. Network transfer. Every reduce task pulls from every map task, so the fetch pattern is all-to-all.
  4. Sorting or merging. Depending on the writer path, records are sorted on the map side and merged on the reduce side.
  5. The stage barrier. No reduce task can start before every map task has finished, so the slowest map task sets the floor.

The barrier is the one people forget, and it is why skew is so damaging: a single slow map task holds up the entire downstream stage, wasting every other slot.

There is also a scaling cost: with m map and r reduce tasks there are m x r blocks to track and fetch. Very high parallelism on both sides turns into metadata and connection overhead, which is one reason “just raise the partition count” has a limit.

What to do about it, in order: eliminate the shuffle (broadcast, or bucket the tables); reduce the data entering it (filter and project first); reduce the number of shuffles (aggregate and join on the same key); then size its partitions sensibly.

88. What are Spark's shuffle writers?

SortShuffleManager selects between three writer paths, checked cheapest first:

Writer Chosen when Behaviour
Bypass merge sort No map-side combine, and partition count below spark.shuffle.sort.bypassMergeThreshold Writes one file per partition then concatenates. No sorting
Unsafe (serialized) sort No map-side combine, and the serializer supports relocating serialized objects Sorts serialized records in binary form, so no deserialization
Sort Everything else Sorts deserialized records, spilling as needed. Most general, most expensive

The condition that disqualifies both fast paths is the same: a map-side combine. And a map-side combine is exactly what reduceByKey, aggregateByKey and a partial aggregation do.

So there is a real tension worth articulating:

reduceByKey   ->  less data shuffled, through the slowest writer
groupByKey    ->  more data shuffled, possibly through a faster writer

reduceByKey still usually wins, because network volume dominates writer efficiency, and because groupByKey risks an OutOfMemoryError on a hot key. But knowing why the trade exists, rather than reciting “always use reduceByKey”, is what a deeper question is probing for.

For DataFrames you rarely choose directly: Catalyst inserts partial aggregation automatically, which is a map-side combine, so aggregations take the sort path.

89. What is data skew?

An uneven distribution of rows across partitions, so some tasks have far more work than others.

The cause is almost always the data itself: one key dominates. A city_id column where 60% of trips are in one city, a customer_id where a single enterprise account has millions of rows, or a null key that every unmatched row shares.

flowchart TB
  subgraph E["even"]
    A1["task 0: 1M rows"] --- A2["task 1: 1M"] --- A3["task 2: 1M"] --- A4["task 3: 1M"]
  end
  subgraph S["skewed"]
    B1["task 0: 100K"] --- B2["task 1: 100K"] --- B3["task 2: <b>3.7M</b>"] --- B4["task 3: 100K"]
  end

Why it is so costly: a task is the unit of work and cannot be split. One task with ten times the rows takes roughly ten times as long, and because of the stage barrier, every other slot sits idle waiting for it. A 200-task stage can be limited by one task.

The signatures to recognise:

  • One or two tasks with far longer duration than the median, in the same stage.
  • Those same tasks showing much larger shuffle read.
  • A stage that is 95% complete for a long time.
  • An OutOfMemoryError on one executor while others are fine.

Skew is a data problem that presents as a performance problem, which is why config tuning rarely fixes it.

90. How do you detect skew?

Two places: the UI, and the data.

In the UI, open the stage page and read the task metrics summary, which gives min, 25th percentile, median, 75th and max for duration and shuffle read. A median of 3 seconds with a max of 40 minutes is skew, and the shuffle read column confirms it is data volume rather than a slow machine.

That distinction matters:

Duration Shuffle read Diagnosis
Max » median Max » median Skew. Fix the distribution
Max » median Roughly even Environment. A slow node, GC, or a bad disk

In the data, count rows per key and look at the head of the distribution:

(df.groupBy("city_id").count()
   .orderBy(F.desc("count"))
   .show(20, truncate=False))

For an RDD, count rows per partition directly:

print(rdd.mapPartitionsWithIndex(
    lambda i, rows: [(i, sum(1 for _ in rows))]).collect())

And for a join, remember to check both sides plus the null count on the key, since nulls all hash together:

df.filter(F.col("city_id").isNull()).count()
91. How do you fix skew?

In order of preference, cheapest and least invasive first.

1. Let AQE handle it. spark.sql.adaptive.skewJoin.enabled defaults to true and splits skewed partitions automatically. Check the thresholds first, because they are strict: it needs a partition larger than 5x the median and larger than 256 MB.

2. Broadcast the small side, if there is one. No shuffle means no skew.

trips.join(F.broadcast(cities), "city_id")

3. Salt the hot key. Spread it across N synthetic sub-keys, join, then aggregate the salt away. This is the general answer when the hot key is genuinely large.

4. Separate the hot keys. Process them with a different strategy and union the results. More code, but optimal when there are only a handful.

hot = ["sf"]
a = trips.filter(F.col("city_id").isin(hot)).join(F.broadcast(cities), "city_id")
b = trips.filter(~F.col("city_id").isin(hot)).join(cities, "city_id")
result = a.unionByName(b)

5. Filter nulls out of the join key if they carry no meaning, since they all land in one partition.

What not to do: enable speculation. The duplicate task has the same rows and is equally slow, so it burns a second slot for nothing.

92. How does salting work in practice?

Add a random suffix to the hot key on the fact side, and replicate the dimension side once per possible suffix so every match still exists.

from pyspark.sql import functions as F

N = 16

facts = trips.withColumn("salt", (F.rand() * N).cast("int"))

dims = (cities
    .withColumn("salt", F.explode(F.array([F.lit(i) for i in range(N)]))))

joined = facts.join(dims, ["city_id", "salt"]).drop("salt")
flowchart TB
  A["all 'sf' rows<br/>one partition"] --> B["salt 0..15"]
  B --> C["16 partitions:<br/>sf-0, sf-1, ... sf-15"]
  D["cities row 'sf'"] --> E["replicated x16<br/>sf-0 ... sf-15"]
  C --> F["join, now spread<br/>across 16 tasks"]
  E --> F

The trade is explicit: the dimension side is N times larger. With N=16 and a 1 MB dimension that is 16 MB, which is nothing. With a 3 GB dimension it is 48 GB, which is not. So N stays small, and salting suits a large-fact-to-small-dimension shape.

For an aggregation rather than a join, salting is two-phase and cheaper, because nothing is replicated:

(trips
  .withColumn("salt", (F.rand() * N).cast("int"))
  .groupBy("city_id", "salt").agg(F.sum("fare_amount").alias("part"))
  .groupBy("city_id").agg(F.sum("part").alias("total")))

Aggregate per salted key, then aggregate the partials. The first pass spreads the hot key over N tasks; the second is tiny.

93. What are the adaptive skew join settings and their defaults?
Config Default
spark.sql.adaptive.enabled true
spark.sql.adaptive.skewJoin.enabled true
spark.sql.adaptive.skewJoin.skewedPartitionFactor 5.0
spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes 256MB

Both conditions must hold. A partition is treated as skewed only when it is larger than 5 times the median and larger than 256 MB.

That conjunction is the answer to “AQE skew handling is on but nothing happened”. On a job whose partitions are all under 256 MB, skew handling never fires no matter how uneven they are. A stage with a 2 MB median and a 100 MB outlier is 50x skewed and completely untouched.

# make it fire on a smaller job
spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes", "32MB")

When it does fire, Spark splits the skewed partition into several smaller reduce tasks and replicates the matching partition from the other side so each split still finds its matches. That is salting, done automatically at runtime.

Two limits worth knowing: it applies to joins, so a skewed groupBy is not covered, and it needs the shuffle to have actually happened, so it cannot help a skewed read.

94. What is adaptive query execution?

Re-planning parts of a query at runtime, using real statistics from stages that have already finished, instead of relying on planning-time estimates.

It is on by default since Spark 3.2: spark.sql.adaptive.enabled is true.

Three things it does:

Optimisation What it fixes
Coalescing shuffle partitions A static 200 that is wrong for this data
Switching join strategy An estimate that said “large” when the side is actually small
Splitting skewed partitions One task with far more rows than its peers
flowchart TB
  S["a shuffle stage completes:<br/>real sizes are now known"]
  S --> A["<b>coalesce partitions</b><br/>merge the small ones"]
  S --> B["<b>switch join strategy</b><br/>sort-merge to broadcast"]
  S --> C["<b>split skewed partitions</b><br/>one big task into several"]

Why it works so well: planning-time estimates on derived data are usually wrong. Spark can estimate the size of a table it has statistics for, but after three filters and a join, the estimate is a guess. After the shuffle, it is a measurement.

Reading it: before execution the plan shows AdaptiveSparkPlan isFinalPlan=false, so explain() gives you the planned shape. To see what actually ran, look at the SQL tab after the query finishes.

The main caveat is that AQE works at shuffle boundaries, so a query with no shuffle gets nothing from it, and the first stage is always planned blind.

95. What does AQE partition coalescing do?

After a shuffle, it merges small output partitions into fewer, larger tasks.

The relevant settings:

Config Default
spark.sql.adaptive.coalescePartitions.enabled true
spark.sql.adaptive.advisoryPartitionSizeInBytes 64MB
spark.sql.adaptive.coalescePartitions.minPartitionSize 1MB

Spark writes the full spark.sql.shuffle.partitions worth of output, looks at the actual size of each, and then assigns contiguous runs of them to single reduce tasks so each task gets roughly the advisory size.

200 shuffle partitions, 800 MB total
   -> without coalescing: 200 tasks x 4 MB
   -> with coalescing:     13 tasks x ~64 MB

This is what makes the static default tolerable. Before AQE, 200 had to be tuned per job or you paid for it; now a high value is mostly harmless, because coalescing brings it down.

The practical consequence for tuning: set the advisory size, not the partition count.

spark.conf.set("spark.sql.adaptive.advisoryPartitionSizeInBytes", "128MB")

That expresses the intent (“I want partitions about this big”) and lets Spark do the arithmetic per stage, which is better than one number for a query whose stages differ in size.

Note it only merges; it never splits. A stage whose partitions are too large needs a higher spark.sql.shuffle.partitions, or skew handling.

96. How do you choose the number of shuffle partitions?

Work from two constraints and take whichever is larger.

Constraint 1, partition size. Aim for partitions in the low hundreds of megabytes of shuffled data. Too large and you spill; too small and scheduling dominates.

partitions = total shuffle bytes / target partition size
           = 600 GB / 128 MB
           = ~4,800

Constraint 2, core count. You want a small multiple of your total cores, so every core gets several tasks and stragglers average out.

cores = num_executors x spark.executor.cores = 50 x 4 = 200
partitions = 200 x 3 = 600      (a floor, not a target)

With AQE, the better move is to express the first constraint directly and let coalescing handle the rest:

spark.conf.set("spark.sql.shuffle.partitions", "2000")            # generous
spark.conf.set("spark.sql.adaptive.advisoryPartitionSizeInBytes", "128MB")

Set the count high enough that no partition is too big, and let AQE merge the small ones down. That handles a query whose stages differ in volume, which a single static number cannot.

The diagnostic loop: if the stage page shows spill, partitions are too large. If it shows thousands of tasks lasting milliseconds, they are too small. Those two symptoms are the feedback signal, not a formula.

97. What is the small files problem, and how do you avoid it?

Many tiny files make every later read slow, for reasons that are not about total bytes:

  • Listing cost. Planning must enumerate the files, which on object storage means many API calls.
  • Per-file open cost. Spark charges spark.sql.files.openCostInBytes, default 4MB, per file when computing splits, precisely because opening one is not free.
  • Metadata overhead. Each Parquet file has a footer to read.
  • Compression loss. Small files compress worse, because dictionaries and encodings have less to work with.

Where they come from, and the fix for each:

Cause Fix
High write parallelism coalesce(n) before the write
partitionBy on many values repartition(same cols) first
Streaming with a short trigger Longer trigger, or periodic compaction
Over-partitioning Partition on a lower-cardinality column
(df.repartition("city_id")
   .write.partitionBy("city_id")
   .mode("overwrite").parquet(path))

You can also cap file size rather than count:

spark.conf.set("spark.sql.files.maxRecordsPerFile", "5000000")

which defaults to 0, meaning unlimited.

The structural fix for a continuously written table is a lakehouse format, because Hudi, Iceberg and Delta all provide compaction as a managed table service rather than something you bolt on.

98. What is the difference between `repartition(n)` and `repartition(col)`?

repartition(n) distributes rows round-robin into n partitions, producing even sizes and no relationship between a row’s value and its partition. repartition(col) hash-partitions by that column, so all rows sharing a value land together.

df.repartition(400)                 # even, value-agnostic
df.repartition("city_id")           # by value, into spark.sql.shuffle.partitions
df.repartition(400, "city_id")      # by value, into 400
  repartition(n) repartition(col)
Distribution Even by construction Even only if the key is uniform
Co-locates by value No Yes
Use before Heavy per-row work needing balance A write partitioned by that column, or a join or aggregate on it

The catch with repartition(col) is skew. Hash partitioning a column where one value dominates puts all those rows in one partition, by design. So it fixes file layout and can remove a later shuffle, while potentially creating the very skew you were trying to avoid. Check the distribution before reaching for it.

The most common correct use is immediately before a partitioned write, so each output directory is produced by one task rather than all of them. The second is before several operations that key on the same column, so one exchange serves them all instead of each operator shuffling separately.

99. Does `orderBy` shuffle, and how?

Yes, and it is more expensive than most shuffles because of how it partitions.

A global sort needs range partitioning: partition 0 gets the lowest values, partition 1 the next, and so on, so that concatenating the sorted partitions yields a globally sorted result. But Spark does not know the value distribution in advance, so it must sample the data first to pick range boundaries.

flowchart LR
  A["sample the data<br/><i>an extra job</i>"] --> B["compute range boundaries"]
  B --> C["shuffle by range"]
  C --> D["sort within each partition"]

That sampling job is why orderBy can show an unexpected extra job in the UI before the real work starts.

If you only need ordering within partitions, which is often enough (for example to make a file’s contents ordered, or for a window already partitioned by the same key), use the cheaper form:

df.sortWithinPartitions("started_at")     # no range partitioning, no sampling

Two related notes:

  • A RangePartitioner handles skew badly when many rows share a value, since boundaries cannot split a single value.
  • orderBy before a limit is a special case Spark optimises with TakeOrderedAndProject, which avoids a full global sort by taking the top n per partition and merging.
100. How can you join without a shuffle?

Three ways, and each removes the shuffle for a different reason.

1. Broadcast the small side. Every executor gets the whole small table, so each large-side partition can join locally. No redistribution needed.

trips.join(F.broadcast(cities), "city_id")

Automatic when the estimated size is below spark.sql.autoBroadcastJoinThreshold, default 10485760 (10 MB).

2. Bucket both tables the same way. If both were written with bucketBy on the join key with the same bucket count, matching rows are already co-located in corresponding files, and the metastore records that fact.

(df.write.bucketBy(64, "rider_id").sortBy("rider_id")
   .saveAsTable("lakehouse.trips_bucketed"))

This is the only shuffle-free option for a genuinely large-to-large join.

3. Reuse an existing partitioning. If a DataFrame is already hash-partitioned on the join key, from a previous repartition(col) or an earlier shuffle on the same key, Catalyst can satisfy the join’s requirement without a new exchange. With RDDs this is explicit co-partitioning plus cache.

flowchart TB
  A["shuffle needed?"] --> B{"one side small?"}
  B -->|"yes"| C["broadcast: no shuffle"]
  B -->|"no"| D{"both bucketed<br/>same key, same count?"}
  D -->|"yes"| E["bucketed join: no shuffle"]
  D -->|"no"| F{"already partitioned<br/>on the key?"}
  F -->|"yes"| G["reuse: no shuffle"]
  F -->|"no"| H["sort-merge join: shuffle"]
101. What is the shuffle fetch side, and which knobs matter?

The reduce side pulling its blocks from every map output. It is the half that crosses the network, and the knobs govern how aggressively it does so.

Config Governs
spark.reducer.maxSizeInFlight How much data may be in flight per reduce task
spark.reducer.maxReqsInFlight Concurrent fetch requests
spark.shuffle.io.maxRetries Retries per failed fetch
spark.shuffle.io.retryWait Wait between retries
spark.network.timeout The general network timeout, default 120s

Raising the in-flight limits increases throughput and memory pressure; lowering them is gentler on a loaded cluster.

But the important answer is that fetch-side tuning is rarely the fix. A FetchFailedException almost always means the executor holding those blocks is gone, was killed for exceeding memory, or is too overwhelmed to serve. Raising retries makes the job take longer to fail rather than succeed.

So the diagnostic order is:

  1. Find the executor that was serving those blocks and why it died. The executors tab shows loss reasons.
  2. If it was killed for memory, fix that: overhead, skew, or partition size.
  3. Consider an external shuffle service, so shuffle output survives the executor.
  4. Only then touch fetch settings, and only with evidence of genuinely slow but healthy transfers.
102. What causes `FetchFailedException`?

The reduce side could not get shuffle blocks it needs. The block was written, but whatever was serving it is no longer able to.

The causes, most common first:

  1. The executor died. Killed by the cluster manager for exceeding memory limits, or crashed with an OutOfMemoryError. Its shuffle files went with it.
  2. The executor is alive but overwhelmed, so fetch requests time out while it is busy with GC or computation.
  3. The node died, taking its local disk with it.
  4. Disk full, so the map side could not write or the service cannot read.
  5. Genuine network problems, which are the rarest in practice.
flowchart LR
  A["FetchFailedException"] --> B["look at the EXECUTORS tab,<br/>not the network config"]
  B --> C{"why did it die?"}
  C -->|"killed for memory"| D["overhead, skew,<br/>or partition size"]
  C -->|"OOM"| E["partition too large"]
  C -->|"still alive"| F["GC pressure or overload"]

What Spark does automatically: DAGScheduler treats it as lost map output and resubmits the upstream stage to regenerate it. So the job often recovers, just slowly, which is why you may see a stage run twice.

What to do: fix the executor deaths. An external shuffle service also helps directly, because shuffle blocks are then served by a node-level daemon that outlives the executor.

103. What is shuffle spill, and what does it tell you?

Spill is data written to disk because it did not fit in execution memory. The stage page reports it as two numbers: spill (memory), the in-memory size of what was spilled, and spill (disk), the serialized size actually written.

It happens on either side: the map side while sorting and buffering output, or the reduce side while merging fetched blocks and aggregating.

What it means: the per-task working set exceeds the execution memory each task gets. That memory is roughly the unified pool divided by the number of concurrently running tasks on that executor.

executor memory 16g
  - 300 MB reserved
  x spark.memory.fraction 0.6          -> ~9.6 GB unified pool
  / spark.executor.cores 4             -> ~2.4 GB per running task

So the fixes, in order:

  1. More partitions, so each task’s share is smaller. This is usually the right answer and costs nothing but a config.
  2. Fewer cores per executor, which gives each running task a larger slice of the same pool.
  3. Check for skew, because one oversized partition spills while its peers do not.
  4. More memory, last, because it treats the symptom.

A little spill is not an emergency; Spark is designed to spill rather than fail. Heavy spill relative to shuffle size means you are paying disk I/O for work that should have been in memory.

104. Why can `coalesce(1)` be dangerous?

Because coalesce is a narrow dependency, so it creates no stage boundary, and the reduced parallelism propagates backwards into the stage that feeds it.

flowchart TB
  subgraph C["coalesce(1): ONE stage, 1 task total"]
    R1["read 200 partitions"] --> M1["expensive transform"] --> Z1["coalesce(1)"] --> W1["write 1 file"]
  end
  subgraph R["repartition(1): TWO stages"]
    R2["read 200 partitions"] --> M2["expensive transform<br/><b>200 tasks</b>"] --> S2[("shuffle")] --> W2["write 1 file<br/>1 task"]
  end

With coalesce(1), the read and the expensive transform are in the same stage as the single-partition write, so the whole pipeline runs in one task on one core. A job that should take minutes takes hours, and the cluster sits idle.

repartition(1) inserts a shuffle, which costs a full redistribution but keeps the upstream stage at full parallelism.

So the rule:

Situation Use
Reducing partitions after cheap work coalesce(n)
Reducing to a very small number after expensive work repartition(n)
Needing balance, or more partitions repartition(n)

A middle path is often best: coalesce(200) to trim file count without a shuffle, rather than collapsing to one. And if you need exactly one file, ask whether you actually do, because a single-file requirement usually comes from a downstream tool that could read a directory instead.

Memory, caching and persistence

Where the errors are most alarming and the fixes least obvious. Interviewers use this area to find out whether you understand the memory model or just raise numbers until it works.

105. How is executor memory divided?

UnifiedMemoryManager does the arithmetic in a fixed order:

  1. Subtract 300 MB of reserved memory, a constant for Spark’s own overhead.
  2. Take spark.memory.fraction of what remains, default 0.6, as the unified pool shared by execution and storage.
  3. What is left, 40% of the post-reserved heap, is user memory: your own data structures, UDF state, and Spark internal metadata.
--executor-memory 16g  (16,384 MB)
  - 300 MB reserved              -> 16,084 MB usable
  x 0.6 memory.fraction          -> 9,650 MB unified pool
      of which storageFraction 0.5 -> 4,825 MB storage floor
  remaining 40%                  -> 6,434 MB user memory
  + memoryOverhead               -> ~18,022 MB requested from the cluster
flowchart TB
  C["container: heap + overhead"] --> H["heap 16g"]
  C --> O["overhead, off-heap and native"]
  H --> R["300 MB reserved"]
  H --> U["unified pool, 60%<br/>execution + storage"]
  H --> M["user memory, 40%"]
  U --> E["execution: shuffles, joins, sorts"]
  U --> S["storage: cached blocks, broadcasts"]

Two things people miss: the pool is shared between execution and storage combined, which is usually less than expected; and the container asks the cluster manager for heap plus spark.executor.memoryOverhead, which is why a 16g executor consumes about 18g of cluster capacity.

Off-heap is separate and off by default: spark.memory.offHeap.enabled is false.

106. What is the difference between execution and storage memory?

They are two uses of one pool.

  Execution memory Storage memory
Used for Shuffles, joins, sorts, aggregations Cached blocks, broadcast data
Lifetime The duration of a task Until evicted or unpersisted
If it runs short Spills to disk Evicts blocks, or spills if the level allows
Floor None spark.memory.storageFraction, default 0.5

They borrow from each other, which is what “unified” means. If nothing is cached, execution can use the whole pool. If no shuffle is running, cached blocks can fill it.

spark.memory.storageFraction is not a partition; it is the floor below which cached blocks will not be evicted. Cached data can grow past it when execution is idle, but cannot be pushed below it.

Before Spark 1.6 these were static pools, and the tuning pain of that is why unification happened: a job that cached nothing still had memory reserved for caching it could not use.

The consequence for tuning: spark.memory.fraction is rarely worth changing, because raising it steals from user memory, which UDFs and your own objects need. If you are tempted to raise it, the real problem is usually partition size.

107. Can execution and storage evict each other?

No, and the asymmetry is the whole answer.

Execution can evict cached blocks, down to the storage floor. Storage can never evict execution.

flowchart LR
  E["execution needs memory"] -->|"evicts cached blocks<br/>down to storageFraction"| S["storage"]
  S -.->|"never evicts"| E

The reason is correctness and progress: a task that cannot get execution memory can only spill or fail, whereas a cached block that is evicted can simply be recomputed from lineage. So Spark prioritises the thing that cannot recover.

The practical consequences, and these are the interview-worthy part:

  • Caching a large DataFrame cannot fail a shuffle. Execution will take the memory it needs.
  • But it can make a shuffle spill, because the cache occupied the pool first and execution had to reclaim it, and because whatever stays above the floor is unavailable.
  • So adding a cache() can make a job slower, which is a counterintuitive result worth being able to explain.
  • And an under-populated cache is the worst outcome: you paid to write it, execution evicted most of it, and now every use partly recomputes.

Check the storage tab. A cached dataset showing 40% in memory is costing you without paying you back, and either needs MEMORY_AND_DISK, a serialized level, or not to be cached at all.

108. What is the difference between heap and off-heap memory?

On-heap memory is allocated inside the JVM heap and managed by the garbage collector. Off-heap is allocated outside it, through sun.misc.Unsafe, and is not subject to GC.

  On-heap Off-heap
Managed by JVM garbage collector Spark, explicitly
Config spark.executor.memory spark.memory.offHeap.size
Enabled Always spark.memory.offHeap.enabled, default false
GC pressure Yes No
Risk Long GC pauses Leaks are not reclaimed automatically

Why off-heap exists: a large heap full of small objects makes full garbage collections slow, and GC pauses stall the executor. Storing data off-heap in a compact binary form removes it from the collector’s view entirely.

This is what Tungsten is built on. Spark’s binary row format stores rows as contiguous bytes with a known layout, so it can be read field by field without constructing JVM objects, and can live off-heap.

--conf spark.memory.offHeap.enabled=true --conf spark.memory.offHeap.size=4g

When to bother: very large executors where GC time is a visible fraction of task time, or heavy shuffle and aggregation workloads. Note that off-heap size counts toward the container total, so it must fit inside what the cluster manager granted, which usually means accounting for it in memoryOverhead.

109. What is executor memory overhead?

Memory requested from the cluster manager in addition to the JVM heap, to cover everything that is not heap:

  • Off-heap allocations and Tungsten buffers.
  • The Python worker process, for PySpark.
  • Native libraries, and the JVM’s own non-heap regions such as metaspace, thread stacks and code cache.
  • Netty’s direct buffers for shuffle transfer.

spark.executor.memoryOverhead defaults to a fraction of executor memory with a minimum floor, so a 16g executor requests roughly 18g in total.

--conf spark.executor.memory=16g --conf spark.executor.memoryOverhead=3g

Why it is the answer to a specific failure. When the cluster manager kills a container for exceeding its memory limit, that is the total being exceeded, not the heap. A JVM OutOfMemoryError means the heap; a container kill means the total. They have different fixes, and telling them apart is the skill:

Symptom Cause Fix
java.lang.OutOfMemoryError: Java heap space Heap exhausted More partitions, or more executor.memory
Container killed by YARN or Kubernetes for exceeding limits Total exceeded Raise memoryOverhead

PySpark jobs need noticeably more overhead than JVM ones, because every executor runs Python workers whose memory is entirely outside the heap. A PySpark job being killed while the heap looks healthy is the classic case.

110. What causes a driver `OutOfMemoryError`?

Something is pulling data into the driver, or the driver is tracking too much metadata.

The causes, most common first:

  1. collect() or toPandas() on a large result. This is by far the most common.
  2. A broadcast that is not small. The driver must assemble the whole table before sending it, so broadcast() on a multi-gigabyte side fails on the driver.
  3. A huge number of tasks. The driver tracks metrics and state per task, so a job with hundreds of thousands of tasks needs real memory just for bookkeeping.
  4. A deep or wide query plan. Analysis and optimisation hold the plan tree in driver memory; thousands of chained withColumn calls or very wide schemas add up.
  5. Accumulating results in a loop, appending to a driver-side list across iterations.

spark.driver.memory defaults to 1g, which is a coordinator’s budget.

The distinction from spark.driver.maxResultSize, default 1g, is worth making: that config aborts the job with a clear named error before the heap is exhausted. So a maxResultSize error is Spark protecting you, and an OutOfMemoryError is what happens when the data arrived some other way, such as through a broadcast.

The fix is usually structural, not a larger driver: write instead of collect, aggregate instead of transfer, and check the plan for an accidental broadcast.

111. What causes an executor `OutOfMemoryError`?

A single task needing more heap than it can get. The usual sources:

  1. A partition too large. The task must hold its working set, and if that exceeds its share of the unified pool plus what can spill, it fails.
  2. Skew. One key concentrating rows into one task, so that task alone is oversized.
  3. A high-cardinality aggregation, where the hash map of groups grows large.
  4. groupByKey on a hot key, because all values for one key must fit on the reduce side.
  5. A UDF materialising too much, for example list(rows) inside mapPartitions, or an applyInPandas group that does not fit.
  6. Too many cores per executor, so the pool is divided too many ways.

The arithmetic that explains most cases:

per-task execution memory ~= unified pool / concurrent tasks on the executor
                          ~= 9.6 GB / 4 cores
                          ~= 2.4 GB

So the fixes in order:

  1. More partitions. Halving partition size halves the working set, and this is usually the correct answer.
  2. Check and fix skew, since one bad partition is invisible in an average.
  3. Fewer cores per executor, which gives each task a bigger slice of the same pool.
  4. A serialized storage level if caching is competing.
  5. More executor memory, last.

Raising memory first often just moves the failure to a larger dataset next month.

112. When should you cache?

When both conditions hold: the dataset is used more than once, and producing it was expensive.

base = (spark.read.parquet(path)
    .filter(F.col("started_at") >= "2026-09-01")
    .join(dims, "city_id"))        # expensive: a scan and a shuffle

base.cache()
base.count()                        # populates

by_city = base.groupBy("city_id").count()
by_rider = base.groupBy("rider_id").count()     # reuses the cache

Without the cache, base is recomputed for each aggregation, including the scan and the join.

Good candidates:

  • A cleaned, filtered base dataset feeding several outputs.
  • An iterative algorithm reading the same data each pass.
  • A dataset used by both branches of a self-join or a union.
  • Anything you will query interactively in a notebook.

And cache deliberately, then verify. Call an action to populate it, and check the storage tab: if it is not fully in memory, the cache may be costing more than it saves. Release it when done:

base.unpersist()

One subtlety: caching is by plan, not by variable. Two DataFrames with identical plans share the cache entry, which is occasionally surprising and occasionally very useful.

113. When should you not cache?

Four cases where it costs more than it saves.

  1. Used once. You pay to write the cache and gain nothing. This is the most common mistake, and in a notebook it is everywhere.
  2. Cheap to recompute. A plain scan of a small Parquet file is faster to redo than to serialise into memory.
  3. Memory already tight. The cache occupies the unified pool, execution evicts it, and you end up partly recomputing anyway while having made every shuffle more likely to spill.
  4. Immediately before one big shuffle. The cache and the shuffle want the same pool, and the shuffle is the thing that cannot recover.

There is also a correctness-adjacent case: caching a non-deterministic DataFrame pins one particular result, which may be what you want or may hide a bug. A rand() or a current_timestamp() in the plan behaves differently cached than uncached.

The anti-pattern to name explicitly:

# caching everything, which is caching nothing usefully
df1.cache(); df2.cache(); df3.cache(); df4.cache()

Each competes for the same pool, most get evicted, and every stage now spills more. Cache the one dataset that is genuinely reused, populate it, verify it in the storage tab, and unpersist it when the reuse is over.

114. What are the DataFrame storage levels worth knowing?

Four, and the choice is a memory-versus-CPU trade.

Level When
MEMORY_AND_DISK The default for df.cache(). Spills rather than losing blocks
MEMORY_AND_DISK_SER Memory is tight. More fits, at deserialization cost
DISK_ONLY Recomputation is more expensive than reading back
MEMORY_ONLY Rarely for DataFrames; the RDD default, and it silently drops what does not fit
from pyspark import StorageLevel
df.persist(StorageLevel.MEMORY_AND_DISK_SER)

The serialized variants are the useful lever. A deserialized cache stores rows close to their in-memory representation, which is fast to read and large. A serialized cache is far more compact, so more of the dataset stays in memory, at the cost of CPU to decode on every read. When the storage tab shows a cache only partly resident, serialized is usually the better trade.

The _2 replicated variants put a copy on a second node so an executor loss does not force recomputation. They double the memory cost, so they are for the case where recomputation is genuinely very expensive.

Note DataFrame caching already stores data in a compressed columnar form internally, controlled by spark.sql.inMemoryColumnarStorage.compressed, so the gap between deserialized and serialized is smaller than for RDDs.

115. How do you uncache, and why does it matter?
df.unpersist()                    # one dataset, blocking=False by default
spark.catalog.clearCache()        # everything
spark.sql("UNCACHE TABLE trips")  # a named table

It matters because cached blocks hold storage memory for the lifetime of the application, not the lifetime of the variable. Python’s garbage collector reclaiming a DataFrame object does not release the blocks on the executors: Spark’s block manager still holds them.

So in a long-running application, a notebook, a streaming job, or a multi-stage pipeline, forgotten caches accumulate. Every later stage then has less execution memory, spills more, and runs slower, with no obvious cause. The symptom is a pipeline that gets slower toward the end for no reason visible in the plan.

The habit worth building:

base.cache()
base.count()
# ... the reuse that justified the cache ...
base.unpersist()

Check what is actually held in the storage tab of the UI, which lists every cached dataset with its size and memory fraction. Anything there that you no longer read is pure cost.

One note on unpersist(blocking=True): it waits for the blocks to actually be removed, which is worth using when you need the memory freed before the next stage rather than eventually.

116. Is `cache()` an action?

No. cache() and persist() are lazy: they mark the dataset as one to store, and the blocks are written during the next action that computes it.

df.cache()          # nothing happens yet
df.count()          # computes, and populates the cache on the way
df.groupBy(...)     # now reads from the cache

This matters in two ways:

  • The first action pays the cost. It computes the dataset and writes the cache, so it is slower than the uncached version. The saving starts with the second use.
  • If you never call an action, nothing is cached, which is a silent no-op people hit when they cache and then only inspect the schema.

The idiom is to force population deliberately with a cheap action, usually count(), so the cost lands where you expect rather than inside the first real query.

The exception is SQL. CACHE TABLE is eager by default:

CACHE TABLE trips;          -- populates immediately
CACHE LAZY TABLE trips;     -- behaves like df.cache()

That asymmetry between the DataFrame API and SQL is a favourite detail to probe, because it is easy to assume they behave the same.

117. What is checkpointing, and how does it differ from caching?

Checkpointing writes a dataset to reliable storage and truncates its lineage. Caching stores a copy while keeping the lineage intact.

spark.sparkContext.setCheckpointDir("hdfs:///checkpoints")
df = df.checkpoint()          # eager by default for DataFrames
  Cache Checkpoint
Stored where Executor memory or local disk Reliable storage, HDFS or S3
Lineage Kept Truncated
Survives executor loss Only with _2 levels Yes
Survives application end No The files do
Cost Cheap A full write to reliable storage

The reason lineage truncation matters: with a long chain of transformations, a single lost partition means replaying the whole chain, and the plan itself grows until analysis is slow or overflows the stack. Checkpointing cuts the chain, so recovery reads a file instead of recomputing a hundred steps.

Where you need it: iterative algorithms that build lineage per iteration, and Structured Streaming, where checkpointing is mandatory and stores offsets and state.

localCheckpoint() truncates lineage using executor storage instead of reliable storage: much faster, but lost if the executor is lost, so it trades durability for speed.

A common pattern is to do both, cache so the write is not recomputed, then checkpoint.

118. Why would a long lineage be a problem even if the data fits?

Because lineage is a data structure held on the driver, and it grows whether or not the data does.

Three separate problems:

  1. Planning cost. Every action re-analyses and re-optimises the plan. Catalyst walks the tree repeatedly, so a tree with thousands of nodes makes planning a measurable fraction of the job. In a loop, planning time grows each iteration.
  2. Stack depth. Many Catalyst rules are recursive over the tree. A deep enough plan overflows the thread stack, producing a java.lang.StackOverflowError that looks nothing like a data problem.
  3. Recovery cost. One lost partition replays the entire chain that produced it.

The classic trigger is a loop:

for i in range(200):
    df = df.withColumn(f"c{i}", F.col("x") + i)      # 200 levels deep

or an iterative algorithm where each pass depends on the last.

The remedies:

  • Checkpoint periodically, every few iterations, which truncates the chain.
  • Collapse the operations. One select with 200 expressions instead of 200 withColumn calls.
  • Raise -Xss on the side that threw, as a stopgap for the stack overflow specifically.

The tell that distinguishes this from other slowness: the job spends its time before any task starts. If the stages are quick but the gap between them is long, you are looking at planning, not execution.

119. What is the default serializer, and should you change it?

spark.serializer defaults to org.apache.spark.serializer.JavaSerializer.

Java serialization is general and slow, and produces large output. Kryo is substantially more compact and faster:

--conf spark.serializer=org.apache.spark.serializer.KryoSerializer --conf spark.kryo.registrationRequired=true
conf.registerKryoClasses(Array(classOf[Trip], classOf[City]))

Registering your classes matters: unregistered classes are written with their full class name per record, which throws away much of the benefit. registrationRequired=true turns a missed registration into an error rather than silent bloat.

Where it does and does not help is the important half:

Workload Effect of Kryo
RDDs of custom objects Large
Shuffle of RDD data Large
Caching with a _SER level Large
Broadcast variables Moderate
DataFrame and SQL operations Small

DataFrames mostly use Tungsten encoders, which are generated per schema and already compact, so they bypass spark.serializer for row data. That is why “set Kryo” is standard advice for RDD-era code and much less impactful for modern DataFrame pipelines.

So: set it if you have RDD-heavy code or serialized caching. If your job is pure DataFrame, expect little and look elsewhere.

120. What is Tungsten?

Spark’s execution backend, aimed at making the JVM behave more like hand-written memory-managing code. Three parts:

  1. Binary row format. UnsafeRow stores a row as a contiguous block of bytes with a fixed-size region for primitives and offsets into a variable-length region. Fields are read by offset arithmetic, with no objects constructed.
  2. Explicit memory management. Spark allocates and tracks memory itself, on or off heap, rather than relying on the garbage collector for data.
  3. Whole-stage code generation. Operators in a stage are fused into one generated Java method, so there is no per-row virtual call chain.

Why it matters, concretely: a JVM object per row costs a header, pointer indirection, and GC attention. A million-row partition becomes a million objects. UnsafeRow makes it one byte array plus offsets.

JVM object row   : header + pointers + boxed fields, GC-visible
UnsafeRow        : contiguous bytes, read by offset, GC-invisible

This is the technical answer to “why do DataFrames beat RDD code”. It is not that the DataFrame API is better written; it is that RDDs must materialise your objects because Spark does not know their shape, while a DataFrame’s schema lets Tungsten use a compact layout and generate code against it.

It also explains why a Python UDF is expensive: it forces rows out of that binary format, across a process boundary, and back.

121. How do you diagnose GC pressure?

Read the GC Time column in the executors tab of the UI, and compare it to task time. A rough guide: consistently above 10% of task time is worth investigating, and above 20% is a problem.

The causes, and the fix for each:

Cause Fix
Many small objects, typically RDDs or Python UDFs Move to DataFrames; serialized storage levels
Too much data per executor More, smaller executors rather than few huge ones
A very large heap Full collections take longer; consider several executors instead
Deserialized caching MEMORY_AND_DISK_SER, which stores far more compactly
Long-lived cached data plus heavy execution Unpersist what is not reused

To see what the collector is actually doing, log it:

--conf "spark.executor.extraJavaOptions=-verbose:gc -XX:+PrintGCDetails -XX:+PrintGCTimeStamps"

Two structural points worth making:

  • Off-heap storage removes data from the collector’s view entirely, which is the real fix for GC-bound shuffle and aggregation work.
  • Very large heaps are not free. A 64g executor can spend a long time in a full collection, which is one reason the usual advice is several medium executors rather than one enormous one.

And distinguish GC pressure from spill: GC is CPU spent on memory management, spill is disk I/O from insufficient memory. They often appear together, but the remedies differ.

122. What does the storage tab of the UI tell you?

For every cached dataset: its name, storage level, the number of cached partitions, the fraction held in memory, and the sizes in memory and on disk.

The single most useful number is the memory fraction. What it tells you:

Fraction in memory Meaning
100% The cache is working as intended
Partial, with a MEMORY_AND_DISK level The remainder is being read from disk each time
Partial, with MEMORY_ONLY The remainder is being recomputed each time
Not listed at all Nothing was cached, usually because no action ran after cache()

That third row is the trap: an RDD cached MEMORY_ONLY that only half fits silently recomputes half the partitions on every use, which looks like the cache having no effect at all.

What to do with the information:

  • Partial residency means either switch to a serialized level so more fits, accept MEMORY_AND_DISK, cache a smaller projection, or do not cache.
  • Several datasets listed when you only meant to cache one is the “cache everything” anti-pattern, and they are competing for the same pool.
  • An empty tab on a job you expected to use a cache means the cache() call never populated.

It is also the place to confirm an unpersist() actually released the blocks rather than just dropping your reference.

Joins, Catalyst and adaptive execution

Join strategy is where most real tuning happens, and where an interviewer can tell in one answer whether you have read a physical plan before.

123. What join strategies does Spark have?

Five, and Catalyst picks one per join based on size estimates, join type and hints.

Strategy How it works Shuffle Good when
Broadcast hash join Ship the small side to every executor, build a hash table locally None One side is small
Shuffle hash join Shuffle both, build a hash table per partition from the smaller side Both sides One side much smaller, but not broadcastable
Sort-merge join Shuffle both, sort each partition, merge Both sides Both sides large
Broadcast nested loop Broadcast one side, compare every pair One side Tiny input, or a non-equi condition
Cartesian product Full cross product Both sides Explicit cross joins only

The selection order is roughly: can it broadcast, then can it shuffle-hash, then sort-merge, and for non-equi conditions only the nested-loop options are available at all.

flowchart TB
  A["equi-join?"] -->|"no"| N["BroadcastNestedLoop<br/>or Cartesian"]
  A -->|"yes"| B{"one side <<br/>autoBroadcastJoinThreshold?"}
  B -->|"yes"| BH["BroadcastHashJoin"]
  B -->|"no"| C{"one side small enough<br/>to hash per partition?"}
  C -->|"yes"| SH["ShuffleHashJoin"]
  C -->|"no"| SM["SortMergeJoin"]

The last two rows of the table are the ones to notice in a plan: seeing CartesianProduct or BroadcastNestedLoopJoin where you expected an equi-join almost always means the join condition is not what you think.

124. What is a broadcast hash join and when does Spark choose it?

Spark sends the entire smaller side to every executor, where each builds a hash table in memory. Each partition of the large side then joins locally, so the large side is never shuffled.

flowchart TB
  D["driver collects<br/>the small side"] --> E1["executor 1<br/>hash table"]
  D --> E2["executor 2<br/>hash table"]
  D --> E3["executor 3<br/>hash table"]
  L1["large partition"] --> E1
  L2["large partition"] --> E2
  L3["large partition"] --> E3

Spark chooses it automatically when the estimated size of one side is below spark.sql.autoBroadcastJoinThreshold, default 10485760, which is 10 MB. Set it to -1 to disable automatic broadcasting.

spark.conf.get("spark.sql.autoBroadcastJoinThreshold")   # 10485760

Why it is the fastest option: no shuffle at all on the large side, which removes the network transfer, the disk write and the stage barrier.

The costs, which a good answer names:

  • The driver assembles it first. The small side is collected to the driver before being broadcast, so a too-large broadcast fails on the driver, not the executors.
  • Memory per executor. Every executor holds a full copy, so 200 MB broadcast across 50 executors is 10 GB of cluster memory.
  • It is bounded by join type. A full outer join cannot be broadcast, and for a left outer join only the right side can be.
125. How do you force or prevent a broadcast?
from pyspark.sql.functions import broadcast
joined = trips.join(broadcast(cities), "city_id")
SELECT /*+ BROADCAST(c) */ *
FROM trips t JOIN cities c ON t.city_id = c.city_id

The available hints, and what each forces:

Hint Effect
BROADCAST Broadcast that side, regardless of the estimate
MERGE Force a sort-merge join
SHUFFLE_HASH Force a shuffle hash join
SHUFFLE_REPLICATE_NL Force a replicated nested loop

To prevent broadcasting globally:

spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1")

When forcing is legitimate: Spark’s estimate is wrong, usually because statistics are missing or the side is derived through several transformations, and you know it is small.

When it backfires: hinting BROADCAST on something genuinely large. The driver must collect it first, so you move the failure from a slow join to a driver OutOfMemoryError, which is harder to diagnose. A hint overrides the safety check, so it is a claim you are making about the data.

A better habit than hinting is often to fix the estimate: ANALYZE TABLE ... COMPUTE STATISTICS, or filter and project the small side first so it is genuinely below the threshold. And with AQE enabled, a runtime switch to a broadcast often happens without a hint at all.

126. What is `spark.sql.broadcastTimeout`?

How long Spark will wait for a broadcast to be built and distributed before failing the query. The default is 300 seconds.

The error is explicit about the config, which is a clue people sometimes miss: it tells you to either raise the timeout or disable broadcast joins.

spark.conf.get("spark.sql.broadcastTimeout")   # 300

Why raising it is usually the wrong move. Hitting a five-minute timeout to build what should be a small hash table means the side is not small. The sequence of events is: Spark estimated it as broadcastable, began collecting it to the driver, and that collection is taking minutes.

So the diagnosis order:

  1. Check the actual size of the side being broadcast. If it is hundreds of megabytes, the estimate was wrong.
  2. Check whether a hint forced it. A BROADCAST hint overrides the size check entirely.
  3. Check for a missing filter. A side that should be small after filtering is being broadcast before the filter, or the filter cannot push down.
  4. Then consider the timeout, legitimately, if the table is genuinely small but the cluster is slow or contended.

The alternative to raising it is to remove the broadcast and let the join shuffle, which is slower per-query but predictable:

spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1")
127. Why is sort-merge join the default for large joins?

Because its memory requirement does not scale with the size of either side.

A hash join must hold one side’s hash table in memory per partition. A sort-merge join sorts both sides, then walks them together with two cursors, so it only needs buffers for the merge, not the whole partition.

left  (sorted):  a a b c d ...
right (sorted):  a b b d e ...
                 ^ ^            advance the smaller cursor, emit matches

And the sort can spill to disk, so a partition larger than memory still completes rather than failing.

spark.sql.join.preferSortMergeJoin defaults to true, which biases the planner toward it over shuffle hash join.

  Sort-merge Shuffle hash
Memory Merge buffers, spillable A whole hash table per partition
Extra work Sorting both sides None
Fails on a huge partition No, spills Possibly

So the trade is CPU for robustness: you pay the sort in exchange for bounded memory. For two large tables that is the right trade, which is why it is the default.

A useful corollary: if both sides are already sorted on the join key, for example from bucketing with sortBy, the sort is skipped and sort-merge becomes very cheap.

128. What is a shuffle hash join, and when is it better than sort-merge?

Both sides are shuffled on the join key, then for each partition Spark builds a hash table from the smaller side and probes it with the larger.

It beats sort-merge when one side is small enough to hash per partition but too large to broadcast, because it skips the sort entirely.

sort-merge   : shuffle both, sort both, merge
shuffle hash : shuffle both, hash the small side, probe

The conditions Spark requires before choosing it:

  • The join is an equi-join.
  • One side is estimated small enough that a per-partition hash table fits, judged against spark.sql.autoBroadcastJoinThreshold multiplied by the number of shuffle partitions.
  • spark.sql.join.preferSortMergeJoin is false, or the cost model prefers it.

Because preferSortMergeJoin defaults to true, Spark picks it less often than you might expect, and the usual reason is exactly that conservatism: sort-merge cannot run out of memory, and shuffle hash can.

To force it when you know the shape:

SELECT /*+ SHUFFLE_HASH(c) */ * FROM trips t JOIN cities c ON t.city_id = c.city_id

The honest summary for an interview: it occupies a narrow band between broadcast and sort-merge, and in modern Spark, AQE converting a planned sort-merge into a broadcast at runtime covers much of that band more safely.

129. What is the difference between a bucketed join and a broadcast join?

Both avoid a shuffle, for entirely different reasons and at different scales.

  Broadcast join Bucketed join
Avoids the shuffle by Copying the small side everywhere Pre-arranging both sides on disk
Size limit One side must be small None
Set up Nothing, automatic Write both tables bucketed, in advance
Needs a metastore No Yes, so saveAsTable
Cost Memory per executor A one-off write, plus rigidity

A broadcast join is a runtime decision about a small side. A bucketed join is a storage layout decision that pays off on every future join.

(trips.write.bucketBy(64, "rider_id").sortBy("rider_id")
   .mode("overwrite").saveAsTable("lakehouse.trips_b"))
(riders.write.bucketBy(64, "rider_id").sortBy("rider_id")
   .mode("overwrite").saveAsTable("lakehouse.riders_b"))

spark.table("lakehouse.trips_b").join(spark.table("lakehouse.riders_b"), "rider_id")

The conditions are strict, and naming them is the good answer: both sides bucketed, on the same column, with the same number of buckets. A mismatch in bucket count means Spark shuffles anyway, silently, and you have paid the write cost for nothing.

So bucketing is the only shuffle-free route for large-to-large joins, and its drawback is inflexibility: the bucket count is baked in, and a join on a different key gets no benefit.

130. Why might a broadcast join not happen even though the table is small?

Because Spark broadcasts on an estimate, and the estimate can be wrong or unavailable.

The reasons, most common first:

  1. No statistics. For a file-based table Spark falls back to the total file size, which can be far larger than the in-memory size after filtering, or unavailable for some sources.
  2. The side is derived. After a filter and a join, the estimate is a computed guess. Spark multiplies selectivity factors, and those compound badly.
  3. The filter cannot push down, so the estimate is of the unfiltered table.
  4. Compression. A 50 MB Parquet file may be 500 MB in memory, and Spark’s estimate accounts for that, so a “small” file can exceed the threshold.
  5. The join type forbids it. A full outer join cannot be broadcast, and for outer joins only the nullable side can be.

The fixes, in order of preference:

ANALYZE TABLE lakehouse.cities COMPUTE STATISTICS FOR ALL COLUMNS;
cities_small = cities.filter(F.col("active")).select("city_id", "name")
trips.join(F.broadcast(cities_small), "city_id")

And the modern answer: adaptive query execution largely solves this. Because AQE re-plans after the shuffle stage completes, it knows the actual size and can convert a planned sort-merge join into a broadcast. So on Spark 3.2 and later, a missing broadcast at plan time often becomes a broadcast at runtime, which is why the SQL tab shows a different plan than explain() did.

131. What is cost-based optimisation and is it on?

Cost-based optimisation uses table and column statistics to make planning decisions that rule-based optimisation cannot: which join order is cheapest, and which physical strategy to pick.

spark.sql.cbo.enabled defaults to false.

Even when enabled it needs statistics to exist:

ANALYZE TABLE lakehouse.trips COMPUTE STATISTICS FOR ALL COLUMNS;
spark.conf.set("spark.sql.cbo.enabled", "true")
spark.conf.set("spark.sql.cbo.joinReorder.enabled", "true")

What it can do that rules cannot: reorder a multi-way join so the most selective join happens first, which changes the size of every intermediate result.

Why it is off by default, and why that is defensible: CBO depends on statistics being present, current and representative. In a lakehouse where tables are written continuously, statistics go stale quickly, and a confidently wrong estimate is worse than no estimate. Collecting column statistics on a large table is also itself an expensive scan.

And adaptive execution has taken over much of its job using real runtime numbers rather than stored estimates, which do not go stale. So the modern answer to “should I enable CBO” is usually: make sure AQE is on, and reach for CBO only for complex multi-way joins over stable tables where join order is the problem.

132. What does AQE do about join strategy?

After a shuffle stage completes, AQE knows the actual size of each side, so it can replace the planned strategy with a better one. The main conversion is demoting a sort-merge join to a broadcast hash join.

flowchart LR
  A["plan time:<br/>estimate says 800 MB<br/>-> SortMergeJoin"] --> B["stage runs"]
  B --> C["actual: 6 MB"]
  C --> D["re-plan:<br/>BroadcastHashJoin"]

This is the single most valuable thing AQE does, and the reason is structural: planning-time estimates on derived data are usually wrong. Spark can estimate a base table it has statistics for, but after two filters and a join the number is a product of guessed selectivities. After the shuffle, it is a measurement.

The relevant setting is spark.sql.adaptive.autoBroadcastJoinThreshold, which if unset falls back to the static spark.sql.autoBroadcastJoinThreshold.

Two consequences worth knowing:

  • explain() lies about the final plan. Before execution the tree is AdaptiveSparkPlan isFinalPlan=false, so what you read is the planned shape. The SQL tab after the query shows what ran.
  • There must be a shuffle for AQE to act at. A query with no exchange gets no adaptive help, and the first stage is always planned blind.

A related optimisation follows automatically: once the join becomes a broadcast, the already-written shuffle can be read locally rather than over the network.

133. What is a local shuffle reader?

An AQE optimisation that avoids network traffic when a join strategy changes at runtime.

The situation: Spark planned a sort-merge join, so both sides were shuffled. Then AQE discovered one side is small and converted the join to a broadcast. The shuffle files already exist, and with a broadcast join each task no longer needs a specific partition from every map output; it just needs the rows. So instead of fetching them over the network in the shuffle pattern, each executor reads the shuffle blocks it wrote locally.

spark.sql.adaptive.localShuffleReader.enabled defaults to true.

without: every reduce task fetches from every map output    (network)
with   : each executor reads its own local shuffle blocks   (local disk)

In the plan it appears as LocalShuffleReader above the exchange.

The reason to know it exists is mainly for reading plans: seeing LocalShuffleReader tells you AQE changed a join strategy after the shuffle had already been written, which is a useful signal that your plan-time estimates were wrong. That in turn tells you statistics might be worth collecting, or that a filter should be pushed earlier so the plan is right the first time.

134. How do you read whether AQE actually changed the plan?

Compare the planned tree with what actually ran, because they are different objects.

Before execution, explain() shows the planned shape wrapped in an adaptive node:

AdaptiveSparkPlan isFinalPlan=false
+- SortMergeJoin [city_id#21], [city_id#40], Inner
   :- Sort [city_id#21 ASC], false, 0
   :  +- Exchange hashpartitioning(city_id#21, 200)
   ...

After execution, the same query’s entry in the SQL tab of the UI shows isFinalPlan=true and the strategies AQE settled on. You can also re-run explain() on an already-executed DataFrame and get the final plan.

What to look for as evidence AQE acted:

Sign Meaning
isFinalPlan=true You are looking at what ran
BroadcastHashJoin where the plan said SortMergeJoin A runtime strategy switch
LocalShuffleReader Follows such a switch
CustomShuffleReader or coalesced partition counts Partition coalescing happened
A stage with far fewer tasks than spark.sql.shuffle.partitions Coalescing again

The practical habit: never conclude anything about a modern Spark query from explain() alone. It tells you the starting point. The SQL tab tells you the result, and when the two differ, the difference is exactly what AQE contributed.

135. What is the difference between a semi join and an anti join?

Both filter the left side by whether a match exists on the right, and neither returns any right-side columns.

trips.join(cities, "city_id", "left_semi")    # rows WITH a match
trips.join(cities, "city_id", "left_anti")    # rows WITHOUT a match
  Left semi Left anti
Returns Left rows that match Left rows that do not match
Columns Left only Left only
Duplicates One left row per left row, regardless of how many right matches Same
SQL equivalent WHERE EXISTS or IN WHERE NOT EXISTS

Why they are better than the alternatives, which is the point of the question:

  • Against an inner join plus select: a semi join stops at the first match, so a left row matching 100 right rows is emitted once. An inner join emits it 100 times, and then you have to deduplicate.
  • Against a left_outer plus isNull filter for the anti case: the outer join carries all the right-side columns through the shuffle before you discard them.

So semi and anti joins are both correct where the alternatives need a distinct, and cheaper, because less data crosses the shuffle.

The null behaviour is worth knowing: for a left anti join, a null key on the left never matches, so those rows are returned, which is usually what you want but occasionally surprising.

136. How do you handle a join key with nulls?

Null never equals null in SQL semantics, so a null key matches nothing. Three consequences and their handling:

1. Inner joins silently drop null keys. If the null means “not applicable”, that is correct. If it means “unknown but should match”, you must say so:

a.join(b, a.city_id.eqNullSafe(b.city_id))      # null == null is true here

2. Null keys are a skew source. Every null hashes to the same partition, so a column that is 30% null puts 30% of the rows in one task. If they carry no meaning, drop them before the join:

a.filter(F.col("city_id").isNotNull()).join(b, "city_id")

Or replace them with a distributed sentinel, which is salting for nulls:

a.withColumn("join_key", F.coalesce("city_id",
    F.concat(F.lit("__null_"), (F.rand() * 16).cast("int"))))

3. Outer joins produce nulls too, so after a left outer join you cannot distinguish “matched with a null value” from “did not match” by looking at the column. Add a marker before joining if you need to tell them apart:

b_marked = b.withColumn("_matched", F.lit(True))
a.join(b_marked, "city_id", "left").withColumn(
    "had_match", F.coalesce("_matched", F.lit(False)))

The habit: check the null count on every join key before writing the join. It answers both the correctness and the skew question at once.

137. What causes a Cartesian product by accident?

A join condition Spark cannot treat as an equality, so it has no key to partition on and must compare every pair.

The usual causes:

  1. A non-equi condition. ON a.ts BETWEEN b.start AND b.end has no equality, so there is nothing to hash.
  2. The condition moved into a where after a cross join. Semantically equivalent SQL, but if Spark cannot push it into the join it evaluates the product first.
  3. A join on an expression that Catalyst cannot match to a key, for example ON cast(a.id AS string) = b.id.
  4. A missing condition entirely, such as a.join(b) with no key, or a multi-column join where one predicate was forgotten.

In the plan it shows as CartesianProduct or BroadcastNestedLoopJoin:

CartesianProduct
:- FileScan parquet [...]
+- FileScan parquet [...]

How to catch it: read the plan whenever a join is unexpectedly slow. A Cartesian on two 1-million-row tables is a trillion comparisons, so the symptom is a stage that never finishes rather than one that is merely slow.

There is a guard: spark.sql.crossJoin.enabled historically had to be set for implicit cross joins, forcing you to be explicit with crossJoin().

For a legitimate range join, the usual technique is to add an equality that buckets both sides, for example joining on a truncated date as well as the range condition, so Spark has a key to partition on and the nested-loop comparison happens only within each bucket.

138. What is dynamic partition pruning?

Using the result of a filter on one side of a join to prune partitions of the other side at runtime, rather than scanning all of them.

The classic star-schema shape: a large partitioned fact table, joined to a small filtered dimension table.

SELECT t.*
FROM trips t JOIN cities c ON t.city_id = c.city_id
WHERE c.region = 'west'

Static partition pruning cannot help, because the filter is on c.region while trips is partitioned by city_id. Without DPP, Spark scans every partition of trips and discards most of them after the join.

With DPP, Spark runs the dimension side first, collects the matching city_id values, and injects them as a partition filter on the fact scan.

flowchart LR
  A["filter cities<br/>region = 'west'"] --> B["city_ids: sf, la, sea"]
  B --> C["inject as a partition filter"]
  C --> D["scan only those<br/>trips partitions"]

The conditions: the fact table must be partitioned on the join key, the join must be an equi-join, and the dimension side must be filtered and broadcastable.

In the plan it appears as a dynamicpruningexpression in the fact-side PartitionFilters.

The payoff is often the largest single win available on a star schema, because it changes how much data is read rather than how efficiently it is processed.

139. How do you join a large table to a large table efficiently?

Broadcasting is off the table, so the work is in reducing what enters the shuffle and making the shuffle well-shaped.

In order:

1. Reduce both sides first. Filter and project before joining. This is the largest lever and the most often skipped:

t = trips.filter(F.col("started_at") >= "2026-09-01").select("trip_id", "rider_id", "fare_amount")
r = riders.filter(F.col("active")).select("rider_id", "segment")

Sometimes this alone makes one side broadcastable, which changes the problem entirely.

2. Check whether one side can become small. An aggregation before the join often shrinks it by orders of magnitude.

3. Bucket both tables on the join key with the same bucket count, if this join is a recurring pattern. That removes the shuffle permanently.

4. Size the shuffle. Aim for partitions in the low hundreds of megabytes, and set spark.sql.adaptive.advisoryPartitionSizeInBytes rather than guessing a count.

5. Handle skew on the join key, since a large-to-large join is where skew hurts most. AQE skew handling first, then salting.

6. Avoid re-shuffling. If you then aggregate on the same key, one exchange serves both. If you aggregate on a different key, you pay a second shuffle, so consider whether the order can change.

The one thing not to do is hint BROADCAST and hope.

140. What is the `Exchange` node in a plan?

A shuffle. Every Exchange in a physical plan is one redistribution of data, and therefore one stage boundary.

Its description carries the useful detail:

Exchange hashpartitioning(city_id#21, 200), ENSURE_REQUIREMENTS, [plan_id=141]
Exchange rangepartitioning(started_at#30 ASC NULLS FIRST, 200), ...
Exchange SinglePartition, ...
Partitioning Means
hashpartitioning(cols, n) Hash on those columns into n partitions. The usual case
rangepartitioning(cols, n) Range partitioning, from a global sort
SinglePartition Everything into one partition. A red flag unless intended

The ENSURE_REQUIREMENTS tag means the exchange was inserted because an operator above it required that partitioning, which is the normal reason.

How to use it:

  • Count them to count your shuffles, and therefore your stages.
  • Read the key to see whether consecutive operations shuffle on the same column, in which case one exchange can serve both.
  • Watch for SinglePartition, which usually comes from a global orderBy, a coalesce(1), or a window with no partitionBy. That last one is a common accidental performance disaster: a window function without partitionBy sends the entire dataset through one task.
141. What is `ReusedExchange`?

A shuffle whose output is read by more than one branch of the plan, rather than being computed twice.

It appears when the same subtree occurs in two places, for example in a self-join, a union of two aggregations over the same base, or a query referencing the same CTE twice.

:- Exchange hashpartitioning(city_id#21, 200)
:  +- HashAggregate(...)
+- ReusedExchange [city_id#40], Exchange hashpartitioning(city_id#21, 200)

Seeing it is good news: Catalyst recognised the common subtree and is reusing the shuffle files rather than recomputing the stage. In the UI this is also what produces “skipped stages”, shown in grey, which are not a problem but a saving.

The interesting case is its absence. If you expected reuse and do not see it, the two branches differ in some way you did not intend:

  • A different filter, even a logically equivalent one written differently.
  • A different column projection, so the plans are not identical.
  • A non-deterministic expression such as rand() or current_timestamp(), which makes the subtrees non-equivalent by definition.

When reuse does not happen and you want it, cache() on the shared base is the explicit way to force one computation, at the cost of storage memory. Checking for ReusedExchange first tells you whether the cache is even necessary.

142. How do you reduce the number of shuffles in a query?

Six techniques, roughly in order of payoff.

1. Aggregate and join on the same key. If both operations need the same partitioning, one exchange satisfies both. Reordering a pipeline so the keys line up can halve the shuffles.

2. Repartition once, deliberately. If several operators key on the same column, one repartition(col) up front means each of them finds the partitioning already satisfied.

3. Replace distinct plus join with a semi join. distinct is a shuffle; a semi join achieves the same filtering in the join’s own shuffle.

trips.join(active_riders, "rider_id", "left_semi")     # not distinct + inner join

4. Broadcast what can be broadcast. A broadcast join has no exchange at all, so reducing a dimension until it fits under the threshold removes a shuffle outright.

5. Avoid unnecessary global ordering. orderBy costs a range-partitioning shuffle plus a sampling job. If you only need order within files, use sortWithinPartitions.

6. Bucket recurring join keys, so the shuffle is paid once at write time rather than on every read.

Then verify, because this is all measurable:

df.explain("formatted")     # count the Exchange nodes before and after

The mindset worth conveying: do not think of shuffles as something to tune, think of them as something to count, and then ask whether each one is earning its place.

Performance tuning and troubleshooting

What an interviewer really wants to know: can you find the bottleneck, or do you guess at knobs.

143. How do you approach a slow Spark job?

Find the bottleneck before changing anything. The order matters more than any individual technique.

  1. Open the SQL tab and find the slowest operator. The per-node timing tells you where the work is, which is often not where you assumed.
  2. Open the stage page for that stage and read the task duration distribution: min, median, max.
  3. Classify it. Almost every slow stage is one of four things:
Signal Diagnosis
Max task time » median, with larger shuffle read Skew
Large shuffle read and write across all tasks Too much shuffle
Spill reported Partitions too large
Thousands of tasks, each milliseconds long Too many partitions
  1. Check the plan for structural mistakes before tuning: a missing broadcast, an accidental Cartesian, a filter that did not push down, columns you did not need.
  2. Change one thing, re-run, and compare the same stage.

The answer that impresses is the negative one: say what you would not do first. You would not raise executor memory, not enable speculation, and not change spark.sql.shuffle.partitions until you knew which of the four it was. Tuning without a diagnosis is how jobs acquire a page of configs that nobody can justify.

144. What does the Spark UI tell you, tab by tab?
Tab What it answers
Jobs Which action is slow, how many stages, which stages were skipped
Stages Task duration and shuffle distribution, spill, the DAG for that stage
SQL The final physical plan and per-operator row counts and time
Storage What is cached, and what fraction is actually in memory
Executors Per-executor task counts, GC time, failures, and loss reasons
Environment The configs actually in effect

How to use them together, which is the real skill:

  • Start in SQL for a DataFrame job, because it maps operators to time. Starting in Stages tells you where but not what.
  • Stages gives you the distribution. The summary metrics table with percentiles is the single most valuable view in the UI, because it separates skew from uniform slowness.
  • Executors answers “why did it fail”. Loss reasons live here: killed for memory, heartbeat timeout, or exited.
  • Environment settles arguments. If someone insists a config is set, this is the authoritative answer.
  • Storage is where caches are exposed as a cost rather than assumed to be a benefit.

And remember the UI dies with the application, which is why event logging plus a history server is close to mandatory: most investigations happen after the job.

145. A few tasks take far longer than the median. What is it?

Skew, in almost every case. Confirm it and then fix the distribution.

Confirm: on the stage page, compare the max and median for both duration and shuffle read size. The second column is what distinguishes the two possible causes:

Duration Shuffle read Cause
Max » median Max » median Skew: the task has more data
Max » median Roughly even Environment: slow node, GC, bad disk

If it is the second, look at the executors tab for GC time and at which host the slow tasks ran on; a single bad machine shows up as all slow tasks on one executor.

If it is skew, fix it in this order:

  1. Check whether AQE skew handling could fire. It needs a partition over 5x the median and over 256 MB, so on a smaller job it silently never triggers.
  2. Broadcast the small side if there is one.
  3. Salt the hot key.
  4. Split the hot keys out and union.

And do not enable speculation. The duplicate task has the same rows and is equally slow, so it consumes a second slot for no benefit and masks the real cause by making it look like flaky infrastructure.

146. The stage shows large spill. What do you change?

Spill means the per-task working set exceeded the execution memory that task could get. So reduce the working set, or increase the share.

The arithmetic that frames the fix:

per-task execution memory ~= unified pool / concurrent tasks per executor
                          ~= (heap - 300MB) x 0.6 / spark.executor.cores

Changes in order of preference:

  1. More partitions, so each task handles less. This is usually the right answer, costs nothing, and scales.
spark.conf.set("spark.sql.adaptive.advisoryPartitionSizeInBytes", "64MB")
spark.conf.set("spark.sql.shuffle.partitions", "2000")
  1. Check for skew. One oversized partition spills while its peers do not, and no amount of memory fixes a distribution problem.
  2. Fewer cores per executor, which gives each running task a larger slice of the same pool. Going from 5 cores to 3 raises per-task memory by two thirds.
  3. More executor memory, last, because it treats the symptom and costs cluster capacity.

Context worth adding: some spill is normal and by design. Spark spills rather than failing, which is a feature. The question is proportion: spill of a few percent of shuffle size is fine, spill several times the shuffle size means you are paying disk I/O for work that should have been in memory.

147. Planning takes longer than execution. What is wrong?

The job is spending its time on the driver before tasks start. Three usual causes:

1. Too many input files. Planning must list them and compute splits, which on object storage is many API calls. Tens of thousands of small files can take minutes before a single task runs.

print(spark.read.parquet(path).inputFiles().__len__())

2. A deep plan. Hundreds of chained withColumn calls, or a loop that rebuilds a DataFrame, produce a very deep tree that Catalyst walks repeatedly. Collapse them into one select, or checkpoint to truncate.

3. Very wide schemas. Thousands of columns, or deeply nested structs, make every analysis and optimisation pass expensive.

How to tell it apart from execution slowness: in the Jobs tab, look at the gap between the job being submitted and its first stage starting. In the SQL tab, a long total duration with short stage durations points the same way.

The extreme version is a StackOverflowError from the recursive tree walk, which is the same problem at a larger scale.

The related structural fix for the file-count case is a lakehouse format, because Hudi, Iceberg and Delta keep file metadata in their own manifests rather than requiring a directory listing, which is one of the main reasons they exist.

148. There are thousands of tasks each lasting milliseconds. What is wrong?

Too many partitions, so per-task overhead dominates the actual work.

Each task carries real fixed costs: serialization of the task description, scheduling on the driver, JVM thread setup, reading a partition’s metadata, and reporting results. At a few milliseconds of work per task, those costs are the job.

It also loads the driver, which must track every task’s state and metrics. A job with hundreds of thousands of tasks can exhaust driver memory on bookkeeping alone.

The fixes depend on where the partitions came from:

Source Fix
A shuffle Lower spark.sql.shuffle.partitions, or let AQE coalesce
Many small input files Compact them, or raise spark.sql.files.maxPartitionBytes
An explicit repartition(n) Lower n
Before a write coalesce(n)

With AQE on, coalescing should already be merging small shuffle partitions toward spark.sql.adaptive.advisoryPartitionSizeInBytes, default 64MB. So if you are seeing thousands of tiny tasks after a shuffle, check that AQE is enabled and that coalescing is not disabled.

The target to state: partitions in the low hundreds of megabytes, and a total that is a small multiple of your core count. Tasks lasting well under a second are a signal you have overshot.

149. How do you size executors?

Favour several medium executors over a few enormous ones or many tiny ones.

A common starting point is 4 to 5 cores per executor, and there are two reasons:

  1. HDFS client contention. Concurrent tasks in one JVM share the same client, and throughput stops improving past roughly that point.
  2. Garbage collection. A very large heap makes full collections slow, and a GC pause stalls every task on that executor.

Then memory follows from partition size:

per-task execution memory ~= (heap - 300MB) x 0.6 / cores

Size the heap so that is comfortably larger than your largest partition’s working set, and add spark.executor.memoryOverhead for off-heap, Python workers and native libraries.

--executor-cores 4 --executor-memory 16g --conf spark.executor.memoryOverhead=3g

The failure modes at each extreme:

Shape Problem
Few huge executors (e.g. 32 cores, 200g) Slow full GC, HDFS contention, coarse scheduling granularity
Many tiny executors (1 core, 2g) No memory sharing between tasks, more overhead copies, more broadcast copies

And state the caveat: these are starting points, not laws. The right answer depends on your partition sizes and workload, and the honest thing to say is that you would start here, measure spill and GC, and adjust.

150. How do you decide `--num-executors` versus dynamic allocation?

Pick one. Setting both is contradictory, and Spark resolves it by treating the static number as merely the initial size.

  Static sizing Dynamic allocation
Config --num-executors spark.dynamicAllocation.enabled=true
Suits Predictable batch jobs Bursty, interactive, multi-tenant
Cost model Reserves capacity for the whole run Releases idle capacity
Startup All at once Ramps up
Needs Nothing Shuffle tracking or an external shuffle service

Use static sizing when the job’s parallelism is roughly constant, you know the shape from experience, and predictable performance matters more than cluster efficiency. A nightly batch job is the archetype.

Use dynamic allocation when parallelism varies a lot between stages, the application is long-lived and often idle (a notebook, a Thrift server), or the cluster is shared and holding idle executors is antisocial.

--conf spark.dynamicAllocation.enabled=true --conf spark.dynamicAllocation.minExecutors=2 --conf spark.dynamicAllocation.maxExecutors=200

The requirement to name: an idle executor may still hold shuffle blocks a later stage needs. So dynamic allocation needs either an external shuffle service, or spark.dynamicAllocation.shuffleTracking.enabled, which defaults to true and keeps executors holding live shuffle data.

The case against dynamic allocation for short jobs: the ramp-up costs more than the elasticity saves.

151. What does a container killed for exceeding memory limits mean?

The cluster manager killed the JVM because the process total exceeded what was requested, which is heap plus overhead. It is not a JVM OutOfMemoryError, and the distinction drives the fix.

Symptom Means Fix
java.lang.OutOfMemoryError: Java heap space The heap filled More partitions, then more executor.memory
Container killed by YARN or Kubernetes for exceeding limits The total exceeded More memoryOverhead

What lives in overhead: off-heap buffers, Netty’s direct buffers for shuffle, the JVM’s metaspace and thread stacks, native libraries, and the Python worker processes for PySpark.

That last one is why PySpark jobs hit this far more often than JVM jobs: Python memory is entirely outside the heap, so a job with heavy Python UDFs can have a perfectly healthy heap and still be killed.

The sequence:

  1. Confirm which it is from the message. This determines everything.
  2. Raise spark.executor.memoryOverhead, generously for PySpark.
  3. Check for skew, because one oversized partition is a common root cause that presents as a memory shortage.
  4. Raise partition count so each task’s footprint shrinks.
  5. Raise executor memory, last.

The point to make: raising memory first often just moves the failure to a larger dataset next month, whereas fixing the partition size or the skew fixes it structurally.

152. What is `spark.network.timeout` and when do you raise it?

The default timeout for network interactions between Spark components, 120s. It acts as the default for several more specific timeouts, including RPC lookups and shuffle fetch waits.

spark.conf.get("spark.network.timeout")   # 120s

When raising it is legitimate: genuinely long operations on healthy infrastructure. Very large shuffle fetches, a heavily contended cluster, or cross-availability-zone traffic.

When raising it is the wrong answer, which is most of the time: the timeout is reporting that something is dead or overwhelmed, not causing the problem. An executor stuck in a long garbage collection cannot answer heartbeats; an executor that was killed cannot serve shuffle blocks. Raising the timeout makes the job take longer to notice, not longer to succeed.

So the diagnosis order:

  1. Look at the executors tab. Did an executor die, and why?
  2. Check GC time. An executor spending 40% of its time in GC will miss heartbeats, and the fix is memory or partition size, not the timeout.
  3. Check for container kills, which point at overhead.
  4. Only then consider the timeout, and pair it with spark.executor.heartbeatInterval, which must stay well below it.

A related note: spark.executor.heartbeatInterval defaults to 10s, and it must be considerably smaller than the network timeout or executors get declared lost spuriously.

153. How do you debug a `NoClassDefFoundError` or a version conflict?

Find out which JAR each class is actually loaded from, rather than reasoning about what should be on the classpath.

--conf "spark.driver.extraJavaOptions=-verbose:class" --conf "spark.executor.extraJavaOptions=-verbose:class"

That prints a line per class load with its source JAR, turning a guess about shading into a direct answer:

[Loaded com.fasterxml.jackson.databind.ObjectMapper from file:/opt/spark/jars/jackson-databind-2.15.2.jar]

The usual causes:

  1. Two versions on the classpath, one from Spark’s jars/ and one from your assembly. Spark’s usually wins, which is why your newer version appears absent.
  2. A provided dependency not actually provided at runtime.
  3. A transitive dependency you did not know you had.

The fixes, in order:

  • Shade the conflicting dependency in your assembly, which is the robust answer: relocate the package so there is no collision.
  • Align versions with what Spark ships, which is the simplest when possible.
  • spark.driver.userClassPathFirst and spark.executor.userClassPathFirst, which put your jars ahead of Spark’s. Effective but risky, because it can break Spark’s own dependencies.
  • --jars or --packages for genuinely missing classes.

The habit to convey: -verbose:class converts an argument into a fact. It is noisy, so filter it, but it answers the question definitively.

154. What causes a `java.lang.StackOverflowError` in Spark?

A recursive walk over a structure deeper than the thread’s stack. In Spark that structure is almost always the query plan or the RDD lineage, not your code’s own recursion.

The plan gets deep from:

  • Hundreds of chained withColumn or filter calls.
  • A loop that rebuilds a DataFrame each iteration.
  • Deeply nested structs or very wide schemas.
  • A long iterative lineage with no checkpoint.

The stack trace is the tell: if the frames are Catalyst classes, an Avro or Parquet schema walker, or TreeNode methods, it is the plan.

Two fixes, and the right one depends on which side threw:

1. Raise the stack on the side that threw, as a direct remedy:

--conf "spark.driver.extraJavaOptions=-Xss64m"
--conf "spark.executor.extraJavaOptions=-Xss64m"

Note in client mode the driver JVM already exists, so use --driver-java-options.

2. Make the plan shallower, which is the structural fix:

  • Collapse many withColumn calls into one select.
  • checkpoint() periodically inside an iterative loop, truncating lineage.

The judgement to express: if you find yourself needing hundreds of megabytes of stack, the plan itself is the problem and raising -Xss is buying time rather than fixing it.

155. How do you use custom logging to see what a job did?

Write a Log4j 2 config, ship it with --files, and point both sides at it.

rootLogger.level = warn
rootLogger.appenderRef.stdout.ref = console
appender.console.type = Console
appender.console.name = console
appender.console.layout.type = PatternLayout
appender.console.layout.pattern = %d{HH:mm:ss} %-5level %logger{1} - %msg%n

logger.sql.name = org.apache.spark.sql.execution
logger.sql.level = debug
spark-submit \
  --files /etc/spark/conf/log4j2.properties \
  --conf "spark.driver.extraJavaOptions=-Dlog4j.configurationFile=log4j2.properties" \
  --conf "spark.executor.extraJavaOptions=-Dlog4j.configurationFile=log4j2.properties" \
  app.jar

Three details that make it work:

  • The value is a bare filename, not a path, because --files places the file in each container’s working directory.
  • Set both sides. The driver logs planning and scheduling; executors log task-level work. Setting one and forgetting the other looks like the feature not working.
  • Spark 3.3.0 and later is Log4j 2. A config in the old log4j.rootLogger= syntax is silently ignored, not rejected.

Useful namespaces to raise selectively: org.apache.spark.sql.execution for plan and codegen detail, org.apache.spark.scheduler for stage and task decisions, org.apache.spark.storage for block and cache activity.

Keep the root at warn and raise one namespace, otherwise the signal is buried. And never leave DEBUG on in production.

156. How do you confirm a config actually took effect?

Three ways, in descending authority:

1. The Environment tab of the UI. It lists the Spark properties actually in force for the running application. This is the authoritative answer and settles most arguments.

2. At runtime, for SQL configs:

print(spark.conf.get("spark.sql.shuffle.partitions"))

Note this works for spark.sql.* configs but raises for many core ones, which are read from SparkConf at startup:

print(spark.sparkContext.getConf().get("spark.executor.memory", "(default)"))

3. --verbose on submit, which prints the resolved configuration and where each value came from.

The cases where a config silently does not apply, which is the real content of this question:

Case Why
Driver JVM options in client mode The JVM already started
spark.executor.* set after the session exists Executors are already running
spark.default.parallelism on a DataFrame job Wrong config; you want spark.sql.shuffle.partitions
A --conf placed after the application jar Passed to your app, not to Spark

That last one is worth emphasising: it is not an error, it is an argument to your main, so Spark never sees it and nothing warns you.

157. What is the difference between `spark.sql.files.maxPartitionBytes` and `spark.sql.shuffle.partitions`?

They control parallelism at two different points, and reaching for the wrong one is a common waste of time.

  spark.sql.files.maxPartitionBytes spark.sql.shuffle.partitions
Controls Input partitions when reading files Output partitions of a shuffle
Default 128MB 200
Affects The first stage Every stage after a shuffle
Units Bytes per partition A count
flowchart LR
  A["files on storage"] -->|"maxPartitionBytes<br/>decides split size"| B["stage 0 tasks"]
  B -->|"shuffle"| C["spark.sql.shuffle.partitions<br/>decides task count"]
  C --> D["stage 1 tasks"]

So if the read stage has too few tasks, lowering maxPartitionBytes gives you more. If the stage after a shuffle has the wrong number, spark.sql.shuffle.partitions or AQE coalescing is the lever.

Note the units differ in a way that matters: one is a size and the other a count. The size one adapts automatically as data grows; the count one does not, which is exactly why a static 200 ages badly and why AQE coalescing exists.

And remember maxPartitionBytes is only one input to the split calculation. spark.sql.files.openCostInBytes, default 4MB, is charged per file, so with many small files the packing behaviour dominates and lowering maxPartitionBytes changes nothing.

158. How do you speed up a job that writes many partitions?

Repartition by the same columns you partition by, so each output directory is written by one task instead of all of them.

(df.repartition("city_id", "dt")
   .write.partitionBy("city_id", "dt")
   .mode("overwrite")
   .parquet(path))

The arithmetic that explains why. Without the repartition, every task may hold rows for every partition value, so it opens a writer per value:

files = tasks x partition values      (worst case)
      = 200 x 250
      = 50,000 small files

With the repartition, each value lives in one task:

files = partition values = 250

That also fixes the write itself, not just later reads: 50,000 open file handles and 50,000 small footers is slow to produce, and on object storage it is thousands of PUT requests.

Other levers:

  • spark.sql.files.maxRecordsPerFile, default 0 (unlimited), caps rows per file so one huge partition splits into several properly sized files rather than one enormous one.
  • Dynamic partition overwrite, so a job writing one day does not rewrite the table:
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")
  • Reconsider the partition column. If it has thousands of values, directory-per-value is the wrong layout and a lower-cardinality column, or bucketing, is better.
159. What are the highest-value things to check before tuning configs?

Six checks, each of which typically outweighs any config change.

  1. Are you reading more than you need? Columns you never select, rows you filter after the join instead of before, partitions you could have pruned. Check ReadSchema and PartitionFilters in the plan.
  2. Is the plan what you think? An accidental CartesianProduct, a SortMergeJoin where a broadcast was possible, a SinglePartition exchange from a window without partitionBy.
  3. Are there small files? Thousands of tiny inputs make planning dominate execution.
  4. Is there skew? Max versus median task time on the slowest stage.
  5. Is something cached that should not be? The storage tab: a partly resident cache is pure cost.
  6. Is a UDF doing work a built-in could do? Every UDF is an optimisation barrier.
flowchart TB
  A["slow job"] --> B["read less?"] --> C["plan correct?"] --> D["file count?"]
  D --> E["skew?"] --> F["caches?"] --> G["UDFs?"] --> H["only now: configs"]

The reason to put it this way in an interview: it shows you treat tuning as diagnosis rather than ritual. A job with a page of carefully tuned configs and a missing broadcast is slower than a job with default configs and a correct plan. Configs are the last 20%, and the plan is the first 80%.

160. How do you benchmark a change honestly?

Change one thing, run at representative scale, and compare the stage you expected to change rather than total wall clock.

The rules that keep a benchmark honest:

  • One variable at a time. Two changes together tell you nothing about either.
  • Representative data. A change that helps on 1 GB may do nothing at 1 TB, and skew often only appears at scale.
  • Repeat runs. A single run on a shared cluster measures your neighbours as much as your change.
  • Compare the right metric. If you changed shuffle partitions, compare that stage’s duration and spill, not the whole job, which includes unrelated stages.
  • Account for warm caches. The second read of the same data is faster for reasons unrelated to your change: page cache, warm JVMs, and Spark’s own caches.

And the thing not to do, which is worth saying explicitly: do not publish or rely on wall-clock timings from a laptop or a demo dataset. They are dominated by JVM warm-up and OS page cache, and they do not transfer to a cluster. A timing harness is fine; presenting its local output as a result is not.

What to report instead: the mechanism, and the structural numbers that are deterministic. “Spill dropped from 40 GB to zero and the stage went from 200 tasks to 2,000” is reproducible and explains itself. “It got 3x faster on my machine” does not.

Structured Streaming

Enough to show you understand the model, not just the API. Watermarks and state are where these questions go.

161. What is Structured Streaming?

A stream processing model where a stream is treated as an unbounded table that rows are appended to, and a query over it is re-executed incrementally as data arrives.

flowchart TB
  A["input stream"] --> B["unbounded input table<br/>rows appended"]
  B --> C["your query,<br/>the same as batch"]
  C --> D["result table,<br/>updated each trigger"]
  D --> E["sink, per output mode"]

The consequence that makes it valuable: you write the same DataFrame operations as for batch. The engine works out how to maintain the result incrementally, keeping whatever state that requires.

events = (spark.readStream.format("kafka")
    .option("kafka.bootstrap.servers", "kafka-broker1:9092")
    .option("subscribe", "trips")
    .load())

(events.selectExpr("CAST(value AS STRING)")
    .writeStream.format("parquet")
    .option("checkpointLocation", "s3a://lakehouse-prod/checkpoints/trips")
    .option("path", "s3a://lakehouse-prod/warehouse/trips_raw")
    .outputMode("append")
    .start())

By default it executes as micro-batches: the engine runs a small batch job per trigger, which is why the same Catalyst machinery and the same UI apply.

The model is what you should lead with in an interview, because it explains everything else: output modes exist because the result table can change in different ways, and watermarks exist because maintaining the result incrementally requires bounded state.

162. What is the difference between Structured Streaming and the old DStream API?

DStreams are the original RDD-based API; Structured Streaming is the DataFrame-based replacement.

  DStreams Structured Streaming
Built on RDDs DataFrames and Catalyst
Optimisation None Full Catalyst and Tungsten
Time semantics Processing time only Event time, with watermarks
Late data Manual handling Watermarks and allowed lateness
End-to-end guarantee At-least-once in practice Exactly-once with a replayable source and idempotent sink
Batch and stream code Different APIs The same API
Status Legacy Current

The two differences that matter most in practice are event time and the unified API. DStreams only knew when the engine saw a record, so correct windowing over event time was the application’s problem. And because DStream code was a separate API, a batch backfill of the same logic had to be written twice.

New work should use Structured Streaming. The right way to say this in an interview is that DStreams are not merely older but architecturally unable to offer event-time correctness or plan optimisation, because they operate on opaque RDDs.

The one thing DStreams retained for a while was some lower-level control over batch composition, which the newer API has largely absorbed through triggers and foreachBatch.

163. What are the output modes?

Three, and which are legal depends on your query.

Mode Writes each trigger Legal for
append Only new rows that will never change Queries with no aggregation, or aggregations with a watermark
update Only rows whose value changed Most aggregations
complete The entire result table Aggregations only
.outputMode("append")

The logic behind the restrictions is worth explaining, because it shows you understand the model:

  • append requires that a row, once emitted, is final. For a plain projection that is always true. For a windowed aggregation it is only true once the watermark has passed the window’s end, which is why append on an aggregation requires a watermark.
  • update emits changed rows, so it works for aggregations without needing finality, but the sink must be able to handle updates to keys it has seen.
  • complete rewrites everything every trigger, so it requires keeping the whole result table in state forever. It does not scale, and it is only appropriate for small result sets such as a dashboard aggregate.

The common mistake is choosing complete for convenience, then finding state grows without bound. The other is expecting append to emit aggregation results immediately, when in fact nothing appears until the watermark closes the window.

164. What is a trigger, and what are the options?

A trigger controls when the engine runs the next micro-batch.

Trigger Behaviour
Default, unspecified Start the next batch as soon as the previous finishes
processingTime="1 minute" Start a batch on a fixed interval
availableNow=True Process all available data in as many batches as needed, then stop
continuous="1 second" Experimental continuous processing, with limited operation support
.trigger(processingTime="1 minute")
.trigger(availableNow=True)

How to choose:

  • The default gives lowest latency and continuous resource use. It also produces many small files if the sink is file-based, because it writes per batch.
  • A fixed interval is the usual production choice, because it batches work into sensibly sized writes and makes resource use predictable. If a batch takes longer than the interval, the next starts immediately, so the interval is a floor, not a guarantee.
  • availableNow is how you run a streaming query as a scheduled batch job: it consumes everything outstanding using the streaming checkpoint for offsets, then exits. That gives you exactly-once bookkeeping with batch economics, and it replaced the older once=True.

The practical trade to name: shorter triggers mean lower latency and more small files; longer triggers mean fewer, larger files and higher latency. For a lakehouse sink, the file-count consequence usually dominates the latency preference.

165. What is event time versus processing time?

Event time is when the thing happened, carried in the record. Processing time is when the engine saw it.

event time      10:00:05   (the trip started)
ingest time     10:00:07   (Kafka received it)
processing time 10:04:30   (Spark's batch picked it up)

They diverge for ordinary reasons: a mobile client was offline, a broker rebalanced, a consumer lagged, or a batch was delayed.

Correctness depends on event time. A trip that happened at 09:59 belongs in the 09:55-10:00 window no matter when it arrives, and aggregating by processing time would attribute it to the wrong window and produce a different answer on every replay.

The mistake to watch for, and interviewers do ask about it: using the source’s own timestamp as if it were event time. Kafka’s timestamp column is broker ingest time, not when the event occurred:

# wrong: this is ingest time
events.withWatermark("timestamp", "10 minutes")

# right: event time from the payload
(events
  .select(F.from_json(F.col("value").cast("string"), schema).alias("t"))
  .select("t.*")
  .withWatermark("started_at", "10 minutes"))

The practical test is replayability: if you reprocess yesterday’s data today, an event-time query gives the same answer and a processing-time query does not.

166. What is a watermark and what does it actually do?

A watermark is a moving threshold on event time that declares how late data may be. It is computed as the maximum event time seen so far, minus the delay you allow.

(events
  .withWatermark("started_at", "10 minutes")
  .groupBy(F.window("started_at", "5 minutes"), "city_id")
  .agg(F.count("*").alias("trips")))

It does two things, and naming both is the strong answer:

  1. Finalises windows, so their results can be emitted in append mode. Until the watermark passes a window’s end, that window might still change, so nothing can be appended.
  2. Bounds state. Once a window can no longer change, the engine drops the state it was holding for it. Without that, state grows for the life of the query.
flowchart TB
  A["max event time seen<br/>10:30"] --> B["minus 10 min delay"]
  B --> C["watermark 10:20"]
  C --> D["windows ending before 10:20:<br/>emit and drop state"]
  C --> E["records older than 10:20:<br/>dropped as too late"]

The trade is explicit: a longer watermark tolerates later data and holds more state; a shorter one is cheaper and drops more. That is a business decision about how much lateness matters, not a tuning parameter.

The signature of getting it wrong is state size climbing without plateau in lastProgress, which is the number to monitor.

167. Does a watermark do anything on a non-aggregating stream?

Essentially nothing, and recognising that is a good sign.

withWatermark gives the engine two capabilities: finalise results so they can be appended, and drop state that can no longer change. A pass-through stream, a projection, a filter, a from_json, keeps no state and has no results to finalise. So there is nothing for the watermark to do.

# the watermark here is inert
(events
  .withWatermark("started_at", "10 minutes")
  .select("trip_id", "city_id")
  .writeStream.format("parquet").outputMode("append").start())

Seeing one there usually means it was copied from an aggregation example.

Where it does matter is any stateful operation:

Operation Watermark needed
Windowed aggregation Yes, for append mode and to bound state
dropDuplicates on a stream Yes, or the dedup set grows forever
Stream-stream join Yes, on both sides, plus a time condition
flatMapGroupsWithState Yes, to expire state
Projection, filter, from_json No

The one subtle exception is worth mentioning: a watermark also drops individual records older than the threshold, so on a non-aggregating stream it could silently discard late rows without giving you anything in return. So it is not merely useless there, it can be actively harmful.

168. What is checkpointing in Structured Streaming?

A required storage location where the engine records what it has processed and the state it holds, so a restarted query resumes exactly where it stopped.

.option("checkpointLocation", "s3a://lakehouse-prod/checkpoints/trips_by_window")

What it contains:

  • Offsets, a write-ahead log of which source offsets each batch covered, written before the batch runs.
  • Commits, marking which batches completed.
  • State, the keyed state for stateful operations, versioned per batch.
  • Metadata, including the query id.

The offsets-before-execution ordering is what makes recovery deterministic: on restart the engine knows which batch was in flight and re-runs it over exactly the same input range, which is the basis of the exactly-once story.

The rules:

  • Every streaming query needs one. Without it the query cannot recover, and file sinks will refuse to run.
  • One per query. Two queries sharing a location will corrupt each other.
  • It must be durable and shared. Local disk fails when the driver moves, so use HDFS or object storage.
  • It is not disposable. Deleting it means reprocessing from the source’s earliest offset, or losing data, depending on startingOffsets.

That last point is the operational trap: people delete a checkpoint to “reset” a query and either reprocess everything or skip data.

169. How does Structured Streaming achieve exactly-once?

By combining three things, and it is a property of the whole pipeline, not of the engine alone.

  1. A replayable source. One that can be re-read from a recorded offset: Kafka, or files. A source you cannot rewind, such as a plain socket, cannot give exactly-once.
  2. Checkpointed offsets. The engine writes which offsets a batch covers before running it, so a restart re-runs the identical input range.
  3. An idempotent sink. One where writing the same batch twice has the same effect as writing it once, usually by keying the write on the batch id.
flowchart TB
  A["replayable source<br/>Kafka offsets"] --> B["checkpoint offsets<br/>before the batch"]
  B --> C["run the batch"]
  C --> D["idempotent sink<br/>dedupes on batch id"]
  D --> E["exactly-once end to end"]

The honest framing: if any of the three is missing, you get at-least-once. A sink that appends without deduplication will duplicate on retry no matter what the engine does. Spark’s file sinks handle this with a commit log, which is why they are exactly-once, and foreachBatch puts the responsibility on you.

The interview-grade version of the answer is exactly that distinction: the engine guarantees it will re-present the same batch, and the sink must guarantee that re-presenting is harmless. Saying “Spark gives exactly-once” without the sink condition is the answer that gets probed.

170. What is `foreachBatch` and why is it useful?

A sink that hands you each micro-batch as an ordinary DataFrame, letting you do anything batch code can do.

def upsert(batch_df, batch_id):
    (batch_df.write.format("hudi")
       .options(**hudi_options)
       .mode("append")
       .save(table_path))

(events.writeStream
   .foreachBatch(upsert)
   .option("checkpointLocation", ckpt)
   .start())

Why it matters: streaming sinks support a limited set of operations, while a batch DataFrame supports everything. So foreachBatch is how you reach:

  • MERGE into a lakehouse table, which is the standard way to drive upserts from a stream into Hudi, Iceberg or Delta.
  • Multiple sinks from one stream, writing the same batch twice without reading the source twice.
  • JDBC writes, and any sink with no streaming implementation.
  • Batch-only operations on the micro-batch, such as a join to a static table with a strategy that streaming would not choose.

Two responsibilities it transfers to you:

  1. Idempotency. The same batch_id can be re-presented after a failure, so your write must tolerate that. The batch_id argument exists precisely so you can deduplicate on it.
  2. Caching if you use the batch twice. Writing batch_df to two sinks without caching recomputes it, which for a Kafka source means re-reading.

So foreachBatch is the escape hatch that makes streaming practical for lakehouse ingestion, at the cost of owning the exactly-once contract yourself.

171. What happens if you change a streaming query's code and restart it?

It depends on whether the change is compatible with the checkpointed state. The engine validates on restart and fails rather than producing wrong results.

Usually safe:

  • Adding or removing a projection column.
  • Changing a filter predicate.
  • Changing the sink’s output path.
  • Changing the trigger interval.
  • Adding a new stateless operation.

Usually not safe:

  • Changing the aggregation keys, since the state is keyed by them.
  • Changing the output mode.
  • Adding or removing a stateful operation, which changes the state layout.
  • Changing the watermark column.
  • Changing the source in a way that changes offset semantics, such as a different topic.
incompatible change -> restart fails, or requires a new checkpoint
new checkpoint      -> reprocess from startingOffsets, or skip data

The operational consequence, which is the real content of this question: for an incompatible change you need a migration plan, not just a deploy. The usual options are to start a new query with a new checkpoint and backfill the gap with a batch job, or to run both old and new in parallel and cut over.

This is why streaming state should be treated as a schema: changing it is a migration. Teams that treat a streaming job as stateless code get caught by this the first time they need to add a grouping key.

172. How do you handle late data?

Decide how late is tolerable, express that as the watermark, and then decide what happens to anything later.

.withWatermark("started_at", "30 minutes")

Data arriving within 30 minutes of the maximum event time seen is included in its correct window. Data arriving later is dropped from the stateful aggregation, because the state for its window has been released.

The trade is direct:

Watermark Late data State
Short, minutes More dropped Small, cheap
Long, hours or days More included Large, expensive

So the choice is a business question: how much does a late record matter, and is a slightly wrong aggregate acceptable in exchange for bounded cost?

If you cannot afford to drop it, the standard pattern is two paths:

  1. The streaming aggregation with a modest watermark, for timely approximate results.
  2. The raw stream written as-is, with no watermark and no aggregation, so nothing is lost.

Then a periodic batch job recomputes the affected windows from the raw data and corrects the serving table. That is the lambda-shaped compromise, and in a lakehouse it is usually a MERGE that overwrites the recomputed windows.

You can also monitor what you are losing: the difference between raw record counts and aggregated counts over the same period tells you whether the watermark is set sensibly.

173. What are stateful operations, and which are they?

Any operation whose result depends on data seen across triggers, so the engine must carry something between batches.

Operation State held
Windowed aggregation A partial aggregate per window per key
groupBy aggregation without windows A running aggregate per key, forever
dropDuplicates The set of keys already seen
Stream-stream join Buffered rows from both sides awaiting matches
flatMapGroupsWithState, transformWithState Whatever you choose to keep

Stateless by contrast: select, filter, withColumn, from_json, and a stream-to-static join.

All of the stateful ones need a watermark to bound state, except a non-windowed groupBy, which is unbounded by construction: a running count per key keeps one entry per key forever, and no watermark can help because the key never expires. That is worth knowing, because it is the one case where the answer is “restructure the query”, not “add a watermark”.

State lives in a state store, checkpointed per batch, and it is the main scaling dimension of a streaming job. lastProgress reports stateOperators with row counts and memory used, which is the number to watch:

print(query.lastProgress["stateOperators"])

Plateauing state means the watermark is doing its job. Monotonically growing state is the signature of a missing watermark, an over-generous one, or an unbounded aggregation.

174. What is special about a stream-stream join?

Both sides are unbounded, so the engine must buffer rows from each side while waiting for a match that may not have arrived yet. That makes it the most state-hungry operation in Structured Streaming.

Two requirements follow:

  1. Watermarks on both sides, so the engine knows when a row can no longer find a match and its state can be dropped.
  2. A time constraint in the join condition, bounding how far apart matching events can be.
(impressions.withWatermark("imp_time", "2 hours")
 .join(
   clicks.withWatermark("click_time", "3 hours"),
   F.expr(
     "imp_id = click_imp_id AND "
     "click_time >= imp_time AND "
     "click_time <= imp_time + interval 1 hour")))

Without the time condition the engine must keep every row forever, because any future row might match.

Outer joins add a further subtlety worth volunteering: a non-match cannot be emitted immediately, because a match might still arrive. The engine must wait until the watermark proves no match is possible, so outer-join results appear delayed by the watermark, not at the trigger where the row arrived.

The practical advice: prefer a stream-to-static join where you can, joining the stream to a dimension table. It is stateless, needs no watermark, and is usually what the requirement actually calls for.

175. How do you monitor a streaming query?

Through the query’s own progress objects, and through the Streaming tab of the UI.

print(query.status)          # what it is doing right now
print(query.lastProgress)    # metrics for the last completed batch

The numbers that matter, in priority order:

Metric What it tells you
inputRowsPerSecond vs processedRowsPerSecond Whether you are keeping up
stateOperators[].numRowsTotal Whether state is bounded
batchDuration against your trigger interval Whether batches fit in their window
durationMs breakdown Where the time goes: addBatch, getBatch, walCommit
Source endOffset minus committed Consumer lag

The two alerts worth having:

  1. Processing rate below input rate, sustained. That is falling behind, and lag grows without bound until you add resources or reduce work.
  2. State size growing without plateau. That is the watermark signature, and it ends in an executor OutOfMemoryError hours or days later.

For programmatic alerting, attach a listener rather than polling:

spark.streams.addListener(my_listener)   # onQueryProgress, onQueryTerminated

The thing to say about alerting on duration: batch duration alone is a poor signal, because it moves with input volume for innocent reasons. The ratio of processing rate to input rate is the structural signal, and state size is the one that predicts a failure before it happens.

Deployment, cluster managers and operations

The practical half that candidates often cannot answer, which makes it a good differentiator.

176. How do you run Spark on YARN?

--master yarn, with the deploy mode choosing where the driver lives.

export HADOOP_CONF_DIR=/etc/hadoop/conf
spark-submit \
  --master yarn \
  --deploy-mode cluster \
  --queue analytics \
  --num-executors 50 --executor-cores 4 --executor-memory 16g \
  --conf spark.executor.memoryOverhead=3g \
  --class com.example.TripsIngest trips.jar

What happens: Spark starts an ApplicationMaster in the cluster, which requests containers from the ResourceManager through YarnAllocator. Each granted container becomes an executor. In cluster mode the driver runs inside the ApplicationMaster container; in client mode the AM only handles resource requests while the driver stays on your machine.

Things that are YARN-specific:

  • HADOOP_CONF_DIR or YARN_CONF_DIR must be set, or Spark cannot find the ResourceManager.
  • --queue places the application in a scheduler queue, which is how YARN arbitrates between applications.
  • Container kills for exceeding memory are reported by the NodeManager, and the fix is memoryOverhead.
  • The external shuffle service runs as a NodeManager auxiliary service.
  • Kerberos delegation tokens are obtained by the driver and renewed for long-running jobs, which needs --principal and --keytab.

The tracking URL in the ResourceManager UI links to the Spark UI while running, and to the history server afterwards.

177. How do you run Spark on Kubernetes?

--master k8s://https://<api-server>:<port>, with a container image that has Spark in it.

spark-submit \
  --master k8s://https://k8s.internal:6443 \
  --deploy-mode cluster \
  --name trips-ingest \
  --conf spark.kubernetes.container.image=registry.internal/spark:4.1.3 \
  --conf spark.kubernetes.authenticate.driver.serviceAccountName=spark \
  --conf spark.executor.instances=50 \
  --class com.example.TripsIngest \
  local:///opt/app/trips.jar

The architecture is distinctive: the driver runs in a pod, and the driver’s own ExecutorPodsAllocator creates executor pods by talking to the API server. There is no Spark cluster-manager daemon at all.

What that implies:

  • A service account with permission to create pods is required, which is what the serviceAccountName config is for.
  • The image is the dependency mechanism. local:// paths refer to files inside the image, which is the repeatable way to ship code.
  • Pod templates let you set anything Kubernetes supports that Spark has no config for: node selectors, tolerations, volumes, sidecars.
  • Dynamic allocation works via shuffle tracking, since there is no external shuffle service on Kubernetes in the YARN sense.
flowchart TB
  S["spark-submit"] --> A["API server"]
  A --> D["driver pod"]
  D -->|"ExecutorPodsAllocator<br/>creates pods"| A
  A --> E1["executor pod"]
  A --> E2["executor pod"]
178. What is the difference between Standalone, YARN and Kubernetes for Spark?
  Standalone YARN Kubernetes
Provided by Spark itself Hadoop Kubernetes
Grants containers Master daemon ResourceManager, via ApplicationMaster API server, via the driver
Isolation Process-level Containers, cgroups Pods, full container isolation
Multi-tenancy Limited Queues, mature Namespaces, quotas
Dependency shipping Jars on a shared path Jars, --packages Container images
Best for Dedicated Spark clusters Existing Hadoop estates Cloud-native, mixed workloads

How to choose, honestly:

  • Standalone when the cluster runs only Spark and you want the least machinery. It is simple and it works, but it has the weakest multi-tenancy.
  • YARN when you already have Hadoop. It is the most mature for queue-based sharing and for Kerberos, and the tooling around it is well understood.
  • Kubernetes when you are already running Kubernetes, want per-job container images, and want Spark to share a cluster with non-Spark workloads.

And note Mesos is gone. It was removed in the 4.x line, so a candidate listing four options is working from older material.

The interview-relevant point beyond the table: the choice changes deployment and sharing, but not how Spark itself schedules. The driver, the stages and the tasks behave identically, because only SchedulerBackend differs. That is worth saying, because it reframes the question as an operations decision rather than a performance one.

179. How do you share a cluster between Spark applications?

At two levels, and keeping them distinct is the answer.

Between applications, the cluster manager arbitrates:

  • YARN: queues with capacity or fair scheduling, so each team or pipeline gets a guaranteed share.
  • Kubernetes: namespaces with resource quotas, and optionally a scheduler that understands gang scheduling.
  • Standalone: limited, mostly spark.cores.max per application.

Dynamic allocation matters here, because an application holding idle executors is holding capacity another application could use.

Within one application, spark.scheduler.mode governs how concurrent jobs share slots:

--conf spark.scheduler.mode=FAIR --conf spark.scheduler.allocation.file=/etc/spark/fairscheduler.xml
sc.setLocalProperty("spark.scheduler.pool", "interactive")

Default is FIFO, under which a long job submitted first starves a short one submitted second. That is the wrong behaviour for a notebook server or a Thrift server serving many users from one SparkSession, which is exactly where FAIR and pools belong.

The distinction to state plainly: one application cannot share executors with another, because executors are per-application. So between-application sharing is about who gets containers, and within-application sharing is about who gets task slots in the containers you already hold.

180. How do you pass configuration to a Spark application, and what is the precedence?

Four mechanisms, and the precedence, highest first:

  1. SparkConf in code, set before the session is created.
  2. spark-submit flags, including --conf and the named ones like --executor-memory.
  3. A properties file named with --properties-file.
  4. spark-defaults.conf in SPARK_HOME/conf.
spark = (SparkSession.builder
    .config("spark.sql.shuffle.partitions", "400")     # beats --conf
    .getOrCreate())
spark-submit --properties-file /etc/spark/job.conf              --load-spark-defaults              --conf spark.sql.shuffle.partitions=200 app.jar

--load-spark-defaults makes Spark read spark-defaults.conf as well as the file given by --properties-file, which otherwise replaces it. That is useful when you want cluster-wide defaults plus job-specific overrides.

The traps worth naming:

  • A --conf after the application jar is passed to your application, not to Spark. No warning.
  • Options that only apply at JVM start, notably spark.driver.extraJavaOptions in client mode, cannot be set from code.
  • Executor configs set after the session exists do not retroactively change running executors.

To see what actually applied, --verbose prints the resolved configuration and its source, and the Environment tab of the UI is authoritative at runtime.

181. How do you add third-party libraries?

Four mechanisms, each for a different situation:

Mechanism For
--packages Maven coordinates, resolved with transitive dependencies at submit time
--jars Local or remote jars, added to driver and executor classpaths
--py-files Python .zip, .egg or .py
--files Arbitrary files placed in each working directory
spark-submit \
  --packages org.apache.hudi:hudi-spark3.5-bundle_2.12:1.2.0 \
  --jars /opt/lib/custom-udfs.jar \
  --files /etc/spark/conf/log4j2.properties \
  app.py

For production, prefer not resolving at submit time. --packages reaches out to Maven Central on every run, which makes your job depend on network availability and on a repository staying up, and makes the resolved set non-deterministic if a version range is involved. Bake dependencies into an assembly jar, or into a container image on Kubernetes.

Two related notes:

  • Conflicts are the usual pain. If your version collides with one Spark ships, Spark’s usually wins. -verbose:class tells you which JAR won, and shading in your assembly is the robust fix.
  • --jars does not expand directories, and URLs must be comma-separated. The local: scheme means “already present on every node”, which avoids copying large jars.

For PySpark specifically, a virtual environment archive shipped with --archives is the cleanest way to get consistent Python dependencies on every executor.

182. What is the history server and why do you need it?

The Spark UI is served by the driver, so it dies when the application ends. The history server replays event logs to provide the same UI for completed applications.

Two parts. The application writes a log:

--conf spark.eventLog.enabled=true --conf spark.eventLog.dir=hdfs:///spark-events --conf spark.eventLog.compress=true

And a separately running server reads that directory:

export SPARK_HISTORY_OPTS="-Dspark.history.fs.logDirectory=hdfs:///spark-events"
$SPARK_HOME/sbin/start-history-server.sh

Why you need it is simple: almost all investigation happens after the fact. A job that failed at 03:00 leaves you with only the driver’s stdout unless event logging was on. Enabling it is close to mandatory in production, and it is the single cheapest operational improvement for a Spark platform.

Operational details worth knowing:

  • Logs grow, so enable compression and set retention (spark.history.fs.cleaner.enabled).
  • A very large log is slow to replay, which happens with jobs that had hundreds of thousands of tasks.
  • On YARN, the ResourceManager’s tracking URL redirects to the history server once an application finishes, so links keep working.

The thing that goes wrong most often is enabling spark.eventLog.dir on a path the history server cannot read, so logs are written and never shown.

183. How do you secure a Spark cluster?

Four layers, and a good answer covers more than one.

1. Authentication between Spark components.

--conf spark.authenticate=true --conf spark.authenticate.secret=<secret>

On YARN and Kubernetes the secret can be distributed by the cluster manager rather than being passed on a command line.

2. Encryption in transit. RPC encryption with spark.network.crypto.enabled, and separately encryption of shuffle and cached data spilled to disk with spark.io.encryption.enabled. TLS for the UI and the history server.

3. Authorisation and identity. Kerberos on Hadoop, with --principal and --keytab so long-running jobs can renew delegation tokens. On Kubernetes, service accounts and RBAC.

4. Isolation. YARN queues or Kubernetes namespaces, so one tenant cannot exhaust another’s capacity.

Then the mistakes to avoid, which is where practical knowledge shows:

  • Secrets in extraJavaOptions or --conf are visible in the Environment tab of the UI and in process listings. Use a secrets mechanism instead.
  • The UI has no authentication by default. It should not be reachable from an untrusted network; put it behind a proxy with authentication, using spark.ui.filters.
  • Event logs contain SQL text and configs, so the log directory needs the same access control as the data.
184. How do you route Spark traffic through an HTTP proxy?

JVM proxy system properties, passed through the same extraJavaOptions mechanism as any other JVM flag.

spark-submit \
  --conf "spark.driver.extraJavaOptions=-Dhttp.proxyHost=proxy.internal -Dhttp.proxyPort=3128 -Dhttps.proxyHost=proxy.internal -Dhttps.proxyPort=3128 -Dhttp.nonProxyHosts=localhost|127.0.0.1|*.internal" \
  --conf "spark.executor.extraJavaOptions=-Dhttp.proxyHost=proxy.internal -Dhttp.proxyPort=3128 -Dhttps.proxyHost=proxy.internal -Dhttps.proxyPort=3128 -Dhttp.nonProxyHosts=localhost|127.0.0.1|*.internal" \
  app.jar

Three practical points:

  • Set it on both sides. Driver-side covers dependency resolution and metadata calls; executor-side covers whatever your code fetches from tasks.
  • http.nonProxyHosts matters as much as the proxy itself. Without it, intra-cluster traffic tries to go through the proxy, which is slow at best and broken at worst. The separator is |, not a comma.
  • https.nonProxyHosts does not exist; the HTTP one covers both.

The reason this question appears in a Spark interview is that it is really a question about whether you understand the delivery mechanism. Proxy settings, logging configuration, stack size, GC flags and heap dump options all arrive the same way: through spark.driver.extraJavaOptions and spark.executor.extraJavaOptions. Learning one teaches the rest.

And the client-mode caveat applies here too: the driver JVM already exists, so use --driver-java-options if you are setting it for a client-mode driver.

185. What should you monitor in production?

Monitor structural signals, not total job duration, because duration moves for many innocent reasons.

For batch jobs:

Signal Why
Task duration skew, max versus median per stage Detects skew before it becomes a failure
Shuffle spill The early sign that partitions are too large
GC time per executor Predicts heartbeat timeouts and slowdowns
Executor loss count and reason Distinguishes container kills from crashes
Input file count Rising file counts predict planning slowdowns
Stage retry count A stage running twice means lost shuffle output

For streaming jobs:

Signal Why
Processing rate versus input rate The definition of keeping up
State store row count Growth without plateau means a watermark problem
Consumer lag The business-visible measure of freshness
Batch duration versus trigger interval Batches overrunning their window

How to collect it: event logs plus the history server for post-mortem, and Spark’s metrics system (spark.metrics.conf) exporting to Prometheus or StatsD for live dashboards. StreamingQueryListener for streaming.

The reason to frame it this way in an interview: anyone can alert on “job took too long”. Alerting on spill, skew ratio and state growth means you find out before the job fails, and you know what to do when the alert fires. That is the difference between monitoring and having a dashboard.

Scenario-based questions

These separate people who have run Spark in production from people who have read about it. Each is a symptom; the answer is a way of looking, not a fact to recall. Practise saying these out loud.

186. A job that took 20 minutes now takes 2 hours, and nothing was deployed. How do you investigate?

No code changed, so something changed in the data or the cluster. Work through it in this order, because it is cheapest first.

1. Compare the physical plan against a previous run. This finds the most common cause by far: a broadcast join that no longer broadcasts, because the dimension table grew past spark.sql.autoBroadcastJoinThreshold. One join strategy change can account for the entire regression.

df.explain("formatted")     # and the SQL tab of the history server for the old run

2. Check input size and file count separately. Ten times the files with the same bytes changes planning cost and partition packing, which looks like slowness with no plan change.

3. Check for new skew. On the slowest stage, max versus median task duration and shuffle read. A newly hot key, or a surge of null keys, concentrates work.

4. Check the executors tab for loss and retries. A stage that now runs twice because an executor keeps dying doubles the time with no plan difference.

5. Check cluster contention. A noisy neighbour on a shared cluster, or fewer executors granted than before.

6. Only then look at configs, which did not change, but might now be wrong for the new data volume.

The answer that scores: name the lost broadcast first and say why you would look there before anything else. It is the single most likely cause of a step-change regression with unchanged code.

187. A stage has 200 tasks; 198 finish in seconds and 2 run for an hour. What do you do?

That is skew, and the first thing to do is confirm it rather than assume it.

Confirm. On the stage page compare max and median for shuffle read size, not just duration:

Duration Shuffle read Cause
Max » median Max » median Skew. Fix the distribution
Max » median Roughly even Environment. Slow node, GC, bad disk

If shuffle read is even, stop treating it as skew and look at the executors tab for that host.

Then, if it is skew:

  1. Check whether AQE could fire. spark.sql.adaptive.skewJoin.enabled is true, but a partition must be over 5x the median and over 256 MB. Two tasks running an hour on a job whose partitions are 80 MB will never trigger it. Lowering skewedPartitionThresholdInBytes may be the whole fix.
  2. Broadcast the small side, if there is one. No shuffle, no skew.
  3. Salt the hot key, which is the general answer.
  4. Split the hot keys out and union, if there are only a few.

And say what you would not do: enable speculation. The duplicate task has the same rows and is equally slow, so it burns a second slot for nothing and disguises a data problem as flaky infrastructure.

Also check the join key’s null count. Nulls all hash together, and a column that is 40% null produces exactly this picture.

188. Your job writes 50,000 tiny Parquet files. Why, and how do you fix it?

It is arithmetic. Every task that holds rows for a partition value writes a file into that value’s directory:

files = tasks x partition values
      = 200 shuffle partitions x 250 cities
      = up to 50,000

Each task has a few rows for almost every city, so it opens a writer for almost every city.

The fix is to make each value live in one task, by repartitioning on the same columns you partition by:

(df.repartition("city_id")
   .write.partitionBy("city_id")
   .mode("overwrite")
   .parquet(path))

That gives 250 files instead of 50,000.

Then check three related things:

  1. Is the partition column the right cardinality? 250 is fine; 50,000 distinct values is not, and directory-per-value is the wrong layout. Consider a coarser column, or bucketing.
  2. Cap file size for the big partitions. spark.sql.files.maxRecordsPerFile, default 0, splits one enormous partition into several sensible files.
  3. Watch for skew from the repartition itself. repartition("city_id") hash-partitions, so a dominant city becomes one large task. That may be an acceptable trade for the file layout, or may need salting.

Why it matters beyond tidiness: every later read pays listing cost, per-file open cost (openCostInBytes, default 4MB each), and worse compression. The write is slow too, because 50,000 files on object storage is 50,000 PUTs.

189. The driver runs out of memory on a job that only aggregates. What happened?

“Only aggregates” means the result should be small, so something is pulling data back or the driver is tracking too much.

Check these in order:

  1. Is there a collect(), toPandas(), or a wide show()? Even a show(1000) on a very wide row is a lot of data. This is the most common answer.
  2. Is something being broadcast? The driver must assemble a broadcast side before sending it, so a broadcast() hint on a large table fails on the driver, not the executors. Check the plan for BroadcastHashJoin or BroadcastExchange and the size of what feeds it.
  3. How many tasks? The driver holds metrics and state per task. A job with hundreds of thousands of tasks, from tiny input files or an enormous shuffle partition count, needs real memory just for bookkeeping.
  4. Is the plan deep or very wide? Analysis holds the tree in driver memory; a loop of withColumn calls or a 5,000-column schema adds up.
  5. Is there a driver-side loop accumulating results? Appending each iteration’s output to a Python list is a classic.
spark.conf.get("spark.driver.maxResultSize")     # 1g

The distinction worth drawing: a maxResultSize error is Spark protecting you with a clear named failure, while an OutOfMemoryError means the data arrived some other way, usually a broadcast. So the error type narrows the search.

And note spark.driver.memory defaults to 1g, which is a coordinator’s budget.

190. Executors keep being killed for exceeding memory limits. What is your sequence?

1. Confirm which failure it is, because the two have different fixes:

Message Means Fix
Container killed by YARN/Kubernetes for exceeding limits The process total exceeded heap plus overhead memoryOverhead
java.lang.OutOfMemoryError: Java heap space The heap filled Partitions, then executor.memory

A container kill is not a JVM error, and reading the message rather than assuming is the first step.

2. Raise spark.executor.memoryOverhead. Overhead holds off-heap buffers, Netty’s direct buffers, metaspace, thread stacks, and the Python workers. PySpark jobs hit this far more often for exactly that reason, and can be killed with a perfectly healthy heap.

--conf spark.executor.memoryOverhead=4g

3. Check for skew, because one oversized partition is a very common root cause that presents as a memory shortage. Fixing the distribution fixes the memory.

4. Raise the partition count so each task’s footprint shrinks. This scales with your data, whereas a memory bump does not.

5. Reduce cores per executor, which gives each running task a larger share of the same pool.

6. Raise spark.executor.memory, last.

The judgement to voice: raising memory first usually just moves the failure to next month’s data volume. Steps 3 and 4 fix it structurally.

191. A `FetchFailedException` keeps failing your job. What is the real problem?

Almost never the fetch settings. The blocks were written; whatever was serving them can no longer do so.

Look at the executors tab, not the network config. The question is which executor was serving those blocks and why it stopped:

Finding Root cause Fix
Killed for exceeding memory Overhead, skew, or partition size Those, not fetch retries
OutOfMemoryError Partition too large More partitions
Alive but very high GC time Too busy to serve Memory, or fewer cores
Node lost Hardware or preemption Retries, and an external shuffle service
Disk full Spill filled local disk More disk, or less spill

What Spark does on its own: DAGScheduler treats it as lost map output and resubmits the upstream stage to regenerate it. So the job often recovers slowly, and the tell is a stage that runs twice.

The structural fix is the external shuffle service, or shuffle tracking, so shuffle output is served by a node-level daemon that outlives the executor:

--conf spark.shuffle.service.enabled=true

What not to do: raise spark.shuffle.io.maxRetries and spark.network.timeout and call it fixed. That makes the job take longer to fail, and hides the executor deaths that are the actual problem.

192. You have a 2 TB fact table and a 3 GB dimension table to join. How?

3 GB is 300 times the 10 MB broadcast threshold, so a broadcast is out as things stand. The work is to change that, or to make the shuffle good.

1. Try to make the dimension small. This is the highest-value step and the most often skipped. Filter and project it first:

dim = (dimension
    .filter(F.col("active"))
    .select("city_id", "region"))        # 3 GB may become 40 MB

A 3 GB dimension is usually 3 GB because of columns and rows the query never needs. If this gets it under the threshold, the problem disappears.

2. If it cannot shrink, accept a sort-merge join and make it well-shaped: filter the fact side first, size the shuffle with spark.sql.adaptive.advisoryPartitionSizeInBytes, and handle skew on the join key.

3. If this join is recurring, bucket both tables on the join key with the same bucket count. That removes the shuffle permanently, at the cost of a one-off write and some rigidity.

4. If the fact table is partitioned on the join key and the dimension is filtered, dynamic partition pruning may already be limiting the fact scan. Check for dynamicpruningexpression in the plan.

And the thing not to do: hint BROADCAST on 3 GB. The driver must collect it first, so you convert a slow join into a driver OutOfMemoryError, which is harder to diagnose.

193. Your streaming query's state grows until the job dies. Why?

A stateful operation whose state is never released. In practice one of four causes.

1. No watermark on a stateful operation. Without one, no window is ever final, so nothing is dropped.

2. The watermark column is not event time. The classic mistake is using the source’s ingest timestamp:

# wrong: Kafka's timestamp is broker ingest time
events.withWatermark("timestamp", "10 minutes")

# right: event time from the payload
events.withWatermark("started_at", "10 minutes")

If the watermark column advances with processing time rather than event time it still moves, so this can look like it is working while behaving wrongly for late data.

3. The watermark is attached in the wrong place. It must be applied to the stream before the aggregation that needs it. Applied after, it does nothing.

4. The aggregation is unbounded by construction. A non-windowed groupBy keeps one entry per key forever, and no watermark can expire a key. That one needs a restructure, usually into a windowed aggregation, not a config.

Confirm which it is from the progress metrics:

print(query.lastProgress["stateOperators"])   # numRowsTotal over time

State that plateaus is healthy. State that climbs monotonically is the signature. Then check the event-time gap in lastProgress, which tells you whether the watermark is advancing at all.

194. A nightly job silently produced half the expected rows. How do you find out why?

Silent wrong results point at semantics, not performance, and Spark has a handful of well-known ways to lose rows quietly.

The candidates, most common first:

  1. An inner join dropping non-matches, including rows whose key is null, because null never equals null. Half the rows is a very plausible join-coverage figure.
  2. union matching by position where the two sides’ column orders differ. This corrupts rather than drops, but produces nonsense that a later filter then drops.
  3. A filter against a sometimes-null column. col != 'x' is null, not true, for null rows, so they are excluded.
  4. A changed upstream schema, so a select picks up a renamed column as null and a downstream filter removes it.
  5. ANSI mode on a Spark 4 upgrade, now raising or changing results where invalid casts previously yielded null.
  6. Dynamic partition overwrite replacing more or fewer partitions than intended.

The method, which matters more than the list: instrument row counts at each stage and find where they fall off.

for name, d in [("raw", raw), ("filtered", filtered), ("joined", joined)]:
    print(name, d.count())

Then for the join specifically:

raw.join(dim, "city_id", "left_anti").count()      # rows that found no match
raw.filter(F.col("city_id").isNull()).count()      # null keys

The habit to convey: for a silent correctness problem, reach for counts and anti-joins, not the UI. The UI tells you about speed, not about semantics.

195. You must deduplicate a CDC feed to the latest row per key. How?

A window function is the idiomatic batch answer.

from pyspark.sql import Window
from pyspark.sql.functions import row_number, desc, col

w = Window.partitionBy("trip_id").orderBy(desc("updated_at"))

latest = (changes
    .withColumn("rn", row_number().over(w))
    .where(col("rn") == 1)
    .drop("rn"))

That shuffles once, by trip_id.

The things to check, which is what the question is really about:

  1. Is the key skewed? The window partitions by trip_id, so one enormous key is one task. Check the distribution first.
  2. Is the ordering column unique enough? Ties in updated_at make the winner arbitrary. Add a tiebreaker such as an offset or a sequence number, or the result is non-deterministic between runs.
  3. Are deletes represented? A CDC feed usually carries an operation column, and the latest row for a key may be a delete, which should remove the row rather than keep it.
latest.filter(col("op") != "D")

Alternatives worth naming:

  • dropDuplicates(["trip_id"]) is cheaper but keeps an arbitrary row, not the latest. Only correct if the feed is already ordered, which it usually is not.
  • For a stream, dropDuplicates with a watermark bounds the state, which is the streaming equivalent.
  • In a lakehouse table, this is what an upsert does natively: Hudi with a record key and an ordering field applies exactly this rule on write, so the deduplication becomes the table’s job rather than the query’s.
196. A query works on a sample but fails at full scale with spill and no progress. What now?

The shape that works at one scale rarely fails at another for a config reason. It fails because partition sizes scaled with the data and the per-task working set outgrew its memory.

The arithmetic to check first:

per-task execution memory ~= (heap - 300MB) x spark.memory.fraction / cores
                          ~= (16g - 300MB) x 0.6 / 4
                          ~= 2.4 GB

partition size = total shuffle bytes / partition count
               = 4 TB / 200
               = 20 GB per task        <- far too large

200 partitions was fine for the sample and is absurd at full scale.

Then, in order:

  1. Raise the partition count, or lower the advisory size, so partitions land in the low hundreds of megabytes.
spark.conf.set("spark.sql.shuffle.partitions", "4000")
spark.conf.set("spark.sql.adaptive.advisoryPartitionSizeInBytes", "128MB")
  1. Look for skew that the sample hid. A sample of 1% of rows may contain none of the hot keys, so a job that looked balanced is not.
  2. Check for an aggregation whose cardinality scaled. A groupBy with 1,000 groups in the sample and 50 million at full scale is a different operation.
  3. Check whether something cached no longer fits, and is now evicting execution memory on every stage.

The framing that scores: say explicitly that a sample validates logic, not shape, and that partition sizing is the first thing to recompute at a new scale.

197. You need to read 10 TB but only one day of it. The job reads everything. Why?

The filter is not pruning, and there are only a few possible reasons.

Check the plan first. Two fields answer it:

PartitionFilters: [isnotnull(dt#42), (dt#42 = 2026-09-17)]     <- pruning works
PushedFilters:    [IsNotNull(started_at), GreaterThan(...)]     <- source filtering works

If your date predicate appears in neither, nothing is being skipped.

The causes:

  1. The table is not partitioned on the column you filter. Filtering started_at on a table partitioned by dt prunes nothing, because Spark cannot know the relationship. Filter the partition column.
  2. The predicate is wrapped in a function. This is the most common, and it is subtle:
# cannot prune: the column is inside a function
df.filter(F.to_date("started_at") == "2026-09-17")

# can prune
df.filter((F.col("started_at") >= "2026-09-17") &
          (F.col("started_at") <  "2026-09-18"))
  1. The type does not match. Comparing a string partition column to a date literal, or vice versa, can prevent pruning depending on the cast direction.
  2. The filter is applied after a join or aggregation, so it is no longer adjacent to the scan and cannot be pushed through.
  3. A UDF in the predicate, which is opaque to Catalyst.

And the structural note: for a table where this keeps happening, a lakehouse format prunes on column statistics rather than only directory names, which makes the layout less brittle.

198. Two jobs writing the same table produce corrupt or missing data. What is wrong?

A directory of Parquet files has no commit protocol, so concurrent writers have no way to coordinate and a failed writer has no way to roll back.

What goes wrong concretely:

  • Interleaved output. Two jobs writing the same partition both add files, and readers see the union, which is neither job’s intended result.
  • Partial output on failure. A job that dies halfway leaves the files it already wrote, with nothing marking them as incomplete.
  • overwrite racing append. The overwrite deletes files the appender is adding, or vice versa.
  • Readers seeing a moving target, because there is no snapshot isolation: a scan that lists files, then reads them, can have the list change underneath it.

The two fixes, and they are different in kind:

  1. Serialise the writers. An orchestration-level lock, or a single job. Correct, and limiting.
  2. Use a table format with atomic commits. Hudi, Iceberg or Delta add a commit protocol over the same Parquet files: a reader sees a consistent snapshot, a failed write commits nothing, and concurrent writers are arbitrated by the format rather than by luck.
(df.write.format("iceberg").mode("append").save("prod.db.trips"))

This is the structural answer and worth stating as such: it is not a Spark tuning problem, it is a missing transaction boundary. Lakehouse formats exist precisely for this, and saying so shows you understand why they were built rather than just that they exist.

199. A `groupBy` on a high-cardinality column is unbearably slow. Options?

First ask whether the aggregation is needed at that cardinality, because the cheapest fix is usually to not do it.

Then, in order:

  1. Is the output actually consumed at that granularity? An aggregation to 50 million groups that then feeds a join on a coarser key can often be reordered: join first, aggregate second.

  2. Is an approximation acceptable? Exact countDistinct requires a second shuffle; the approximate version does not:

F.approx_count_distinct("rider_id", rsd=0.02)     # instead of countDistinct
  1. Is the shuffle sized for the cardinality? 50 million groups across 200 partitions is 250,000 groups per task, and the hash map for that may not fit. Raise the partition count so each task’s map is smaller.

  2. Is there skew within the groups? High cardinality and skew often coexist: millions of small groups plus a handful of enormous ones. Two-phase aggregation with a salt handles the large ones:

(df.withColumn("salt", (F.rand() * 32).cast("int"))
   .groupBy("rider_id", "salt").agg(F.sum("fare").alias("p"))
   .groupBy("rider_id").agg(F.sum("p").alias("total")))
  1. Is the output written partitioned by that column? Then high cardinality is also a small-files problem, and the layout needs rethinking.

The framing: high cardinality is not itself a problem; the problem is a per-task hash map that does not fit, and the lever is partition count before memory.

200. Your team wants to move a Spark job's source table from Hudi to Iceberg. What do you consider?

First separate the layers, because the question conflates them. Spark is the compute engine; Hudi and Iceberg are table formats. The job’s Spark logic barely changes. What changes is write semantics, read reach and operations.

The questions that actually decide it:

  1. What is the write pattern? High-frequency keyed upserts are what Hudi’s record index is built for: a key lookup finds its file group directly. Iceberg has no record index, so a MERGE plans a join whose cost scales with how much of the table it touches. If the job is CDC ingestion, this is the dominant consideration.

  2. Which engines must read it? Iceberg has the widest warehouse support, and if a warehouse that only reads Iceberg is the reason for the move, that is a legitimate and sufficient reason.

  3. Who operates the table services? Both need compaction and cleanup, but they are exposed differently: Hudi ships them as configurable services that can run inline or async, Iceberg as procedures you schedule (rewrite_data_files, expire_snapshots, remove_orphan_files).

  4. What is the migration mechanism? Iceberg has snapshot and migrate procedures for in-place adoption of existing Parquet. A full rewrite of a large table is a real cost.

  5. What breaks? Incremental reads, time travel syntax and concurrency configuration all differ, so downstream consumers need changes even though the Spark job does not.

And the option people forget: if the answer is “both, for different readers”, Apache XTable converts metadata in place without copying data, so one set of Parquet files can be read as either. That reframes the decision from a migration to a conversion.

Conclusion

Two hundred questions is a lot of surface, and it collapses into four ideas.

Work is lazy until an action asks for a result, which is what allows the whole plan to be optimised rather than each call executed as it arrives. Partitions are the unit of parallelism, so nearly every performance question is really about how many you have and how evenly the data sits across them. Shuffles are the expensive thing, and stage boundaries fall exactly where they occur, which is why counting Exchange nodes tells you so much. And the driver holds all the state that cannot be recomputed, which is why it is the one process with no redundancy.

Derive an answer from those four and you will handle questions this list does not contain, which is the actual goal. The scenario section is where that shows: nobody can memorise the answer to “it was 20 minutes and now it is 2 hours”, because the answer is a way of looking rather than a fact.

References

Trademarks

Apache Spark, Apache Hadoop, Apache Hive, Apache Parquet, Apache ORC, Apache Avro, Apache Kafka, Apache Arrow, Apache Hudi, Apache Iceberg, Apache XTable (incubating), Apache Mesos, Apache Log4j and Apache are either registered trademarks or trademarks of The Apache Software Foundation in the United States and other countries. Delta Lake is a trademark of the Linux Foundation. Kubernetes is a registered trademark of the Linux Foundation.

Found this useful?

These posts and tools are free. If one saved you an afternoon, you can buy me a coffee.

Buy me a coffee