Interview/Spark

Spark senior interview questions

Skew, shuffle, join strategy, Catalyst, and Spark UI evidence.

Lesson

A dashboard job has 1 job, 4 stages, 80k tasks. What do you inspect first?

Answer it out loud, then reveal. Play steps through like the simulators.

All questions on this page

Indexed as FAQ. Open any item if you prefer a list to Play.

architect

A dashboard job has 1 job, 4 stages, 80k tasks. What do you inspect first?

What happens when Spark executes a job? · tap to open the answer

Short: Task count and the slowest stage — not executor count.

Detailed: 80k tasks usually means tiny files or over-partitioning. Stage time is max(task). Find the stage with shuffle or scan skew, then file layout, then AQE/coalesce.

Senior: Hundreds to low thousands of healthy tasks, not tens of thousands of 50 ms tasks.

Common mistake: Doubling the cluster because 'there are many tasks'.

Follow-up: What number of tasks would you aim for on a 1 TB nightly?

Lesson · Simulation

senior

A job is 'Spark' but the SLA miss is storage. How do you tell?

What is Apache Spark? · tap to open the answer

Short: If input is tiny files or a full scan, the cluster is waiting on the lake, not the CPU.

Detailed: Check Spark UI: many tiny tasks, huge task overhead, or a scan that ignores partition filters. Scaling executors will not fix a thousand 2 MB files.

Senior: Fix the files or the plan. Then scale. Never the reverse as a first move.

Common mistake: Adding workers before looking at file sizes and partition pruning.

Follow-up: What would you change first — cluster size or the table layout?

Lesson · Simulation

senior

How do you prove a failure is driver-side vs executor-side in five minutes?

Spark cluster architecture · tap to open the answer

Short: Driver OOM: notebook/kernel dies, executors still green. Executor OOM: one worker lost, driver alive, task failed with ExecutorLostFailure.

Detailed: Spark UI Executors tab: look at driver vs executor memory. Driver heap used near max plus a collect/broadcast in the plan → driver. One executor with huge spill or a killed container → executor.

Senior: collect, toPandas, and broadcasting a large dimension. Write the hypothesis, then the UI evidence.

Common mistake: Restarting the cluster without reading the error class.

Follow-up: Which action is the classic driver killer?

Lesson · Simulation

senior

You need a custom partitioner for a skewed key. Do you drop to RDD?

RDDs: lineage and partitions · tap to open the answer

Short: Sometimes — RDD partitionBy with a custom Partitioner, or stay in DataFrame with salting.

Detailed: Prefer salting + AQE in DataFrame land. Custom partitioners are powerful and easy to get wrong (breaking join co-partitioning). Measure first.

Common mistake: Custom partitioner as the first skew fix.

Follow-up: What join property do you lose if only one side is custom-partitioned?

Lesson · Simulation

senior

createDataFrame(huge_python_list) OOMs before any executor works. Why?

DataFrames and schemas · tap to open the answer

Short: The list lives on the driver; Spark then has to ship it.

Detailed: Parallelize/createDataFrame from driver memory. Read from storage or a Spark table instead. Same class of bug as collect() in reverse.

Common mistake: Raising executor memory for a driver-side Python list.

Follow-up: How should a 50 GB CSV enter Spark?

Lesson · Simulation

senior

A managed table and an external table both show up in SHOW TABLES. What changes the moment someone runs DROP TABLE?

Spark SQL: the same engine, a different keyboard · tap to open the answer

Short: Dropping a managed table deletes the data files; dropping an external table only removes the metadata.

Detailed: A managed table's location is owned by the catalog, so DROP removes the metastore entry and the underlying directory. An external table keeps its files at the LOCATION you declared — DROP leaves the data and you can re-register it. DESCRIBE EXTENDED shows Type: MANAGED vs EXTERNAL along with the Location.

Senior: DESCRIBE EXTENDED before any DROP in a runbook. Managed is fine for data the pipeline owns; external is for data other systems also write.

Common mistake: Saying DROP TABLE never deletes data because 'Spark doesn't own storage'.

Follow-up: How do you verify which kind you are about to drop?

Lesson · Simulation

architect

