Ch.2: The Systems Behind a BigQuery Query
Outline
- 0:00 Four systems, one query
- 0:35 Not one server with disks
- 1:15 Map the pieces
- 1:55 Dremel fans SQL out
- 2:18 Summaries move up
- 2:42 Colossus holds the table
- 3:15 The network is in the query path
- 3:59 Borg places the work
- 4:39 Submission is not execution
- 5:17 The plan shapes the work
- 5:56 Read the columns you need
- 6:33 Partial results move up
- 7:11 Serverless is not physics-free
- 7:51 Rules become consequences
- 8:32 One product boundary
Transcript
0:00 Welcome to Learning Podcasts. BigQuery for Any Software Engineer: Four Systems Working as One. A BigQuery query feels like one product. You paste SQL, press run, and a result table comes back. Underneath, it is not one machine and not even one kind of system. Right. The useful mental model is four systems working together: Dremel for query execution, Colossus for storage, Jupiter for the network, and Borg for scheduling hardware. Once that map is in your head, BigQuery stops feeling like a black-box database and starts feeling like a distributed query service.
0:34 The trap is picturing a database server with disks attached. BigQuery was built around the opposite move. Separate storage from compute, then connect them with enough network bandwidth that the separation still performs. So the table is not sitting on the same worker that runs my query. The worker is temporary, but the storage is durable. Exactly. The query workers can appear for one job, read only the needed data from distributed storage, and disappear when the job finishes. That is why there is no cluster for you to resize before every analysis.
1:08 Right. Capacity exists, but it is managed behind the service boundary. Put the four pieces on the map. Dremel is the query engine. It turns SQL into a tree of parallel work. Colossus is the distributed storage layer, the place BigQuery table data actually lives. Jupiter is the data-center network that lets compute and storage stay separate without making every query crawl. And Borg is the cluster manager that allocates work onto machines. BigQuery is the product boundary around those systems. So when a query runs, it is not "the BigQuery box" doing work.
1:48 It is the engine, storage, network, and scheduler coordinating. Start with Dremel. Dremel is the query execution model behind BigQuery, the part that makes SQL fan out instead of crawl. Its key shape is a serving tree. A root coordinates the query, intermediate mixers aggregate partial results, and leaf workers scan pieces of the data. That tree matters because analytical SQL is usually reducible. Many workers can scan different chunks, compute partial results, and send smaller summaries upward.
2:21 So the root is not reading every row itself. Exactly. The root coordinates. The leaf level does the parallel scan. The mixers reduce fan-in so the final result can come back fast. That is the first big architecture move. Colossus is the storage side. BigQuery table data lives in Google's distributed storage layer, not on a private disk attached to one query worker. That is the storage-compute separation in physical form. Yes. It means storage can scale separately from query execution. A huge table can sit in storage even when no one is querying it.
2:58 And when someone does query it, compute can be assigned for that moment. Exactly. The hard part is making remote storage feel close enough for analytics, which is where the network stops being an implementation detail. If compute and storage are separate, the network becomes part of the query engine. Right. Jupiter is Google's high-bandwidth data-center network. For BigQuery, that network is what lets many workers stream column data from Colossus at the same time. Without that, storage separation would be elegant on a diagram and painful in production.
3:34 It would become a very expensive remote-disk story. Exactly. The network has to carry the scan fan-out, shuffle traffic, and aggregation flow without becoming the bottleneck every time. So when BigQuery feels serverless, it is partly because the network is absorbing the distance between storage and compute. That is the hidden cost of the magic. The 4th piece is Borg, Google's cluster manager. The ancestor idea behind Kubernetes. Right. Borg schedules work across large shared fleets. In the BigQuery mental model, it is the reason query work can be placed onto available machines instead of onto a user-managed cluster.
4:14 So when I press run, I am not leasing one fixed warehouse that sits idle afterward. Exactly. The service allocates execution resources for the query shape it needs. That is a different operational contract. You manage queries and data layout; the platform manages machine placement. Now trace a query. It starts as SQL sent through the console, API, client library, scheduler, or pipeline. Same service boundary, many entry points. BigQuery receives the job, checks metadata, validates permissions, and starts planning the execution.
4:51 The important part is that submission is not execution yet. Exactly. At this point, the service knows what tables are involved, which columns are referenced, what filters exist, and what output shape the query wants. That is already enough to avoid some work before scanning anything. Planning determines which parts of the table are even candidates for reading. The planner turns the SQL into stages: scan, filter, join, aggregate, sort, whatever the query requires. And Dremel maps that into parallel work across the serving tree.
5:24 Yes. Leaf workers read column chunks. Mixers combine partial results. The root assembles the final answer. Here is where it gets really interesting. The amount of work is shaped by the columns and filters, not by the number of rows you eventually display. Exactly. A query can return 10 rows and still scan a huge amount if the plan has to read huge columns. That connects directly to the cost model from the first video. During execution, workers stream the needed data from Colossus. Not whole rows by default.
5:56 BigQuery's storage format is columnar, so analytical queries can read selected columns across many rows. The details of that format come later, but the architecture implication matters now. Dremel wants many leaf workers reading relevant column chunks. Colossus stores those chunks. Jupiter moves them fast enough that the workers can stay fed. So storage, network, and query engine are all in the critical path. If any one of them is weak, the serverless illusion breaks. After workers scan their slices, the result does not come back as one giant unstructured pile.
6:34 The tree shape matters again. Leaf workers compute partial results. Mixers combine those partials closer to where the work happened. The root receives a smaller, more organized stream. That is why aggregation-heavy queries fit the model so well. Count, sum, average, group by. Each leaf can do part of the job. Joins and shuffles are harder, but the same idea holds: split work, move intermediate data, combine results. The service is constantly trying to reduce what has to flow upward. The operational payoff is that you do not capacity-plan BigQuery the way you capacity-plan a classic database cluster.
7:14 Storage can grow without keeping compute hot. Compute can spike for a query without permanently attaching machines to that table. That is powerful, but it moves responsibility. You still have to write queries that scan the right columns, filter the right partitions, and avoid wasteful stages. Exactly. Serverless does not mean physics-free. Oh, wow. That is the line people need before the bill arrives, especially because the interface makes expensive work feel frictionless. BigQuery hides machine management.
7:47 It does not hide the cost of a bad execution shape. This four-system model explains several BigQuery behaviors that otherwise feel arbitrary. Like why selecting fewer columns matters so much. Yes. Columnar storage plus distributed scans. It explains why huge jobs can start without you provisioning a cluster. Borg and the shared fleet. It explains why network and shuffle show up in performance conversations, even though the product feels like SQL. And it explains why query shape matters more than the size of the result table.
8:21 I love that. The model turns random rules into consequences. Exactly. You are not memorizing tips. You are seeing the architecture leak in useful ways. BigQuery is one product boundary around four major pieces: Dremel, Colossus, Jupiter, and Borg. Dremel plans and executes the query tree. Colossus stores the data. Jupiter moves data fast enough for separation to work. Borg places the work on shared compute. The user sees SQL and a result table. The platform is coordinating a distributed execution system underneath.
8:58 Next, we look at slots: the compute currency behind query execution, why adding resources is not a manual knob, and how pricing models change the way teams think about capacity. Thanks for listening to Learning Podcasts.