Roadmap To Be A Data Engineer / Lesson 18

Lesson 18 Skew and shuffle Fundamentals §5 Orchestration / execution (+§2, §6) About 20 min read

One Task Left Running

One big seller turns a nine-minute job into four hours.

At 06:40 on Friday the seller report is four hours late, 199 of 200 tasks finished before 02:17, and no one has touched the code since April. The job is not too big for the cluster. It is too lopsided for it.

Model seller_performance_daily
Started 02:15 CEST · normally done 02:24
Window 5 Aug – 3 Sep 2026 · 30 days
Rows read 1,202,993,410
Cluster 64 cores · 200 shuffle partitions
Status 1 task running · 199 finished

106:40, Friday

The nightly model that ranks sellers starts at 02:15 and is normally on the warehouse by 02:24. This morning it is still running. The Spark UI shows one stage, 200 tasks: 199 of them finished within two minutes of starting, and one has been going for four hours and nine minutes.

The three things everyone reaches for first are all checkable, and all three are wrong. Volume is flat — the table took in −0.07% week on week and −0.19% against yesterday. The cluster is the same 64 cores it was on Wednesday. And seller_performance_daily has not been edited since April.

Nothing changed in the job. Two things changed around it, five months apart, and neither of them was wrong.


2One task holds 28% of the work

A shuffle is the moment a distributed engine decides which machine will handle each row. It does it by hashing the key and taking the remainder: hash(seller_id) mod 200. Rows with the same key always get the same answer, which is the entire point — a join or a group-by can only work if every row for a key meets every other row for that key.

It is also the entire problem. A partition can hold many keys, but a key cannot be split across partitions. So the largest key sets a floor under the largest task, and the largest task sets a floor under the whole stage.

Figure 1 — where 1.2 billion rows went when they were hashed into 200 buckets. Partition 170 receives 337,026,699 rows — 28.02% of the stage, 106.1× the median partition (3,175,510 rows) and 56× what a perfectly even split would give it (6,014,967). Partition 25 holds a second large key (110,442,371 rows — the unknown-member row). Every other partition is within a factor of five of the median. Exact — hashed row counts, no model.

Largest partition 337,026,699 rows — 28.02% of the stage
Versus the median 106×
Core utilisation 5.58% of 64 cores while that task runs
Hardware stops helping at 2.6 cores — a 3-core box finishes this stage when a 6,400-core one does

Those last two are arithmetic, not opinion. While the hot task runs, the other 63 cores have nothing left to do: the remaining 199 partitions together hold less work than the one. The stage therefore takes 17.9× longer than the same rows would take spread evenly, and 94.4% of the core-seconds on the invoice buy nothing. Doubling the cluster to 128 cores does not move the finish time by one second; it moves utilisation to 2.79% and the bill to double.

Figure 2 — measured, not modelled. Reproduced on Spark 4.2.0, local[2], 8,000,000 rows, identical totals in both runs. With one key holding 27.74% of the rows the biggest task read 2,248,550 records against a median of 29,455 — 76× — and ran 81× longer than the median task (3,896 ms against 48 ms). Spread the same rows evenly and the worst task is 1.16× the median.

An honest note on that laptop. The skewed run took 10.40 s and the even one 8.32 s — only 25% apart, because with two cores there is almost no parallelism to waste. That is the rule in one line: skew is a tax on the parallelism you have already paid for. It costs nothing on a laptop and 94% of a 64-core cluster.


3Seventy sellers

Spark will not shuffle a table it can simply copy. If one side of a join is small enough, the planner ships a whole copy to every executor and the join happens in place — a broadcast join, with no network partitioning at all. “Small enough” is a single number, spark.sql.autoBroadcastJoinThreshold, and its default has been 10,485,760 bytes for years. We read the constant out of the Spark source at both v3.1.2 and v4.2.0; it is unchanged.

What the planner compares against it is the size of the columns it will actually read, on disk. dim_seller gains about 71 rows a day, because about 71 people a day sign up to sell clothes.

  • Wednesday night: 84,539 rows → 10,482,836 bytes. 2,924 bytes of headroom.
  • Thursday night: 84,609 rows → 10,491,516 bytes. 5,756 bytes over.

Figure 3 — a continuous input crossing a step function. 124 bytes a row is not an estimate: a six-column seller dimension written as Parquet measured 123.74 bytes a row in the lab. The comparison is exact to the byte — holding a table fixed at 10,469,607 bytes and moving the threshold one byte either side flips the plan: threshold = size gives a BroadcastHashJoin, threshold = size − 1 gives a SortMergeJoin.

