BigQuery Architecture Deep Dive: How Google's Serverless Data Warehouse Really Works
BigQuery isn't a database in the traditional sense — it's a separation of storage and compute taken to its logical extreme. This deep dive explains BigQuery's Dremel execution engine, Colossus storage layer, and why understanding the architecture makes you dramatically better at writing efficient queries.
If you've been using BigQuery for a while, you've probably noticed something strange: sometimes a query over 1 TB finishes in 3 seconds, and sometimes a query over 100 GB takes 45 seconds. The price is the same per byte scanned, but the performance is wildly inconsistent. Understanding why requires understanding what BigQuery actually is under the hood — and it's not what most people think.
BigQuery is not a managed database. It's not Redshift or Snowflake with a GCP logo. It's built on two internal Google systems — Dremel (the execution engine) and Colossus (the distributed file system) — that existed long before BigQuery was a product, and the architecture shapes everything about how you use it effectively.
This article covers the full architecture: how queries execute, what slots actually are and why they matter, how storage is structured, why partitioning and clustering aren't just optional hints, and how to think about query performance from first principles. Once you see the shape of the machine, the "why" behind every best practice becomes obvious.
The Fundamental Architecture: Separated Storage and Compute
Traditional databases co-locate storage and compute. An RDS instance has a disk, and the CPU on that machine runs queries against that disk. BigQuery separates them completely.
Storage lives in Colossus, Google's distributed file system. Your tables are sharded across thousands of machines in a Google data center. When you write data to BigQuery, it goes into Capacitor format — a column-oriented file format designed specifically for analytical queries. Think of it like Parquet, but internal to Google.
Compute lives in a pool of workers called Dremel nodes. When you run a query, BigQuery allocates some of these nodes to your query, has them read the relevant data from Colossus, process it, and return the result.
The path looks like this:
Your query
↓
Query planner (figures out which data to read and how)
↓
Job scheduler (allocates slots to your query)
↓
Leaf nodes (read data from Colossus, apply filters/projections)
↓
Intermediate nodes (aggregate, join, sort)
↓
Root node (final result)
↓
Your client
This tree-shaped execution is Dremel. The query fans out from root to leaves, data flows back up to the root, and the result lands in a table or gets returned to you.
What Slots Actually Are
"Slots" is BigQuery's unit of compute. A slot is a unit of CPU, memory, and network capacity. When Google says your on-demand project gets 2,000 slots for free queries, they mean your queries can occupy up to 2,000 units of this compute resource simultaneously.
What does one slot actually do? It's roughly one vCPU worth of execution capacity. A query that reads 1 billion rows and aggregates them might use 1,000 slots for 5 seconds — that's the slot-seconds cost.
The on-demand pricing model is a sleight of hand: you pay per TB scanned, not per slot-second. Under the hood, BigQuery still uses slots — it just lets you borrow from a shared pool without paying for them explicitly. When the pool is congested, your query waits. This is why on-demand queries can be fast on Tuesday at 2am and slow on Tuesday at 2pm.
Slot Reservation
Under BigQuery Editions (Standard, Enterprise, Enterprise Plus), you can reserve a dedicated slot pool. Your queries always get those slots, regardless of what other projects are doing.
-- Check current slot utilization from your billing export
SELECT
job_type,
SUM(total_slot_ms) / 1000 / 3600 AS total_slot_hours,
AVG(total_slot_ms / TIMESTAMP_DIFF(end_time, start_time, MILLISECOND)) AS avg_slots_per_job
FROM `region-us.INFORMATION_SCHEMA.JOBS`
WHERE DATE(creation_time) = CURRENT_DATE()
AND state = 'DONE'
GROUP BY job_type
ORDER BY total_slot_hours DESC;
This query (which you can run in any project with the right IAM) tells you how many slots your jobs are consuming. If your average job uses 200 slots and you have 2,000 in your pool, you can run 10 concurrent queries at full speed. If you need to run 100 concurrent queries, you need more slots or a queue.
Slot Efficiency
Not all queries use slots efficiently. Common slot wastes:
Skewed data in joins. If one join key has 10 million rows and all others have 10, one leaf node gets 10 million rows and all others get 10. The query finishes when the slowest worker finishes. 10 million vs. 10 is extreme; 10,000 vs. 5 happens in real data.
Shuffle operations. When a query needs to move data between workers (e.g., for a GROUP BY on a column that isn't the partition column), it shuffles data. Shuffle is expensive in both slots and time.
Unnecessary columns. Reading 50 columns when you need 3. BigQuery is columnar; you only pay for the columns you read. SELECT * costs significantly more than SELECT id, name, amount.
The Columnar Storage Model
Capacitor stores data in columns, not rows. This is the single most important architectural decision for understanding BigQuery performance.
In a row-based store, if you have a table with 100 columns and want the average of one column, you read all 100 columns for every row and discard 99% of what you read.
In a column-based store, you read only the columns in your query. A table with 100 columns and 1 billion rows costs the same to query on 3 columns as a 3-column table would.
The practical implication: SELECT * is genuinely harmful in BigQuery. Not a style concern — a cost and performance concern. Every column scanned costs money and uses more slots.
-- Bad: reads every column
SELECT * FROM `my_project.analytics.events` WHERE date = '2025-01-15';
-- Good: reads only what you need
SELECT user_id, event_type, created_at
FROM `my_project.analytics.events`
WHERE date = '2025-01-15';
The second query might scan 1/50th as much data if your table has 50 columns.
Partitioning: How BigQuery Prunes Data
Partitioning divides a table into separate segments — usually by date. When you query with a filter on the partition column, BigQuery reads only the relevant partitions rather than the whole table.
-- Create a partitioned table
CREATE TABLE `analytics.events`
PARTITION BY DATE(created_at)
CLUSTER BY user_id, event_type
OPTIONS (require_partition_filter = false) AS
SELECT * FROM `analytics.events_raw`;
The difference between a partitioned and unpartitioned query:
-- Scans entire 10 TB table
SELECT COUNT(*) FROM `analytics.events` WHERE user_id = '12345';
-- Scans only 1 day's partition (~27 GB in this example)
SELECT COUNT(*) FROM `analytics.events`
WHERE DATE(created_at) = '2025-01-15' AND user_id = '12345';
This is not a performance hint — BigQuery physically reads different files. The partition filter makes it a physically different operation.
Partition Pruning in Practice
Partition pruning only applies when the filter is on the partition column and uses a literal value or a constant expression that can be evaluated at query planning time.
-- Partition pruning works: literal date
WHERE DATE(created_at) = '2025-01-15'
-- Partition pruning works: date function
WHERE DATE(created_at) = DATE_SUB(CURRENT_DATE(), INTERVAL 1 DAY)
-- Partition pruning does NOT work: column comparison
WHERE DATE(created_at) = DATE(some_other_column)
-- Partition pruning does NOT work: UDF in predicate
WHERE DATE(created_at) = my_custom_date_function('some_input')
When in doubt, run EXPLAIN or check the bytes processed estimate before running the query. A query plan that shows 0 B processed after filtering is partition-pruned.
Clustering: Sorted Storage for Range Scans
Clustering is complementary to partitioning. Where partitioning segments the table by date, clustering sorts data within each partition by one or more columns.
-- Well-clustered table for common query patterns
CREATE TABLE `analytics.user_events`
PARTITION BY DATE(event_date)
CLUSTER BY user_id, event_type, country;
If you frequently filter or aggregate by user_id, BigQuery can skip entire blocks within a partition where user_id doesn't match your filter. This is called "block pruning."
-- With clustering on user_id, this reads only the blocks containing this user
SELECT event_type, COUNT(*)
FROM `analytics.user_events`
WHERE DATE(event_date) = '2025-01-15'
AND user_id = 'usr_12345'
GROUP BY event_type;
Clustering is automatic (no syntax to specify block sizes) and maintained on write. BigQuery reclusters data periodically as new data arrives — you don't manage this.
The efficiency of clustering depends on cardinality and selectivity. A column with 100,000 unique values clustered well will prune much more than a column with 5 unique values. Gender as a clustering column helps almost nothing; user_id helps enormously.
The Query Execution Model: From SQL to Dremel
When you submit a SQL query, BigQuery's query planner converts it to a Dremel execution plan. This happens in milliseconds. Understanding the execution tree explains why some queries are fast and others aren't.
Execution Stages
A query execution breaks into stages:
- Scan stage: leaf workers read data from Colossus, applying filters and projections.
- Shuffle stage (if needed): workers redistribute data for joins, aggregations, or sorts.
- Aggregate stage: intermediate workers aggregate results.
- Output stage: final worker writes or returns results.
Every shuffle is expensive because it moves data between workers over the network. Queries with multiple GROUP BYs, joins on multiple keys, or window functions may have multiple shuffle stages.
-- Use EXPLAIN to see the query plan
SELECT * FROM `project.dataset.table`
WHERE column = 'value';
-- EXPLAIN output shows bytes processed per stage and slot estimates
You can view query plans in the BigQuery console under the "Execution details" tab of a completed job. The most useful metrics:
- Bytes read from storage: how much physical data the leaf workers read
- Bytes shuffled: data movement between workers (high numbers indicate shuffle-heavy queries)
- Max slot-ms per stage: which stage was the bottleneck
Join Strategies
BigQuery uses two join strategies: hash joins and broadcast joins.
Broadcast joins happen when one side of a join is small enough to fit in memory on every worker. The small table is broadcast to all workers; the large table is scanned locally. No shuffle needed.
Hash joins happen when both sides are large. BigQuery shuffles both tables by the join key into buckets, ensures matching keys land on the same worker, and executes the join locally. Shuffle is the expensive part.
You can encourage broadcast joins by filtering before joining:
-- Less efficient: full scan of both tables before join
SELECT o.*, u.name
FROM `orders` o
JOIN `users` u ON o.user_id = u.id;
-- More efficient: filter first, join after (optimizer usually handles this)
WITH active_users AS (
SELECT id, name FROM `users` WHERE status = 'active'
)
SELECT o.*, u.name
FROM `orders` o
JOIN active_users u ON o.user_id = u.id;
Modern BigQuery's optimizer is quite good at this, but complex multi-join queries benefit from explicit filtering.
Cost Model: Bytes Scanned, Not Rows
BigQuery's on-demand pricing is $5-6.25 per TB scanned (varies by region). The meter measures bytes read from Colossus before filtering — partition pruning and clustering reduce what's scanned, but WHERE clauses on non-partition, non-clustered columns don't.
-- The meter runs before this filter executes
SELECT * FROM `analytics.events` WHERE country = 'NL';
-- Scans the entire table even if only 1% of rows are from NL
-- unless country is a partition or cluster key
This surprises people from traditional database backgrounds. In PostgreSQL, a WHERE clause with an index skips rows. In BigQuery (unpartitioned/unclustered), you always pay for the full scan.
The practical implication: structure your tables around your most expensive query patterns. The partition column and clustering columns should match the filters your highest-volume queries use.
For more strategies on controlling BigQuery costs, see our BigQuery cost optimization guide.
Materialized Views and Results Caching
Results Caching
BigQuery caches query results for 24 hours. If you run the exact same query twice in 24 hours and the underlying data hasn't changed, the second run is free and instantaneous.
-- This query costs $0 on the second run if nothing changed
SELECT region, SUM(revenue) FROM `sales.orders`
WHERE year = 2024
GROUP BY region;
Cache invalidation happens automatically when underlying tables are updated. You can disable caching per job if you need fresh results for benchmarking.
Materialized Views
Materialized views maintain precomputed aggregates that BigQuery automatically refreshes as base tables change. For common aggregation patterns, they dramatically reduce scan costs.
CREATE MATERIALIZED VIEW `analytics.daily_revenue_mv`
OPTIONS (enable_refresh = true, refresh_interval_minutes = 60)
AS SELECT
DATE(order_date) AS order_day,
region,
SUM(amount) AS total_revenue,
COUNT(*) AS order_count
FROM `sales.orders`
GROUP BY order_day, region;
When a query matches the materialized view pattern, BigQuery rewrites the query to scan the view instead of the base table. The view might be 1,000x smaller than the base table, cutting cost and time proportionally.
Materialized views are not free: they incur storage costs for the precomputed data and slot usage during refresh. For frequently-run aggregation queries over large tables, the tradeoff is almost always worth it.
BigQuery Storage API and External Tables
Storage Read API
The BigQuery Storage Read API lets you read BigQuery tables in parallel from external tools — Spark, Pandas, Arrow — without going through the normal query interface. This enables high-throughput ML training pipelines, ETL, and data export.
from google.cloud import bigquery_storage
client = bigquery_storage.BigQueryReadClient()
table = "projects/my-project/datasets/analytics/tables/events"
requested_session = bigquery_storage.types.ReadSession()
requested_session.table = table
requested_session.data_format = bigquery_storage.types.DataFormat.ARROW
requested_session.read_options.selected_fields = ["user_id", "event_type", "created_at"]
# Parallel streams for distributed reading
session = client.create_read_session(
parent="projects/my-project",
read_session=requested_session,
max_stream_count=10,
)
The Storage Read API is used under the hood by BigQuery DataFrames, the Spark BigQuery connector, and Vertex AI datasets.
External Tables
BigQuery can query files directly in GCS (Parquet, ORC, CSV, JSON, Avro) without loading them. External tables are useful for ad-hoc analysis or as a staging step before deciding whether data deserves to live in native BigQuery storage.
CREATE EXTERNAL TABLE `analytics.events_external`
OPTIONS (
format = 'PARQUET',
uris = ['gs://my-data-lake/events/2025/*/*.parquet']
);
-- Now query it like a regular table
SELECT event_type, COUNT(*) FROM `analytics.events_external`
WHERE _FILE_NAME LIKE '%2025-01%'
GROUP BY event_type;
External table queries are slower than native storage queries and don't benefit from BigQuery's columnar optimization. Use them for occasional access to data that doesn't justify the cost of loading into native storage.
Performance Anti-Patterns to Avoid
After years of BigQuery optimization, here are the patterns that reliably cause problems.
Self-joins on large tables. A self-join — joining a table to itself to do a running total, for example — can be rewritten with window functions at a fraction of the cost.
-- Terrible: self-join for running total
SELECT a.date, a.revenue, SUM(b.revenue)
FROM `sales` a JOIN `sales` b ON b.date <= a.date
GROUP BY a.date, a.revenue;
-- Good: window function
SELECT date, revenue, SUM(revenue) OVER (ORDER BY date ROWS UNBOUNDED PRECEDING)
FROM `sales`;
UNION without UNION ALL. UNION deduplicates rows, requiring a sort/hash operation. If you know your sources are non-overlapping, UNION ALL is much faster and cheaper.
Overly wide CTEs. Common Table Expressions (CTEs) are materialized in BigQuery, not inlined as in some databases. A CTE that produces a billion rows and is joined in two subsequent CTEs will materialize that billion rows twice. Sometimes it's better to use a temp table or restructure.
DATE_TRUNC without partition column. DATE_TRUNC(created_at, MONTH) does not enable partition pruning even if created_at is the partition column. Use DATE(created_at) BETWEEN ... AND ... for partition pruning.
For SQL-based machine learning, see our BigQuery ML guide.
Connecting BigQuery to Looker Studio and BI Tools
BigQuery's BI Engine is an in-memory analysis service that accelerates queries from BI tools like Looker Studio, Tableau, and others. Queries that hit BI Engine skip the Dremel scan entirely — they run against cached, in-memory data.
# Enable BI Engine reservation
gcloud bigquery bi-engine update --reservation-size=10 --location=US --project=my-project
10 GB of BI Engine keeps your most-queried data warm. Dashboard queries that used to take 2-5 seconds now return in under 100ms.
For Looker Studio specifically, enable "Connect to BI Engine" when creating a data source. Looker Studio dashboards that previously scanned GBs per dashboard load become free after BI Engine absorbs the cache hit.
See our Looker Studio enterprise BI guide for patterns we use to build financial and operational dashboards on top of BigQuery.
Putting It Together: Architecture-Aware Query Design
The mental model for fast, cheap BigQuery queries:
-
Start with schema design. Choose partition and cluster columns based on your most frequent, most expensive queries. Change them at the beginning — repartitioning a multi-TB table later is expensive.
-
Select only what you need. Never SELECT * in production queries. List columns explicitly.
-
Partition-prune with every query. Every query that touches a partitioned table should filter on the partition column.
-
Push filters early. Filter before joining. The optimizer usually handles this, but explicit subqueries make it explicit.
-
Use clustering for high-cardinality filter columns. If you query by user_id or session_id frequently, cluster on them.
-
Prefer aggregation over row-level work. BigQuery is built for aggregation over billions of rows. Row-level operations with UDFs or complex per-row logic fight the architecture.
-
Materialize frequently-aggregated patterns. If your dashboard runs the same GROUP BY every hour, a materialized view pays back its overhead in days.
-
Profile before optimizing. Use the Execution Details panel to find the expensive stage, not the expensive query.
BigQuery's architecture is genuinely well-designed for the workload it targets. The more your query patterns align with how Dremel processes data — columnar scans, aggregate over large sets, partition-aware filters — the faster and cheaper it gets.
For cost optimization tactics built on top of this architecture, see our BigQuery cost optimization strategies. For real-time streaming into BigQuery from Pub/Sub, see our BigQuery streaming analytics guide.