skip to content

Which aggregates return exactly the same answer after a heavy key is spread over suffixes and recombined, and which do not?

level: middleimportance: must knowfreq 52%

answer

  1. not every aggregate survives two passes
  2. the same operation folds the partials
  3. carry sum and count, divide last
  4. distinct values overlap across random suffixes
  5. suffix by the value being counted

basics

~20 s

Aggregates whose partials combine with the same operation - sum, count, minimum, maximum - come back identical. A mean survives only if each partial carries a sum and a count. Exact distinct counts and percentiles do not combine that way at all.

solid answer

~50 s

Spreading a heavy key means appending a small integer to it so one value's records reach several destinations, then folding the partial results under the original key. It is only legitimate if that fold rebuilds exactly the single-pass answer. Three classes. **Directly decomposable**: sum, count, minimum, maximum - the second pass applies the same operation to the partials, which is all associativity with an identity buys you. **Decomposable with a richer partial**: a mean carried as a sum and a count; a variance carried as a count, a sum and a sum of squares; any ratio carried as its numerator and denominator, divided once at the end. **Not decomposable by a random suffix**: an exact distinct count, because the same value can appear under several suffixes and the partial counts overlap; and a median or percentile, because there is no combining operation over partial medians.

go deeper

for a junior

Remember that sums, counts, minimums and maximums come back unchanged from a two-pass fold, and that a mean needs its sum and its count carried separately rather than averaged.

for a middle

Explain why unequal group sizes make averaging averages wrong, what each partial must carry for a ratio or a variance, and why overlapping values break an added distinct count.

for a senior

Show that the partial's size matters as much as its combinability, and pick the spread to suit the aggregate - by the counted value for exact distinct counts, only where those values are varied.

for a principal

Decide what the team accepts when an aggregate does not decompose: a slow exact run, a different formulation of the metric, or a documented approximation with a stated error, and who owns that choice.

## What the rewrite is asking of the aggregate A grouping cannot be computed from what one worker (one process on one machine) already holds, so it opens with a **redistributing step**: every record is sent to the destination that owns its key, chosen by a function of the key value alone. A **heavy key** - one value covering a large share of the records - therefore lands whole on one destination. The standard remedy is to append a small integer suffix to that key so its records reach several destinations, aggregate them there, and fold the partial results under the original key in a second pass. That is a rewrite of the job, not of the data, so it carries an obligation: the second pass must reconstruct exactly the answer the single-pass grouping would have produced. Two conditions decide whether it can. 1. **A combining operation exists.** Pass two must fold up to `W` partials into one value, so the fold has to be associative and have an identity - the general law governing any parallel reduction. 2. **The partial is small.** If a partial is as large as the group it summarises, pass two rebuilds the heavy group on one worker and the imbalance returns at the combining step. This is the condition people forget. ## The three classes | aggregate | partial per `(key, suffix)` | pass two | identical answer? | |---|---|---|---| | sum, count | partial sum, partial count | add | yes | | minimum, maximum | the partial extreme | extreme of the extremes | yes | | mean | sum and count | add both, divide once | yes | | variance, deviation | count, sum, sum of squares | add the three, finish once | yes, with a numerical caveat | | latest value per key | value plus its ordering field | keep the extreme ordering field | yes, if ties break deterministically | | exact distinct count, random suffix | the distinct values themselves | union, then size | only with a partial as big as the group | | exact distinct count, suffix from the value | a partial distinct count | add | yes | | median, percentile | no summary that combines | no combining operation | no | | collect the group's records | the partial list | concatenate | yes, but the partial is the whole group | ## The two failures worth being able to explain - **Averaging the averages.** Each suffix holds a different number of records, so the partial means carry different weights. Averaging them weights every destination equally and the answer is wrong by however unevenly the spread fell. The repair is to carry the numerator and the denominator separately and divide once, which also makes the aggregate associative in the process. - **Adding distinct counts.** If the suffix is random, one value under the heavy key can appear beneath several suffixes, so the partial counts overlap and adding them over-counts. Carrying the actual distinct values fixes the arithmetic and destroys the point of the rewrite, because the set that has to be rebuilt per key is the very thing making the key heavy. ## The distinct count has a different spread For an exact distinct count, do not suffix at random - derive the suffix from the value being counted. Each distinct value then reaches exactly one destination, the partial counts are disjoint, and they simply add. The heavy key's records still scatter, because it is the varied values beneath it that drive the spread. This carries its own condition, and it is the honest limit of the trick: it works when the heavy key covers many different values. If the key is heavy because one value repeats millions of times, spreading by value sends all those repeats to one destination again and nothing has been gained. The measurement to take before choosing is not just records per key but distinct values per key. ## Deviations, ordering, and what nothing checks A variance combines from a count, a sum and a sum of squares, which is exact in arithmetic but can lose precision when the values are large and the deviations small; pairwise combination of a running mean and a running sum of squared deviations is the sturdier form and combines just as well across suffixes. Aggregates that depend on the order records are seen in were already unreliable: a grouping makes no promise about arrival order at a destination, on any engine of this class. The rewrite does not cause that, it exposes it - if the answer changed when the key was spread, order-dependence was there before. Finally, the rewrite is a hand transformation. No runtime checks that the two-pass form still returns the original answer, so the check belongs to the author: run both forms over a slice small enough to finish quickly and compare per key, including the keys that were never spread.

  • How do you make an exact per-key distinct count work under this rewrite?
    Derive the suffix from the value being counted rather than at random. Each distinct value then reaches exactly one destination, so the partial counts are disjoint and add. It only helps when the heavy key covers many different values; if it is heavy because one value repeats, those repeats still converge on one destination.
  • What must a partial carry for a variance or a standard deviation?
    Either a count, a sum and a sum of squares, added component-wise and finished once, or a count with a running mean and a running sum of squared deviations combined pairwise. The second is less prone to precision loss when the values are large and the spread around the mean is small.
  • Collecting every record of a group into one list combines by concatenation, so why is it a poor fit?
    Because the combining operation is not the problem - the size of the partial is. Concatenation is perfectly associative, but the second pass rebuilds the heavy key's entire collection on one worker, so the imbalance simply moves from the first pass to the second and the rewrite buys nothing.

saying these in an interview costs you the question

  • Average the per-suffix averages to get the key's mean
  • Add the per-suffix distinct counts for an exact total
  • Take the median of the partial medians
  • Any aggregate works as long as you group twice
  • Close enough counts as equivalent after the rewrite
  • Carrying the values themselves makes distinct counting cheap again