skip to content

Apache Spark

Spark is the default distributed compute engine for batch and streaming work, and the one interviewers assume you have used. Expect questions across the whole stack: DataFrame APIs, the Catalyst optimizer, shuffles, and why a job that works on sample data melts on production volume.

on this pageshow

explore

questions

61 · 6 sections

In Spark, what is the difference between a transformation and an action on a DataFrame?

level: juniorimportance: must knowfreq 82%
basics
~20 s

A transformation such as filter or select only records intent and returns a new DataFrame; an action such as count, collect, show or write executes the accumulated plan on the cluster and returns a result outside Spark.

open as a page

Why does a Spark DataFrame filter usually outperform the equivalent RDD filter?

level: middleimportance: must knowfreq 72%
basics
~20 s

A DataFrame filter is an expression Spark understands, so it can prune columns, push the predicate into the file scan and generate compiled code over compact binary rows. An RDD lambda is an opaque closure Spark must feed deserialized objects, one row at a time.

open as a page

In Spark, what is RDD lineage and how is it used when an executor is lost?

level: middleimportance: should knowfreq 58%
basics
~20 s

Lineage is the recorded graph of parent datasets and the deterministic functions that produced each partition. When an executor dies, Spark does not restore a replica: it re-runs that graph to recompute only the lost partitions.

open as a page

What does schema inference cost when Spark reads a large CSV or JSON directory?

level: middleimportance: should knowfreq 45%
basics
~20 s

It costs an extra full pass over the data before your query even starts, and it makes the schema depend on today's values, so a column can silently change type between runs. Supplying an explicit StructType removes both problems.

open as a page

When is dropping from the DataFrame API to Spark's RDD API still justified?

level: seniorimportance: should knowfreq 48%
basics
~20 s

Rarely, and only for work the DataFrame API cannot express: non-tabular or custom-parsed data, explicit control of physical partition placement, position-based operators like zipWithIndex, or interop with RDD-based libraries. Everything expressible as columns should stay in DataFrames.

open as a page

How does Spark's planner choose between BroadcastHashJoin and SortMergeJoin?

level: middleimportance: must knowfreq 64%
basics
~20 s

Spark broadcasts when a hint says so or when one side's estimated size fits under spark.sql.autoBroadcastJoinThreshold, which defaults to 10 MB; otherwise an equi-join on sortable keys becomes a SortMergeJoin that shuffles and sorts both sides.

open as a page

In Spark, what does Catalyst do between your DataFrame code and the physical plan?

level: middleimportance: must knowfreq 70%
basics
~10 s

Catalyst runs four stages: analysis resolves names and types against the catalog, rule-based logical optimization rewrites the plan, physical planning picks operators and inserts shuffles, and Tungsten generates JVM code for fused operator chains.

open as a page

In Spark, does writing SQL instead of DataFrame code change how fast a query runs?

level: juniorimportance: should knowfreq 55%
basics
~20 s

No. A SQL string and the equivalent DataFrame code build the same logical plan, and Catalyst analyzes, optimizes and compiles both identically. Pick whichever reads better. The real exceptions are the RDD API and Python UDFs, which Catalyst cannot optimize through.

open as a page

In Spark, how do you read the physical plan that df.explain() prints?

level: middleimportance: should knowfreq 66%
basics
~20 s

Read it bottom-up: leaves are scans, indentation shows the tree. Each Exchange is a shuffle and a stage boundary, a *(n) prefix marks whole-stage code generation, and the scan line's PushedFilters and ReadSchema show what was pushed into the source.

open as a page

Why does a Python UDF in a Spark query defeat Catalyst optimization?

level: seniorimportance: should knowfreq 52%
basics
~20 s

Catalyst cannot see inside a Python UDF, so it cannot push it into the data source, fold it, or compile it into generated code. The plan gains a Python-evaluation node that breaks the fused loop and ships every row into a separate Python process and back.

open as a page

In Spark, what is a partition and what decides how many partitions a DataFrame starts with?

level: juniorimportance: must knowfreq 78%
basics
~20 s

A partition is Spark's unit of parallelism: one partition is processed by one task on one core. The starting count comes from the input — file splits packed to spark.sql.files.maxPartitionBytes (128 MB by default) — not from a fixed number.

open as a page

In Spark, how does repartition differ from coalesce, and when is coalesce the wrong choice?

level: middleimportance: must knowfreq 85%
basics
~20 s

repartition performs a full shuffle, so it can raise or lower the partition count and rebalances rows evenly. coalesce only merges existing partitions without a shuffle, can only lower the count, and caps the parallelism of the whole upstream stage.

open as a page

In Spark, what files does a map task write during a sort-based shuffle?

level: middleimportance: must knowfreq 68%
basics
~20 s

