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 pageshowhide
explore
- RDDs and DataFrames6 questions
- Spark SQL and Catalyst6 questions
- Partitioning and Shuffles19 questions
- Partitions and Partitioners6 questions
- Shuffle Internals6 questions
- Skew and Adaptive Execution7 questions
- Structured Streaming6 questions
- Cluster Execution and Tuning18 questions
- Jobs, Stages and Tasks6 questions
- Memory Management and Caching6 questions
- Submission and Deploy Modes6 questions
- MLlib Pipelines6 questions
questions
61 · 6 sectionsIn Spark, what is the difference between a transformation and an action on a DataFrame?
basics
~20 sA 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.
Why does a Spark DataFrame filter usually outperform the equivalent RDD filter?
basics
~20 sA 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.
In Spark, what is RDD lineage and how is it used when an executor is lost?
basics
~20 sLineage 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.
What does schema inference cost when Spark reads a large CSV or JSON directory?
basics
~20 sIt 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.
When is dropping from the DataFrame API to Spark's RDD API still justified?
basics
~20 sRarely, 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.
How does Spark's planner choose between BroadcastHashJoin and SortMergeJoin?
basics
~20 sSpark 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.
In Spark, what does Catalyst do between your DataFrame code and the physical plan?
basics
~10 sCatalyst 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.
In Spark, does writing SQL instead of DataFrame code change how fast a query runs?
basics
~20 sNo. 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.
In Spark, how do you read the physical plan that df.explain() prints?
basics
~20 sRead 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.
Why does a Python UDF in a Spark query defeat Catalyst optimization?
basics
~20 sCatalyst 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.
In Spark, what is a partition and what decides how many partitions a DataFrame starts with?
basics
~20 sA 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.
In Spark, how does repartition differ from coalesce, and when is coalesce the wrong choice?
basics
~20 srepartition 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.
In Spark, what files does a map task write during a sort-based shuffle?
basics
~20 sEach 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.
In a Spark stage, which task metrics separate data skew from a slow executor?
basics
~20 sCompare 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.
A Spark stage fails with FetchFailedException — what happened, and what does Spark do next?
basics
~20 sA 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.
In Spark Structured Streaming, what does each output mode — append, update, complete — write?
basics
~10 sAppend 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.
In Spark Structured Streaming, what does withWatermark do to state and to late rows?
basics
~20 swithWatermark 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.
How does a Spark Structured Streaming query achieve end-to-end exactly-once output?
basics
~20 sSpark 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.
A Spark Structured Streaming windowed count in append mode emits nothing for an hour — why?
basics
~20 sAppend 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.
When would you schedule a Spark Structured Streaming query with Trigger.AvailableNow instead of running it continuously?
basics
~20 sUse 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.
In spark-submit, what do the --master, --class and application-jar arguments specify?
basics
~20 sIn 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.
In Spark, what is the difference between a job, a stage and a task?
basics
~10 sIn 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.
In Spark, how do the MEMORY_ONLY and MEMORY_AND_DISK persist levels differ when a partition will not fit?
basics
~20 sWith 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.
In spark-submit, what is the difference between --deploy-mode client and --deploy-mode cluster?
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.
In Spark, what determines where the DAG scheduler places a stage boundary?
basics
~10 sSpark 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.
In Spark MLlib, what is the difference between a Transformer and an Estimator?
basics
~20 sA 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.
What does calling fit() on a Spark MLlib Pipeline produce, and how does each stage run?
basics
~20 sPipeline.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.
Why does a Spark MLlib StringIndexer fail on a category it never saw during fit()?
basics
~20 sStringIndexer 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.
Why should a Spark MLlib CrossValidator wrap the whole Pipeline instead of just the final estimator?
basics
~20 sA 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.
How do you save and reload a fitted Spark MLlib PipelineModel, and what are the limits?
basics
~10 sWrite 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.