Apache Spark

Apache Spark

Definition: A distributed data processing engine that splits large computations across many machines, able to handle batch processing, real-time streaming, SQL analytics, and machine learning through one unified API. It originated around 2009 as a research project at UC Berkeley’s AMPLab, aiming to outperform Hadoop MapReduce on iterative workloads like machine learning training, and was open-sourced before being donated to the Apache Software Foundation in 2013. Spark’s defining idea was keeping intermediate data in memory across computation steps instead of writing every intermediate result to disk, which on many workloads made it roughly 10-100x faster than the MapReduce jobs it displaced. It’s now the default processing engine underneath most modern big-data platforms, including Databricks, which its original creators founded to commercialize it.

Architecture

  • RDDs and DataFrames: RDDs (Resilient Distributed Datasets) are Spark’s original core abstraction — an immutable, partitioned collection of records recoverable after failure by recomputing lost partitions from lineage; DataFrames are a higher-level, schema-aware API built on top of RDDs that lets Spark apply query optimizations RDDs’ opaque function calls can’t expose
  • Driver and executor model: a single driver process runs the user’s main program, builds the execution plan, and schedules work; a set of executor processes on worker machines actually run tasks and hold data partitions in memory or on local disk
  • Lazy evaluation: transformations like map, filter, and select don’t run immediately — Spark just records each one onto a growing logical plan, deferring real computation until an action forces it to execute
  • Transformations vs. actions: transformations (filter, groupBy, join) return a new, still-unevaluated dataset; actions (count, collect, write) trigger real execution across the cluster and return a concrete result
  • In-memory caching: .cache() or .persist() keeps a DataFrame that’s reused across multiple downstream operations resident in executor memory, avoiding recomputing its entire lineage from scratch each time it’s referenced
  • DAG scheduler: the driver compiles the chain of transformations into a directed acyclic graph of stages, splitting the graph at every point where data must be shuffled across the network
  • Catalyst optimizer: DataFrame and Spark SQL queries are parsed into a logical plan, rewritten through rule-based and cost-based optimizations — predicate pushdown, column pruning, join reordering — then compiled into an optimized physical execution plan
  • Cluster manager options: the driver requests resources from a pluggable cluster manager — Spark’s own standalone manager, Hadoop YARN, or Kubernetes (K8s) — which allocates and monitors executor processes on worker nodes

The driver, cluster manager, and executor layers map onto physical machines roughly like this:

Why It Matters

  • Made large-scale iterative computation practical — before Spark, algorithms that reread the same data many times (machine learning training, graph algorithms) paid Hadoop MapReduce’s full disk read/write cost on every single pass
  • Unified a previously fragmented stack — one engine now commonly handles batch ETL, streaming, SQL analytics, and ML training, where teams once stitched together several specialized systems for each
  • Language-agnostic access to distributed computing — APIs in Python, Scala, Java, R, and SQL let data scientists and analysts use a language they already know instead of learning a distributed-systems framework from scratch
  • Powers most modern managed big-data platforms — Databricks, AWS EMR, Google Dataproc, and Azure Synapse all ship Spark as a core or optional processing engine
  • Scales down as easily as it scales up — the same code that runs in local mode on a laptop for development runs unmodified on a production cluster of thousands of cores, lowering the barrier to testing and iterating on real logic

Under the Hood: Lazy Evaluation and the DAG

Spark deliberately avoids running a transformation the moment it’s called. Writing .filter() or .groupBy() on a DataFrame doesn’t touch any data — it just appends a node to a logical plan the driver is silently building. Only an action, like .count() or .write(), forces real computation to happen. This laziness exists for a concrete reason: it lets Spark see the entire chain of operations before running any of them, rather than executing each line blindly and in isolation the instant it’s written.

That full visibility is what makes optimization possible, roughly in this order:

  1. Build the logical plan — each transformation appends to an abstract, unevaluated description of the computation
  2. Catalyst analyzes and rewrites it — resolving column references against the schema, then applying rule-based rewrites such as pushing filters down close to the data source and pruning columns nothing downstream reads
  3. Cost-based optimization picks a physical strategy — for example, choosing a broadcast join over a shuffle join once it’s known one side of a join is small enough to fit in executor memory
  4. The physical plan compiles into a DAG of stages — the scheduler walks the optimized plan and cuts it into stages at every point a shuffle is unavoidable
  5. Stages break into tasks — each stage runs one parallel task per partition, scheduled onto whatever executor cores are free

