BigQuery Ch.1: Columnar Storage

Outline

Transcript

0:00 Welcome to the Deep Dive. Today we are digging into a really fascinating stack of sources. Yeah, specifically we're pulling from chapter one of the series, BigQuery for any software engineer. Right, and our mission for this Deep Dive is to, well, basically solve a very specific, very painful mystery that a lot of you might have experienced. Oh, absolutely. The classic BigQuery shock. Exactly. So I want you, the listener, to picture this. You're a senior backend software engineer. You know your way around PostgreSQL.

0:30 You understand how to architect distributed systems. You know what you're doing, basically. Right, you know what you're doing, and you're debugging a pipeline, so you just write a basic everyday query. You know, select from events where user 1, 2, 3, 4, 5. Just does the totally standard debugging query. Totally standard. You hit execute, and a few seconds later you get 10 rows back. Exactly what you expected. So far, so good. Yeah, until you happen to glance over at your Google Cloud billing dashboard, and your stomach just completely drops.

1:00 Because the query engine just scanned, what, 800 gigabytes of data? 800 gigabytes. I mean, 800 gigabytes to retrieve 10 rows. Yeah, that is the exact moment. Most engineers who are, you know, used to traditional relational databases, they just want to unplug their router and walk away. I would. Right. Because the ratio of data retrieved to data scan there, it just feels fundamentally broken. It really does. Because, okay, if I have a B tree index on that user ID column in Postgres, I'm doing a quick index lookup, right?

1:31 I grab a few pages from the disk, and I'm done. The I.O. is almost nothing. Exactly. The I.O. is minimal. So, if BigQuery is scanning nearly a terabyte of data just to find 10 records, it means, well, everything you know about indexing and row retrieval just, it doesn't apply here. No, it doesn't. It feels like a completely different paradigm. Because it is a completely different beast. I mean, the reason that 800 gigabyte bill feels so deeply wrong to you is because you're applying a mental model built for transactional systems to, well, a massive analytical processing engine.

2:03 So we have to unlearn some things. You do. To understand why BigQuery behaves the way it does, and crucially, why it costs what it costs, we need to completely tear down that row-oriented mindset. Okay, let's get right into the physical mechanics of this then. In a traditional system like Postgres or MySQL, data is row-oriented. We all know how this works, right? The entire row is stored contiguously on the disk. Right, all bundled together. Yeah. When I fetch a user record, I pull the whole record, the ID, the email, the timestamp, the status into memory all at once.

2:39 Because it's sitting together physically on the storage drive. And that architecture is, you know, perfectly optimized for transactional workloads where you're frequently reading or inserting or updating entire individual records. Right. But BigQuery is a distributed query engine built on columnar storage, which means it flips that physical layout completely sideways. Sideways. Okay. So how should we visualize that? Well, think of a spreadsheet. In Postgres, you read across the rows. But BigQuery weeds down the columns.

3:06 Instead of storing all the fields for one event together, it takes a single column, say the event type, and it stores all the values for that column across millions of rows contiguously. Oh, wow. So if I'm visualizing this as a distributed systems engineer, it's almost like a microservices architecture applied to storage. That's a great way to put it, actually. Because instead of a monolithic row where all the data is bound together, every single column is effectively its own independent physical file or, you know, series of files on the disk.

3:37 That is a much more accurate way to model it. So if your events table has 50 columns, BigQuery essentially treats those as 50 separate data vectors. Okay. When you write a query asking for just three of those columns, the engine completely ignores the other 47. It goes to the storage layer, locates the specific files for just those three columns, and streams them from top to bottom. Wait, from top to bottom. So it's doing a full table scan on those specific columns. Yeah. It's not like doing an index lookup to find the 10 rows first.

4:06 In a purely unoptimized table, yes, it reads the entire column. That sounds horribly inefficient, though. It sounds that way, but this is the critical architectural tradeoff Google made. You have to remember, when you're operating at petabyte scale, doing random disk seeks via a B-tree index across distributed storage nodes, that becomes a massive latency bottleneck. Oh, because of the network overhead. Exactly. It's actually faster to brute force read a massive contiguous block of column data using thousands of parallel workers than it is to coordinate millions of tiny random reads over a network.

