Why does a top-N-per-group window function cost more than the N rows it returns?
answer
- ranking happens before the filter
- result size is not the cost driver
- aggregates pre-reduce, windows do not
- a local top-N per worker would suffice
- check rows entering the exchange
basics
~20 sBecause 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.
solid answer
~50 sA window plus a filter on the row number is evaluated in that order: rank first, filter second. The plan therefore moves all rows through the exchange and sorts every partition end to end, even though only N rows per key survive. Nothing about the query lets the engine reduce data before the shuffle, because a window emits one row per input row. Some optimisers detect the shape — a ranking window immediately filtered to a small N — and push a per-partition top-N below the exchange, which changes the cost dramatically; check the plan rather than assume it. When the engine does not, rewrite the query as an aggregate: `GROUP BY key` with `max(ts)` (or an argmax-style function where the dialect offers one) is partially reducible, so each worker sends one row per key instead of all of them, and the payload is recovered with a join. Watch for ties, which the aggregate form can duplicate.
code
sql · 19 lines-- Before: every row is shuffled and every partition fully sorted
SELECT * FROM (
SELECT o.*,
ROW_NUMBER() OVER (PARTITION BY customer_id ORDER BY order_ts DESC) AS rn
FROM orders o
) t
WHERE rn = 1;
-- After: the aggregate pre-reduces to one row per customer per worker
WITH latest AS (
SELECT customer_id, max(order_ts) AS order_ts
FROM orders
GROUP BY customer_id
)
SELECT o.*
FROM orders o
JOIN latest l
ON o.customer_id = l.customer_id
AND o.order_ts = l.order_ts;go deeper
Recall that the ranking is computed for every row before anything is filtered, so a query returning one row per customer may still process the entire table. Small output does not mean small work.
Explain why a window cannot pre-reduce before the shuffle while an aggregation can, and describe the group-by-plus-join rewrite along with the tie behaviour that differs between the two forms.
Read the plan for rows entering the exchange, know that some engines push a per-partition top-N below it, and pick between the window and the aggregate form on evidence rather than habit. Add the time filter that shrinks the input in either case.
Decide when "latest row per key" should stop being a query at all. A maintained current-state table updated on ingest turns a recurring full-table sort into a point lookup, and that is a modelling decision with a cost curve attached.
## The order of operations is the whole story The idiomatic top-N-per-group query computes a ranking in a subquery and filters it outside: ```sql SELECT * FROM ( SELECT o.*, ROW_NUMBER() OVER (PARTITION BY customer_id ORDER BY order_ts DESC) AS rn FROM orders o ) t WHERE rn = 1; ``` The filter on `rn` cannot be applied before the ranking exists, and the ranking is defined over the full ordered partition. So the physical plan does the expensive things first: hash every row on `customer_id`, ship it across the network, sort each partition by `order_ts`, number the rows, and only then discard all but one per customer. On a billion-row table returning ten million customers, roughly 99% of the shuffled and sorted bytes are thrown away. This is why the query "feels" cheap — the result is small — and is not. Result size is a terrible predictor of window cost; input size and partition size are the predictors. ## Why the engine cannot reduce early on its own Aggregation is *partially reducible*: each worker can compute a local partial result per group and send only that, so the exchange carries groups rather than rows. A general window function is not, because it produces one output row per input row; there is no smaller intermediate that preserves the answer. The top-N shape is a special case where reduction *is* possible in principle: a worker only needs its own local top N per key before the exchange, since a row that is not in the local top N of its partition-fragment cannot be in the global top N. Some optimisers implement exactly this — often called a partial or pushed-down top-N — recognising a `ROW_NUMBER`/`RANK` window whose only consumer is a filter such as `rn <= 3`. When it fires, the exchange carries N rows per key per worker instead of everything, and the query becomes cheap. When it does not fire, or when the filter is written in a way the optimiser does not recognise (buried behind other operators, combined with unrelated predicates, or applied to a saved intermediate), you pay the full price. The only reliable move is to read the plan and look for the row count entering the exchange. ## The aggregate rewrite Where the optimisation is absent, express the intent as an aggregation: ```sql WITH latest AS ( SELECT customer_id, max(order_ts) AS order_ts FROM orders GROUP BY customer_id ) SELECT o.* FROM orders o JOIN latest l ON o.customer_id = l.customer_id AND o.order_ts = l.order_ts; ``` The aggregation pre-reduces on every worker before the shuffle, so the network sees one row per customer per worker rather than every order. The join then fetches the payload. Three caveats: - **Ties duplicate.** If two orders share the maximum timestamp, both come back. `ROW_NUMBER` picks exactly one arbitrarily; the aggregate form does not. Add a deterministic tiebreaker, or accept and document the duplicate. - **It reads the table twice** — once for the aggregate, once for the join. That is usually still far cheaper than a full sort of every partition, but it is not free, and for a very narrow table the window form may win. - **Dialects with an argmax-style aggregate** (a function returning the value of one column at the row where another is maximal) collapse this into a single pass with no self-join. Where the dialect offers one, it is the best form; where it does not, the join is the portable answer. For N greater than one, the aggregate rewrite becomes clumsier — you need per-key top-N collection — and a well-optimised window with pushed-down top-N is usually the better tool. That is exactly the case where confirming the optimisation in the plan is worth the minute it takes. ## Shrink the input regardless Orthogonal to the operator choice: most "latest row per key" queries have a time bound in reality. If the newest order per customer is always within the last 90 days for active customers, filtering to that range before the window cuts the input by whatever fraction history represents, and any remaining customers can be handled by a second, small query. Partition pruning on a time-partitioned table makes this nearly free. ## Reading the plan The diagnostic is one number: rows entering the exchange below the window. If it equals the table's row count, no reduction happened and the rewrite is worth doing. If it is close to keys times workers times N, the optimiser has already done the work for you and the query is as good as it gets.
- What does the aggregate rewrite do differently when two rows tie on the ordering column?It returns both. Taking max(order_ts) per customer and joining back matches every row at that timestamp, whereas ROW_NUMBER assigns 1 to exactly one of the tied rows and discards the rest. If exactly one row per key is required, add a deterministic tiebreaker such as the primary key to both the aggregate and the join, or deduplicate after the join.
- How would you confirm from a plan that a per-partition top-N was pushed below the exchange?Compare the row count entering the exchange with the table's row count. If the exchange carries roughly keys times workers times N rows, the reduction fired; if it carries the full table, it did not. Some plans also name a partial or limited top-N operator beneath the exchange. Do not infer it from query shape alone — it depends on the engine and on how the filter is written.
- Does the same cost argument apply to a top-N over the whole table rather than per group?It applies even more clearly, and the fix is simpler: ORDER BY with a LIMIT lets every worker keep only its own top N and merge a tiny result, while a global ROW_NUMBER window typically gathers everything onto one worker first. Use the window form only when you need the number itself in the output or a per-group result.
saying these in an interview costs you the question
- Assumes a small result set implies a cheap query
- Thinks the filter on the row number prunes rows before the sort
- Believes every engine pushes a partial top-N below the exchange
- Rewrites to an aggregate without handling ties
- Ignores an available time filter that would shrink the input first