Losing the broadcast does not make the join a little slower. It inserts a shuffle of the fact table, keyed on seller_id, in front of everything else: 1,202,993,410 rows and 15.96 GB across the network that were not crossing it the night before. And the key it partitions on is the one key in this warehouse that is wildly unevenly distributed.

Figure 4 — the same model, two nights running. Both plans shuffle twice for the count(distinct); only the second one shuffles the fact table itself. That single extra exchange multiplied the largest task by 32.3×, from 10,445,577 rows to 337,026,699.


4Why the same key is harmless in one operator and fatal in another

The old plan was already lopsided. Its worst task took 10,445,577 rows — 10.6× the median — and nobody ever noticed, because ten million rows in one task is a couple of minutes. The house account has been enormous since April. It only became expensive when an operator arrived that could not fold it up first.

Operator What crosses the network for the hot key Rows in one task vs. the join
count(*), sum(price) one partial row per key per map task — the engine adds them up before anything moves ~1,920 173,787× less
count(distinct buyer_id) two shuffles: the first on (seller_id, buyer_id), which spreads perfectly; the second carries one row per distinct pair 9,399,856 35× less
the join on seller_id every single row — a row can only meet its match if it is physically there 333,670,656 —

The difference is whether the operator is algebraic — whether a partial answer can be combined with another partial answer. A count can. A join cannot: there is no partial join. An exact distinct sits in between, which is why Spark plans it as two shuffles rather than one. This is the plan it printed, unedited:

*(3) HashAggregate(keys=[seller_id], functions=[count(1), count(distinct buyer_id)])
+- Exchange hashpartitioning(seller_id, 200), ENSURE_REQUIREMENTS
   +- *(2) HashAggregate(keys=[seller_id], functions=[merge_count(1), partial_count(distinct buyer_id)])
      +- *(2) HashAggregate(keys=[seller_id, buyer_id], functions=[merge_count(1)])
         +- Exchange hashpartitioning(seller_id, buyer_id, 200), ENSURE_REQUIREMENTS
            +- *(1) HashAggregate(keys=[seller_id, buyer_id], functions=[partial_count(1)])

Read it from the bottom. The first aggregate collapses duplicate (seller, buyer) pairs while the rows are still where they were read; the first exchange partitions on the pair, so even a seller with nine million buyers is spread across all 200 partitions; only the second exchange keys on the seller alone, and by then the house account has been reduced from 333,670,656 rows to 9,399,856. The engine is already doing, for the distinct, exactly the thing it cannot do for the join.


5Detected in ten minutes, undiagnosed for a day

This failure is not subtle. A job that normally takes nine minutes taking four hours trips every alarm anyone has ever written. Detection was never the problem. Diagnosis was.

Every number this company stores about that table is a sum or an average: rows per day, GMV per day, cost per day, mean task duration. The quantity that governs the job’s behaviour is a maximum — the largest key. Those two move independently, and for five months they did.

Figure 5 — the metric everyone watches, and the metric that decides how long the job takes. Panel A is a business growing quietly: no step, no spike, nothing on any single day that would make anyone look twice. In panel B the largest single key goes from 8.63% of the day (12 April) to 28.35% (3 September) — and the total is untouched by that, because a row moving from one seller to a bigger one is still one row.

The ninth failure shape in this series — the extremal

Events, constants, drifts, the unwatched, the reversible, the unreproduced, the bundled and the transient all describe failures that are hard to see. This one is loud. It is hard to explain, because the statistic that governs the system is an extreme and every statistic anyone keeps is a total. When behaviour is set by a maximum, storing the mean tells you nothing about it.

The same shape turns up everywhere once you look for it: p99 latency against average latency, the biggest customer against total customers, the largest file in a bucket against the bucket’s size, the one query that scans a year against the median query. In each case the mean is comfortable and the maximum is what breaks.

And it explains why the change review found nothing. Two changes, five months apart, in two different systems, each individually harmless: a house account grew, and a dimension table gained seventy rows. No reviewer of either had any reason to think about the other.


6What the engine will do about this by itself

Adaptive Query Execution watches the real size of each shuffle partition and, when one is much bigger than the rest, splits it into pieces and replicates the matching rows on the other side. It is a real feature that really works. Three things stop it being the answer here.

First, it is switched off. At v3.1.2, spark.sql.adaptive.enabled is createWithDefault(false). At v3.2.0 the same line reads createWithDefault(true). Our runtime image is from 2021. Nobody turned AQE off; the platform turned it on after we stopped upgrading.

Second, the optimiser has two absolute gates. A partition counts as skewed only if it is more than five times the median and bigger than 256 MiB. The factor is relative; the size is not. Any job whose worst partition is under about 22,755,791 rows of this shape is invisible to the feature no matter how lopsided it is.

