What you will be able to do
- Explain why shuffles are the most expensive operations in a Spark workload
- Recognize spill, skew and garbage-collection pressure as the main memory-related challenges and describe their symptoms
- Identify the limits that come from Spark's cluster architecture, such as application isolation, driver memory and cluster cost
- Decide when a workload fits Spark poorly, for example task-parallel work or Python UDFs that replace native functions
1.The shuffle: the price of distributing data
Spark is fast because each task works on one partition independently. That breaks down when an operation needs data that lives in many partitions. To combine all the values for one key, for example, 'Spark needs to perform an all-to-all operation'. It reads from every partition and brings matching values together. That movement of data is called the shuffle.
A shuffle 'typically involves copying data across executors and machines, making the shuffle a complex and costly operation'. Join operations, most ByKey operations and repartitioning can all trigger one. Shuffles also stress memory. Databricks notes that the memory used for shuffles, joins, sorts and aggregations is where trouble usually starts. When that memory runs out, data moves to disk, and this 'is most common during data shuffling'. The next section covers that memory pressure.
Shuffles also add a tuning burden. The right number of shuffle partitions depends on the data, and Databricks notes that data sizes 'may differ vastly from stage to stage, query to query, making this number hard to tune'. The Spark guide likewise says shuffle behavior is tuned by adjusting a variety of configuration parameters.
Checkpoint 1 of 6· Check yourself
A colleague says a groupByKey is cheap because Spark is an in-memory engine. What is the best correction?
ByKey operations other than counting can trigger a shuffle, and a shuffle moves data between machines at the cost of disk, serialization and network I/O.
“This typically involves copying data across executors and machines, making the shuffle a complex and costly operation.”Source: spark.apache.org
2.Memory pressure: spill, skew and garbage collection
Running in memory is an advantage only while the data fits. Databricks defines the first failure mode directly: 'Spill is what happens when Spark runs low on execution memory.' Spark then moves data from memory to disk, which can be expensive. The RDD guide describes the same effect for shuffle data structures: 'Spark will spill these tables to disk, incurring the additional overhead of disk I/O and increased garbage collection.' So memory pressure costs you twice, once in disk I/O and again in JVM garbage collection.
Shuffles also leave intermediate files on disk. Spark keeps them until the related RDDs are garbage collected, and that can take a long time if the application still holds references. The guide warns that 'long-running Spark jobs may consume a large amount of disk space'.
The second failure mode is skew. 'Skew is when one or just a few tasks take much longer than the rest.' Because a stage isn't finished until its slowest task is, skew 'results in poor cluster utilization and longer jobs': most executors sit idle while one works through an oversized partition.
| Challenge | What is happening | What you see |
|---|---|---|
| Spill | Spark runs low on execution memory and moves data to disk | Spill stats appear in the stage details |
| Skew | One or a few tasks take much longer than the rest | Max task duration well above the 75th percentile |
| Garbage collection | Spilled shuffle tables add JVM collection overhead | The RDD guide lists increased garbage collection as an extra cost of spilling, on top of disk I/O |
Checkpoint 2 of 6· Check yourself
In a stage's summary metrics, the 75th-percentile task duration is 40 seconds and the Max is 3 minutes. What is the most likely problem?
Databricks uses a Max duration more than 50% above the 75th percentile as the sign of skew. Spill is confirmed by spill stats in the stage details, not by the spread of task durations.
“If the Max duration is 50% more than the 75th percentile, you may be suffering from skew.”Source: docs.databricks.com
The advantage that memory pressure undermines is in-memory reuse. Spark lets you persist a dataset in memory, 'allowing it to be reused efficiently across parallel operations', and that 'allows future actions to be much faster (often by more than 10x)'. Caching has a cost, though. Each node keeps the partitions it computes in memory, and Spark drops old partitions in least-recently-used fashion when the cache fills. Databricks also warns that manually caching data in production pipelines 'can interrupt these optimizations and lead to increases in cost and latency'. Caching is a trade-off, not a free speedup.
Checkpoint 3 of 6· Exam question
An engineer migrates a nightly aggregation pipeline from disk-based MapReduce to Spark and observes a large runtime improvement on a job that repeatedly scans and re-scans the same intermediate dataset across several stages. Which characteristic of Spark's design explains this speedup?
Correct answer: A — Spark can hold intermediate results in executor memory across stages, avoiding the repeated disk reads and writes that MapReduce performs between each map and reduce phase.
- A. This is correct: Spark's in-memory processing lets it cache intermediate DataFrames or RDDs between stages, so a pipeline that reuses the same dataset avoids the round-trip to disk that MapReduce requires between phases. This is the core mechanism behind Spark's speed advantage on iterative or reused-data workloads.
- B. Spark does not automatically compress every dataset with a columnar codec between stages; compression is format- and configuration-dependent, not an automatic universal behavior. This is not the reason repeated scans get faster.
- C. Stages within a job generally have data dependencies, so later stages wait on the shuffle output of earlier ones rather than running fully in parallel. Describing all stages as simultaneous misrepresents Spark's DAG-based scheduling.
- D. Spark still runs on the JVM and is subject to its garbage collector; it does not replace GC with a pause-free custom manager. Memory tuning reduces GC impact but does not eliminate it.
- E. Spark does not compile job logic ahead of time into native machine code; execution still goes through the JVM (or Python worker processes for PySpark UDFs). Native compilation is not the mechanism behind the observed speedup.
3.Limits that come from the cluster architecture
Some challenges come from how a Spark application is built. Each application gets its own executor processes for its whole lifetime. 'This has the benefit of isolating applications from each other', since tasks from different applications run in different JVMs. The trade-off is that 'data cannot be shared across different Spark applications (instances of SparkContext) without writing it to an external storage system'. Cluster managers such as Standalone, YARN and Kubernetes allocate resources across applications, so concurrent applications compete for the same cluster resources and can't hand data to each other in memory.
The driver brings constraints of its own. It schedules every task, so it must stay reachable: 'the driver program must be network addressable from the worker nodes', and it should run close to them, preferably on the same local area network. Its memory is also limited. Databricks warns that pulling results back from the cluster has to stay small: 'Only collect small amounts of data back to R data frames, or the Spark driver will run out of memory.'
Not necessarily. You pay for workers for as long as the workload runs. If a job scales linearly, four workers for half an hour costs the same as two workers for an hour, and finishes sooner. Databricks adds that if cost matters most and the SLA is flexible, an autoscaling cluster is usually the cheapest, but not necessarily the fastest. Sizing is therefore a real decision: VM families differ in RAM, cores, network bandwidth and local storage.
Checkpoint 4 of 6· Check yourself
Two separate Spark applications on the same cluster need to use the same intermediate result. What must happen?
Each application has its own isolated executors, so data passes between applications only through external storage.
“data cannot be shared across different Spark applications (instances of SparkContext) without writing it to an external storage system.”Source: spark.apache.org
4.When Spark is the wrong tool
The last challenge is choosing Spark for work it doesn't suit. Spark is built for data parallelism and is recommended for large-scale data processing such as joins, filtering and aggregation. Work that consists of many independent computations is task parallelism, and Databricks points to Ray for it: 'Ray is designed for task parallelism, where multiple tasks run concurrently and independently.'
| Workload | Better fit |
|---|---|
| ETL, analytics reporting, feature engineering, data preprocessing | Spark |
| Large-scale machine learning with MLlib | Spark |
| Reinforcement learning, simulation modeling, hyperparameter search | Ray |
| Deep learning training and high-performance computing (HPC) | Ray |
Even inside Spark, some code works against the engine. Python UDFs that replace a native function are the usual example. 'Serialization is required to transfer data between Python and Spark. This significantly slows down queries.' Use native functions where they exist. When a Python UDF is unavoidable, Databricks recommends Pandas UDFs, where Apache Arrow moves data efficiently between Spark and Python.
What Spark does handle for you on suitable workloads is recovery from failures. The RDD guide says 'RDDs automatically recover from node failures', and that a lost cached partition 'will automatically be recomputed using the transformations that originally created it'. Replicated storage levels let tasks keep running without waiting for that recomputation.
Checkpoint 5 of 6· Match them up
Match each situation to the recommendation the documentation gives
Tap a term, then the definition that fits it.
Spark fits data-parallel work and native operations. Task-parallel workloads and row-by-row Python serialization are where it struggles.
“Use Ray for workloads where Spark is less optimized, such as reinforcement learning, hierarchical time series forecasting, simulation modeling, hyperparameter search”Source: docs.databricks.com
Checkpoint 6 of 6· Exam question
During a long-running Spark job, one executor is terminated after the underlying VM is reclaimed by the cloud provider. The job continues and eventually finishes with correct results. Which two statements correctly describe how Spark achieves this fault tolerance? (Choose 2 answers)(Select 2)
Correct answers: A, B — Spark tracks the lineage of transformations that produced each partition, so lost partitions on the failed executor can be recomputed from the original data and transformations.; The driver detects the missing executor's heartbeat and reschedules its pending and lost tasks onto other executors that remain part of the cluster.
- A. This is correct: Spark's RDD/DataFrame lineage graph records the sequence of transformations, so when a partition is lost it can be recomputed deterministically from source data rather than requiring a live replica. This lineage-based recovery is central to Spark's fault tolerance model.
- B. This is correct: the driver monitors executor heartbeats and, on detecting a failure, reassigns the affected tasks to healthy executors so the job can continue to completion. This rescheduling is what lets the job finish despite the lost VM.
- C. Spark does not silently drop lost partitions and return partial results; it recomputes them from lineage so the final output remains complete and correct. Returning incomplete data would defeat the purpose of fault tolerance.
- D. Spark does not synchronously mirror every partition to a second executor; that would double memory and network cost for every job. Recovery instead relies on recomputation from lineage, not live replication.
- E. A restarted executor process does not retain its prior in-memory state, since that state was lost when the VM was reclaimed. Spark must resubmit the lost tasks rather than assume the old state survives.
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.Spark is an in-memory engine, so a Spark job never touches disk during processing.Why is that wrong?
When Spark runs low on execution memory it spills data to disk, most often during shuffles, which adds disk I/O and extra garbage collection.
Covered in Memory pressure: spill, skew and garbage collection
2.Applications on the same Spark cluster share executors, so one application can reuse another's in-memory data.Why is that wrong?
Each application has its own isolated executor processes. Data passes between applications only through external storage.
3.Spark is the best choice for every distributed Python workload.Why is that wrong?
Spark excels at data parallelism. Task-parallel workloads such as hyperparameter search or simulation suit Ray better.
Covered in When Spark is the wrong tool
Sources
Every claim above is drawn from one of these pages, quoted as it was written on the date shown.
- 1.https://spark.apache.org/docs/latest/rdd-programming-guide.htmlSecondary source
“The Shuffle is an expensive operation since it involves disk I/O, data serialization, and network I/O.”
↩︎ The shuffle: the price of distributing data“This typically involves copying data across executors and machines, making the shuffle a complex and costly operation.”
↩︎ The shuffle: the price of distributing data“Shuffle behavior can be tuned by adjusting a variety of configuration parameters.”
↩︎ The shuffle: the price of distributing data“Spark will spill these tables to disk, incurring the additional overhead of disk I/O and increased garbage collection.”
↩︎ Memory pressure: spill, skew and garbage collection“This means that long-running Spark jobs may consume a large amount of disk space.”
↩︎ Memory pressure: spill, skew and garbage collection“allowing it to be reused efficiently across parallel operations”
↩︎ Memory pressure: spill, skew and garbage collection“This allows future actions to be much faster (often by more than 10x).”
↩︎ Memory pressure: spill, skew and garbage collection“Finally, RDDs automatically recover from node failures.”
↩︎ When Spark is the wrong tool“if any partition of an RDD is lost, it will automatically be recomputed using the transformations that originally created it.”
↩︎ When Spark is the wrong tool - 2.
“It starts to move data from memory to disk, which can be expensive. It is most common during data shuffling.”
↩︎ The shuffle: the price of distributing data“Spill is what happens when Spark runs low on execution memory.”
↩︎ Memory pressure: spill, skew and garbage collection“Skew is when one or just a few tasks take much longer than the rest.”
↩︎ Memory pressure: spill, skew and garbage collection“This results in poor cluster utilization and longer jobs.”
↩︎ Memory pressure: spill, skew and garbage collection“Spill is what happens when Spark runs low on execution memory.”
↩︎ Exam trap 1“If the Max duration is 50% more than the 75th percentile, you may be suffering from skew.”
↩︎ Checkpoint - 3.https://www.databricks.com/blog/2020/05/29/adaptive-query-execution-speeding-up-spark-sql-at-runtime.htmlSecondary source
“data sizes may differ vastly from stage to stage, query to query, making this number hard to tune”
↩︎ The shuffle: the price of distributing data - 4.https://docs.databricks.com/aws/en/spark/faqOfficial docs
“Manually caching data or returning preview results in production pipelines can interrupt these optimizations and lead to increases in cost and latency.”
↩︎ Memory pressure: spill, skew and garbage collection - 5.https://spark.apache.org/docs/latest/cluster-overview.htmlSecondary source
“This has the benefit of isolating applications from each other”
↩︎ Limits that come from the cluster architecture“the driver program must be network addressable from the worker nodes.”
↩︎ Limits that come from the cluster architecture“data cannot be shared across different Spark applications (instances of SparkContext) without writing it to an external storage system.”
↩︎ Exam trap 2“data cannot be shared across different Spark applications (instances of SparkContext) without writing it to an external storage system.”
↩︎ Checkpoint - 6.
“Only collect small amounts of data back to R data frames, or the Spark driver will run out of memory.”
↩︎ Limits that come from the cluster architecture - 7.https://docs.databricks.com/aws/en/lakehouse-architecture/performance-efficiency/best-practicesOfficial docs
“So, if you spin up two worker clusters and it takes an hour, you are paying for those workers for the full hour.”
↩︎ Limits that come from the cluster architecture“an autoscaling cluster is usually the cheapest, but not necessarily the fastest.”
↩︎ Limits that come from the cluster architecture“Serialization is required to transfer data between Python and Spark. This significantly slows down queries.”
↩︎ When Spark is the wrong tool - 8.
“Ray is designed for task parallelism, where multiple tasks run concurrently and independently.”
↩︎ When Spark is the wrong tool“Ray is designed for task parallelism, where multiple tasks run concurrently and independently.”
↩︎ Exam trap 3“Use Ray for workloads where Spark is less optimized, such as reinforcement learning, hierarchical time series forecasting, simulation modeling, hyperparameter search”
↩︎ Checkpoint