The stage boundary in step 4 is the single most consequential detail in Spark performance work. Within a stage, every operation is “narrow” — each output partition depends only on a small, known set of input partitions, so no data has to move between machines. The moment an operation needs rows with the same key to land on the same partition — groupBy, join, distinct, repartition — Spark has no choice but to shuffle: write intermediate data to disk, transfer it across the network, and read it back on the other side. A concrete pipeline — reading a CSV, filtering, grouping, aggregating, and writing the result — splits into exactly two stages around that one shuffle:

Every task in Stage 0 can run entirely independently, on whichever executor already holds that partition. Stage 1 can’t start until Stage 0’s shuffle output is fully written, because any task in Stage 1 might need rows that were scattered across many different Stage 0 partitions.

DimensionApache SparkHadoop MapReduceApache Flink
Execution modelIn-memory DAG of stages, lazy evaluationRigid two-phase map-then-reduce, disk-basedNative streaming dataflow graph of continuous operators
SpeedRoughly 10-100x faster than MapReduce on iterative workloadsSlowest of the three — reads/writes disk between every stageComparable to or faster than Spark specifically on streaming workloads
Streaming supportMicro-batch via Structured Streaming, near-real-timeNone natively — batch-only by designTrue event-at-a-time streaming, sub-second latency
Typical latencySeconds, per micro-batchMinutes to hours per jobMilliseconds, continuous
Typical use caseGeneral-purpose batch ETL, ML pipelines, ad-hoc SQLLegacy large-scale batch jobs, largely supersededLow-latency streaming analytics, event-driven applications

Common Pitfalls

  • Using Spark for small datasets — the distributed-computing overhead (serialization, network shuffles, JVM startup) makes it slower than Pandas or plain SQL once data genuinely fits on one machine
  • Data skew — a poorly chosen partition or join key sends a disproportionate share of rows to a few partitions, so a handful of tasks run far longer than the rest while other executors sit idle
  • Too many small output files — an over-partitioned write, common in streaming jobs, produces thousands of tiny files that overwhelm the driver and downstream readers with file-listing overhead
  • Not caching a reused DataFrame — without .cache() or .persist(), Spark recomputes an entire lineage from scratch every time an already-computed DataFrame is referenced again in a new action
  • Collecting too much data to the driver — calling .collect() on a large DataFrame pulls every partition into the driver’s single JVM heap, a common and entirely avoidable cause of driver out-of-memory crashes
  • Confusing narrow and wide transformations — assuming an operation is as cheap as a filter or map when it’s actually a groupBy, join, or distinct that triggers a full shuffle

Worked Example: Tracing a Job Through Its Stages

Consider a job that reads a CSV of orders, filters to one region, groups by customer, sums order totals, and writes the result:

StageOperationShuffle?What happens
0spark.read.csv(...)NoReads the file and splits it into partitions across executors
0.filter(...) then .select(...)NoNarrow transformations — each partition filters and trims its own rows independently
1.groupBy("customer_id")YesWide transformation — matching keys must land on the same partition, forcing a shuffle
1.agg(sum("amount"))NoRuns after the shuffle, once matching keys are already co-located
1.write.parquet(...)NoAction — triggers execution of the entire DAG built above

Nothing in Stage 0 or Stage 1 actually ran until that final .write() call — up to that point, every line above only extended the logical plan.

Real-World Use

  • ETL at scale — nightly batch jobs clean, join, and reshape terabytes of raw data as one stage in a larger Data Pipeline, feeding a Data Lake vs Data Warehouse and commonly orchestrated by Apache Airflow
  • Machine learning pipelines — MLlib provides distributed implementations of common algorithms (logistic regression, random forests, ALS recommendation) that train directly on data too large for one machine
  • Streaming analytics — Structured Streaming applies the same DataFrame API to unbounded data arriving continuously from sources like Apache Kafka, unifying batch and streaming code in one job
  • Ad-hoc SQL analytics on huge datasets — analysts query billions of rows interactively through Spark SQL without writing any Scala or Python
  • Graph processing — GraphFrames runs algorithms like PageRank or connected-components across graphs with billions of edges

