skip to content

How does BigQuery execute a query across Dremel, Colossus and the Jupiter network?

level: middleimportance: must knowfreq 68%

answer

  1. four names: engine, storage, format, network
  2. a plan is a graph of stages, not one scan
  3. stages hand data over through a middle tier
  4. the tree shape became a DAG
  5. bytes read and slot time are different costs

basics

~20 s

BigQuery compiles SQL into a directed graph of stages executed by Dremel. Slots read only the referenced columns from Capacitor files on Colossus over the Jupiter network, pass intermediate rows through an in-memory shuffle tier, and a final stage writes the result.

solid answer

~50 s

A query becomes a **plan of stages**, not a single scan. The planner uses table metadata to decide what to read and how to partition the work. The leaf stage's slots read only the referenced columns from BigQuery's Capacitor columnar files stored on **Colossus**, over Google's **Jupiter** datacenter network — storage is not on the worker's disk, and the network is fast enough that this is workable. Each stage writes its output into an **in-memory shuffle tier** rather than handing rows directly to a peer; the next stage reads from shuffle. That decoupling is what lets BigQuery restart a failed worker, repartition data mid-query, and adjust the plan while it runs based on observed data volumes. The final stage writes to a temporary result table the client pages through. The original Dremel serving tree — root, mixers, leaves — is the ancestor of this design; today the shape is a general DAG.

code

text · 11 lines
text
S00: Input       READ      events.user_id, events.amount
                 FILTER    event_date >= '2026-01-01'
                 AGGREGATE GROUP BY user_id
                 WRITE     __stage00_output   (shuffle)

S01: Aggregate   READ      __stage00_output
                 AGGREGATE SUM(amount)
                 WRITE     __stage01_output   (shuffle)

S02: Output      READ      __stage01_output
                 WRITE     __result

go deeper

for a junior

Be able to name the pieces and their jobs: Dremel executes, Colossus stores, Capacitor is the columnar file format, Jupiter is the network that connects them.

for a middle

Walk the lifecycle end to end — plan into stages, leaf slots read only the needed columns, stages exchange rows through an in-memory shuffle, the last stage writes a result table — and say why a columnar format makes remote reads affordable.

for a senior

Use the model diagnostically: point at a stage graph and say whether the cost is in the scan, in the shuffle, or in skew, and explain why bytes scanned and slot time can diverge by an order of magnitude.

for a principal

Be ready to argue when a disaggregated, shuffle-based engine is the right platform at all: it trades per-query predictability and machine-level tuning for elasticity, statelessness and zero capacity planning.

