Apache Spark interview questionsQuestion 5 of 5
Apache Spark interview question · Question 5 of 5
What is data skew and how can you mitigate it?
Short answer
Data skew means some keys have far more rows than others, so after a shuffle a few partitions are huge and their tasks run much longer while the rest of the cluster waits. I confirm it in the Spark UI by comparing maximum and median task time and shuffle read size in the slow stage. Then I fix it in order: let AQE's skew-join handling split partitions, broadcast the small side if possible, filter or separately process hot or NULL keys, and as a last resort salt the key.
Detailed explanation
Symptoms: one stage where most tasks finish in seconds and a handful take minutes; large shuffle read on those tasks; spills or out-of-memory errors on a single executor.
Common causes: NULL or default keys (-1, "unknown"), a dominant customer or product, low-cardinality join keys.
Fixes
- AQE skew join (on by default) splits oversized partitions in sort-merge joins.
- Broadcast the small side so the skewed large side is never shuffled.
- Handle hot keys separately: filter
NULLkeys that cannot match, or process the top key in its own branch. - Salting: spread a hot key across N sub-keys.
Salting a join
from pyspark.sql import SparkSession, functions as F
spark = SparkSession.builder.master("local[2]").appName("example").getOrCreate()
N = 8
facts = spark.range(0, 50_000).withColumn("key", F.when(F.col("id") % 10 < 8, F.lit("hot")).otherwise(F.col("id").cast("string")))
dims = spark.createDataFrame([("hot", "big customer")], ["key", "label"])
salted_facts = facts.withColumn("salt", (F.rand(seed=1) * N).cast("int"))
salted_dims = dims.crossJoin(spark.range(N).withColumnRenamed("id", "salt"))
joined = salted_facts.join(salted_dims, ["key", "salt"], "left")
print(joined.filter(F.col("label").isNotNull()).count())
40000
The 40,000 “hot” rows are now spread across 8 salt values instead of one partition, and every one still finds its match because the dimension row was replicated for each salt.
Skew in aggregations
For groupBy on a hot key, aggregate in two phases: group by (key, salt) first, then group the partial results by key.
Common mistakes
- Salting without confirming skew in the UI.
- Forgetting to replicate the other side when salting a join, which drops matches.
- Assuming AQE handles aggregation skew.
Progress is saved in this browser only. No account needed.