skip to content

Why does BigQuery keep table data in Colossus instead of on the query workers' disks?

level: middleimportance: should knowfreq 56%

answer

  1. nodes that own data are hard to resize
  2. workers keep nothing durable
  3. the network stands in for local disk
  4. columnar reads move a small slice of the table
  5. storage bills even when nothing queries

basics

~20 s

Keeping data in Colossus makes compute stateless: any worker can serve any query, capacity can change mid-query without rebalancing data, and storage persists and is billed independently. Google's Jupiter network plus columnar files make the remote reads fast enough.

solid answer

~50 s

In a classic shared-nothing warehouse, each node owns a slice of the data, so adding or removing nodes means physically redistributing rows, and idle nodes still hold data you cannot detach from. BigQuery breaks that coupling: tables live in **Colossus**, Google's distributed file system, and query workers own nothing. That makes compute fungible — the scheduler can hand your query any free slots, grow a stage's parallelism on the fly, and reclaim everything afterwards. It also means storage survives and is billed whether or not you ever run a query, and a huge table costs nothing extra in compute until you touch it. The design only works because two things offset the loss of data locality: the **Jupiter** datacenter network provides enough bandwidth to keep workers fed, and the columnar format means a query typically pulls a small fraction of the table's bytes across it.

go deeper

for a junior

Know the headline: table data lives in Google's storage service, not on the machines that run your query, so storage and query capacity are separate things you think about separately.

for a middle

Explain why statelessness matters — no data to rebalance when capacity changes, any worker can run any stage, storage persists and bills on its own — and why a fast network plus columnar files make remote reads affordable.

for a senior

Bring the operational consequences: no local block cache means repeated similar scans keep paying, cost tracks bytes read rather than table size, and idle projects still carry storage charges that migrations routinely forget.

for a principal

Be able to compare disaggregation strategies across platforms and argue the trade you are buying: elasticity, multi-tenancy and independent cost axes in exchange for lost locality and less predictable per-query latency.

## The coupling this design removes A traditional MPP warehouse assigns each table's rows to specific nodes. That gives excellent locality — a scan reads local disk — but three things get welded together: 1. **Capacity and data placement.** Adding a node means moving data onto it. Removing one means moving data off. Resize becomes a data-movement operation measured in hours. 2. **Idle cost and durability.** You cannot switch off the compute without losing access to the data, because the data is on the compute. 3. **Concurrency and hardware.** More concurrent users means more nodes, which means yet more data redistribution. BigQuery separates the two axes entirely. Colossus holds the columnar files; the execution service holds nothing durable. Every query reads what it needs over the network and forgets it. ## What statelessness buys - **Elastic parallelism per stage.** Because no worker owns a partition of the table, the number of workers on a stage is a scheduling decision, not a data-placement fact. A stage can run 40 slots wide or 4,000 depending on input splits and availability. - **No resize operation.** There is nothing to rebalance, so there is no maintenance window, no vacuum-after-resize, no skewed-node aftermath. - **Fault tolerance without replication of compute state.** A dead worker's work unit is simply re-run somewhere else; its input is still sitting in Colossus or in shuffle. - **Multi-tenancy.** The same physical fleet serves many customers, because a machine is not tied to any customer's data. That is what makes the on-demand model economically possible. - **Independent cost axes.** Storage is charged for holding bytes; query execution is charged for the work done. A rarely queried archive costs storage only. Long-untouched data automatically moves to a cheaper long-term storage rate, with no action from you. ## Why remote storage isn't a disaster for performance The usual objection is that a network read is slower than a local disk read. Two properties change the arithmetic: **The network is not the 1 Gb link you're imagining.** Google's Jupiter fabric provides very high bisection bandwidth inside a datacenter — enough that a worker pulling columnar data from Colossus stays CPU-bound rather than IO-starved. Disaggregation is only a good idea when the network is engineered for it. **Columnar layout shrinks the transfer.** A query that touches 3 columns of a 200-column table moves roughly 1.5% of the row bytes, before compression. Partition metadata removes whole date ranges before a byte is read. So the bytes actually crossing the network are a fraction of the logical table size. What you *do* give up is a local block cache. In a node-owning warehouse, hot blocks sit on the node's SSD and a repeat scan is nearly free. BigQuery instead offers reuse at a different level: an identical query over unchanged tables can be answered from the **cached result** without re-executing, and there is a separate in-memory acceleration layer for BI-style repeated queries. Neither is a block cache, so a *slightly different* query re-reads the data. ## The practical consequences to state - **A query's cost is driven by bytes it must read, not by how big the table is.** Column pruning and partition filters are your first-order levers precisely because every read is a fresh remote read. - **Nothing is warm.** There is no advantage to "keeping the cluster hot" because there is no cluster and no resident data. Repeated identical queries benefit from the results cache; repeated *similar* queries do not. - **Storage keeps costing money when compute is silent.** Teams migrating from a node-based warehouse often forget this and are surprised by a storage line item on a project nobody queries. - **Copying and sharing are cheap.** Because data lives in one storage service rather than on someone's nodes, giving another project read access to a dataset does not require moving or duplicating anything, and the reader's queries are executed with the reader's own compute. ## Comparing honestly with other designs Other modern warehouses reach the same destination by different routes: some run stateless compute clusters over object storage with a local SSD cache per cluster; some keep a managed storage layer behind node-based compute. The shared idea — durable data in a service, transient compute above it — is the defining architectural move of the cloud-warehouse generation. BigQuery's version is the most extreme, because there is no per-customer compute object at all, and the network is the substitute for local disk. If an interviewer pushes on the weakness, name it: without local data caching, a workload of many slightly different scans over the same hot table pays the read repeatedly, and per-query latency depends on shared slot availability rather than on hardware you control.

  • What does BigQuery cache, given that workers have no local copy of the table data?
    Finished query results. An identical query text over unchanged tables can be served from the cached result without executing or scanning. There is also a separate in-memory acceleration layer aimed at repeated BI queries. Neither is a block cache, so a query that differs even slightly re-reads the columns it needs from storage.
  • What is the main weakness of putting all table data behind the network like this?
    No data locality and no local block cache, so repeated but non-identical scans over the same hot table pay the read every time. Latency also depends on shared slot availability rather than on hardware you own, which makes per-query timing less predictable than a dedicated cluster until you buy capacity.
  • How does this design change what happens when you share a dataset with another team?
    Nothing is copied or moved. The data stays in one place in storage and the other project is granted read access; their queries run on their own compute allocation. Because storage and compute are separate services, the reader pays for their own execution while the owner keeps paying for storage.

saying these in an interview costs you the question

  • Believes each worker owns a shard of the table
  • Assumes a local SSD cache makes repeat scans free
  • Says storage is free when no queries run
  • Thinks scaling compute requires redistributing data
  • Claims remote storage is always slower than local disk here

context