4:42 So if I need to calculate the average of an N64 column across 10 billion rows? Contiguous column or storage is basically the only way to do it efficiently. Okay. That completely redefines the cost equation for me? Right. Because if it's reading the columns top to bottom, then the fact that my where clause eventually filtered the result set down to 10 rows, that's totally irrelevant. Completely irrelevant to the billing. The cost isn't based on the output. The cost is purely based on the total number of bytes the engine had to physically read from those column files before it applied the filter.

5:13 You've hit the nail on the head. I mean, you could return 10 rows or 10 million rows. The financial cost is identical because the physical disk I.O. is identical. That is wild. But here's the silver lining. If you only actually need userid, event timestamp, and event type, querying just those three columns drops your scanned bytes by like 90% compared to reading all 50 columns. And therefore drops your bill by 90%. Exactly. Which brings us to probably the most dangerous muscle memory a backend engineer has, the select trap.

5:46 Oh, yeah. The silent killer. Because in a row-oriented database, select is a lazy habit. Sure. We know we shouldn't use it in production because it sends unnecessary data over the wire to the application. But from a disk I.O. perspective, it's effectively free, right? Right. Because the engine has to pull the whole row block off the disk into memory anyway. Grabbing the extra columns doesn't trigger a massive penalty. But think about what select commands a columnar engine to do. It's a disaster.

6:14 You're explicitly telling BigQuery, go find the physical files for every single one of my 50 columns and read all of them from top to bottom. You're forcing it to scan the maximum possible amount of data. Which means you are paying the maximum possible price. It is the single most expensive habit you can have in this environment. I have to push back here, though, just thinking about developer velocity. Okay, lay it on me. Let's say I have a table with 60 columns. Yeah. And I actually need 40 of them for a complex debugging session.

6:44 Explicitly typing out 40 column names in a fast-moving, agile environment, that's a nightmare for readability and developer time. It is tedious. I'll give you that. Right. And isn't developer time usually more expensive than a few extra bytes scanned on a cloud bill? Why punish the user for just exploring the data? That's a very fair point regarding the developer experience. And, you know, there are ways around it, like creating focused views or using explicit EXCEPT clauses in your SQL. But the financial penalty isn't just arbitrary.

7:14 It's a direct reflection of the physical hardware utilization. Okay, break that down for me. Let's use an analogy for the financial mechanism here. Traditional databases are like an all-you-can-eat buffet where you pay a flat entry fee at the door. Like you provisioned the AWS EC2 instance you paid for. Exactly. Once you're inside, filling your plate with every dish, doing a select a fin, it doesn't cost you extra money, even if you only take two bytes. Because the server is already running. The overhead is a sunk cost.

7:44 Right. But BigQuery's pay-per-scan model is like a buffet where you pay strictly by the physical weight of the food you put on your plate at the scale. Oh, man. That's a great analogy. So select is the equivalent of dumping the entire buffet line onto your tray, taking it to the scale, paying for 100 pounds of food, and then eating a single dinner roll. The engine had to do the heavy lifting to move those bytes over the network. Google is charging you for the lifting. That buffet analogy makes the financial pain very, very real.

8:15 And there's a specific quirk mentioned in the source material that makes it even more insidious. Oh, the wide column problem. Yes. Let's say my team has a highly optimized BigQuery table. We query it every day. Then, an engineer working on a totally different feature decides to alter the table and add a new, incredibly wide column to store giant JSON blobs of raw HTTP request payloads. Yeah, that happens all the time. In a row-oriented system, adding that column might cause some page fragmentation.

8:46 But if existing queries are doing index lookups, the cost impact is minimal. Almost unnoticeable. But in BigQuery, because of column or storage, adding that heavy column retroactively increases the financial cost of every single future select query run against that table. Every single one. Even if the person running the query has absolutely no idea the new column exists, they're just running the same select dashboard they run every morning, and suddenly the bill spikes because they're unwittingly scanning gigabytes of JSON blobs they aren't even using.

