skip to content

One MPP node runs at 100% while the rest idle during a large GROUP BY — what is happening and how do you fix it?

level: seniorimportance: must knowfreq 66%

answer

  1. all rows with one value must land together
  2. look at rows received per node, not total
  3. NULLs and 'unknown' are one value to a hash
  4. runtime equals the slowest worker
  5. split the hot value into synthetic buckets

basics

~20 s

The shuffle key is skewed: a few values, often NULL or a placeholder, own most of the rows, so one node receives that whole bucket. Runtime equals the slowest node. Fix by isolating or salting the hot values, or by pre-aggregating before the shuffle.

solid answer

~50 s

A hash shuffle sends every row with the same key to the same node, so the row distribution across nodes mirrors the *value* distribution of the key. If 40% of rows carry `user_id IS NULL`, or a sentinel like `-1` or `'unknown'`, or one whale tenant, that single bucket lands on one node while the others finish and wait. Confirm it by measuring rows or bytes received per node for the failing stage, then run a `SELECT key, COUNT(*) ... ORDER BY 2 DESC LIMIT 20` to see the offenders. Fixes, in order of preference: filter or exclude junk sentinel values; make sure the engine pre-aggregates locally before the exchange so only partials cross; **salt** the hot key by adding a random bucket to the grouping key and aggregating twice; or split the query, handling hot keys separately with a broadcast join. Adding nodes does not help — the hot bucket is still indivisible.

code

sql · 6 lines
sql
-- 1. Confirm the skew before changing anything
SELECT customer_id, COUNT(*) AS n
FROM orders
GROUP BY customer_id
ORDER BY n DESC
LIMIT 20;

go deeper

for a junior

Recall that rows sharing a key value all go to one machine, so a value that appears far more often than others makes one machine do most of the work.

for a middle

Explain that hash placement distributes values rather than rows, name the usual offenders — NULL, sentinels, whale entities, low cardinality — and describe two-phase aggregation as the first defence.

for a senior

Show the diagnostic path: rows received per worker, then a top-N count by key, then the targeted fix. Distinguish skew from a straggler and justify salting or splitting hot keys.

for a principal

Treat recurring skew as a modelling and data-quality defect, not a query bug: upstream sentinel defaults, grain choices and distribution keys that make skew structural across every workload.

