Relational engines can implement grouping and duplicate elimination either by hashing on the grouping key or by sorting the input and collapsing adjacent equal rows. Compare the two strategies and explain when a planner picks each.
answer
- Hash: memory ∝ distinct groups, no order needed
- Sort: memory ∝ sort, output ordered, free if pre-sorted
- Distinct = grouping with no aggregate
- Planner compares estimated groups vs memory budget
- Under-estimated groups → hash agg spill
basics
~20 sHashing builds one in-memory entry per distinct key and needs memory proportional to the number of groups, but no ordering. Sorting needs no per-group memory but costs an O(n log n) sort, yet is free when input is already ordered and it delivers ordered output. Planners choose by estimated group count versus memory budget.
solid answer
~60 s**Hash-based grouping** builds a hash table keyed by the grouping columns, updating an accumulator per group as rows stream in. Cost is roughly one hash and probe per input row; memory scales with the number of **distinct groups** and the size of each group's state. It emits nothing until input ends, and its output has no useful order. If the table exceeds the memory budget it partitions and spills. **Sort-based grouping** sorts the input on the grouping key, then walks it emitting one result per run of equal keys. Memory is O(1) per group; the price is the sort. It shines when the input is *already* ordered — via an index scan, a merge join, or an earlier sort — in which case it is nearly free and streams with tiny memory. Planners choose using the estimated distinct-group count: few groups relative to the memory budget favours hashing; ordered input, very high cardinality (grouping nearly every row), or a downstream requirement for ordered output favours sorting. A badly under-estimated group count is the classic cause of a hash aggregate that spills and collapses.
code
text · 12 lines-- hash strategy
HASH AGGREGATE (group key=customer_id, est_groups=120,000)
-> SEQ SCAN orders (rows=8,000,000)
-- sort strategy (output stays ordered on customer_id)
GROUP AGGREGATE (group key=customer_id)
-> SORT (key=customer_id)
-> SEQ SCAN orders (rows=8,000,000)
-- sort strategy with pre-ordered input: no sort node at all
GROUP AGGREGATE (group key=customer_id)
-> INDEX SCAN orders_customer_id_idx (ordered)go deeper
Know that there are two ways — build a hash table keyed by the grouping columns, or sort and collapse neighbours — and that hashing needs memory for the groups.
Give the memory and CPU drivers for each, and name the conditions that decide: estimated group count, pre-existing order, and whether ordered output is needed downstream.
Discuss the under-estimated-group-count failure, how it shows in a plan as estimated-vs-actual skew with temp-file activity, and the remediation order from fixing statistics to supplying an ordered access path.
Frame it as a memory-versus-order design question across a workload: which indexes exist to supply interesting orders, how much per-operation memory the concurrency budget allows, and whether pre-aggregation upstream removes the choice entirely.
## Two ways to answer the same question Grouping (aggregate per key) and duplicate elimination (distinct) are the same underlying problem: partition rows into equivalence classes by a key and produce one output per class. There are exactly two mainstream physical strategies. ### Hash aggregation Build a hash table whose key is the grouping columns. For each input row: hash the key, look it up, and either create an entry with an initialized accumulator or update the existing accumulator (`count+1`, `sum+=x`, and so on). When the input ends, walk the table and emit one row per entry. - **CPU:** ~O(rows) — one hash and one probe per row, plus the aggregate update. - **Memory:** proportional to **distinct groups × per-group state size**. Not to input rows. Grouping a billion rows into 12 status values costs 12 entries. - **Ordering:** input order irrelevant; output order arbitrary. - **Blocking:** yes — nothing can be emitted until the last row is seen, because a new row could still join any group. - **Overflow:** partition by hash of the key, keep some partitions resident, write the rest to temp files, and process spilled partitions after the input pass. Correct because equal keys always land in the same partition. A sharp edge: aggregates with **unbounded per-group state** — collecting arrays, string aggregation, exact distinct counts inside each group — make each entry large, so the hash table can be huge even at modest group counts. ### Sort-based grouping Sort the input on the grouping key, then scan it once: rows with equal keys are adjacent, so you accumulate until the key changes, emit, and reset. - **CPU:** the sort, ~O(rows · log rows) comparisons, plus a cheap linear pass. - **Memory:** the sort's memory; the grouping step itself holds only the current group's accumulator. - **Ordering:** requires ordered input, produces ordered output — an *interesting order* the optimizer can reuse for a downstream ORDER BY, merge join, or a later grouping on the same prefix. - **Blocking:** the sort blocks; if input already arrives ordered, the grouping step is fully streaming and can emit each group as its key boundary passes. ### The decisive comparison | | Hash | Sort | |---|---|---| | Memory driver | number of distinct groups | rows being sorted | | Cost when input already ordered | same as always | nearly free | | Output order | arbitrary | ordered on the key | | Behaviour with huge group counts | large table, likely spill | sort cost, graceful | | Behaviour with few groups | ideal | wasteful (sorts everything) | | First row latency | after all input | after all input, unless pre-ordered | ## How the planner decides The optimizer estimates the number of distinct groups from column statistics (distinct-value counts, possibly multi-column statistics for composite keys), multiplies by an estimated per-group state size, and compares that with the operator's memory budget. Roughly: - **Estimated groups small relative to the budget** → hash aggregation, because it avoids the sort entirely. - **Input already ordered on the grouping key** (index scan, merge join output, prior sort) → sort/stream aggregation, because the expensive part is already paid. - **Estimated groups approaching the row count** (grouping by something near-unique) → the hash table would be as big as the data; sorting is safer. - **The query needs ordered output anyway** → sorting amortizes across both requirements. - **Some aggregates are not hashable** in a given engine (certain ordered-set or user-defined aggregates) → forced sort path. Duplicate elimination follows exactly the same fork, since distinct is grouping with no aggregate function. ## The failure mode you will be asked about The dangerous case is a **hash aggregate whose group count was under-estimated** — commonly because the grouping key is composite and the columns are correlated, or because statistics are stale after a data shift. The planner sizes the plan for, say, 50,000 groups; reality is 40 million. The hash table blows its budget and either spills heavily or, in engines/versions where the aggregate could not spill, drives memory far past the intended limit. The symptom is a plan whose estimated rows out of the aggregate are orders of magnitude below the actual, with large temp-file activity. Remedies, in order: fix the estimate (refresh statistics, add multi-column statistics for correlated grouping keys); reduce per-group state (avoid array/string accumulation over huge group counts); pre-aggregate or pre-filter so fewer rows and fewer groups reach the operator; provide an ordered input via an index so the engine can stream-aggregate with near-zero memory; and last, raise the memory budget for that specific workload. ## What a strong answer sounds like Say that the two strategies trade **memory proportional to group count** against **an O(n log n) sort plus the gift of ordered output**, that the planner's choice hinges on the estimated distinct-group count versus the memory budget, and that pre-existing order flips the decision decisively toward sorting. Mentioning that hash aggregation's memory depends on distinct groups rather than input rows is the single line that separates candidates who have read plans from those who have not.
- Grouping a billion-row table by a column with only 8 distinct values — which strategy do you expect, and why?Hash aggregation. Its memory need is driven by distinct groups, so eight entries plus accumulators is negligible, and it costs a single streaming pass over the input. Sorting a billion rows to collapse them into eight groups would pay an enormous O(n log n) cost for nothing, unless the input already arrives ordered on that column or the query needs ordered output anyway.
- How does DISTINCT relate to this choice?Duplicate elimination is grouping with no aggregate function, so the same two strategies apply: a hash table keyed by the whole projected row, or a sort followed by collapsing adjacent equal rows. The same trade-offs hold — hashing needs memory proportional to the number of distinct rows, sorting is preferable when the input is already ordered or when the distinct count approaches the row count.
saying these in an interview costs you the question
- Saying hash aggregation's memory scales with input rows rather than with distinct group count
- Claiming sorting is always slower, ignoring that pre-ordered input makes the grouping step nearly free
- Forgetting that sort-based grouping yields an interesting order the optimizer can reuse downstream
- Ignoring per-group state size — array or string accumulation can blow up a hash table at modest group counts
- Treating DISTINCT as a fundamentally different operation from GROUP BY at the execution level