skip to content

How do you find which stage of a slow BigQuery query consumed the slot time?

level: seniorimportance: should knowfreq 44%

answer

  1. the plan is stored, not just drawn
  2. rank the stages by their slot time
  3. rows in versus rows out per stage
  4. average against maximum tells you about skew
  5. waiting is contention, not a bad query

basics

~20 s

Read the job's per-stage statistics: the execution details in the console, or the job_stages array in INFORMATION_SCHEMA.JOBS_BY_PROJECT. Rank stages by slot_ms, then look at records read, shuffle bytes, and wait/read/compute/write timings to see why that stage was expensive.

solid answer

~40 s

Every BigQuery job exposes its plan as a list of stages with per-stage statistics. Query `INFORMATION_SCHEMA.JOBS_BY_PROJECT`, unnest `job_stages`, and sort by `slot_ms` — that names the expensive stage immediately. Then read that stage's detail: `records_read` and `records_written` show whether it exploded or reduced rows; `shuffle_output_bytes` shows how much it redistributed; and the four timing families tell you where its time went. High `wait_ms` means it sat waiting for slots, which is contention, not a query problem. High `read_ms` means the scan is too big — fix with partition filters, clustering, or fewer columns. High `compute_ms` means real CPU work. And when `compute_ms_max` is far above `compute_ms_avg`, one worker got a disproportionate share, which is skew on the shuffle key. Always filter the view by `creation_time` so the query itself stays cheap.

code

sql · 16 lines
sql
SELECT
  stage.name,
  stage.slot_ms,
  stage.records_read,
  stage.records_written,
  stage.shuffle_output_bytes,
  stage.shuffle_output_bytes_spilled,
  stage.wait_ms_avg,
  stage.read_ms_avg,
  stage.compute_ms_avg,
  stage.compute_ms_max
FROM `region-us`.INFORMATION_SCHEMA.JOBS_BY_PROJECT AS j,
  UNNEST(j.job_stages) AS stage
WHERE j.job_id = 'bquxjob_0123456789abcdef'
  AND j.creation_time > TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 1 DAY)
ORDER BY stage.slot_ms DESC

go deeper

for a junior

Know that BigQuery records an execution plan for every job and that the console shows stage-by-stage details rather than a single opaque timing.

for a middle

Be able to read a stage: rows in versus rows out, shuffle bytes, and where its time went across waiting, reading, computing and writing — and say what each pattern implies.

for a senior

Demonstrate the loop end to end on a real job: rank stages by slot time, classify the cause, use max-versus-average to detect skew and spilled bytes to detect memory pressure, then name the specific fix.

for a principal

Turn the same data into a programme — aggregate slot consumption by team, label and destination table over time, rank recurring spend, and drive structural fixes such as pre-aggregation and workload separation rather than one-off tuning.