Best Practices

  • Size partitions deliberately — a rough rule of thumb targets somewhere around 100-200MB per partition; too small wastes scheduling overhead, too large risks executor memory pressure
  • Cache strategically, not reflexively — only .persist() a DataFrame genuinely reused across multiple actions, and .unpersist() it once it’s no longer needed to free executor memory
  • Prefer narrow transformations where possible — push filtering and column pruning before a shuffle rather than after, so less data has to move across the network
  • Use broadcast joins for small tables — explicitly broadcasting a small lookup table avoids shuffling a much larger table across the network just to join against it
  • Monitor jobs through the Spark UI — the stages, tasks, and storage tabs surface skew, disk spill, and garbage-collection pressure long before a job fails outright
  • Avoid .collect() on large results — write large outputs to storage, or sample explicitly, rather than pulling an entire DataFrame back to the driver

FAQ

Is Spark a database? No — it’s a processing engine, not a storage system. It reads from and writes to external storage like a data lake or warehouse but doesn’t persist data itself between jobs.

Does Spark replace Hadoop entirely? Partially — Spark replaced MapReduce as the dominant processing engine, but it commonly still runs on Hadoop’s YARN resource manager and reads from HDFS, so the wider Hadoop ecosystem often persists even where MapReduce itself doesn’t.

Why is Spark faster than MapReduce? Mainly because it keeps intermediate results in memory across the steps of a pipeline instead of writing every intermediate result to disk, and because its DAG scheduler optimizes a whole chain of operations rather than treating each step as an isolated job.

What’s the difference between an RDD and a DataFrame? An RDD is a low-level, unstructured collection of arbitrary objects with no schema; a DataFrame is a structured, schema-aware, table-like abstraction built on top of RDDs that lets Catalyst apply optimizations RDDs’ opaque code can’t expose.

Can Spark run on a single machine? Yes — local mode runs the driver and all executors as threads inside a single JVM, commonly used for development and testing before the same code is deployed to a real cluster.

History

  • Started in 2009 as a research project at UC Berkeley’s AMPLab, led by Matei Zaharia, aiming to outperform MapReduce on iterative workloads like machine learning
  • Open-sourced in 2010, donated to the Apache Software Foundation in 2013, and became a top-level Apache project in February 2014 — one of the fastest-growing open-source data projects of its era
  • Matei Zaharia and several other original creators founded Databricks in 2013 to commercialize and support Spark
  • Spark 1.x introduced Spark SQL and an early DataFrame API, moving the primary interface beyond raw RDDs
  • Spark 2.x (2016) unified the DataFrame and Dataset APIs and introduced Structured Streaming
  • Spark 3.x (2020) added Adaptive Query Execution, letting the engine re-optimize a query plan mid-job using real statistics rather than only upfront estimates

Common Interview Questions

  • “Explain the difference between a transformation and an action” — expect an answer centered on laziness: transformations build up a plan, only an action triggers real computation across the cluster
  • “What causes a shuffle, and why is it expensive?” — expect an explanation that wide transformations require moving data between partitions across the network and disk, unlike narrow transformations that operate purely locally
  • “How does Spark achieve fault tolerance without replicating data?” — expect an answer about RDD lineage: any lost partition can be recomputed from the transformations that produced it, rather than needing a standby replica
  • “When would you choose a broadcast join over a shuffle join?” — expect a discussion of table size: broadcasting a small table to every executor avoids shuffling a much larger table just to align matching keys

Example

An e-commerce company runs a nightly Spark job that reads several terabytes of raw clickstream logs from cloud storage, filters out bot traffic, joins page-view events against a much smaller product catalog table using a broadcast join, groups the result by user session, and aggregates it into per-session summaries — time on site, pages viewed, items added to cart. Because every step above is a lazy transformation, none of it actually runs until the final .write() call, by which point Catalyst has already rewritten the whole chain into an optimized physical plan and the DAG scheduler has split it into stages around the one unavoidable shuffle, the groupBy on session ID. The entire job, spanning billions of rows, finishes in minutes on a cluster of a few dozen machines — work that would take many hours, or simply exhaust available memory, on a single machine.

Dig deeper