One task ran for forty minutes. Data skew in Databricks, explained

How to find a skewed key, and the five fixes to try in order

The job used to finish in nine minutes. On Monday it took fifty-one.

Nothing had changed in the code. No new columns, no new logic, no new tables. The only thing that changed was the data. One customer had started sending far more events than everyone else.

Open the Spark UI on a run like that and you see the same picture every time. Two hundred tasks in the stage. A hundred and ninety-nine of them finish in seconds. One of them sits there, alone, for forty minutes, chewing through gigabytes while every other core on the cluster does nothing.

That is data skew. It is one of the most common reasons a healthy pipeline slowly turns into an expensive one, and it is one of the few performance problems that adding a bigger cluster does not fix.

Why one task gets all the work

Spark does most of its heavy work by shuffling. When you join two tables, or group by a column, Spark has to bring all rows with the same key to the same place. It hashes the key, and every row with that hash lands in the same partition, and that partition is processed by exactly one task.

That design is what makes Spark scale. It is also the source of the problem.

If your keys are spread evenly, each task gets a similar slice of the data and they all finish around the same time. If one key holds a large share of the rows, the task that owns that key has to process all of it by itself.

A stage finishes when its slowest task finishes. One heavy key is enough to make a hundred idle cores wait.

Real data is rarely even. A few examples you will meet in practice:

One enterprise customer sends more traffic than the next five hundred combined. A null in the join key stands in for "unknown", and half the rows are unknown. A default country code gets used whenever the app cannot detect location. A single popular product appears in most orders. A backfill loads one very large day alongside dozens of small ones.

None of these are mistakes. They are just how the world looks when you write it down.

First, prove it before you fix it

Skew is easy to guess at and easy to get wrong. Before changing any code, look at the data.

The cheapest check is a count by key, sorted descending. In Databricks Free Edition you can run this on any table you already have.

from pyspark.sql import functions as F

orders_df = spark.read.table("workspace.default.orders")

key_counts_df = (
    orders_df
    .groupBy("customer_id")
    .agg(F.count("*").alias("row_count"))
    .orderBy(F.desc("row_count"))
)

key_counts_df.show(10, truncate=False)

Then compare the top key against the typical key. If the largest key holds many times more rows than the median key, you have skew on that column.

stats_df = key_counts_df.agg(
    F.max("row_count").alias("largest_key"),
    F.expr("percentile_approx(row_count, 0.5)").alias("median_key"),
    F.count("*").alias("distinct_keys"),
    F.sum("row_count").alias("total_rows"),
)

stats_df.show(truncate=False)

A ratio of two or three between largest and median is normal. A ratio of a few hundred is a slow task waiting to happen.

Also check nulls separately, because they are the most common hidden culprit.

null_share_df = orders_df.select(
    F.count("*").alias("total_rows"),
    F.sum(F.when(F.col("customer_id").isNull(), 1).otherwise(0)).alias("null_keys"),
)

null_share_df.show(truncate=False)

The second place to look is the Spark UI. Open the job, click into the slow stage, and read the task summary table. When the max duration and max shuffle read are far larger than the 75th percentile, the stage is skewed. This is exactly the kind of reading we practise in Debugging and Monitoring, and it is the difference between fixing a real problem and guessing.

What to actually do about it

There is no single fix. There is an order to try things in, from cheapest to most invasive.

1. Let adaptive query execution handle it

Spark's adaptive query execution can detect a skewed partition at runtime and split it into smaller pieces automatically. On modern Databricks runtimes this is on by default, and it quietly solves a large share of ordinary skew.

print(spark.conf.get("spark.sql.adaptive.enabled"))
print(spark.conf.get("spark.sql.adaptive.skewJoin.enabled"))

If both are true and the stage is still lopsided, adaptive execution is not enough for your case. That usually means the skew is inside an aggregation rather than a join, or a single key is so large that even split pieces are heavy.

2. Remove the keys that should not be joining at all

Half of the skew I have seen in real pipelines was never real. It was null joining to null, or a placeholder value like -1 or UNKNOWN matching thousands of rows on both sides.

