skip to content

OLAP Query Patterns

The SQL shapes analytics really runs — huge star-schema joins, windowed metrics, pre-aggregated views and approximate counts — and what each costs on a distributed columnar engine. Interviewers use it to check I can write analytics SQL that scales, not merely SQL that is correct.

on this pageshow

questions

23

How does APPROX_COUNT_DISTINCT differ from COUNT(DISTINCT) in an analytical engine?

level: juniorimportance: must knowfreq 60%

answer

  1. one returns an estimate, the other the truth
  2. memory grows with distinct values, not rows
  3. a fixed-size summary replaces the full value set
  4. hashing plus registers instead of a hash table
  5. typical error is a low single-digit percentage

basics

~20 s

APPROX_COUNT_DISTINCT returns an estimated number of distinct values computed from a small fixed-size sketch, rather than the exact number computed by tracking every value seen. It uses far less memory and time, at a few percent typical error.

solid answer

~50 s

`COUNT(DISTINCT col)` has to remember every distinct value it has seen, so its memory grows with the cardinality of the column and, in a distributed engine, all copies of a value must be brought together before they can be deduplicated. `APPROX_COUNT_DISTINCT` (most analytical engines expose it under that name or a close variant) instead hashes each value into a small fixed-size probabilistic sketch — usually HyperLogLog — and estimates the cardinality from the shape of that sketch. The sketch is a few kilobytes regardless of whether the column holds a thousand or a billion distinct values, so the query needs no large hash table and no reshuffling of the data. The price is that the answer is an estimate, typically within a low single-digit percentage of the truth. Use it for dashboards, trends and exploration; use the exact function when the number is billed, audited or reconciled.

code

sql · 9 lines
sql
-- exact: memory grows with the number of distinct user_ids
SELECT event_date, COUNT(DISTINCT user_id) AS dau
FROM events
GROUP BY event_date;

-- approximate: fixed-size sketch per group, a few percent error
SELECT event_date, APPROX_COUNT_DISTINCT(user_id) AS dau_est
FROM events
GROUP BY event_date;

go deeper

for a junior

Be ready to say plainly that one function is exact and the other is an estimate from a compact sketch, and that the estimate is normally within a few percent. Naming HyperLogLog as the usual mechanism is a bonus.

for a middle

Explain why the exact version is expensive — it must remember every distinct value, so memory tracks cardinality — and why the sketch is a fixed size no matter how many distinct values arrive.

for a senior

Show judgment about which reported metrics may be approximate. Interviewers want to hear you separate indicator numbers from billed or audited numbers, and hear you name the consistency oddities estimates introduce.

for a principal

Own this as a platform policy rather than a per-query choice: which metrics in the semantic layer are declared approximate, how the error bound is published to consumers, and why a metric must never silently switch implementations.

## The two functions side by side `COUNT(DISTINCT user_id)` is defined to return the exact number of different non-null values in the column. To honour that definition, an engine must be able to answer "have I seen this value already?" for every row, which means it must keep the set of values it has already seen. The natural implementation is a hash table keyed by the value. Its size is driven by the **cardinality** of the column — the number of distinct values — not by the number of rows. Ten billion rows over fifty distinct countries costs almost nothing; ten billion rows over five hundred million user ids costs a large hash table. `APPROX_COUNT_DISTINCT(user_id)` answers a deliberately weaker question: roughly how many distinct values are there? It replaces the exact set with a **sketch** — a small, fixed-size summary that supports the operation you actually need (an estimate of cardinality) while throwing away the ability to answer membership. The near-universal choice is HyperLogLog. ## What the sketch actually stores HyperLogLog hashes each input value to a uniformly distributed bit string. A fixed number of leading bits selects one of `m` registers; the remainder is inspected for its run of leading zeros, and the register keeps the maximum run length it has ever seen. A long run of zeros is rare, so seeing one is evidence that many distinct values were hashed. Averaging that evidence across all `m` registers (with a harmonic mean and a bias correction) yields the estimate. Two consequences follow directly from this design. First, the state is **fixed size**: `m` small registers, typically a few to a few tens of kilobytes, no matter how many distinct values arrive. Second, feeding the same value in twice changes nothing, because the register only ever takes a maximum — the sketch is naturally idempotent, which is exactly what deduplication needs. ## The cost difference in practice In a single-node engine the difference shows up as memory: the exact hash table for hundreds of millions of ids can exceed the operator's memory budget and spill to disk, turning a scan-bound query into an I/O-bound one. In a distributed engine the difference is larger still, because exactness forces data movement — every occurrence of a given value has to reach the same worker before duplicates can be collapsed. The approximate version builds one small sketch per worker locally, and the workers exchange only the sketches. ## The accuracy you can expect HyperLogLog's relative standard error is about `1.04 / sqrt(m)`. With 16384 registers that is roughly 0.8%: on a true value of 50,000,000 you would typically land within a few hundred thousand, and occasionally further. Most engines expose a precision or accuracy parameter that increases `m` — the error shrinks with the square root of the sketch size, so halving the error costs four times the memory. Note that this is a *relative* error: the absolute miss grows with the number you are estimating. Many implementations also switch to an exact or sparse representation at low cardinality, so small groups often come back exactly right. That is an implementation convenience, not a guarantee you should build on. ## Where it belongs and where it does not Approximate distinct counts are right whenever the decision the number feeds is insensitive to a percent or two: daily and monthly active users on a trend chart, unique visitors per page, distinct error codes in a monitoring view, cardinality checks while exploring a new dataset, funnel steps. They are wrong whenever the number is the product rather than an indicator: anything invoiced, anything reported to a regulator, anything reconciled against an upstream system of record, and any threshold decision where the boundary sits inside the error band. One subtle trap: because each estimate carries its own independent error, two separately estimated results need not be consistent with each other. A filtered subset can come back with a slightly larger estimate than the unfiltered table, which looks like a data bug and is not. ## The rest of the family The same trade appears for other aggregates that are expensive because they need to remember detail. Approximate percentile functions replace a full sort with a quantile summary. Approximate top-K functions replace a full frequency table with a heavy-hitters sketch. In every case the question to ask is the same one: is this number a decision input where a small, bounded error is invisible, or is it the answer itself?

  • Does the approximate function get slower as the number of distinct values grows?
    Barely. The per-row work is a hash plus a register update regardless of cardinality, and the sketch stays the same size, so runtime tracks the number of rows scanned rather than the number of distinct values. The exact version is the one that degrades with cardinality, because its hash table grows and eventually spills.
  • If a group only has a few hundred distinct values, is the approximate function still worth using?
    Not really — at low cardinality the exact count is cheap, and many implementations fall back to an exact or sparse representation anyway, so you gain nothing and give up a guarantee. The approximation earns its keep at high cardinality: millions or billions of distinct values, especially when many groups are computed at once.
  • Can you make the approximate answer more accurate?
    Yes, by raising the precision or accuracy parameter, which increases the number of registers in the sketch. Error falls with the square root of the register count, so cutting the error in half costs roughly four times the sketch memory. Beyond a point, computing exactly is cheaper than chasing accuracy.

