Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. BigQuery architecture: Dremel, Colossus, Jupiter

Snowflake, BigQuery & Databricks · BigQuery

BigQuery architecture: Dremel, Colossus, Jupiter

Easywarehouses-20
bigqueryarchitecturedremelcolossusserverless

Question

How is BigQuery's architecture different from a traditional warehouse?

Solution

BigQuery is serverless. You do not create clusters, size nodes or manage indexes. Behind that, storage and compute are separate systems joined by a very fast network, and compute is a pool of workers that BigQuery assigns to your query for the few seconds it runs.

The pieces

Query -> Dremel execution tree (slots do the work)
            |  Jupiter network + in-memory shuffle
         Colossus (distributed file system)  <- tables in Capacitor columnar format

What each piece does:

  • Colossus: Google's distributed file system. Table data lives there in Capacitor, a columnar format with compression, so a query reads only the columns it needs.
  • Dremel: the execution engine. It splits a query into stages and runs each stage in parallel on many workers, organised like a tree. Leaf workers scan and filter data, and higher levels combine partial results.
  • Slots: the units of compute, a share of CPU and memory, that Dremel uses.
  • Jupiter: Google's data centre network, fast enough that storage and compute can sit apart without the usual penalty.
  • Shuffle: a distributed in-memory service moves intermediate data between stages for joins and aggregations, so workers do not keep it on local disks.

How it differs from a classic warehouse

In a traditional MPP warehouse (older Redshift node types, Teradata), data is spread over nodes, and the nodes also run queries. To add capacity you add nodes, and data has to be redistributed. You also tune distribution keys, sort keys, vacuum and indexes.

In BigQuery, storage and compute scale independently. A query can use thousands of slots for a few seconds and then release them. There are no indexes to build. The levers you do control are table design (partitioning, clustering), what columns and rows you scan, and how you pay (on-demand per bytes scanned, or reserved slots).

What this means for you

Cost and speed are mostly about how much data you read. A simple query on a well-partitioned table is fast and cheap. A careless SELECT * on a huge table is slow and expensive, and no amount of cluster tuning applies. Say this in an interview, because it shows you understand the model.

PreviousNext