A scheduled SQL pipeline fails with an AnalysisException and Spark UI shows zero jobs for that run. Where do you look, and what is it definitely not?

Spark SQL: the same engine, a different keyboard · tap to open the answer

Short: It failed in Catalyst on the driver — name and type resolution — so it is not an executor or data-volume problem.

Detailed: Analysis resolves tables, columns, and functions against the catalog before a single stage is submitted, so a resolution failure means no job, no stage, and nothing in executor logs. spark.sql() builds and analyzes its plan when called, so the guilty statement is the first one referencing the renamed column or dropped table — not the write at the bottom of the notebook.

Senior: File permissions, corrupt Parquet footers, and Delta schema-on-write mismatches — those need a task to run. Resolution failures never get that far.

Common mistake: Digging through executor stderr for a query that never submitted a task.

Follow-up: Which failures do get deferred until execution time?

Lesson · Simulation

senior

You chained 30 withColumn calls. The job is slow but the plan looks simple. What happened?

Narrow vs wide transformations · tap to open the answer

Short: Each withColumn can add projection overhead; UDFs in those columns serialize every row.

Detailed: Prefer select with expressions. If columns are Python UDFs, you lost Tungsten/Photon. Spark UI: time in Python vs CPU.

Common mistake: More executors for a UDF pipeline.

Follow-up: How do you rewrite 30 withColumns cleanly?

Lesson · Simulation

senior

display() in Databricks feels harmless. When is it collect in disguise?

Actions that trigger jobs · tap to open the answer

Short: When the notebook pulls a large unaggregated result to render.

Detailed: UI limits help, but a wide table with huge strings still hammers the driver. Prefer aggregations or LIMIT. Job clusters should not display fact tables.

Common mistake: Leaving display(df) in a scheduled notebook.

Follow-up: How do you QA a 2 TB write without pulling it?

Lesson · Simulation

senior

When is laziness a production footgun?

Why Spark is lazy · tap to open the answer

Short: When invalid tables, bad schemas, or huge plans only fail at the action — in a job cluster at 2am.

Detailed: Validate schema and table existence early. Unit-test explain() on representative data. Don't leave 15 debug actions that each rerun a shuffle.

Common mistake: Adding count() everywhere as 'validation'.

Follow-up: What is a cheap eager check that isn't a full shuffle?

Lesson · Simulation

senior

How do you size partitions for a shuffle?

Partitions: the unit of parallelism · tap to open the answer

Short: Aim for ~100–256 MB per shuffle partition, then look at the histogram.

Detailed: Too few: fat tasks, spill, OOM. Too many: scheduler overhead. AQE coalesce helps small data. Skew means one partition will still be fat — salting, not more partitions, fixes that.

Senior: You want enough partitions for cores × duration, sized by data bytes, not by machine count.

Common mistake: Setting shuffle partitions to executor_count.

Follow-up: Why is executor_count the wrong formula?

Lesson · Simulation

senior

You have a job that joins on customer_id and then windows on customer_id. When is repartition('customer_id') worth paying for?

Repartition vs coalesce vs partitionBy · tap to open the answer

Short: When one shuffle on that key gets reused by more than one downstream operator.

Detailed: repartition('customer_id') hash-partitions rows so matching keys are co-located, and a following join or window on the same key can then skip its own Exchange. If only one operator needs the key, it would have shuffled once anyway — you have just paid for two shuffles instead of one.

Senior: explain('formatted') and count Exchange nodes before and after the change. If the count did not drop, the repartition bought you nothing.

Common mistake: Adding repartition(col) before every join out of habit and doubling the shuffle count.

Follow-up: How do you verify the second shuffle actually disappeared?

Lesson · Simulation

architect

spark.sql.shuffle.partitions is 200 and the code also calls repartition(1000). Which one wins, and where?

Repartition vs coalesce vs partitionBy · tap to open the answer

Short: They have different scopes: the config sizes the engine's own Exchanges, repartition(1000) sizes that one explicit shuffle.

Detailed: Every groupBy or sort-merge join Exchange targets spark.sql.shuffle.partitions unless AQE coalesces it at runtime; repartition(1000) is its own Exchange with its own target, and a later groupBy re-shuffles back to the configured number. df.write.partitionBy is a third, unrelated thing — it controls the on-disk directory layout, not in-memory partition count.