Counting distinct values exactly is like keeping a guest list at the door; the approximate version is like glancing at the crowd and estimating attendance — far cheaper, close enough to plan for, useless for billing each guest.

saying these in an interview costs you the question

  • Claims the approximate function samples rows instead of hashing all of them
  • Says the error is unbounded or unpredictable rather than a known relative error
  • Thinks approximate memory grows with the number of distinct values
  • Uses the approximate count for invoiced or audited numbers
  • Believes the two functions differ only in speed, not in the answer

context

open as a page

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

level: middleimportance: must knowfreq 65%

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.

open as a page

Why can a rollup table answer SUM queries directly but not AVG or COUNT(DISTINCT user_id)?

level: middleimportance: must knowfreq 58%

basics

~20 s

SUM, COUNT, MIN and MAX re-aggregate correctly because combining group results is the same operation again. AVG is a ratio, so store SUM and COUNT and divide at read time. Distinct counts and medians cannot be recombined at all from stored group results.

open as a page

What determines how much a pre-aggregated rollup table reduces the data an analytical query scans?

level: middleimportance: must knowfreq 62%

basics

~20 s

A rollup's win is base rows divided by the number of distinct grouping-key combinations that actually occur in the data. Coarse grain and low-cardinality keys win most; adding a near-unique key such as session_id saves almost nothing.

open as a page

When joining a huge fact table to a small dimension, why does an MPP engine broadcast the dimension?

level: middleimportance: must knowfreq 72%

basics

~20 s

Broadcasting sends one full copy of the small dimension to every node so each node joins its local fact rows without moving them. The fact table is orders of magnitude larger, so shipping the dimension costs far less network traffic than redistributing both sides.

open as a page

Why does PARTITION BY in a window function force a shuffle in an MPP warehouse?

level: middleimportance: must knowfreq 70%

basics

~20 s

A window function must see all rows of a partition together, so an MPP engine hashes the PARTITION BY key and redistributes every row across nodes — a full network shuffle — then sorts each partition locally before computing.

open as a page

A star join on a 6 TB fact table runs 40x slower after a dimension grew. How do you diagnose it?

level: seniorimportance: must knowfreq 62%

basics

~20 s

Compare estimated versus actual rows on the join's build side and check whether the plan still broadcasts. A dimension that outgrew the broadcast threshold either replicates far too much data to every node or flips to redistribution, where a skewed join key can pin the work onto one node.

open as a page

A window partitioned by tenant_id pins one node at 100% while others idle — what is happening?

level: seniorimportance: must knowfreq 62%

basics

~20 s

Partition-key skew. Rows are hashed by tenant_id, so one huge tenant lands entirely on one worker, which must sort and buffer that whole partition alone. It spills and becomes a straggler while peers finish early and wait.

open as a page

In a columnar warehouse, why can't you scale a 1% TABLESAMPLE distinct count up by 100?

level: middleimportance: should knowfreq 45%

basics

~20 s

Distinct counts do not scale with the sample rate. If values repeat, a 1% sample already sees nearly all of them and multiplying by 100 overcounts massively; if values are unique the raw sample count undercounts. The correct factor depends on an unknown frequency distribution.