Each Spark map task writes two local files: one data file whose records are grouped and ordered by destination reduce partition, plus a small index file giving each partition's byte offset inside that data file.

open as a page

In a Spark stage, which task metrics separate data skew from a slow executor?

level: middleimportance: must knowfreq 70%
basics
~20 s

Compare the stage's per-task quantiles in the Spark UI. Skew shows one task reading many times the median task's shuffle bytes or records; a slow executor shows normal input sizes with long durations concentrated on one host.

open as a page

A Spark stage fails with FetchFailedException — what happened, and what does Spark do next?

level: seniorimportance: must knowfreq 62%
basics
~20 s

A reduce task could not read shuffle blocks from a peer after retries, almost always because the executor holding those files died or the node is unreachable. Spark marks the lost map output as missing, recomputes it, and retries the stage.

open as a page

In Spark Structured Streaming, what does each output mode — append, update, complete — write?

level: middleimportance: must knowfreq 74%
basics
~10 s

Append writes only rows that are final and never revises them. Update writes only rows whose value changed in this trigger. Complete rewrites the entire result table every trigger and requires an aggregation.

open as a page

In Spark Structured Streaming, what does withWatermark do to state and to late rows?

level: middleimportance: must knowfreq 68%
basics
~20 s

withWatermark declares how late event-time data may arrive. Spark tracks the maximum event time it has seen, subtracts that threshold, and uses the result to finalize windows, evict their state, and drop rows older than it.

open as a page

How does a Spark Structured Streaming query achieve end-to-end exactly-once output?

level: seniorimportance: should knowfreq 60%
basics
~20 s

Spark writes each micro-batch's input offset range to a write-ahead log under checkpointLocation before processing it, so a restart replays exactly that batch. Exactly-once then holds if the source is replayable by offset and the sink is idempotent or transactional.

open as a page

A Spark Structured Streaming windowed count in append mode emits nothing for an hour — why?

level: seniorimportance: should knowfreq 46%
basics
~20 s

Append mode holds each window until the watermark passes its end, and the watermark only advances from event times in data actually read. A stalled watermark, an oversized delay threshold, or a misbound withWatermark all produce a healthy query that emits nothing.

open as a page

When would you schedule a Spark Structured Streaming query with Trigger.AvailableNow instead of running it continuously?

level: principalimportance: should knowfreq 36%
basics
~20 s

Use Trigger.AvailableNow when the latency budget is minutes to hours and cluster cost dominates: it processes everything currently available in several rate-limited micro-batches, then stops, keeping checkpointed offsets and exactly-once semantics while the cluster runs only briefly.

open as a page

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

level: juniorimportance: must knowfreq 55%
basics
~20 s

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

open as a page

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

level: juniorimportance: must knowfreq 82%
basics
~10 s

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

open as a page

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

level: juniorimportance: must knowfreq 70%
basics
~20 s

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

open as a page

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

level: middleimportance: must knowfreq 78%
basics
~20 s

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

open as a page

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

level: middleimportance: must knowfreq 68%
basics
~10 s

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

open as a page

In Spark MLlib, what is the difference between a Transformer and an Estimator?

level: juniorimportance: must knowfreq 55%
basics
~20 s

A Transformer implements transform() and converts one DataFrame into another, usually by appending columns. An Estimator implements fit(), which learns from a DataFrame and returns a Model — and that Model is itself a Transformer.

open as a page

What does calling fit() on a Spark MLlib Pipeline produce, and how does each stage run?

level: middleimportance: must knowfreq 45%
basics
~20 s

Pipeline.fit() returns a PipelineModel. The stages run in the declared order: Transformer stages transform the DataFrame, and each Estimator stage is fitted and replaced by the Model it produced, so the result contains only Transformers.

open as a page

Why does a Spark MLlib StringIndexer fail on a category it never saw during fit()?

level: middleimportance: should knowfreq 38%
basics
~20 s

StringIndexer learns its label vocabulary during fit(), and its handleInvalid parameter defaults to error, so any category absent from the training data throws. Set it to keep, which assigns the extra index numLabels, or skip, which drops the row.

open as a page

Why should a Spark MLlib CrossValidator wrap the whole Pipeline instead of just the final estimator?

level: seniorimportance: should knowfreq 34%
basics
~20 s

A Spark Pipeline is itself an Estimator, so CrossValidator can re-fit every feature stage inside each fold's training split. Fitting scalers or indexers once over the full dataset leaks held-out statistics into training and inflates the score.

open as a page

How do you save and reload a fitted Spark MLlib PipelineModel, and what are the limits?

level: seniorimportance: should knowfreq 30%
basics
~10 s

Write a fitted pipeline with model.write().overwrite().save(path), which produces a directory on any Hadoop-compatible filesystem, and read it back with PipelineModel.load(path). Loading is guaranteed across minor and patch Spark versions, best-effort across majors.

open as a page