Senior: repartition on the same columns you pass to write.partitionBy, so each task owns one output directory and writes one decent-sized file instead of every task dropping a fragment into every directory.

Common mistake: Conflating write.partitionBy with repartition and expecting one to control the other.

Follow-up: How do you get a sane directory layout and reasonable file sizes at the same time?

Lesson · Simulation

senior

How do you remove a shuffle from a join?

What a shuffle actually does · tap to open the answer

Short: Broadcast the small side, or co-partition/bucket both sides on the join key.

Detailed: BroadcastHashJoin keeps the fact table in place. SortMergeJoin shuffles both unless already partitioned the same way. Prove with Exchange absence in explain().

Senior: When both sides are large and unsorted. Then make it balanced — not broadcast a 30 GB dimension.

Common mistake: spark.sql.shuffle.partitions as the join strategy.

Follow-up: When is a shuffle the right answer?

Lesson · Simulation

senior

How do you debug a join that 'exploded' row counts?

Join strategies in Spark · tap to open the answer

Short: It's a many-to-many — duplicates on the key, not Spark magic.

Detailed: Count distinct keys vs rows on both sides before the join. Look for duplicate dimension keys. Spark UI won't show 'duplicates'; data profiling will.

Senior: SELECT key, COUNT(*) FROM dim GROUP BY 1 HAVING COUNT(*) > 1 — then the same on the fact.

Common mistake: Raising shuffle partitions to fix a row explosion.

Follow-up: What's the SQL test you'd run in the interview whiteboard?

Lesson · Simulation

senior

Cached DataFrame, then an executor dies. What happens?

Cache and persist · tap to open the answer

Short: Lost partitions recompute from lineage — or fail if checkpoint wasn't used and shuffle files are gone.

Detailed: Cache is not a durable store. Storage tab shows cached partitions per executor. Executor loss → recompute. That's why checkpoint or a Delta write is the durable form.

Senior: Storage tab + stage that skips the scan. If the scan is back, cache missed or was evicted.

Common mistake: Treating cache as a table.

Follow-up: How do you see cache hit vs recompute in the UI?

Lesson · Simulation

senior

Your predicate is on a non-partition column and the job still only touches 3% of the files. What made that work, and when does it stop working?

File formats: why Parquet reads less · tap to open the answer

Short: Parquet row-group min/max statistics let Spark skip groups whose value range cannot match the predicate.

Detailed: Every row group stores per-column min/max; with pushdown the reader drops groups outside the predicate range, which shows up as a small Input Size / Records relative to the table. It only works when data is physically clustered on that column — if values are scattered, every row group's min/max spans the whole domain and nothing is pruned.

Senior: Sort or cluster on the filter column at write time (Z-ORDER or liquid clustering on Delta). Statistics are only as useful as the physical layout — this is a layout decision, not a config flag.

Common mistake: Calling row-group skipping 'an index' and expecting it to work on randomly ordered data.

Follow-up: How do you make the data cooperate?

Lesson · Simulation

architect

The team wants gzip because it compresses the smallest. What do you tell them?

File formats: why Parquet reads less · tap to open the answer

Short: A gzip text file is not splittable, so one file becomes one task no matter how large it is.

Detailed: Gzip must be decompressed from the beginning, so Spark cannot split a gzip CSV or JSON file — a 20 GB file is a single task and the rest of the cluster idles. Snappy and zstd inside Parquet are applied per page and row group, so the file stays splittable; zstd trades some CPU for noticeably smaller files than snappy.

Senior: It is the opposite failure: thousands of tiny Parquet files make per-file open and metadata cost dominate, with the task duration histogram piled near zero. Compact on write or with OPTIMIZE and aim for files in the hundreds of MB.

Common mistake: Comparing codecs only on compressed size and ignoring splittability and decode cost.

Follow-up: Where does the small-file problem fit into this conversation?

Lesson · Simulation

senior

Your window job has one task running for 40 minutes while 199 finished in seconds. What did you write?

Window functions: keep the row, add the answer · tap to open the answer

Short: A Window.orderBy(...) with no partitionBy — every row was shuffled into a single partition.

