Same code, same cluster, five times slower. A calm, step by step way to read the Spark UI, find the slow stage, and fix the real cause.
Your Databricks job usually finishes in 8 minutes. Today it took 42.
Nobody touched the code. The cluster is the same. The notebook is the same. So what changed?
This is one of the most common questions in data engineering. And the honest answer is that reading the code again will rarely tell you. The code describes what you asked for. The Spark UI shows what actually happened.
This article walks through a calm, repeatable way to find out. It follows one real style of incident from start to finish: the problem, the Spark UI, the root cause, and the fix.
Think of a Spark job as a team of workers sorting mail.
If every worker gets a similar pile, the whole team finishes at about the same time. But if one worker is handed half of all the letters, everyone else finishes early and waits. The job is only done when that last worker is done.
Most sudden slowdowns are some version of this story. The data moved in a way that made one part of the work much bigger than the rest.
So the goal is not to stare at code. The goal is to find where the data moved, and where the work piled up.
Before opening the Spark UI, it helps to know its three levels.
A job is started by an action, like a write or a count.
A job is split into stages. A new stage starts every time Spark has to move data between machines. That movement is called a shuffle.
Each stage is split into tasks. One task processes one partition of data. Tasks run in parallel.
So when something is slow, you go from the outside in. Job, then stage, then task. If you want a gentle refresher on how DataFrames turn into this work, Chapter 2: Spark DataFrames covers it from the start.
Write down what you know before you click anything.
The job usually takes 8 minutes. It now takes 42 minutes. The code and the cluster did not change.
That last sentence matters. If nothing in the code changed, the most likely suspect is the data. Keep that in mind as you look.
Open the job run, then the Spark UI, then the Stages view. Sort by duration.
In our example the job has five stages:
| Stage | Description | Duration | Tasks | Shuffle read | Shuffle write |
|---|---|---|---|---|---|
| 0 | Scan parquet | 21 s | 200 | - | - |
| 1 | Filter | 35 s | 200 | - | - |
| 2 | Project | 1.2 min | 200 | - | - |
| 3 | HashAggregate | 20 min | 200 | 320 GB | 400 GB |
| 4 | Write to Delta | 1.5 min | 200 | - | - |
Stage 3 stands out right away. It takes 20 minutes and it moves a lot of data: 320 GB read from a shuffle and 400 GB written.
A large shuffle is a strong signal. It almost always means a wide transformation, such as a join, groupBy, distinct, orderBy, or repartition.
Click into the slow stage and open the DAG, or the SQL tab, to see what operations it contains.
A simplified view of our job looks like this:
Scan (Parquet) -> Filter (narrow) -> Exchange (shuffle) -> HashAggregate (wide) -> Exchange (shuffle) -> Write (Delta)Each Exchange box is a shuffle. The HashAggregate came from a groupBy in the code. So this stage is doing wide work and moving data across the cluster.
This is where knowing the kinds of transformations really pays off. Here are the seven you will meet most often, in simple words.
1. Narrow. No shuffle. Each output partition depends on one input partition. Examples are filter, select, and withColumn. These are cheap and rarely the cause of a slowdown.
2. Wide. Requires a shuffle. Data moves across partitions so that rows with the same key end up together. Examples are groupBy, join, distinct, and repartition. This is where most performance problems live.
3. Stateless. Each record is processed on its own, without remembering anything from before. A filter on a stream is stateless.
4. Stateful. Spark keeps state across records or across streaming batches. Running counts, windows, and stream deduplication are stateful. State can grow quietly over time, which is why watermarks matter. Chapter 12: Streaming explains this with a small pipeline.
5. Aggregation. Many records become fewer records. groupBy with count, sum, or avg are aggregations. They usually need a shuffle.
6. Join. Two datasets are combined using matching keys. Unless one side is small enough to broadcast, both sides are shuffled by the join key.
7. Map side. Work done locally on each partition before the shuffle, to reduce how much data has to move. Spark does a partial aggregation on each executor first, then shuffles only the partial results. This is why a simple count is often much cheaper than you would expect.
A useful habit: when you see an Exchange in the plan, ask yourself which wide transformation created it, and whether it was needed. Chapter 4: Transformations and Chapter 6: Joins and Aggregations both practice this.
Now open the Tasks view for Stage 3.
This is the most important screen in the whole investigation. Do not just look at the average. Look at the spread.
In our example:
That one task is called a straggler. And a straggler with far more input than its neighbors is the classic sign of data skew.
Remember the mail sorting team. One worker got half the letters. The stage cannot finish until that worker finishes, so the whole job waits.
The summary metrics at the top of the Tasks view help here too. If the max duration is many times larger than the median, you are almost certainly looking at skew.
The Spark UI told you where. Now the data will tell you why.
The aggregation groups by customer_id. So check how evenly that key is spread:
SELECT
customer_id,
COUNT(*) AS record_count
FROM orders
GROUP BY customer_id
ORDER BY record_count DESC
LIMIT 10;The result:
| customer_id | record_count |
|---|---|
| UNKNOWN | 480,000,000 |
| 1001 | 120,000 |
| 1002 | 115,000 |
| 1003 | 110,000 |
There it is. One value, UNKNOWN, has 480 million rows. Every one of those rows has the same key, so Spark sends all of them to the same partition. One task gets all of it.
This also answers the original question. The code did not change. But an upstream system probably started writing UNKNOWN when it could not find a customer. Yesterday that was a small number of rows. Today it is most of the table.
Values like UNKNOWN, NULL, -1, N/A, or an empty string are the most common cause of sudden skew. They are worth checking first, every time.
Choose the fix based on the cause, not on habit. Here are the four families of fixes, in the order we usually try them.
Deal with junk keys first. If UNKNOWN rows do not belong in the aggregation, filter them out. If they do, process them separately and union the result back.
from pyspark.sql import functions as F
known_orders_df = orders_df.filter(F.col("customer_id") != "UNKNOWN")
unknown_orders_df = orders_df.filter(F.col("customer_id") == "UNKNOWN")
known_totals_df = known_orders_df.groupBy("customer_id").agg(F.count("*").alias("record_count"))
unknown_totals_df = unknown_orders_df.agg(F.count("*").alias("record_count")).withColumn("customer_id", F.lit("UNKNOWN"))
customer_totals_df = known_totals_df.unionByName(unknown_totals_df)Salt a real hot key. If one genuine key is huge, add a small random number to it so its rows spread across several partitions. Aggregate twice: first by the salted key, then by the real key.
from pyspark.sql import functions as F
SALT_BUCKETS = 16
salted_df = orders_df.withColumn("salt", (F.rand() * SALT_BUCKETS).cast("int"))
partial_totals_df = salted_df.groupBy("customer_id", "salt").agg(F.count("*").alias("partial_count"))
customer_totals_df = partial_totals_df.groupBy("customer_id").agg(F.sum("partial_count").alias("record_count"))Pre-aggregate large keys before a join, so you join a small summary instead of millions of raw rows.
If the slow stage is a join, check whether one side is small. A small dimension table can be broadcast to every executor, which removes the shuffle completely.
from pyspark.sql import functions as F
orders_with_region_df = orders_df.join(F.broadcast(regions_df), on="region_id", how="left")Also filter early and select only the columns you need before joining. Less data in means less data shuffled.
repartition that is not needed. Use coalesce when you only want fewer output files.Chapter 10: Performance Optimization goes deep on partitions and shuffles, and Chapter 11: Storage and File Optimization covers file layout, OPTIMIZE and small files.
Adaptive Query Execution (AQE) is on by default in Databricks. It looks at real shuffle sizes while the query runs, and it can split skewed join partitions and combine tiny ones. Make sure nobody has turned it off:
spark.conf.get("spark.sql.adaptive.enabled")
spark.conf.get("spark.sql.adaptive.skewJoin.enabled")AQE helps most with skewed joins. For a skewed aggregation like our example, fixing the key itself is still the real solution.
On a paid workspace with classic clusters, you can also tune shuffle partitions, memory and autoscaling. This requires a paid Databricks workspace. The concept is explained here for understanding. In Free Edition, serverless compute manages this for you, which is one more reason to fix the data first.
1. Problem : 8 min became 42 min. Code unchanged. Suspect the data.
2. Stages : Sort by duration. Look for a big shuffle.
3. Plan : Find the Exchange. Name the wide transformation.
4. Tasks : Compare max and median. Find the straggler.
5. Data : GROUP BY the key, ORDER BY count DESC.
6. Fix : Junk keys, salting, broadcast, less shuffle, AQE.Once you start looking at jobs this way, the Spark UI stops being a monitoring screen. It becomes a debugging tool.
You do not need a 42 minute job to learn this. You can create a small skewed dataset in a few seconds and watch the same pattern.
Open a new notebook in Free Edition and run this:
from pyspark.sql import functions as F
# 5 million rows. About 70 percent share one key, the rest are spread out.
events_df = (
spark.range(0, 5_000_000)
.withColumn(
"customer_id",
F.when(F.rand(seed=42) < 0.7, F.lit("UNKNOWN"))
.otherwise((F.col("id") % 50_000).cast("string"))
)
)
# Step 1: check the key distribution
key_counts_df = events_df.groupBy("customer_id").count().orderBy(F.desc("count"))
display(key_counts_df.limit(5))You should see UNKNOWN far above every other key.
Now look at the plan, and run the aggregation:
customer_totals_df = events_df.groupBy("customer_id").agg(F.count("*").alias("record_count"))
customer_totals_df.explain() # look for Exchange and HashAggregate
display(customer_totals_df.orderBy(F.desc("record_count")).limit(5))In the output of explain(), find the Exchange line. That is the shuffle. Notice there are two HashAggregate steps around it. The first one is the map side partial aggregation from transformation type 7.
Finally, try the separated version from Fix 1 and compare the query profile for both runs. In Free Edition, open the query profile from the cell output to see time and rows per step.
Things to notice:
UNKNOWN change the work?Writing down your three answers is the real exercise. That is how the pattern sticks.
Slow jobs cost money, and they cost trust. A dashboard that is late every Monday slowly teaches people not to rely on it.
But the bigger lesson is a way of thinking. When something breaks, start from what happened, not from what you wrote. Find the stage. Find the task. Then ask the data why.
This same thinking appears across many production issues. Skew, small files, slow joins and shuffles are all part of the 50 problems every data engineer faces. And reading logs and the Spark UI calmly is the focus of Chapter 17: Debugging and Monitoring.
If you are new to Spark, start with the free chapters: Start Here, Spark DataFrames and SQL Query Fundamentals.
If you already run pipelines, pick one slow job from last month and walk it through the six steps above. You will probably find the cause in the Tasks view.
And if you want the full path from your first notebook to production performance patterns, Thinking in Data Engineering with Databricks covers it one practical chapter at a time.