9:15 Yeah, it's brutal. The storage model directly shapes the cost model. You cannot decouple the two. This is why strict data governance and avoiding select aren't just best practices in BigQuery. They are mandatory financial controls. Okay, so why does Google charge by the byte scanned in the first place? If I'm renting an EC2 instance with a massive SSD for Postgres, I pay by the hour. I can run a million queries or narrow queries. The hourly rate is the same. Why does BigQuery use this weight-based scale?

9:45 To understand the billing, we have to look under the hood at how the hardware is actually organized. In your Postgres example, storage and compute are inextricably linked. Right. The CPU that executes the query and the SSD that holds the data are sitting in the same physical metal box. Exactly. BigQuery completely severs that tie. Wait. If compute is decoupled, then I'm not paying for an idle CPU at 3 in the morning. Radically decoupled. The data does not sit on the same machines that process your SQL.

10:11 Your data sits in a massive, globally distributed storage system called Colossus. Colossus. So Colossus, Colossus is essentially the hard drive of the operation. It just holds the independent columnar files. Yes. And the compute, the brain parsing your SQL and crunching the numbers, is handled by a completely separate distributed query engine called Dremel. Dremel. Okay, so when I hit execute on my 10-row query, what is Dremel physically doing? Because it has to get the data out of Colossus. Colossus somehow.

10:42 Dremel analyzes your query and dynamically spins up processing units called slots. Slots. Yeah. Think of a slot as a tiny temporary worker node, essentially a fraction of a CPU and some RAM. Depending on the complexity of your query, Dremel might instantly allocate 50 slots, or it might allocate 5,000 slots from a massive shared pool managed by Google. I see. And then those thousands of slots have to reach across Google's data center network to pull the relevant column data out of Colossus, Colossus, in parallel.

11:13 That is the mechanism. They read the data over the network, perform the filtering and aggregations in memory, and return the result. Because these slots are allocated dynamically for the exact millisecond you need them, you never provision servers. Wow. You never manage disk space. And crucially, you never pay for idle time. Which explains the billing model perfectly. If I'm not paying for an idle CPU, I'm literally just renting network bandwidth and disk reads on demand. Google is charging me for the immense network IO required to move data from Colossus, Colossus, to the Dremel slots.

11:51 That is the core realization. The standard on-demand pricing is $6.25 per TB byte scanned. Per tebibyte. tebibyte. Right. And a TB byte is the base 2 equivalent of a terabyte, so it's about 10% larger than a standard metric terabyte. That $6.25 isn't a random number. It is the price of renting an army of parallel workers and the massive network pipes needed to scan petabytes of data in seconds. The scale there is just wild. A small startup query might use 50 slots. A Fortune 500 company running a 10 terabyte aggregation might seamlessly use thousands of slots.

12:28 And the user never configured a single scaling policy. It just happens. But that power is exactly why you have to control how many bytes those slots are pulling over the network. Which brings us to the obvious architectural problem. If I am strictly charged by the bytes scanned over the network, and we know we shouldn't use select, what happens when I do just need three columns? But those columns represent 10 years of company history. Right. The time series problem. Exactly. If I'm querying an events table, how do I avoid reading an entire decade of data just to find what happened last Tuesday?

13:02 We already established there are no traditional B-tree indexes to jump me right to Tuesday. This is where we introduce the mechanisms of partitioning and clustering. Because if you cannot use an index, you have to physically organize the data so the engine knows what files to ignore. Okay, let's start with partitioning. If my columns are already stored as separate files, what does partitioning do to those files? It divides them further based on a specific column, usually a timestamp or date. So instead of one giant contiguous file for your event type column containing 10 years of data, partitioning chops that column into separate, isolated files for every single day.

13:37 Ah, I see. So if I have my table partitioned by day and my where clause explicitly asks for an event timestamp greater than or equal to last Tuesday, Dremel. Dremel looks at that sequel and realizes it doesn't even need to ask Colossus for the files from Monday or Wednesday or the previous nine years. Exactly. It completely prunes them from the execution plan. Those bytes never travel over the network, the slots never process them, and you never pay for them. Partitioning is a coarse-grained mechanism, but it drops your bytes scanned by absolute orders of magnitude.