Detailed: Without partitionBy, Spark must order the whole dataset globally, so it plans an Exchange SinglePartition beneath WindowExec and one task processes everything. The stage shows exactly 1 task carrying the entire Shuffle Read plus a large Spill (Disk) while the rest of the cluster idles.

Senior: Reduce first, then order: filter or pre-aggregate to candidates, or rank within real partitions and stitch them with a small per-partition offset table. A truly global ordering over billions of rows is a design problem, not a tuning knob.

Common mistake: Diagnosing it as skew and salting a key that the window does not even partition by.

Follow-up: How do you compute a global row number without a single-task stage?

Lesson · Simulation

architect

Dedup with row_number() = 1 versus dropDuplicates — which goes in the pipeline and why?

Window functions: keep the row, add the answer · tap to open the answer

Short: row_number() when you must keep a specific survivor; dropDuplicates when any row per key is acceptable.

Detailed: dropDuplicates(keys) dedups via a shuffle and aggregate and hands you an arbitrary row per key. partitionBy(keys).orderBy(updated_at.desc) with a filter on row_number() = 1 pays for a Sort inside WindowExec but deterministically keeps the newest record — which is exactly what CDC and SCD merges require.

Senior: Add a deterministic tiebreaker such as a sequence number or commit ordering. Without one, reruns pick different survivors and downstream diffs look like data loss.

Common mistake: Running dropDuplicates on a CDC feed and shipping a random version of each record.

Follow-up: What breaks if the ordering column has ties?

Lesson · Simulation

senior

Walk me up the fix ladder for a row-at-a-time Python UDF, cheapest first.

UDFs: the box Catalyst cannot see into · tap to open the answer

Short: Built-in expression, then higher-order functions, then pandas_udf, and only then a plain Python UDF.

Detailed: Built-ins keep codegen and pushdown. Higher-order functions (transform, filter, exists, aggregate) handle array and map logic without leaving SQL. pandas_udf ships columnar Arrow batches, so you pay serialization per batch instead of per row and get vectorized pandas or NumPy code. A row-wise Python UDF is the last resort.

Senior: It is still opaque to Catalyst, so nothing pushes through it, and on Databricks it still knocks the query off Photon. Faster, not free.

Common mistake: Jumping straight to pandas_udf when composing two built-ins would have done it.

Follow-up: What does a pandas UDF still not fix?

Lesson · Simulation

architect

A Scala UDF runs inside the JVM, so is it safe to use anywhere?

UDFs: the box Catalyst cannot see into · tap to open the answer

Short: It skips the Python round trip, but it is still a black box to the optimizer.

Detailed: A Scala or Java UDF executes in the executor JVM with no cross-process serialization, so throughput is far better than Python. Catalyst still cannot rewrite or push it, it can break whole-stage codegen for the operators around it, and null handling is entirely your responsibility. A native Expression or a composition of built-ins is still preferable when the logic can be expressed.

Senior: Rerun the query with that column replaced by a literal. If runtime collapses, it is the UDF; if not, you were about to optimize the wrong operator.

Common mistake: Treating 'we made it a Scala UDF' as equivalent to a native expression in a design review.

Follow-up: How do you prove the UDF is or is not your bottleneck?

Lesson · Simulation

senior

You have 200 executors and a 2 GB driver. Which jobs will still fail?

Driver vs executors · tap to open the answer

Short: Any action that pulls the full result to the driver — collect, toPandas, display of an unaggregated frame.

Detailed: Executor count does not grow driver heap. A 2 TB collect still targets one JVM. Broadcast join of a 8 GB dimension also hits the driver, then every executor.

Senior: Aggregate or sample on executors, write the rest, never pull the grain to the notebook.

Common mistake: Scaling workers to fix a driver OOM.

Follow-up: How would you rewrite the notebook so the driver stays small?

Lesson · Simulation

senior

How do you explain a 1:N:M relationship on a whiteboard?

Jobs, stages, and tasks · tap to open the answer

Short: 1 action → N stages (shuffles + 1) → M tasks per stage (partitions).

Detailed: Draw: show() → Job 1 → Stage 0 scan/filter → Exchange → Stage 1 aggregate → few rows to driver. Tasks = partitions in each stage, possibly different counts after shuffle.

