skip to content

A window partitioned by tenant_id pins one node at 100% while others idle — what is happening?

level: seniorimportance: must knowfreq 62%

answer

  1. hashing balances distinct values, not rows
  2. the biggest key decides the wall-clock time
  3. a window partition cannot be split
  4. NULL and sentinel keys share one bucket
  5. salting rescues GROUP BY, not windows

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.

solid answer

~50 s

The window's exchange hashes `tenant_id`, so all rows of one tenant must land on one worker — and a window partition cannot be split. If one tenant holds a large share of the rows (or many rows carry a NULL or default key, which all hash to the same bucket), that worker sorts and buffers a partition orders of magnitude bigger than its peers', spills to local disk, and finishes long after everyone else. Because a window is pipeline-breaking, the whole stage waits on that straggler. Confirm it by counting rows per key and by comparing max versus median rows and spill bytes per worker in the profile. Fixes are about shrinking the hot partition, not about adding nodes: filter and pre-aggregate before the window, add a secondary key such as (tenant_id, day) if the semantics allow, split the hot tenants into their own query and union the results, or replace the window with a partially reducible aggregate.

code

sql · 10 lines
sql
-- Confirm skew before blaming the cluster
SELECT tenant_id, count(*) AS rows_in_partition
FROM events
GROUP BY tenant_id
ORDER BY rows_in_partition DESC
LIMIT 20;

SELECT count(*) AS null_key_rows
FROM events
WHERE tenant_id IS NULL;

go deeper

for a junior

Understand that rows are grouped onto machines by the partition key, and that one very common key value means one very busy machine. Counting rows per key is the first thing to check.

for a middle

Explain hashing on the partition key, the indivisibility of a partition, and why the stage's duration is set by the largest partition rather than by the average. Mention NULL keys collapsing into one bucket.

for a senior

Demonstrate the diagnosis from a runtime profile — max versus median rows and spill per worker — and offer a ranked set of fixes: shrink the input, refine the partition key, isolate hot keys, or change to an aggregate formulation. Say plainly why scaling out does not help.

for a principal

Argue the structural fix: naturally skewed tenants should not share a query with the long tail. Separate schedules, separate physical layout, or a precomputed per-tenant-per-day metric remove the class of failure instead of tuning each instance of it.

## Why the imbalance exists An MPP engine distributes window partitions by hashing the `PARTITION BY` key and assigning hash buckets to workers. That gives balance only if the key's values are roughly uniform in *row count*. Real business keys rarely are: one tenant is ten thousand times bigger than the median, one product SKU dominates the catalogue, one device id belongs to a test harness that emits continuously. Hashing spreads *distinct values* evenly; it does nothing about the rows behind each value. The crucial asymmetry with aggregation is that a window partition is indivisible. A skewed `GROUP BY` can often be salted — compute partial aggregates on `(key, salt)` in parallel, then combine — because aggregation is associative. A window function generally cannot be salted, because it emits one row per input row and its result depends on the ordering of the *entire* partition. Splitting the hot tenant across two workers would produce two independent sequences, not one. ## The straggler effect Once one worker owns a partition many times larger than average, three things compound. It sorts more rows than anyone else, at superlinear cost. It exceeds its memory budget and spills to local disk, adding write plus read-back I/O. And because the window is a pipeline breaker, downstream operators cannot proceed until the last partition is done, so the stage's wall-clock time is set by the worst worker while the rest of the cluster idles. Cost is billed on the whole cluster for the duration, so a skewed window is expensive in credits as well as in latency. NULL keys deserve a special mention: in most engines NULL hashes to a single bucket, so a column that is NULL for a large fraction of rows creates one gigantic partition. The same happens with sentinel values such as `-1`, `0` or `'unknown'` inserted by an upstream job. ## Confirming the diagnosis Two checks settle it quickly. First, look at the data: group by the partition key, count rows, order descending, and inspect the top values plus the count of NULLs. If the top key holds a large multiple of the median, that is the answer. Second, look at the runtime profile: compare rows processed, peak memory and spill bytes across workers for the window stage. Skew shows as one worker with a far higher row count and nonzero spill while peers are near zero. Rule out the alternatives — a genuinely undersized cluster shows *all* workers busy; a hardware straggler shows an imbalance that does not correlate with a data key and moves between runs. ## What actually helps **Shrink the input first.** Filters and pre-aggregation applied before the window reduce the hot partition proportionally. Rows discarded after the window were still shuffled, sorted and possibly spilled. **Refine the partition key** where the business logic permits. Changing `PARTITION BY tenant_id` to `PARTITION BY tenant_id, event_date` turns one enormous partition into many small ones and usually reflects what the report meant anyway (a running total that resets daily). This is the single most effective fix when it is semantically legal, and the question to ask the analyst is precisely whether the sequence is supposed to span all history. **Isolate the hot keys.** Run the query for the handful of giant tenants separately, possibly with a bigger cluster or more memory, and `UNION ALL` the results with the run over the long tail. Ugly but effective when the key list is stable and small. **Change the operator.** If the window is only being used to pick a winner per key — a top-1 or top-N — an aggregate formulation is partially reducible and does not need the whole partition on one node. That converts a skew problem into a normal aggregation problem. **Give the hot worker room.** Raising the per-operator memory budget or moving to nodes with more memory removes the spill even if it does not remove the imbalance. This is the fix that costs money rather than thought, and it is the right one when the skew is inherent and the query is infrequent. **Do not** reach for a larger cluster as the first move. Adding workers gives the hot partition no extra help; it just adds idle machines to the bill. The only sizing change that helps a skewed window is a *bigger* worker, not more of them. ## The design-level answer If a family of reports partitions by a naturally skewed key, the durable fix belongs upstream: keep the giant tenants in separate physical tables or in their own schedule, agree on a coarser grain for the window, or precompute the metric per tenant per day so no single query ever holds one tenant's full history in one operator.

  • Salting fixes a skewed GROUP BY. Why does it not fix a skewed window partition?
    Salting works because aggregation is associative: partial results per salt can be combined into the true group result. A window function emits one row per input row and each value depends on the ordering of the whole partition, so splitting the hot key into salted sub-partitions produces several independent sequences rather than one. Only decomposable cases — a running sum that can be offset by preceding ranges — allow a two-phase equivalent.
  • How would you tell partition skew apart from an undersized cluster in a profile?
    Skew shows one worker with far more rows, peak memory and spill than the others, and the imbalance follows the same key on every run. An undersized cluster shows all workers roughly equally busy and spilling. A hardware straggler shows imbalance that does not line up with any data key and moves between runs. Grouping by the partition key and counting rows confirms the first case in one query.
  • A large fraction of rows have a NULL tenant_id. What does that do to the window?
    In most engines NULL hashes to a single bucket, so all those rows form one enormous partition on one worker — the same straggler pattern, but caused by missing data rather than by a real giant tenant. The fix is usually upstream: reject or route NULL keys at load time, or filter them out before the window when they carry no meaning for the report.

Ten checkout lanes open, but the rule says every item from one customer must go through one lane. If a single customer arrives with a thousand trolleys, nine lanes finish and stand idle while one grinds on.

saying these in an interview costs you the question

  • Suggests scaling the cluster out to fix a skewed partition
  • Salts the window partition key as if it were a GROUP BY
  • Blames the network or a slow node without checking key counts
  • Assumes hash distribution guarantees equal rows per worker
  • Ignores NULL and sentinel keys as a source of one giant partition

context