open as a page

Why does joining a partitioned fact table to a filtered date dimension often still scan every partition?

level: middleimportance: should knowfreq 50%

basics

~20 s

The partition column is filtered only indirectly, through the join, so at planning time no literal date range is known. Unless the engine defers pruning to runtime using the dimension's surviving keys, it must assume every partition might contain matches.

open as a page

Why does a window function with no PARTITION BY collapse onto one worker?

level: middleimportance: should knowfreq 52%

basics

~20 s

With no PARTITION BY there is exactly one partition, so a distributed engine must gather every qualifying row onto a single worker to order and number them. Parallelism drops to one node, which then sorts the entire result and often spills.

open as a page

Why can't daily distinct-user counts be summed, and what should a rollup table store instead?

level: seniorimportance: should knowfreq 50%

basics

~20 s

A user active on several days is counted once per day, so summing daily distinct counts overcounts a longer window and the counts alone cannot be corrected. Store the serialized distinct-count sketch per day and merge sketches at query time.

open as a page

What breaks in an incrementally refreshed rollup when a backfill rewrites last month's fact rows?

level: seniorimportance: should knowfreq 48%

basics

~20 s

Incremental refresh only adds contributions from rows above a watermark, so it never sees a rewrite of old rows and cannot subtract the old contributions. The aggregate keeps the pre-backfill numbers and drifts silently until affected partitions are recomputed or the whole thing is rebuilt.

open as a page

Which conditions let a materialized aggregate answer a warehouse query written against the fact table?

level: seniorimportance: should knowfreq 52%

basics

~20 s

The aggregate's grain must be at least as coarse as the query's and derivable from it, every column the query filters or groups on must exist in the aggregate, its joins and filters must subsume the query's, its measures must be re-aggregatable, and it must be fresh enough for the engine to trust.

open as a page

What is a runtime (bloom) filter in a star join, and how does it cut the fact-table scan?

level: seniorimportance: should knowfreq 55%

basics

~20 s

A runtime filter is built from the join keys surviving the dimension's filter, then pushed into the fact scan so blocks and rows that cannot match are skipped. It turns a dimension-side predicate into pruning on the fact table, which has no such predicate of its own.

open as a page

When is denormalizing a star schema into one wide table worth it in a columnar warehouse?

level: seniorimportance: should knowfreq 58%

basics

~20 s

Denormalizing pays when joins dominate query time and the flattened attributes are small, low-cardinality and rarely restated. Columnar storage compresses the repeated dimension columns heavily and unread columns cost nothing, so the space penalty is far smaller than row-store intuition suggests.

open as a page

Why does a top-N-per-group window function cost more than the N rows it returns?

level: seniorimportance: should knowfreq 56%

basics

~20 s

Because the ranking is computed before the filter. The engine shuffles every row to its partition's worker and fully sorts each partition, then throws away everything except the top N — so the work scales with the table, not with the answer.

open as a page

In a columnar warehouse, what makes a window function spill to disk?

level: seniorimportance: should knowfreq 48%

basics

~20 s

The sort and any row buffering the frame requires. After the shuffle each worker must sort its partitions, carrying every projected column; when a partition plus its frame buffer exceeds the operator's memory budget, the engine writes runs to local disk and merges them back.

open as a page

How do you decide which warehouse metrics may use approximate distinct counts and which may not?

level: principalimportance: should knowfreq 38%

basics

~20 s

Decide by what consumes the number, not by what it costs. Money, regulatory reporting, anything reconciled against a system of record, and threshold decisions inside the error band require exactness; indicators, trends and exploration do not. Declare the choice in the metric layer, not per query.

open as a page

Forty materialized aggregates raised the warehouse bill instead of cutting it — how do you decide which to keep?

level: principalimportance: should knowfreq 38%

basics

~20 s

Build a ledger per aggregate: query cost avoided across its actual hits versus its refresh cost times refresh frequency, plus storage. Kill the zero-hit ones, slow the over-refreshed ones, and consolidate near-duplicates onto the coarsest grain several consumers can share.

open as a page

Why can an approximate percentile over unchanged data return a slightly different p99 each run?

level: seniorimportance: nice to knowfreq 32%

basics

~20 s

Approximate percentile functions build a compact quantile summary whose retained points depend on the order values arrive and how partial summaries merge. Parallel execution varies that order between runs, so the compacted summary and its estimate shift slightly.

open as a page

How do you serve sub-minute-fresh results from a rollup table that only refreshes hourly?

level: seniorimportance: nice to knowfreq 30%

basics

~20 s

Union the rollup for all closed periods with an on-the-fly aggregate of the fact table for the open tail, splitting on an exclusive boundary so nothing is counted twice, then re-aggregate the union. It works only for additive measures and only if the tail scan prunes to recent partitions.

open as a page

A window running-total over full history reruns nightly; how do you make cost grow with new data only?

level: principalimportance: nice to knowfreq 33%

basics

~20 s

Checkpoint the state. Persist the last cumulative value per key at a watermark, then run the window only over rows newer than it and offset each partition by its stored seed. Cost then tracks the daily delta instead of all history.

open as a page