skip to content

Why is exact COUNT(DISTINCT) so expensive across an MPP warehouse, and how does HyperLogLog avoid it?

level: middleimportance: must knowfreq 65%

answer

  1. sums combine locally, distinct counts do not
  2. the same value can live on every worker
  3. memory tracks distinct values, not rows
  4. registers keep a maximum, so duplicates are harmless
  5. error shrinks with the square root of register count

basics

~20 s

Exact distinct counting must bring every copy of a value onto one node and hold a hash table sized by cardinality, so it costs a shuffle plus memory that spills. HyperLogLog builds a fixed-size sketch per node and merges sketches with a per-register maximum.

solid answer

~50 s

A sum is a local operation: each node totals its own rows and the coordinator adds the partials. Distinct counting is not, because the same `user_id` can appear on every node, so the engine must redistribute rows by the distinct key until all copies of a value meet on one worker, then hold a hash table whose size grows with the number of distinct values on that worker. That is a full shuffle plus an operator that spills once cardinality exceeds its memory budget — and a query with several `COUNT(DISTINCT)` over different keys needs the data organised several ways at once. HyperLogLog changes the algebra. Each node hashes its own values into a fixed-size register array locally, no data movement required, and two sketches merge by taking the element-wise maximum of their registers. The merge is commutative and idempotent, so duplicates across nodes take care of themselves and only kilobytes cross the network.

code

sql · 10 lines
sql
-- forces a redistribution by user_id and a cardinality-sized hash table
SELECT country, COUNT(DISTINCT user_id) AS users
FROM events
GROUP BY country;

-- three distinct keys: cannot be satisfied by one redistribution
SELECT COUNT(DISTINCT user_id),
       COUNT(DISTINCT session_id),
       COUNT(DISTINCT device_id)
FROM events;

go deeper

for a junior

Know that counting distinct values is fundamentally harder than summing them because duplicates can be spread across machines, and that engines offer a cheaper estimated version.

for a middle

Explain both halves precisely: the redistribution plus cardinality-sized hash table that exactness demands, and the fixed-size register array whose merge is a maximum. The error formula shape is expected here.

for a senior

Diagnose from a plan or profile — spill on the distinct aggregation, one straggling worker on a skewed key — and choose between raising memory, precomputing sketches and dropping exactness.

for a principal

Frame this as a cost-model decision across the platform: which high-cardinality metrics get a sketch pipeline built once and reused, versus which queries pay for a shuffle every run and why that is acceptable.

