skip to content

Why does a window function with no PARTITION BY collapse onto one worker?

level: middleimportance: should knowfreq 52%

answer

  1. a partition cannot be split across workers
  2. no PARTITION BY means exactly one partition
  3. gathering exchange, not a hash exchange
  4. unordered OVER () can stay distributed
  5. global running totals serialize the final phase

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.

solid answer

~50 s

An empty or absent `PARTITION BY` means the whole result set is one partition, and a partition cannot be split across workers. The exchange below the window therefore routes every row to a single node, which sorts the entire input and scans it once — the rest of the cluster idles. That is why `ROW_NUMBER() OVER (ORDER BY event_ts)` over a billion rows behaves like a serialization point in an otherwise parallel plan, and why it is a common source of spill. There are two escapes. If the window has **no ORDER BY** — a whole-table aggregate such as `SUM(x) OVER ()` — many engines compute it as a normal distributed aggregate and broadcast the single scalar back, so nothing collapses. If a global ordering is genuinely required, partition by something meaningful instead, or accept that the final phase is inherently sequential.

code

sql · 9 lines
sql
-- Collapses: one partition, so one worker sorts the entire table
SELECT user_id, event_ts,
       ROW_NUMBER() OVER (ORDER BY event_ts) AS global_seq
FROM events;

-- Parallel: each worker handles its own share of user_id partitions
SELECT user_id, event_ts,
       ROW_NUMBER() OVER (PARTITION BY user_id ORDER BY event_ts) AS seq_in_user
FROM events;

go deeper

for a junior

Recall that leaving out PARTITION BY makes the whole result one big group, and that on a large table this is slow. Check whether the query actually meant to number rows per user, per account or per day.

for a middle

Explain the mechanism: one partition cannot be split, so the plan gathers every row onto a single worker to sort and scan. Distinguish the unordered whole-table aggregate, which can stay distributed, from the globally ordered form, which cannot.

for a senior

Recognise the shape in a profile — one saturated worker, spill on that worker, idle peers — and know the repertoire of fixes: restore the real partition key, shrink the input first, or replace a numbered global ordering with ORDER BY plus LIMIT.

for a principal

Take a position on whether a global sequence belongs in an interactive query at all. Assigning one at load time, or accepting a per-entity sequence instead, removes a whole class of cluster-wide stalls from a dashboard fleet.

## One partition means one worker A window function is defined over partitions, and a partition is indivisible: correctness requires that every row belonging to it be visible to whatever computes the result. When the `OVER` clause names no `PARTITION BY` columns, there is exactly one partition — the whole (post-`WHERE`, post-join) result set. In a shared-nothing MPP engine the placement rule then degenerates to "send everything to one worker". The plan shows a **gathering exchange** rather than a hash exchange: all nodes stream their rows to a single destination. The consequences are immediate. Cluster parallelism for that stage drops to one. The single worker must hold or sort the entire input, so memory pressure scales with the whole table rather than with one node's share, and spilling to local disk becomes likely. Everything downstream waits, because a window is pipeline-breaking. Adding nodes does not help at all; the stage is bounded by one machine's CPU, memory and disk. ## The case that does not collapse: unordered whole-table aggregates Not every unpartitioned window is a problem. `SUM(amount) OVER ()` or `COUNT(*) OVER ()` — no `PARTITION BY` and no `ORDER BY` — asks for the same scalar on every row. There is no sequential dependency between rows, so many engines recognise the shape and evaluate it as an ordinary distributed aggregate: each node computes a partial sum, the partials are combined, and the single result is broadcast back so each worker can attach it to its own rows. Nothing is gathered. Whether a particular engine performs this rewrite is worth confirming in the plan, but the shape is cheap in principle. The expensive shapes are the ones with a global order: `ROW_NUMBER() OVER (ORDER BY ts)`, `SUM(x) OVER (ORDER BY ts)` as a running total across the whole table, `NTILE(100) OVER (ORDER BY revenue)`. Each output depends on the position of the row in a total order, so the engine needs a globally ordered sequence. ## Can a global order be parallelised at all? Partly, and it depends on the function and the engine. A globally ordered running sum is decomposable in two phases: range-partition the data, compute a local prefix sum on each range in parallel, then add each range's offset (the sum of all preceding ranges) to its rows. Global `ROW_NUMBER` can be produced the same way from per-range counts. Engines that implement this run the heavy pass in parallel and keep only a tiny coordination step sequential. Engines that do not simply gather. You cannot assume the optimisation exists — read the plan for a gathering exchange and a one-worker sort, and treat its presence as the answer. ## What to do instead First, ask whether the global order is real. A surprising share of unpartitioned windows in production SQL want a per-entity sequence and were written without the partition key by accident — the query returns plausible numbers, so nobody notices until the table grows. Adding `PARTITION BY user_id` both fixes the semantics and restores parallelism. Second, shrink the input before the window. A global `ROW_NUMBER` over a filtered, aggregated result of ten thousand rows is trivial; over the raw fact table it is a cluster-wide stall. Pushing the filter and the aggregation below the window is usually the whole fix. Third, if you only need the top rows of a global order, express that as `ORDER BY ... LIMIT` rather than as a window plus a filter on the row number. A top-N with a limit can be evaluated with per-worker partial top-N and a small merge, moving N rows per worker instead of all of them. Fourth, if a global sequence must genuinely be assigned to every row of a large table, treat it as a batch job with room for spill rather than as part of an interactive query, and question whether a monotonic surrogate produced at load time would serve the same purpose. ## Reading it in a plan The tell is a single-destination exchange feeding a sort and a window, with a row count equal to the whole input and a worker count of one. Runtime profiles show one busy worker, high spill bytes on that worker, and near-zero utilisation elsewhere. Once you can recognise that shape, the diagnosis takes seconds.

  • Why is SUM(x) OVER () usually cheap while SUM(x) OVER (ORDER BY ts) over the same table is not?
    The unordered form asks for one scalar shared by all rows, so each node can compute a partial sum and the combined result is broadcast back — a normal distributed aggregate. Adding ORDER BY makes each row's value depend on its position in a total order, which forces either a gather onto one worker or a two-phase prefix-sum scheme, and many engines simply gather.
  • You need the 100 highest-revenue rows overall. Is ROW_NUMBER() OVER (ORDER BY revenue DESC) with a filter a good way to get them?
    Usually not. A plain ORDER BY revenue DESC LIMIT 100 lets each worker keep only its own top 100 and merge a tiny result, whereas the window form assigns a number to every row and typically gathers everything first. Use the window form only when you need per-group top-N or the row number itself in the output.

saying these in an interview costs you the question

  • Thinks an empty OVER () always gathers all rows to one node
  • Says adding workers will speed up a global ordered window
  • Believes a global ROW_NUMBER is as parallel as a partitioned one
  • Uses a window plus row-number filter where ORDER BY with LIMIT would do
  • Misses that a missing partition key is often an outright bug

context