skip to content

Window Functions at Scale

Window functions are the workhorse of analytics, but every PARTITION BY is a shuffle and every ORDER BY is a sort. Interviewers probe whether I see the distributed cost hiding behind one clean-looking statement.

on this pageshow

questions

6

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

level: middleimportance: must knowfreq 70%

answer

  1. placement problem before a math problem
  2. rows must meet before they can rank
  3. the planner inserts an exchange operator
  4. hash the partition key, ship the row
  5. one shuffle per distinct partition key

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.

solid answer

~50 s

A window function is evaluated per partition, so every row carrying the same `PARTITION BY` value has to end up on one worker. In a shared-nothing engine that means the planner inserts an **exchange** operator: each node hashes the partition key of every row it holds and ships the row to the node that owns that hash bucket. After the exchange, each worker sorts its share by (partition key, order key) and scans it sequentially, resetting the accumulator at each partition boundary. So the real cost is the network transfer of the whole intermediate result plus a per-worker sort — not the arithmetic inside the function. Several window functions sharing the same PARTITION BY and ORDER BY normally share one exchange and one sort; a second, different partition key means a second exchange and a second sort.

go deeper

for a junior

Know that PARTITION BY groups rows for the function and that in a distributed warehouse those rows physically have to be brought together first. Being able to say that this movement is the expensive part is enough at this stage.

for a middle

Be ready to describe the mechanics end to end: hash the partition key, exchange rows across workers, sort by partition and order key, scan once resetting state at partition boundaries. Name where the network cost and the sort cost each land.

for a senior

Show that you read plans for exchange operators and count distinct partition keys as a cost proxy. Explain when the exchange disappears because the table is already distributed on that key, and why filtering and projecting before the window is not cosmetic.

for a principal

Own the pipeline-wide argument: choosing one distribution key that serves the join, the aggregation and the window collapses several data movements into one, and that choice is a storage-layout decision made once, not a query rewrite made repeatedly.

## What a window function actually requires A window function computes a value for every input row using a set of *other* rows — its partition, ordered, optionally narrowed to a frame. The defining property is that the answer for one row depends on rows the engine may not have next to it. In a single-machine engine that is only a sorting problem. In a massively parallel (MPP) engine, where the table is spread across many independent nodes that share nothing but a network, it is first a *placement* problem: no worker can compute a correct value for a row until every peer row of that partition is on the same worker. ## The exchange operator The planner solves placement by inserting a repartitioning **exchange** (also called a shuffle) below the window operator. Each node walks the rows it already holds, applies a hash function to the partition-key columns, maps the hash to one of N buckets, and sends each row to the worker that owns that bucket. Every worker is simultaneously a sender and a receiver. When the exchange finishes, the invariant the window operator needs holds: all rows with a given partition key sit on exactly one worker, and no partition is split. This is the same physical mechanism a hash aggregation or a hash-redistribution join uses. The difference matters for cost: an aggregation can pre-reduce on the sending side (compute local partial sums, ship one row per group per node), while a plain window function generally cannot, because it must emit one output row per input row. A window shuffle therefore moves *all* the rows and *all* the columns the query still needs. ## The sort after the shuffle Once rows land, the worker must put each partition into the order the `ORDER BY` inside `OVER` specifies. Engines normally do one combined sort on (partition key, order key) so that partitions come out contiguous and internally ordered, then a single sequential pass computes the function, resetting state whenever the partition key changes. That sort — not the shuffle — is usually the memory-hungry part, because it materialises rows with their full payload. ## Where the cost lands Three separate costs hide behind one clean-looking clause: - **Network**: bytes moved is roughly the width of the projected row times the row count. Selecting columns you do not need inflates this directly. - **Sort**: CPU and memory per worker, proportional to that worker's share of the rows and to row width. If a partition exceeds the operator's memory budget, it spills to local disk. - **Pipeline breaking**: the window cannot emit its first row until its input for that partition is complete, so downstream operators wait, and one slow worker stalls the whole stage. ## When the exchange can be skipped If the table is already distributed on the same key the window partitions by, a good planner recognises that rows are already co-located and omits the exchange — a local sort is all that remains. The same applies to a chain of operators: a join that already redistributed on `customer_id`, followed by a window partitioned by `customer_id`, can keep one distribution for both. Aligning the join key, the group-by key and the window partition key across a pipeline is one of the highest-leverage rewrites in warehouse SQL, because it collapses several exchanges into one. Similarly, if the physical layout already sorts data within each partition on the order key, some engines can cheapen or skip the sort and stream the window instead. That is a property of the storage layout, not of the SQL. ## Multiple windows in one query A query with `sum(x) OVER (PARTITION BY a ORDER BY t)` and `rank() OVER (PARTITION BY a ORDER BY t)` needs one exchange and one sort; the two functions are evaluated in the same pass. Add `sum(x) OVER (PARTITION BY b)` and the plan grows a second exchange plus a second sort, because the placement invariant for `b` is different from the one for `a`. Reading a plan, the count of exchange operators is therefore roughly the count of *distinct* partition keys in the query, and that count is a good first estimate of the query's shuffle cost. ## Practical consequences Filter before the window, not after — rows removed after the window were still shuffled and sorted. Project only the columns you need, since row width multiplies both network and sort cost. Prefer one partition key across a query where the semantics allow. And treat a window over an unfiltered fact table in a dashboard query as an expensive operation by default, because it moves the whole table across the network on every run.

  • Why can a hash aggregation often avoid moving all the rows while a window function usually cannot?
    Aggregation is partially reducible: each node can compute local partial sums per group and ship one row per group, so the exchange carries groups, not rows. A window function emits one output row per input row, so every row must reach its partition's owner. That asymmetry is why rewriting a top-N window as an aggregate plus a join is often dramatically cheaper.
  • How does the physical distribution of the table change the plan for a partitioned window?
    If the table is already distributed on the window's partition key, rows for a partition are already co-located and the planner can drop the exchange, leaving only a local sort. If it is distributed on some other key — or spread evenly with no key at all — the exchange is mandatory. Aligning join, group-by and window keys lets one distribution serve several operators.
  • Two window functions in one query use different PARTITION BY keys. What does the plan look like?
    Two repartitioning stages, executed one after the other: shuffle and sort on the first key, evaluate that window, then shuffle and sort the intermediate result on the second key and evaluate the second window. Each stage is pipeline-breaking, so the query pays two full data movements. Where the business logic allows one shared key, the plan collapses to a single stage.

Think of a deck of cards spread across ten tables. Before anyone can rank the hearts, every heart must be carried to one table — the carrying, not the ranking, is the work.

saying these in an interview costs you the question

  • Claims window functions run locally on whatever node holds the row
  • Thinks ORDER BY inside OVER decides data placement, not PARTITION BY
  • Assumes adding nodes always speeds up a window query
  • Believes several window functions always cost several shuffles
  • Confuses the frame clause with the redistribution key

context

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

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 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

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