## Why distinct is unlike sum Massively parallel engines are fast because most aggregates are **algebraic**: each worker computes a partial result from its own slice, and the partials combine with a cheap function. `SUM` combines by addition, `MIN` by minimum, `COUNT(*)` by addition, `AVG` by keeping a sum and a count. Nothing has to move except the partials. Distinct counting breaks this. Two workers each holding a hundred rows for the same `user_id` cannot produce partial counts that add up, because the correct answer depends on the overlap between their value sets — and neither worker can see the other's set. The count of distinct things is not additive, and no amount of extra partial state short of the sets themselves fixes that. ## What exactness costs An engine has two honest ways to compute an exact distinct count across a cluster. The first is to **redistribute** rows by a hash of the distinct key so that all occurrences of a value land on the same worker. Each worker then deduplicates its slice locally and reports a count; the counts add up because the key spaces are disjoint by construction. The cost is a full shuffle of the column across the network, plus a per-worker hash table whose memory tracks that worker's share of the *cardinality*. When that table exceeds the operator's memory grant it spills to disk, and a scan-bound query becomes an I/O-bound one. If the distinct key is skewed — a handful of values covering a large share of rows — one worker gets a disproportionate slice and the whole query waits on it. The second is to send every distinct value to a single node and deduplicate centrally, which is worse: it serialises the work and concentrates the memory. A further sting: a query with several distinct counts over *different* keys cannot satisfy them with one redistribution. Engines resolve that by repeating the aggregation, or by expanding each row into several tagged rows so all keys can be handled in one pass — either way the work multiplies. ## The HyperLogLog idea Hash each value into a uniformly distributed bit string. Use the first `p` bits to choose one of `m = 2^p` registers, and look at the run of leading zeros in the remaining bits. Under a good hash, seeing a run of `k` zeros has probability `2^-k`, so the longest run observed is evidence of how many distinct items were hashed — see a run of 10 zeros and you have probably hashed on the order of a thousand distinct values. One register would be a wildly noisy estimator, so the sketch keeps `m` of them and combines their evidence with a harmonic mean plus a bias correction. The crucial detail for a distributed engine is what each register stores: the **maximum** run length it has ever seen. Maximum is commutative, associative and idempotent. Hashing the same value a thousand times, on one node or on fifty, leaves the register exactly where a single occurrence would. ## Why it parallelises for free Because of that, every worker can build a complete sketch of its own slice with no coordination, and the sketches merge by taking the element-wise maximum. The merged sketch is bit-for-bit what a single machine would have built from the whole dataset — the distributed answer is not an approximation *of* the single-node answer, it is identical to it. Only the sketches cross the network, a few kilobytes per worker rather than a column of billions of values. There is no shuffle, no spill risk from cardinality, and no skew sensitivity from the distinct key. It also makes the result **deterministic**: because merging is order-independent, the same data returns the same estimate whatever the cluster size or the plan shape. ## Error and precision The relative standard error is approximately `1.04 / sqrt(m)`. With `m = 2^14 = 16384` registers that is roughly 0.8%. Error falls with the square root of the register count, so buying a factor-of-two improvement in accuracy costs a factor of four in sketch size. Engines expose this as a precision or accuracy setting. Because the guarantee is relative, the absolute error scales with the number being estimated: 0.8% of a billion is eight million. ## What this means in practice When a query with `COUNT(DISTINCT high_cardinality_key)` dominates a dashboard's runtime, the fix is rarely a bigger cluster — it is removing the requirement for exactness, or precomputing the sketch once and reusing it. Conversely, a distinct count over a low-cardinality column is cheap even exactly: the hash table is tiny and the shuffle is small, so reaching for a sketch there buys nothing and gives up a guarantee. The diagnostic question is always the cardinality of the counted column and the width of the shuffle it forces, not the number of rows.

  • If HyperLogLog needs no shuffle, why does a grouped approximate distinct count still cost real memory?
    Because you hold one sketch per group. A single sketch is kilobytes, but a `GROUP BY` over a million groups means a million sketches, and the aggregation still has to be partitioned by the grouping key. The saving is on the distinct key's cardinality, not on the number of groups — so a high-cardinality `GROUP BY` remains expensive either way.
  • Why does skew in the distinct key hurt the exact plan but not the sketch plan?
    The exact plan hashes rows by the distinct key, so a value covering a large share of rows sends all those rows to one worker; that worker runs long while the rest idle. The sketch plan never routes by value — each worker summarises whatever slice it already holds — so an uneven value distribution costs nothing extra.
  • Can two sketches built by different systems be merged?
    Only if they agree on the hash function, the precision and the serialised layout. Sketch formats are engine-specific in practice, so a sketch written by one warehouse is generally not readable by another. If you need cross-system merging, standardise on one library's format at the point where the sketches are produced.

Estimating a crowd by asking for the longest run of heads anyone flipped: a very long run is unlikely unless many people flipped, and combining two rooms just means taking the longer run — no need to bring the crowds together.

saying these in an interview costs you the question

  • Says distinct counts add up across nodes like sums do
  • Claims exact distinct memory scales with row count rather than cardinality
  • Thinks HyperLogLog samples rows instead of hashing every one
  • Assumes the merged sketch is less accurate than a single-node one
  • Cannot explain why several distinct keys in one query multiply the work

context