14:12 It's like lopping off massive chunks of the iceberg before you even start digging. That's a good way to look at it. So how does clustering fit in? Because if partitioning isolates the data by day, what if I only want one specific user's events on that specific day? Well, if partitioning chops the data into daily blocks, clustering sorts the data within those daily blocks based on up to four columns you specify. Let's say you cluster by user it. Behind the scenes, BigQuery maintains lightweight block metadata.

14:39 Block metadata. Yeah, essentially keeping track of the minimum and maximum user ID values contained in each chunk of data within that partition. That makes sense. So when Dremel is reading last Tuesday's partition, it checks that block metadata first. If a chunk of data only contains user IDs from 1,000 to 2,000, and I'm looking for user ID 5,000, the Dremel slot just skips that entire chunk. It bypasses the read entirely. So partitioning drops the sledgehammer, cutting out years of data you don't need, and clustering comes in with the scalpel, using block metadata to trim down the exact day's data even further.

15:18 Mastering both mechanisms is how senior engineers keep that $6.25 per tebibyte under strict control. I have to say, even understanding the mechanics now, the separated architecture, the Dremel slots, the network reads from Colossus, Colossus, it still sounds intimidating for an engineer who just wants to poke around and learn. It's definitely a mindset shift. Right, because one missing where clause on an unpartitioned petabyte table, and you've accidentally incurred a massive bill. That is a very rational fear when dealing with serverless architecture, which is why the BigQuery free tier is so critical for engineers transitioning their mental models.

15:57 Yes, we need to highlight this. The source material emphasizes that you get one tebibyte of queries and 10 gigabytes of storage per month, completely free. Yep. And it is an ongoing free tier, not a limited time trial. It resets every single month. It provides a massive safety net. It allows you to load in public datasets, practice your columnar thinking, experiment with partitioning schemas, and actually observe how Dremel, allocates slots without constantly looking over your shoulder at the billing dashboard.

16:28 It's basically a sandbox where you can test the physics of this distributed environment. Exactly, without the financial anxiety. So let's synthesize everything we've uncovered from this first chapter. The fundamental shift you need to make as a backend engineer is profound. You have to stop visualizing your database as a monolith of rows. Goodbye, B-trees. Right. You are no longer navigating via B-Tree indexes. You are working with isolated column vectors stored on a massive distributed file system called Colossus.

16:59 Colossus. And when you write SQL, you aren't just filtering data. You are orchestrating an army of Dremel workers to pull bytes across a network. Which means selecting fewer columns, partitioning your data, and clustering your blocks means moving fewer bytes, which directly translates to a lower bill. Once you internalize the physical reality of the hardware, you start fighting the query engine and start leveraging its actual strength. And we are truly just scratching the surface of this series.

17:26 Oh, definitely. There is so much more. In our upcoming deep dimes into the later chapters of BigQuery for any software engineer, we are going to get into the really heavy machinery. We will be exploring the four Google infrastructure systems that make all of this possible, including the cluster management system called Borg, and the massive data center network called Jupyter. We will also unpack the exact physical layout of the capacitor Kupas EETTOR columnar storage format, dive into the mechanics of the shuffle service used for massive table joins, and break down the complete end-to-end execution of a query graph.

18:01 It is going to be incredibly dense and incredibly fun. But before we wrap up today, I want to leave you, the listener, with a broader architectural thought. We've talked exclusively about how this decoupled system changes the way you write a single SQL query. But there is a much bigger implication here. There really is. If the very act of storing data separately from the compute engine completely flips the financial and performance model of retrieving information, how might that reality fundamentally change the way you architect your application's entire data pipeline from the moment you write your first line of code?

18:33 That's a great question. Right. If storage on Colossus, Colossus, is dirt cheap, but compute is paper scan, how does that change what telemetry you log? How does it change the way you structure your microservices or design your event streaming ingestion APIs? It changes the entire foundation. You aren't just tuning a database anymore. You are designing for a fundamentally different paradigm of distributed computing. So the next time you hit execute on a simple query and you get 10 rows back, take a second, look at the byte scan, remember the columnar files, remember the network pipes, and remember that those 10 rows are just the very tip of an incredibly powerful globally distributed iceberg.

19:13 Thanks for joining us on this deep dive.