## The named components Interviewers ask this question to see whether you can separate four things that often get blurred together: - **Dremel** — the distributed query execution engine. It is the thing that turns SQL into running work. - **Colossus** — Google's distributed file system, where table data physically lives. Durable, replicated, and completely separate from the machines that execute queries. - **Capacitor** — BigQuery's columnar storage format, the file layout written onto Colossus. Each column is stored and compressed separately. - **Jupiter** — Google's datacenter network fabric. It provides the bandwidth that makes reading storage over the network, instead of from a local disk, a reasonable design rather than a bottleneck. A cluster manager schedules the processes Dremel runs on; that is the layer that makes compute fungible across tenants. ## The life of a query 1. **Submission.** The query arrives as a job. It is parsed, type-checked against the table schemas, and planned. 2. **Planning.** The planner reads table metadata — schema, partition and clustering information, size statistics — and produces a **DAG of stages**. Each stage is a set of operations (read, filter, compute, aggregate, join, sort, write) plus a description of how its input is partitioned. 3. **Scheduling.** The scheduler assigns slots to the work units of the runnable stage. A stage's parallel width comes from how many independent input pieces exist and how many slots are currently available. 4. **Leaf reads.** The first stage reads from Colossus. Because Capacitor is columnar, only the columns named in the query are fetched, and partition metadata lets whole ranges of files be skipped. This is why `SELECT *` is so much more expensive than selecting three columns. 5. **Shuffle.** A stage does not hand rows to the next stage directly. It writes its output into BigQuery's in-memory shuffle tier, keyed by whatever the next stage needs to partition on (a join key, a group key). The consuming stage reads its partitions back out. 6. **Repeat.** Aggregations, joins, and analytic functions each become one or more stages, each separated by a shuffle boundary. 7. **Output.** The final stage writes the result to a temporary table; the client fetches pages from it. An identical repeat query against unchanged tables may be served from the cached result without executing at all. ## Why the shuffle boundary matters The original Dremel design was a fixed **serving tree**: a root server received the query, mixer nodes fanned it down, leaf nodes scanned storage, and partial results aggregated back up the tree. That is excellent for selective scan-and-aggregate, and poor for anything that needs data repartitioned by key — large joins in particular. Moving to a shuffle-based execution model changed three things at once: - **Producer and consumer are decoupled in time.** A worker that dies can be restarted and re-read its input; the query does not have to start over. - **Data can be repartitioned arbitrarily** between stages, which is what makes distributed hash joins and large `GROUP BY`s practical. - **The plan can change while it runs.** Having materialized a stage's output, BigQuery knows its real size and can adjust the parallelism or strategy of the next stage instead of relying purely on pre-query estimates. The cost is that shuffling is work. A query that reads few bytes but redistributes billions of rows can burn far more slot time than a query that scans more bytes and shuffles nothing. This is the single most useful practical consequence to state out loud: **bytes scanned and slot time are different currencies**, and the shuffle is where they diverge. ## Why storage-on-the-network is viable The instinct from shared-nothing MPP systems is that compute must sit on top of its data. Two things make disaggregation work here. First, the network: Jupiter's bisection bandwidth means a leaf worker can pull columnar data fast enough to keep its CPU busy. Second, the format: columnar files plus partition metadata mean a query typically reads a small fraction of the table's bytes, so the volume crossing the network is much smaller than the table size suggests. The payoff is that compute is stateless. Any worker can serve any query, workers can be added mid-query, and nothing has to be rebalanced when capacity changes — because no worker owns any data. ## What to take into a plan review When you look at a slow query, map what you see onto this model: an expensive first stage means too many bytes read (fix with fewer columns, partition filters, clustering); a large `shuffle_output_bytes` between stages means too much data is being redistributed (fix by aggregating earlier, filtering earlier, or reducing the join's cardinality); a stage where max compute time dwarfs average means the shuffle key is skewed.

  • Why did BigQuery move from Dremel's original fixed serving tree to a shuffle-based execution model?
    A serving tree aggregates partial results up a fixed hierarchy, which suits selective scan-and-aggregate but cannot repartition data by key. Shuffle decouples producers from consumers: data can be redistributed on a join or group key, a failed worker can be restarted without restarting the query, and because a stage's output is materialized, the engine can re-plan later stages using the real data volume instead of estimates.
  • Two queries scan the same number of bytes but one consumes ten times the slot time. What explains the gap?
    Almost always shuffle and CPU work after the scan. A join that redistributes billions of rows, a high-cardinality GROUP BY, an analytic function, or heavy string and JSON processing all cost slot time without increasing bytes read from storage. Compare per-stage slot_ms and shuffle bytes between the two plans rather than comparing bytes scanned.
  • Where does BigQuery put a query's result once the final stage finishes?
    Into a temporary anonymous table that the client pages through, unless the job specifies a destination table. That temporary result also backs the query cache: an identical query text against unchanged tables can be answered from it without re-executing, which is why a repeated dashboard query can come back instantly and free.

saying these in an interview costs you the question

  • Says data is stored on the query workers' local disks
  • Describes execution as one big parallel table scan
  • Confuses Colossus (storage) with Dremel (execution)
  • Assumes bytes scanned fully determines query cost
  • Thinks the serving tree is still the whole execution model

context