Common mistake: Drawing one box called 'the cluster'.

Follow-up: What changes M without changing N?

Lesson · Simulation

senior

How do you use the DAG to pick a broadcast vs a shuffle without guessing?

The Spark DAG · tap to open the answer

Short: Find the join node and whether an Exchange sits above each child.

Detailed: No Exchange on the fact side → broadcast. Exchange on both → SMJ. That's the DAG, not the SQL text.

Common mistake: Reading only the SQL string.

Follow-up: Where is that visible besides explain()?

Lesson · Simulation

senior

Catalyst picked SMJ. Stats were wrong. What's your move?

Catalyst optimizer · tap to open the answer

Short: ANALYZE TABLE / fix stats, add a broadcast hint if you know the size, or enable AQE join conversion.

Detailed: Don't hint forever as a substitute for stats. Hints rot. Measure with actual size after filters — AQE sees runtime sizes.

Senior: When you've measured the build side after filters and documented why stats cannot see it (UDF, JDBC, etc.).

Common mistake: Broadcast hint on a table that will grow 20×.

Follow-up: When is a hint the right senior answer?

Lesson · Simulation

senior

WholeStageCodegen is off for a stage. How do you find why?

Tungsten execution engine · tap to open the answer

Short: An operator that can't be compiled — often a UDF, an unsupported expression, or a fallback.

Detailed: explain() and SQL UI: look for the break in the * codegen star. Fix the expression. Don't 'tune Tungsten' as a config cult.

Common mistake: Random spark.tungsten flags from a blog.

Follow-up: How does Photon relate to Tungsten on Databricks?

Lesson · Simulation

senior

One task spilled 60 GB to disk and finished. A different task's container was killed. Same memory problem?

Executor memory: execution wins, cache loses · tap to open the answer

Short: No — spill is Spark managing execution memory; a killed container is the cluster manager seeing total process memory exceed its limit.

Detailed: Spark tracks execution memory and spills sorted runs to local disk when it runs short, so the task survives slowly and you see Spill (Memory) and Spill (Disk) on the stage. YARN or Kubernetes knows nothing about that accounting — it sees the whole process: JVM heap plus off-heap buffers plus Python workers. Cross the container limit and it is killed with no Spark-level OutOfMemoryError anywhere.

Senior: spark.executor.memoryOverhead — the non-heap headroom for Python workers, off-heap buffers, and shuffle. On PySpark and pandas UDF workloads that is usually what is actually being exceeded.

Common mistake: Raising spark.executor.memory to fix a container kill, which just asks for a limit it still overshoots.

Follow-up: Which setting do people forget in that second case?

Lesson · Simulation

architect

Why does Tungsten's binary row format belong in a memory conversation, not just a CPU one?

Executor memory: execution wins, cache loses · tap to open the answer

Short: Binary rows outside the Java object model use far less space per row and stop the GC from walking your data.

Detailed: Tungsten keeps rows in compact buffers with explicit offsets instead of JVM objects, so the same records occupy much less memory and the collector ignores them. Turning on spark.memory.offHeap.enabled with a size moves execution memory off the heap entirely, which keeps GC pauses down — but that memory still counts against the container's total limit.

Senior: Python UDFs — rows leave the binary format and become Python objects, so you pay the conversion and the memory back. That is a memory argument for built-ins, not only a CPU one.

Common mistake: Enabling off-heap memory without raising the container's overall allowance.

Follow-up: What destroys the Tungsten advantage in a PySpark pipeline?

Lesson · Simulation

senior

A stateful aggregation ran clean for a week, then started OOMing every few hours. First hypothesis?

Structured Streaming: an unbounded table · tap to open the answer

Short: No watermark, so state retained every key forever and finally outgrew the executors.

Detailed: Without withWatermark the engine cannot know a group will never receive late data, so its state is kept indefinitely. The streaming progress metrics show the state operator's numRowsTotal climbing monotonically and state store commit time growing with it. A watermark lets keys older than the event-time threshold be dropped.

Senior: Events later than the watermark are dropped and cannot be retrofitted — you are signing up to a bounded lateness contract, so route stragglers to a repair batch job if the business needs them.