## Why one node ends up with everything A hash exchange assigns a row to a node by `hash(key) mod N`. This is deterministic and, for a uniformly distributed key, near-perfectly even. But it distributes *values*, not rows: all rows sharing a value are inseparable and must land together, because that is the entire point of the shuffle. So the per-node load is the sum of the frequencies of the values that hash to that node. When the key's value distribution is heavily lopsided, the shuffle faithfully reproduces the lopsidedness. A query's stage completes when its slowest worker completes, so a node holding 40% of the rows makes the stage take roughly 40% of the single-threaded time no matter how many nodes are idle beside it. The cluster's aggregate throughput is irrelevant; the critical path runs through one worker. ## The usual culprits - **NULLs.** In many engines all NULLs hash to a single bucket. A left join or grouping on a mostly-NULL column concentrates that entire population on one node. - **Sentinel values.** `-1`, `0`, `'unknown'`, `'N/A'`, the epoch date, or an `unassigned` foreign key inserted by an upstream ETL default. These are semantically "no value" but syntactically one very frequent value. - **Whale entities.** One tenant, one customer, one bot account, one popular product that legitimately generates orders of magnitude more rows than the median. - **Low-cardinality keys.** Grouping or joining on a column with fewer distinct values than the cluster has workers guarantees that most workers get nothing at all. ## Diagnosing it Two measurements settle it. First, from the query profile or per-node execution statistics, look at rows or bytes *received* per worker for the stage that is slow — a healthy hash exchange shows near-identical counts, a skewed one shows one worker one or two orders of magnitude above the median. Second, in SQL, measure the key directly: ```sql SELECT user_id, COUNT(*) AS n FROM events GROUP BY user_id ORDER BY n DESC LIMIT 20; SELECT COUNT(*) FILTER (WHERE user_id IS NULL) AS nulls, COUNT(*) AS total FROM events; ``` If the top value's count is a meaningful fraction of the total, the shuffle cannot be balanced by any placement of that key. Distinguish this from a **straggler**: a straggler is one node slow for an environmental reason — a degraded disk, a noisy co-tenant, a cold cache, an uneven file assignment — and it moves around between runs. Skew is deterministic: the same node (relative to the key's hash) is slow every single run, and the profile shows it *received* more data, not that it processed the same data slower. ## Fixes, roughly in order **Eliminate junk values.** If the hot key is a sentinel that carries no meaning, filter it out before the shuffle or rewrite the join so those rows do not participate. This is the cheapest fix and often the correct data-quality fix as well. **Make sure pre-aggregation is happening.** For a `GROUP BY`, a two-phase plan aggregates on each node before the exchange, so the hot group crosses the network as one partial row per node rather than as billions of raw rows. Skew then costs almost nothing for decomposable aggregates. If the plan is not doing this — often because the aggregate is not decomposable, such as an exact `COUNT(DISTINCT ...)` — that is the real problem to attack. **Salt the key.** Add a synthetic bucket component to the shuffle key so the hot value is split across many workers, then combine the partial results in a second, much smaller aggregation: ```sql SELECT customer_id, SUM(part) AS total FROM ( SELECT customer_id, MOD(order_id, 32) AS salt, SUM(amount) AS part FROM orders GROUP BY customer_id, MOD(order_id, 32) ) s GROUP BY customer_id; ``` Salting turns one hot bucket into 32 warm ones. It costs an extra aggregation stage, so apply it only to the skewed step, and only when pre-aggregation alone is insufficient. **Split hot from cold.** Run the query twice: once for the handful of hot keys, using a strategy that avoids redistributing them (for a join, broadcasting the other side works well because the hot rows then never move), and once for everything else with the normal plan. `UNION ALL` the results. Some engines do a version of this automatically as skew-aware or adaptive join handling. **Change the layout.** If the same key skews every query, the underlying model may be wrong — a composite key, a different grain, or a denormalized table that removes the join entirely may be the durable answer. ## What does not work Adding nodes does not help: the hot bucket remains indivisible and still lands on exactly one worker, so the critical path is unchanged while you pay for more idle machines. Raising the memory limit only converts a spill into a slow in-memory run. Reordering the query or adding hints without addressing the key distribution moves the symptom rather than the cause.

  • How do you tell data skew apart from a straggler node?
    Skew is deterministic and data-shaped: the same relative worker is slow on every run and the profile shows it received far more rows or bytes than its peers. A straggler receives a normal share but processes it slowly, moves between runs, and traces back to the environment — a degraded disk, a cold cache, an unlucky file assignment, or a noisy neighbour on shared hardware.
  • Why is an exact COUNT(DISTINCT ...) especially vulnerable to skew?
    Distinct counting is not decomposable the way SUM or COUNT is: a node cannot pre-aggregate its share into one value without losing the information needed to deduplicate across nodes. Engines therefore shuffle by the distinct column itself, so a lopsided value distribution concentrates work on one node. Approximate sketch-based counting sidesteps this because sketch state merges cheaply.
  • Does salting help a skewed join as well as a skewed GROUP BY?
    Yes, but it costs more. You salt the large skewed side with a random bucket and expand the other side by replicating each of its rows once per bucket, so the join key becomes key-plus-bucket. That multiplies the smaller side's volume by the salt factor, so it is only worth it when the other side is small and the skew is severe.

saying these in an interview costs you the question

  • Suggesting more nodes will spread a single hot key
  • Blaming the hardware before checking rows received per node
  • Assuming a uniform hash function guarantees uniform node load
  • Forgetting NULLs collapse into a single shuffle bucket
  • Raising the memory limit and calling the skew fixed

context