skip to content

Why does grouping a billion rows by a two-value status column leave most of a hundred-worker cluster idle?

level: middleimportance: nice to knowfreq 44%

answer

  1. count the distinct values first
  2. a value cannot occupy two destinations
  3. cardinality caps the parallelism
  4. even spread needs count and shape both

basics

~20 s

The number of distinct key values is a ceiling on how many destinations can receive anything. Two distinct values means two destinations do all the work, however many workers the cluster has and however many pieces the step runs with.

solid answer

~40 s

A step that groups by key sends each value to exactly one destination. With two distinct values there are at most two non-empty destinations, so 98 of 100 workers get nothing to do and each of the two busy ones handles half a billion rows. Raising the number of pieces makes this worse, not better: it creates more empty destinations. Even spread needs two properties of the key together — enough distinct values to occupy the destinations, and a frequency distribution in which no value dominates. Note that the first does not imply the second, and neither is guaranteed by the mapping: it assigns values to destinations, it does not balance them, so with distinct values roughly equal to the destination count you should expect some destinations doubled up and some empty.

go deeper

for a junior

Recall that each distinct key value goes to one destination, so a key with two values can only ever keep two workers busy.

for a middle

Explain why raising the number of pieces adds empty destinations rather than parallelism, and name the two key properties that even spread actually requires.

for a senior

Show that you ask for the distinct count and the commonest value before touching the cluster, and that you can tell an idle cluster from an overloaded single piece by the shape of the durations.

for a principal

Weigh whether the grouping key should be widened at all: a coarser key may be exactly what the business asked for, and buying parallelism by changing it changes the result being reported.

## Two separate properties of the key Work spreads evenly across a step that groups by key only when the key satisfies **both** of these: 1. **Cardinality** — there are at least as many distinct values as there are destinations, so every destination can be given something. 2. **Frequency shape** — no single value, and no small handful of values, accounts for a large share of the rows. These are independent, and an interviewer will happily test either. A two-value status column fails the first spectacularly. A customer identifier with fifty million distinct values passes the first and can still fail the second if one guest-checkout identifier carries a third of the traffic. ## Why cardinality is a hard ceiling A **redistributing step** — a step that cannot be computed from what one worker already holds, so every record is first sent to the worker that owns its key — assigns destinations with a **key-to-destination rule**, the function from a key value to the number of the destination that receives it. The same value always yields the same destination. Therefore the number of destinations that can receive **any** record is at most the number of distinct key values. With eight distinct values and two hundred destinations, at most eight destinations do work and at least 192 do nothing at all. Raising the number of pieces to four hundred changes the figure 192 to 392 and nothing else. The number of pieces controls how finely the *key space* is divided; it cannot divide a key. ## Key-frequency shapes and what each looks like | Key-frequency shape | Destinations doing work | What you see | | --- | --- | --- | | few distinct values, evenly spread | as many as there are values | most of the cluster idle; the busy pieces are all equally slow | | many distinct values, evenly spread | effectively all of them | the healthy case | | many distinct values, one dominant | all, but one is enormous | one piece runs for hours; the rest finish in seconds | | many distinct values, a handful dominant | all, with several enormous | a few long pieces and a long tail of quick ones | Row one is the case in the question, and it is worth noticing that it *is* a heavy-key case: if a column has two values each covering half the data, both of them are heavy keys by definition. The symptom differs from the single-outlier case — several equally slow pieces rather than one — but the cause is identical, and so is the fact that adding machines does nothing. ## Cardinality above the destination count is still not a guarantee The mapping **assigns**, it does not **balance**. It is chosen so that unrelated values tend to land on different destinations, but nothing arranges for equal counts. With exactly as many distinct values as destinations, some destinations will be assigned two values and others none — the same arithmetic that makes a random allocation of a hundred items into a hundred boxes leave roughly a third of the boxes empty. This matters in practice when the key cardinality is only a small multiple of the destination count: a key with 300 distinct values across 200 destinations will not be evenly spread even if every value has the same number of rows. The practical rule of thumb an interviewer likes to hear is that you want distinct values comfortably exceeding the destination count — a multiple, not a margin — *and* a frequency distribution without a dominant value. ## What varies between engines - Whether an unused destination costs anything at all differs. In a finite job, an empty destination is usually a piece that starts and finishes immediately; in a job over an endless input, the width is typically fixed for the life of the job, so an instance that owns no key simply sits there holding resources until the job is restarted. - Where an engine re-decides the part of a plan it has not run yet from the sizes a finished step actually produced, it may merge many small or empty destinations into fewer pieces after the fact. That trims the overhead of the empty ones; it does not create parallelism that the key's cardinality never permitted. ## What an interviewer is listening for That you check the key before you check the cluster. "How many distinct values does the grouping key have, and what does the commonest one cover?" answers both failure modes in one question, and it is the answer that distinguishes someone who has debugged an idle cluster from someone who has only read about parallelism.

  • If the key has ten thousand distinct values and two hundred destinations, is even spread guaranteed?
    No. Cardinality comfortably above the destination count removes the ceiling, but it says nothing about frequency: if one of those ten thousand values carries a third of the rows, its destination still does a third of the work. Both properties have to hold.
  • Does adding pieces ever make an idle cluster worse rather than merely useless?
    Yes. Every piece carries scheduling and bookkeeping overhead, and each one that produces output may write a file. Going from two hundred to two thousand pieces on an eight-value key adds 1,792 empty units and, in a job that writes per piece, a directory full of empty output — for no change in the busy pieces.

saying these in an interview costs you the question

  • Increasing the number of pieces fixes a key with only a few distinct values.
  • A cluster is idle only when the job was given too few pieces.
  • A high distinct-value count guarantees an even spread of work.
  • Equal-sized groups mean the work is well distributed, whatever their number.
  • The key-to-destination rule balances destinations by row count.