Common mistake: Restarting the stream on a bigger cluster, which just moves the failure a few days out.

Follow-up: What do you give up when you add the watermark?

Lesson · Simulation

architect

Your foreachBatch writes to a transactional store and a restart happens mid-batch. How do you keep the sink correct?

Structured Streaming: an unbounded table · tap to open the answer

Short: Make the write idempotent and keyed, because the same batchId can be re-executed.

Detailed: foreachBatch hands you the DataFrame and the batchId, and on failure the engine may replay a batch that already wrote part of its rows. Use business keys (or the batchId) in a MERGE / upsert, or record the last committed batchId in the sink transactionally so a replay becomes a no-op. Output mode matters too: append emits only new rows, update emits changed keys, complete rewrites the whole result and only works for small aggregates.

Senior: When latency tolerance is hours: AvailableNow consumes everything currently available across several micro-batches and stops, so you keep streaming semantics and checkpointed offsets without paying for a cluster that never sleeps.

Common mistake: Assuming the checkpoint's exactly-once guarantee extends to an arbitrary external sink you wrote yourself.

Follow-up: When would you pick Trigger.AvailableNow over an always-on stream?

Lesson · Simulation

senior

How do you prove skew versus a slow node?

Data skew · tap to open the answer

Short: Skew: same task id always huge shuffle read for a key. Slow node: different keys, one host.

Detailed: Look at shuffle read bytes, not just duration. If bytes are even but duration isn't, suspect disk/CPU/noisy neighbor. If bytes are skewed, it's data.

Senior: Add salt, agg/join on (key, salt), then drop salt and re-agg if needed.

Common mistake: Speculative execution as the root-cause fix.

Follow-up: Walk through salting without ruining the grain.

Lesson · Simulation

senior

Broadcast join killed the driver, not the executors. Walk the path.

Driver OOM · tap to open the answer

Short: Driver collects the build side to ship the broadcast.

Detailed: If the dimension is 6 GB, the driver must hold it. Executors OOM later if they also can't hold the hashed relation — but the first failure can be the driver. Measure broadcast exchange size.

Senior: Error class → driver vs executor → action/broadcast in the plan → rewrite or size the JVM with a documented cap.

Common mistake: Only looking at executor logs.

Follow-up: What's your evidence-ordered debug?

Lesson · Simulation

senior

How do you choose between salting, more partitions, and more executor memory?

Executor OOM · tap to open the answer

Short: If one key is fat: salt or split. If all tasks are fat: more partitions or a smaller row (drop columns). Memory last.

Detailed: UI: shuffle read per task. One bar huge → skew. All bars huge → partition sizing or explode. GC thrash with even bytes → memory/config or too much cache.

Senior: Copy-pasting spark.memory.fraction from OSS blogs onto Photon/DBR without measuring.

Common mistake: A 2× memory ticket with no histogram attached.

Follow-up: What Spark config is a footgun with off-heap and Databricks?

Lesson · Simulation

senior

AQE switched your SMJ to broadcast mid-job. Is that good?

Broadcast joins · tap to open the answer

Short: Yes if the build side was small after filters. Bad if stats were wrong and it now OOMs.

Detailed: AQE join conversion uses runtime size. Watch SQL adaptive UI. If it OOMs, disable conversion for that query or fix the filter so the build side is truly small.

Senior: Don't hint a slowly growing dim forever. Partition/filter it, or SMJ when it crosses a documented threshold.

Common mistake: Disabling AQE globally because one query failed.

Follow-up: How do you cap broadcast size in a lakehouse with growing dims?

Lesson · Simulation

senior

AQE coalesced 2000 shuffle partitions down to 12, and now one task is huge. What happened?

Adaptive Query Execution · tap to open the answer

Short: Coalesce packed leftover data into too few partitions, or skew was hidden.

Detailed: Coalesce is for empty/small partitions. If data is uneven, you can get fat tasks. Inspect post-AQE partition sizes. Disable coalesce for that query or salt.

Senior: spark.sql.adaptive.enabled and compare SQL UI plans — same query, different later stages.

Common mistake: Blaming AQE without looking at the adaptive plan.

Follow-up: How do you show AQE on/off in an interview demo?

Lesson · Simulation

Practice by topic