Gate This job Required Verdict
partition ÷ median partition 106.1× > 5× passes
partition size in bytes 3.98 GB > 256 MiB passes
splitting must not add a shuffle it adds one must not declines
pieces the split can be cut into 60 wanted ≤ map tasks (1,920) available

Third, and this is the one nobody mentions: on this query it declines anyway. We turned AQE on, lowered both gates until the partition qualified several times over, and ran it. Spark left the skew alone. The reason is a fourth setting, spark.sql.adaptive.forceOptimizeSkewedJoin, whose own documentation says “When true, force enable OptimizeSkewedJoin even if it introduces extra shuffle” and whose default is false. Splitting a partition destroys the partitioning the operator above the join depends on, so Spark would have to shuffle again — and by default it would rather leave you with the slow task.

We forced it. The split happened — the plan gained an AQEShuffleRead skewed node, the stage went from 200 tasks to 201, and the hot task’s read halved, from 2,248,550 records to 1,124,491. The job got 3.4× slower (10.62 s → 36.10 s), because of the extra shuffle the default exists to avoid. It also only halved rather than splitting sixty ways, for a reason worth remembering: a skewed partition can only be cut as finely as the number of map tasks that fed it, and our test table was two files.

The lesson is not that AQE is bad — it is excellent, and on a plain two-table join it would have rescued this. The lesson is that an automatic fix is a set of conditions, and you are not entitled to assume you meet them.


7Ten things you could do, and what each is worth

Ranked by what they do to the largest task, which is the only number that moves the finish time. Utilisation is the fraction of 64 cores doing useful work during the worst stage; 78.1% is the best any stage of 200 tasks on 64 cores can do.

What you change Largest task ÷ median Cores working
A As it ran — sort-merge join on seller_id 337,026,699 106.1× 5.58% exact
B Double the cluster (64 → 128 cores) — same finish time, twice the bill 337,026,699 106.1× 2.79% exact
C shuffle.partitions 200 → 2000 — critical path −0.87%, 10× the shuffle files 334,097,283 1523.2× 5.63% exact
D Filter out the unknown member (−8.91% of all rows) — removes 107.2M rows and 0 seconds 337,026,699 106.1× 5.08% exact
E AQE on, Spark defaults — declines: splitting would add a shuffle 337,026,699 106.1× 5.58% measured
F AQE, skew split forced — 3.4× slower in the lab 5,617,112 — — measured
G Broadcast the dimension again — back to last week’s plan, already 10.6× imbalanced 10,445,577 10.6× 37.78% exact
H Salt the hot key 32 ways — needs a replicated dimension and a rewrite 29,937,887 8.9× 57.19% exact
I Aggregate first, join second (exact count-distinct) — the join can never grow back 10,445,577 10.6× 37.78% exact
J Aggregate first + approx_count_distinct — 4.11% error on the hot key 920,931 1.1× 78.12% modelled

Five of the first six rows are what a team reaches for under pressure, and they are worth nothing. The most instructive is D. Dropping the unknown-member rows removes 107,242,026 rows — 8.91% of everything the job reads, rows the dashboard filters out anyway — and changes the finish time by zero, because they were never in the critical partition. It even makes utilisation worse, from 5.58% to 5.08%: less total work, same longest task. That is the extremal shape again, applied to a fix.

G and I land on exactly the same number today and are not the same decision. Raising the broadcast threshold buys time against a table growing 71 rows a day: 32 MiB buys about 7.2 years, 64 MiB about 17.6 — and it applies to every job on the cluster, each of which will now try to copy bigger tables into the driver. Moving the aggregation below the join removes the dependency instead: the join’s left side stops being 1,202,993,410 rows and becomes 84,002 — 14,321× smaller — and no threshold anywhere can make that big again.

Only J removes the skew rather than surviving it, and it costs something real. approx_count_distinct replaces the exact set with a HyperLogLog sketch, which is algebraic, so the query collapses to one shuffle with map-side aggregation. Measured against the exact answer on the same data: mean relative error 1.82% across 60,001 sellers, worst 20% on one of the small ones, and 4.11% on the hot key itself. Fine for an internal ranking; not fine for a number in a seller’s payout statement. That is a business decision, not an engineering one — which is exactly why it should reach you.


8Monday morning