Filter those rows out of the join, handle them explicitly, and bring them back afterwards.

known_orders_df = orders_df.filter(F.col("customer_id").isNotNull())
unknown_orders_df = orders_df.filter(F.col("customer_id").isNull())

joined_df = known_orders_df.join(customers_df, on="customer_id", how="left")

Now the join does real work only, and the unknown rows are visible instead of hiding inside a slow task. That visibility matters more than the speed. A silent placeholder join is one of the ways a pipeline turns green while the numbers go wrong, which is the failure mode we walked through in The pipeline was green. The numbers were wrong.

3. Broadcast the small side

If one side of the join is small, do not shuffle at all. Send a copy of the small table to every executor and join in memory.

from pyspark.sql.functions import broadcast

enriched_df = orders_df.join(broadcast(country_lookup_df), on="country_code", how="left")

Broadcast joins remove the shuffle, and no shuffle means no skewed partition. The limit is size. A dimension table of a few hundred megabytes is fine. A fact table is not. If you are unsure whether a join should be broadcast, Joins and Aggregations walks through how to reason about the two sides.

4. Salt the hot key

When one key genuinely holds a huge share of the data, split it artificially. Add a small random number to the key on the large side, expand the small side to cover every salt value, join on the combined key, then aggregate away the salt.

SALT_BUCKETS = 16

salted_orders_df = orders_df.withColumn(
    "salt", (F.rand() * SALT_BUCKETS).cast("int")
)

expanded_customers_df = (
    customers_df
    .withColumn("salt", F.explode(F.sequence(F.lit(0), F.lit(SALT_BUCKETS - 1))))
)

joined_df = salted_orders_df.join(
    expanded_customers_df,
    on=["customer_id", "salt"],
    how="left",
).drop("salt")

Salting works, and it is ugly. It duplicates the small side sixteen times and it puts a performance trick into your business logic. Use it when the simpler options have failed, comment why it exists, and check whether it is still needed six months later.

5. Pre-aggregate before you join

Often the skew is not needed at all. If you are joining a large event table only to count something per customer, aggregate first and join the small result.

order_totals_df = (
    orders_df
    .groupBy("customer_id")
    .agg(F.sum("amount").alias("total_amount"))
)

customer_summary_df = customers_df.join(order_totals_df, on="customer_id", how="left")

The heavy key still exists in the aggregation, but the aggregation is a much cheaper operation than a wide join, and the join afterwards runs on one row per customer. Thinking about the shape of the data before the shape of the code is the habit we build in Transformations.

Skew in the way the table is stored

Some skew lives on disk rather than in the shuffle. If you partitioned a table by a column where one value dominates, that one directory grows enormous while the others stay tiny. Every read of that value hits one huge set of files.

This is one of several reasons we default to liquid clustering instead of hand-picked partition columns, which we compared in Z-ORDER or liquid clustering? Clustering adapts to the data as it changes. A partition column decided a year ago does not.

Storage skew also has a cousin at the other extreme: thousands of tiny files instead of one giant one. Both are layout problems, and both show up as a job that got slower without the code changing. Partitioning and Performance covers how to think about layout before you reach for tuning knobs.

What good looks like

You have handled skew when three things are true.

The task durations inside your heavy stages look similar to each other, so the cluster finishes together instead of waiting on one core. You know which keys are large in your own data, because you measured them rather than assumed. And the fixes you applied are the simplest ones that worked, with a note explaining any that are not obvious.

One more habit is worth building. Skew changes over time. A key that was ordinary last quarter can become dominant after one big customer signs. Add a small check that records your top key counts on every run, and you will see the change coming instead of discovering it in a fifty-one minute job.

Continue learning

If this is your first time reading a Spark UI, start with Debugging and Monitoring, then Joins and Aggregations for the join mechanics, then Partitioning and Performance for layout.

If you want to feel this rather than read about it, the guided first project gives you deliberately messy data in Databricks Free Edition. Group it, join it, and watch which tasks take the longest.

And if you are hardening a pipeline for production, the companion pieces are safe reruns and backfills without breaking downstream tables. Skew, reruns, and history are the three things that decide whether a pipeline survives its second year.