## Where the data lives BigQuery records a structured execution plan for every query job. You can read it two ways: - **The console's execution-details view**, which draws the stages, their timings, and the steps inside each one. Good for a one-off. - **`INFORMATION_SCHEMA.JOBS_BY_PROJECT`** (and the related job views), which exposes the same information as queryable rows — the right tool for finding the worst offenders across a day of pipeline runs, not just the query you happen to have open. The job row carries whole-query facts, including `total_slot_ms` and `total_bytes_processed`. The `job_stages` array carries one entry per stage with its own name, `slot_ms`, row counts, shuffle bytes, and timing distributions. These views are partitioned on `creation_time`, so always constrain it — otherwise your diagnostic query is itself an expensive scan. ## The two headline numbers, and why they differ **Bytes processed** is what on-demand billing counts and what a dry run predicts. **Slot milliseconds** is how much execution capacity the query actually consumed. These diverge constantly, and understanding the divergence is most of the skill: - A wide scan with no joins burns bytes and little slot time. - A modest scan followed by a huge join, a high-cardinality `GROUP BY`, a regex over every row, or JSON parsing burns slot time far out of proportion to bytes. So *"why is this slow"* and *"why is this expensive"* can have different answers in the same job. Dividing `total_slot_ms` by the query's elapsed milliseconds gives the average number of slots the query held — a useful check on whether it was starved or genuinely wide. ## Reading a stage For the top stage by `slot_ms`, the fields answer specific questions: - **`records_read` / `records_written`.** A stage that reads 10 billion and writes 400 is doing the aggregation you wanted. A stage that reads 2 billion and writes 60 billion is a join blowing up the row count — usually a wrong or missing join predicate, or a fan-out you did not intend. - **`shuffle_output_bytes`.** How much this stage pushed into shuffle for the next one. Big numbers here are the signature of a query that is cheap on bytes-scanned but heavy on slot time. - **`shuffle_output_bytes_spilled`.** Non-zero means the shuffle did not fit in memory. Memory pressure; often precedes an outright failure as data grows. - **`parallel_inputs`.** How many pieces the stage's input was split into — an upper bound on its parallelism. A very small number just before an expensive stage means the work could not be spread. - **The timing families — `wait_ms`, `read_ms`, `compute_ms`, `write_ms`, each reported as an average and a maximum across workers.** ## What each timing family points to **Wait.** Time workers spent waiting to be scheduled. Dominant wait time is a *capacity* story, not a query story: the query was queued behind other work, or the project's slot availability was low at that moment. Rewriting the SQL will not help; scheduling it differently, or running it against reserved capacity, will. **Read.** Time pulling data from storage or from shuffle. Dominant read time in the first stage means the scan is too big — reduce it with a partition filter, clustering on the filter columns, or by selecting fewer columns. Dominant read time in a later stage means it is consuming a lot of shuffle. **Compute.** Real CPU work: expressions, comparisons, hashing, sorting, parsing. Dominant compute time in a stage that reads few rows points at expensive per-row functions — regular expressions, JSON extraction, `UDF`s, or complex `CASE` chains. **Write.** Time producing output into shuffle or the destination. Dominant write time usually means the stage is emitting far more than it consumed. ## Average versus maximum: the skew test The single most valuable comparison in the whole structure is `compute_ms_max` against `compute_ms_avg` within one stage. Similar values mean the work spread evenly. A max many times the average means one worker got a disproportionate share of the shuffle partition — a hot key. That is a data problem, and no amount of extra capacity fixes it; you handle the dominant key separately, salt it, or filter out sentinel values that were never meaningful. ## A workflow that holds up under pressure 1. Pull the job's stages, sorted by `slot_ms` descending. 2. Take the top one or two stages — cost is almost always concentrated. 3. Classify with the timing families: contention (wait), scan (read), CPU (compute), fan-out (write plus a rising row count). 4. Check max-versus-average to rule skew in or out. 5. Check spilled shuffle bytes to see whether memory pressure is building. 6. Fix the specific cause: prune columns and partitions for scan cost; pre-aggregate or fix the join predicate for fan-out; partition or salt for skew; reschedule or reserve capacity for wait. Doing this across a whole day of jobs — group by a query label or by the top expensive stage name — turns a one-query debugging session into a cost programme, which is what senior work on a warehouse actually looks like.

  • Why can bytes processed be small while slot time is enormous?
    Bytes processed counts what the scan reads; slot time counts all execution. A query can read a few gigabytes, then join, redistribute, sort, and run expensive per-row functions over billions of intermediate rows. That work costs slot time and never touches the bytes-scanned figure, which is why on-demand cost and wall-clock slowness are separate diagnoses.
  • A stage's wait time dominates every other timing. What do you change?
    Nothing in the SQL. Dominant wait means workers were queued for slots, so the query was contending with other work rather than doing too much itself. Move it off the peak window, give it reserved capacity or a higher priority within a reservation, or reduce concurrent load. Rewriting a query that is merely waiting wastes effort.
  • How would you use these views to run a cost programme rather than debug one query?
    Query the job views over a time window, filtered by creation_time, and aggregate total_slot_ms and bytes processed by user, by query label, and by destination table. That ranks recurring pipelines and dashboards by real consumption. Then attack the top few with partition filters, pre-aggregated tables, or schedule changes, and re-measure the same aggregate weekly.

saying these in an interview costs you the question

  • Judges query cost by bytes scanned alone
  • Reads the plan only for the query in front of them
  • Ignores the average-versus-maximum skew signal
  • Scans the job views without a creation_time filter
  • Treats wait time as evidence the SQL is bad

context