Three questions worth asking the team

  1. “For each nightly job: what is the biggest single key in the table it shuffles, and what share of the table is it?” A ten-second query nobody has run. Above about 5%, that job has a ceiling on how much hardware can help it.
  2. “Which of our joins are broadcast today, and how many rows of headroom does each have left?” The answer is a row count, it only moves one way, and the day it runs out nothing in our code will have changed.
  3. “In each model, does the join happen before or after the aggregation — and could it happen after?” A join above an aggregation carries every raw row past it. A join below one carries a row per group.

The query

select seller_id,
       count(*)                          as rows,
       count(*) / sum(count(*)) over ()  as share_of_window
from   fct_listing_event
where  event_date between date '2026-08-05' and date '2026-09-03'
group  by 1
order  by 2 desc
limit  10;

Run it against every fact table you shuffle, on the key you shuffle it by. It is the cheapest diagnostic in this lesson and the one that would have turned Friday morning into a two-minute conversation.

The instrument worth adding

One line in the run log per job: the largest key’s row count divided by the median key’s (or, if your engine exposes it, the largest task’s shuffle read divided by the median task’s). On this job it was 10.6 for five months and 106 on Thursday. It costs nothing to record and it is the only number in this incident that was actually moving.

Twenty minutes, hands on

  1. Run the generator in the artifact appendix. It prints the lesson’s figures; check a couple against the text.
  2. Change P from 200 to 2000 — the standard folk remedy — and watch the largest partition fall by 0.87%.
  3. Change SALTS from 32 and find where utilisation stops improving. Ask what you would have to write in SQL to make each salt value work.
  4. Delete the two lines that give OUTLET_ID and UNKNOWN_ID their counts, so every seller is ordinary, and see what a job with no skew looks like.
  5. Then run the SQL above against your own largest table, on the key you actually join by.

Takeaway

A distributed job is only as parallel as its largest key. Before you buy hardware, count the keys — and when you find a hot one, ask in this order: does it belong in the query at all; can the operator fold it up before it moves; and can the aggregation happen before the join. Salting is what you do when all three answers are no.

Vocabulary

  • shuffle / exchange — the step where an engine sends rows across the network so every row with the same key ends up on the same machine. Everything expensive in a distributed query is a shuffle.
  • shuffle partition — one of the buckets a shuffle sorts rows into (200 by default in Spark). A partition can hold many keys; a key can never be split across partitions.
  • skew — an uneven distribution of rows across those buckets. Largest partition ÷ median partition; above about 5 is worth a look.
  • straggler — the one task everyone is waiting for. Skew is the most common cause.
  • broadcast join — the small side is copied whole to every machine, so no shuffle is needed. Fast, and gated on a size threshold you did not choose.
  • sort-merge join — the general join: shuffle both sides on the key, sort each partition, walk them together. What you get when the broadcast threshold says no.
  • algebraic aggregate — one whose partial results can be combined (count, sum, min, max, average). Collapsed on each machine before the shuffle, so nearly immune to skew. An exact distinct is not one; a HyperLogLog sketch is.
  • salting — splitting a hot key into key#0 … key#n so it lands in several partitions, and replicating the other side of the join n times to match. Effective, ugly, last resort.
  • AQE (adaptive query execution) — Spark re-planning a query mid-flight using real sizes instead of estimates. It can re-broadcast a join and split a skewed partition — under conditions worth reading before you rely on them.
  • map task / reduce task — the two halves of a shuffle. The number of map tasks limits how finely a hot bucket can ever be subdivided.

9How these numbers were made

  • Exact. Row counts, key shares, partition loads and every ratio derived from them come out of a deterministic generator (modular arithmetic and closed-form formulas, no random seeds), shipped in full in the artifact appendix. The SQL in §08 was executed against the same key distribution in DuckDB and returns the published partition figures exactly.
  • Measured. Plan shapes, task-level shuffle metrics, the broadcast-threshold flip, AQE’s behaviour and the accuracy of approx_count_distinct come from real Spark 4.2.0 runs on 8,000,000 generated rows, read back out of the Spark event log. Configuration defaults were read out of the Apache Spark source at the tagged versions, not from documentation.
  • Observed, and used for nothing. The wall-clock times in the scene — 02:24, four hours nine minutes — are what the incident showed. No calculation depends on them, because a performance model with invented constants would tell you less than the row counts do.

The one modelled figure is row J of the ladder: the partial-aggregate rows a map-side aggregation emits (1,920 map tasks × 84,002 keys), scaled by the task imbalance measured for that plan shape. Distinct buyers per seller use an occupancy model — a pool of 9,400,000 visitors, each looking at 3.2 listings from the same seller in a month — a stated assumption, not a measurement.

The generator was extracted back out of the rendered page, run in a clean directory, and diffed key by key against the published figures: 35 keys, 0 differences.

Back to top