Why can a BigQuery query fail with "Resources exceeded" if BigQuery scales automatically?
answer
- scaling makes it wider, not deeper
- some stages cannot be split at all
- one worker holding everything is the pattern
- a partitioning key is usually the fix
- more slots do not fix concentration
basics
~20 sBigQuery scales the number of workers, not the memory inside one worker. Operations that funnel all rows into a single worker — a global ORDER BY, a window function with no PARTITION BY, one enormous group — exceed that worker's memory no matter how many slots are free.
solid answer
~50 sAutoscaling adds *width*, and some stages cannot use width. If a stage must see all rows in one place — ranking the entire table with `ROW_NUMBER() OVER (ORDER BY ts)`, a final global sort, an aggregation where one group holds most of the rows, or a `SELECT DISTINCT` over a very wide row — then one worker holds a share of data bounded only by the data itself, and its memory is bounded. BigQuery will spill shuffle data to disk to cope, but past a point the stage fails with a resources-exceeded error. Buying more slots usually does not help, because the problem is concentration rather than total capacity. The fixes are all about breaking the funnel: add a `PARTITION BY` on a high-cardinality key, filter and aggregate before the wide stage, replace an exact global operation with an approximate one, or materialise an intermediate result into a table and continue from there.
code
sql · 7 lines-- Funnels: one worker must rank every row in the table
SELECT *, ROW_NUMBER() OVER (ORDER BY event_ts) AS rn
FROM events;
-- Distributes: each user's rows are ranked independently
SELECT *, ROW_NUMBER() OVER (PARTITION BY user_id ORDER BY event_ts) AS rn
FROM events;go deeper
Recall that a query can still fail even though BigQuery manages the infrastructure, and that a global ORDER BY or an unpartitioned window function over a huge table is the usual trigger.
Explain the mechanism: autoscaling adds workers, but a stage that must see all rows in one place cannot be split, so its per-worker memory grows with the data until it spills and then fails.
Show the diagnosis loop — find the failing stage, read spilled shuffle bytes and the max-versus-average compute time, decide between volume, skew and a non-parallelisable operation — and name the specific rewrite for each.
Own the policy angle: recurring resources-exceeded failures usually mean a modelling problem, and the durable fixes are pre-aggregated tables, approximate aggregation where exactness is not required, and guardrails that keep unbounded global operations out of scheduled pipelines.
## Two different kinds of "more resources" When a query needs more capacity, there are two directions to grow: - **Wider** — run more workers in parallel on more pieces of the input. BigQuery does this automatically, and this is what people mean by autoscaling. - **Deeper** — give one worker more memory to hold a bigger working set. BigQuery does *not* expose this, and there is no per-query memory knob. A resources-exceeded failure is almost always the second kind of problem wearing the first kind's clothes. The stage in question cannot be split further, so its per-worker working set grows with your data, and eventually exceeds what a worker can hold even after spilling. ## The query shapes that funnel **Global ordering without a limit.** `ORDER BY` at the top of a query defines a total order over the result. If the result is huge, the final sort concentrates. A top-*k* pattern (`ORDER BY x DESC LIMIT 100`) is fine, because each worker can locally keep its own top 100 and only those flow up. Ordering millions of output rows is not. **Analytic functions with no partition.** `ROW_NUMBER() OVER (ORDER BY event_ts)` requires one global sequence, so one worker ranks the world. `ROW_NUMBER() OVER (PARTITION BY user_id ORDER BY event_ts)` gives the engine a partitioning key: each user's rows can be handled independently and the stage fans out across as many workers as there are distinct users. The same holds for `SUM(...) OVER ()` versus `SUM(...) OVER (PARTITION BY ...)`. **A skewed grouping or join key.** `GROUP BY country` where 70% of rows share one value means the shuffle sends 70% of the data to one worker. Same for a join on a key where one value dominates — a sentinel like `'unknown'`, an empty string, or a NULL-substitute. The stage is nominally parallel and effectively serial. **Expensive per-row state.** Very wide rows, large arrays or JSON blobs, `ARRAY_AGG` collecting an unbounded group, or `STRING_AGG` over millions of rows all raise the memory a single worker must hold for its share. **Exact distinct over huge cardinality.** Exact `COUNT(DISTINCT ...)` must track every distinct value; over billions of high-cardinality values that state is large. ## How to diagnose it Open the job's stage-level statistics and find the stage that failed or ballooned: - **`shuffle_output_bytes_spilled` greater than zero** means shuffle data did not fit in memory and went to disk — the leading indicator that the stage is under memory pressure. - **`compute_ms_max` far above `compute_ms_avg`** in the same stage means skew: one worker did vastly more work than its peers, which points at a hot key rather than sheer volume. - **A stage with very low parallelism** (few input pieces) right before the failure is the funnel itself. That trio distinguishes the three causes — too much total data through one stage, one hot key, or a genuinely non-parallelisable operation — and each has a different fix. ## The fixes, in the order to try them 1. **Give the operation a partitioning key.** Add `PARTITION BY` to the analytic function. This is by far the most common fix and usually the correct one: the ranking you actually wanted was almost always per-entity, not global. 2. **Reduce before you concentrate.** Filter earlier, project fewer columns, and aggregate before joining. The funnel stage's cost is proportional to what reaches it. 3. **Drop the exactness requirement.** Approximate distinct counting and approximate quantiles use bounded, mergeable state and parallelise cleanly, so they scale where their exact counterparts do not. Use them when a rounding error is acceptable. 4. **Handle the hot key explicitly.** Split the dominant value out and process it separately, or salt the key — add a random suffix to spread it, aggregate, then combine. Filtering out sentinel values that were never meaningful is often enough on its own. 5. **Break the query into steps.** Write an intermediate result to a table and run the heavy stage against the smaller table. This costs an extra write but converts one impossible stage into two possible ones. 6. **Only then, consider capacity.** More slots help when the query is queueing or genuinely under-parallelised — not when a single worker's working set is the constraint. ## The sentence that shows seniority *"Serverless scales the number of workers, not the size of one. This stage can't be split, so the fix is in the query, not the capacity."* Following that with the specific rewrite — `PARTITION BY user_id`, or approximate counting, or pre-aggregation — is what separates a candidate who has debugged this from one who has read about it.
- Why does ORDER BY with a small LIMIT usually succeed where the same ORDER BY without a LIMIT fails?A top-k sort is distributable: each worker sorts its own share and keeps only its best k rows, so just k rows per worker flow to the final stage. Without a limit, the engine must produce a total order over the entire result, so the full result set has to pass through the final sort, and its size is bounded only by the data.
- How would you confirm from a job's statistics that skew rather than sheer volume caused the failure?Compare per-stage timing distributions. If compute_ms_max is many times compute_ms_avg while records_read per worker is wildly uneven, one worker got a disproportionate share — a hot key. If all workers are similarly loaded and shuffle bytes spilled is large, the stage is simply moving more data than fits, which calls for pre-aggregation rather than key handling.
- What is key salting and when is it appropriate here?Salting adds a small random or derived suffix to a hot key so its rows spread across many workers, you aggregate per salted key, then combine the partials in a second pass. It fixes one dominant value skewing a GROUP BY or join. It costs an extra aggregation stage and only makes sense when the skew is a few values, not a broadly uneven distribution.
saying these in an interview costs you the question
- Says buying more slots always fixes resources-exceeded
- Believes serverless engines cannot run out of memory
- Blames table size instead of the concentrating stage
- Adds ORDER BY at the top of every query by habit
- Ignores a hot key and reruns hoping it passes