skip to content

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

level: middleimportance: must knowfreq 62%

answer

  1. ask how many rows the aggregate keeps
  2. distinct key combinations, not multiplied cardinalities
  3. one near-unique key destroys the win
  4. columnar scans already skip unreferenced columns
  5. probe with a DISTINCT count first

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.

solid answer

~50 s

The reduction is essentially **base row count divided by the number of distinct combinations of the grouping keys that actually occur** — not the Cartesian product of their cardinalities, because most combinations never appear. Four billion fact rows grouped by day, country and product might collapse to a few million rows, a several-hundred-fold cut in what a dashboard reads. Grain is the whole decision: every extra dimension multiplies the row count, and one high-cardinality dimension (user id, session id, order id) can push the aggregate to within a few percent of the base table, so you pay storage and refresh for nothing. The secondary effects are compression and layout — the aggregate's repeated low-cardinality keys encode densely and sorting it by the usual filter column improves pruning. Measure the collapse with a DISTINCT-count probe before you build anything.

code

sql · 11 lines
sql
-- Size the aggregate before building it
SELECT count(*) AS base_rows FROM fact_events;

SELECT count(*) AS rollup_rows
FROM (
  SELECT DISTINCT
         CAST(event_ts AS DATE) AS event_day,
         country,
         product_id
  FROM fact_events
) t;

go deeper

for a junior

Know what a rollup table is: the stored result of a GROUP BY over a fact table, read instead of the raw rows. Be able to say that fewer rows means less scanned and therefore lower cost.

for a middle

Explain the arithmetic: reduction equals base rows over distinct occurring key combinations, and each added dimension multiplies that denominator. Show that you would measure the collapse with a distinct-count probe rather than guess.

for a senior

Show judgment on grain selection against real query patterns, know that a columnar scan already skips unreferenced columns, and be able to say when the base table's own clustering already prunes as well as the aggregate would.

for a principal

Own the policy: which grains exist, who may add one, and how the savings are attributed. Be ready to argue that a small set of coarse, widely-shared aggregates beats many narrow ones that each serve a single dashboard.

## What a rollup table is A rollup (or aggregate table) physically stores the result of a `GROUP BY` over a fact table so that repeated queries read pre-computed rows instead of re-scanning the raw facts. It does not matter for the economics whether a scheduled pipeline maintains it or you declare it to the engine as a materialized view: you pay compute once per refresh and save on every read. The whole value proposition is a ratio, and a candidate who cannot estimate that ratio before building the table is guessing. ## The arithmetic of the win The first-order reduction is: ``` reduction ≈ base_row_count / distinct_combinations_of_grouping_keys ``` The denominator is the number of key tuples that **actually occur**, not the product of the individual cardinalities. Real fact data is sparse: not every product sells in every country every day. That is why the honest way to size a rollup is to run the probe rather than multiply cardinalities on a whiteboard. Example: a 4-billion-row event fact grouped by `(day, country, product_id)` might produce 9 million rows. That is roughly a 400x collapse, and since analytical engines bill by bytes scanned or by compute time proportional to bytes read, that is close to a 400x cost cut for every query that can be answered at that grain. ## Grain is the entire decision Grain is the set of grouping keys. Coarser grain means fewer rows, and each dimension you add multiplies the count. Two rules follow: - **Choose the coarsest grain that still serves your queries.** Hourly instead of per-event, day instead of hour, country instead of city — each step is often an order of magnitude. - **One high-cardinality key destroys the win.** Adding `user_id`, `session_id` or `order_id` to the grouping keys makes the aggregate approach the row count of the base table, because those keys are close to unique per fact row. At that point you have built a copy of the fact table with extra refresh cost. If a query genuinely needs per-user detail, the aggregate is the wrong tool for it — serve that query from the well-clustered base table and keep the rollup for the coarse questions. ## Why 'fewer columns' is a weaker argument than it sounds People often justify a rollup by saying it carries only a handful of columns while the fact table has eighty. In a **columnar** engine that argument is mostly already priced in: a scan of the fact table reads only the column chunks the query references, so the eighty unread columns cost nothing. The row-count collapse is the real win; the column reduction matters only for the columns the query would have read anyway. What does add on top of the row collapse is **encoding and layout**. Aggregate rows are dominated by repeated low-cardinality keys and numeric measures, which dictionary- and run-length-encode extremely well, so bytes often shrink faster than rows. And because a rollup is small, you can afford to sort or cluster it on the dimension that dashboards filter on, so the engine prunes most of it away before reading. ## When a rollup does not pay - **The query filters on a dimension the rollup does not carry.** There is no way to select the right subset of pre-aggregated rows, so the engine falls back to the fact table and you get zero benefit. - **The measure is not re-aggregatable.** Distinct counts and medians cannot be recombined from group-level results. - **Refresh cost exceeds query savings.** A rollup rebuilt every fifteen minutes to serve a weekly report is a net loss. - **The base table already prunes to the same bytes.** A highly selective query against a table partitioned and clustered on the same columns may read no more than the rollup would. ## Measure before you build Run a distinct-count probe over the proposed key tuple, compare it with the base row count, and then compare actual bytes scanned for your top few dashboard queries with and without the aggregate. After it ships, track how often it is actually used — the estimate is a prediction, the hit rate is the truth.

  • How would you estimate the rollup's size before committing to build it?
    Run a probe that counts distinct combinations of the proposed grouping keys over a representative window of the fact table, and compare that with the base row count for the same window. Extrapolate to the full retention period, then sanity-check bytes rather than rows, since aggregate rows compress differently. If the collapse is under roughly ten times, question whether the grain is coarse enough to be worth maintaining.
  • Why can a rollup that is only five times smaller still be worth building?
    Because savings multiply by read frequency while refresh cost is fixed per period. A five-fold cut on a dashboard refreshed thousands of times a day still dwarfs one nightly rebuild. A small aggregate can also be sorted on the dashboard's filter column and joined to dimensions cheaply, so latency becomes predictable — which is often the real requirement rather than raw cost.
  • Does a rollup remove the need to partition or cluster the fact table?
    No. Queries that need detail, ad-hoc exploration, and the rollup's own refresh all read the fact table, and the refresh is usually the largest single scan you run. Partitioning on the time column is what lets an incremental refresh touch only recent data. Treat the aggregate as an addition to a good physical layout, not a substitute for one.

saying these in an interview costs you the question

  • Assumes any rollup is automatically much smaller than the fact table
  • Adds a user or session identifier to grouping keys and still expects a win
  • Counts saved columns, forgetting columnar scans read only referenced ones
  • Never measures distinct key combinations before building the aggregate
  • Expects the rollup to help queries filtering on a column it does not store

context