Spark

Mid-level

How do I detect and fix data skew?

Skew is invisible in aggregate metrics and obvious in per-task metrics — you just have to be looking at the right view.

Read in:
Task 1 Task 2 Task 3 Task 4 (the skewed one) Task 5
A skewed stage, illustrated — one task doing far more work than the rest.

A job's average task duration looks fine. Its total duration doesn't. That gap is almost always skew: one or a handful of partitions doing dramatically more work than the rest, with the slowest one dictating the whole stage's wall-clock time.

How to actually see it

Average and even p95 task duration can both look reasonable while skew is still the problem — skew often affects a small enough fraction of tasks that percentile stats hide it. What doesn't hide it: the ratio between max task duration and median task duration within a single stage. A healthy shuffle stage has that ratio well under 3x. A ratio past 10x is skew, not noise.

Fix 1: Salting the skewed key

For a join or aggregation skewed on a specific key value (one customer ID, one region, one status flag dominating volume), salting spreads that key's rows across multiple synthetic sub-keys before the shuffle, then combines results after. This is the most surgical fix when you know which key is the problem, but it requires a code change on both sides of a join.

Fix 2: Adaptive Query Execution (AQE)

If you're on Spark 3.0+, spark.sql.adaptive.enabled=true (on by default in recent versions) includes automatic skew-join handling — Spark detects an overloaded partition at runtime and splits it before the join, no code change required. This is the first thing to check before hand-rolling a salting fix: confirm AQE is actually on and its skew-join threshold isn't set so high that your case doesn't trigger it.

Fix 3: Broadcast the small side instead

If the skew is really a join-size imbalance rather than a within-column skew — one table is small enough to broadcast — switching to a broadcast join sidesteps the shuffle (and the skew that comes with it) entirely. Only viable if the smaller table genuinely fits broadcast-join memory limits; see the OOM article above for what happens when it doesn't.

Where this shows up in opti-pipe: per-task duration variance from your Spark event log is one of the signals the rule engine reads directly — when the max/median task-duration ratio in a shuffle stage crosses the skew threshold, that's flagged as a specific recommendation instead of something you'd have to notice yourself in the Spark UI.

See what this looks like on your own pipeline.

Upload a real Spark event log, dbt run_results.json, or Flink metrics export and get concrete, approve-before-apply recommendations back — not another rule of thumb.