skip to content

Before a step that sends every record to the worker owning its key, how do you check whether the grouping key's values are evenly spread?

level: middleimportance: should knowfreq 52%

answer

  1. measure the data, not the run
  2. group and count before you submit
  3. same expression, same filters, same range
  4. heaviest key as a share of all rows

basics

~20 s

Count records per key over the input, using the exact key expression and filters the job will apply, then compare the heaviest key's count with the total record count and with the typical key's count.

solid answer

~50 s

I measure the data rather than the run. Group the input by the exact expression the job will group or join on, `count(*)` per key, order descending and take the top few. Two readings come out of it: the heaviest key's count as a **share of all the records**, which is what one destination will receive on its own, and the number of distinct keys, which bounds how many destinations can be busy at all. The expression has to match what the job actually does, including any filter applied before the step — counting a raw column when the job groups on a derived value measures a distribution that will never exist. It is a full pass over the input, so it earns its cost on a job you will run repeatedly or one that has already failed; a sample is a cheaper substitute for finding heavy keys, though not for counting rare ones.

code

sql · 19 lines
sql
-- the heaviest keys, on the same expression and filter the job uses
select lower(trim(account_ref)) as grouping_key,
       count(*)                 as records
from   events
where  event_day = date '2026-09-19'
group  by 1
order  by records desc
fetch first 20 rows only;

-- the two readings that decide it
select sum(records)                     as total_records,
       count(*)                         as distinct_keys,
       max(records)                     as heaviest_key_records,
       max(records) * 1.0 / sum(records) as heaviest_key_share
from ( select lower(trim(account_ref)) as grouping_key,
              count(*)                 as records
       from   events
       where  event_day = date '2026-09-19'
       group  by 1 ) per_key;

go deeper

for a junior

Remember that the check is an ordinary group-and-count over the input: how many records fall under each key, with the biggest counts at the top of the list.

for a middle

Explain the two readings the count produces — the heaviest key's share of all records, and the number of distinct keys — and why the count must use the job's own key expression and filters.

for a senior

Show that you weigh the cost of the extra pass, know that sampling preserves a heavy key's share but not the count of rare keys, and can run the same measurement when the input is files rather than a table.

for a principal

Argue about where the measurement belongs. Counting by hand after each incident is a habit; making the key-frequency reading part of what a pipeline reports about its own input is a policy, and it is the cheaper one at scale.

## Why measure the key before the job runs A step that needs records held by other workers — one that cannot be computed from what a single worker already holds, so every record is first sent to the worker that owns its key — divides its work by key value, not by size. The consequence is arithmetical: if one key value covers a large share of the records, the destination that owns it receives that share on its own, and no amount of extra capacity divides it. Waiting for the run to tell you this costs you the run. The distribution is a property of the input, and the input can be measured directly, before anything is submitted. ## The shape of the count The measurement is an ordinary aggregate: group by the key, count, order by the count descending, look at the head of the list. Two derived numbers matter more than the list itself: - **the heaviest key's share of all records** — `heaviest / total`. This is the fraction of the step's work that will land in one place. A key at a few percent is usually unremarkable; a key at a third of the input is the whole story of the run. - **the number of distinct keys** — how many destinations can receive anything at all. Fewer distinct keys than the step has pieces means pieces that receive nothing, regardless of how evenly the rest is spread. A third, cheaper reading is the ratio of the heaviest key's count to the median key's count, which is the key-frequency equivalent of reading a maximum against a median. ## The expression has to match, exactly This is where the check most often lies to you. Three ways to get it wrong: 1. **Counting the wrong expression.** The job groups on a normalised or derived value — lower-cased, trimmed, concatenated from two columns, bucketed by day — and you counted the raw column. Two different expressions are two different distributions, and a value that is heavy in one can be unremarkable in the other. 2. **Counting the wrong rows.** The job applies a filter before the step. Counting the unfiltered input measures a population the step never sees; a filter that removes most of the ordinary traffic can leave a formerly small key dominant. 3. **Counting the wrong time range.** The job reads one day; you counted the history. Heavy keys are often episodic, and a distribution averaged over a year hides the day that broke. The rule is mechanical: write the count against the same expression, the same predicates and the same range the job will use. ## Cost, and when a sample will do The count is a full pass over the input, which is not free — it is a pass you are buying in order to avoid a much more expensive failed one. That trade pays on a job scheduled to run repeatedly, on a job that has already failed or run long, and before a change to the key a step groups on. It rarely pays for a one-off exploration. Where the full pass is too expensive, sample. The useful property is that **a value covering a large fraction of the rows appears at roughly that fraction in a random sample** — a key holding a third of the input will hold about a third of a one-percent sample too, so heavy keys survive sampling well. The property that does *not* survive is the count of rare keys: most values that appear a handful of times will not appear in the sample at all, so a sampled distinct-key count is far too low. Use a sample to find the heavy end, never to count the tail. ## When the input is not a queryable table Often it is not. The equivalent is a small standalone job over the same files the pipeline reads: project the key expression, count by key, write the top counts and the total. It is one pass over one interval of input rather than the full job, and it produces the same two readings. A cheap alternative when even that is too much: sample a handful of input files and count within them, accepting the tail-counting caveat above. ## What the count does not tell you It tells you the shape of the key-frequency distribution, and that is all. It says nothing about machine health, so a job that also suffers from an unhealthy worker will still be slow for a reason this check cannot see. It does not tell you which remedy applies, or whether one is worth applying. And an even count per key does not by itself guarantee even pieces: the rule that maps a key to a destination distributes many keys well in aggregate, but with few keys relative to pieces, collisions are visible. The count is the strongest single piece of evidence available before a run, not a verdict.

  • The input is not a queryable table — how do you count then?
    Run the count as a small standalone job over the same files the pipeline reads: project the key expression, count by key, and write out the top counts and the total. It is one pass over one interval of input rather than the whole job, and it yields the same two readings.
  • Is a sample good enough for this?
    For the reading that matters, usually. A key covering a large fraction of the rows shows up at roughly that fraction in a random sample, so heavy keys survive sampling. The count of distinct keys does not — most rare values never make it into the sample — so use a sample to find the heavy end and not to measure the tail.

saying these in an interview costs you the question

  • Counts distinct keys and calls that an evenness check.
  • Counts the raw column when the job groups on a derived expression.
  • Counts the whole history when the job reads a single day.
  • Treats even record counts per key as a guarantee of even pieces.
  • Believes the runtime already knows the key distribution before the step runs.