Why does PARTITION BY in a window function force a shuffle in an MPP warehouse?
answer
- placement problem before a math problem
- rows must meet before they can rank
- the planner inserts an exchange operator
- hash the partition key, ship the row
- one shuffle per distinct partition key
basics
~20 sA 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 sA 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
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.
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.
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.
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