4.7 Data Storage for Analytics
Data warehouses are usually relational, because SQL fits analytical queries well, and many graphical tools generate SQL and support drill-down and slicing and dicing.
Data warehouses are usually relational, because SQL fits analytical queries well, and many graphical tools generate SQL and support drill-down and slicing and dicing.
On the surface a warehouse and an OLTP database look similar — both have a SQL interface. But the internals look quite different, because they're optimized for very different query patterns. Many vendors now focus on one or the other.
HTAP systems (SQL Server, SAP HANA, SingleStore) support both — but they are increasingly becoming two separate storage and query engines that happen to be accessible through a common SQL interface.
7.1 The unbundled cloud warehouse
Established vendors (Teradata, Vertica, SAP HANA) offer on-prem and cloud. Cloud-only warehouses (BigQuery, Redshift, Snowflake) exploit object storage and serverless computation, integrate with cloud services (automatic log ingestion; Dataflow, Kinesis), and are more elastic because they decouple query computation from the storage layer — data in object storage rather than local disks, so storage capacity and query compute are adjusted independently.
Open source warehouses have broken apart into four separable components — this is one of the most useful mental models in the chapter:
| Layer | Responsibility | Examples |
|---|---|---|
| Data catalog | Which tables are in the database. Usually a standalone service with a REST interface — decoupling it lets data-discovery and governance systems read it too. | Polaris (Snowflake), Unity Catalog (Databricks), Iceberg catalog |
| Table format | Which files constitute a table, plus the schema. Adds row inserts/deletes over immutable files, time travel, garbage collection, even transactions. | Apache Iceberg, Delta Lake |
| Storage format | How rows are encoded as bytes in a file, in object storage or a distributed filesystem — readable by the query engine and by other applications. | Parquet, ORC, Lance, Nimble |
| Query engine | Parse SQL, optimize into a plan, execute. Some have built-in task execution; others use Spark or Flink. | Trino, Apache DataFusion, Presto |
7.2 Column-oriented storage
The motivating query (Example 4-1): are people more inclined to buy fresh fruit or candy depending on the day of the week?
SELECT dim_date.weekday, dim_product.category,
SUM(fact_sales.quantity) AS quantity_sold
FROM fact_sales
JOIN dim_date ON fact_sales.date_key = dim_date.date_key
JOIN dim_product ON fact_sales.product_sk = dim_product.product_sk
WHERE dim_date.year = 2024
AND dim_product.category IN ('Fresh fruit', 'Candy')
GROUP BY dim_date.weekday, dim_product.category;It touches a huge number of rows but only THREE columns of fact_sales: date_key, product_sk, quantity. Fact tables are often over one hundred columns wide, and SELECT * is rarely needed for analytics.
Why row storage fails here: even with indexes on date_key/product_sk, a row-oriented engine still loads all those rows (each with 100+ attributes) from disk into memory, parses them, and filters.
| Row-oriented | Column-oriented |
|---|---|
r1: date|prod|store|promo|cust|qty|price|cost|… 100+ … | date_key: [140,140,140,141,141,…] |
| Reading a row reads all 100 columns ⇒ about 97% of the I/O is wasted. | Reads only the three columns you asked for ⇒ about 3% of the bytes. |
The layout relies on each column storing the rows in THE SAME ORDER. To reassemble a row, take the 23rd entry from each column.
In practice engines don't store a whole column (trillions of rows) in one go. They break the table into blocks of thousands or millions of rows, and within each block store each column separately. Since many queries are restricted to a date range, it's common to make each block contain rows for a particular timestamp range — then the query loads only the needed columns from the blocks overlapping the required range. (This is zone-map / block-skipping territory.)
Where it's used: almost all analytical databases — Snowflake, DuckDB (single-node embedded), Pinot and Druid (product analytics); formats Parquet, ORC, Lance, Nimble; in-memory formats Apache Arrow, Pandas/NumPy; time-series InfluxDB IOx, TimescaleDB.
Note: columnar applies to non-relational data too — Parquet supports a document data model (based on Google's Dremel) using shredding/striping.
⚠️ Do NOT confuse column-oriented databases with the WIDE-COLUMN (column-family) data model, where a row can have thousands of columns and rows needn't share columns. Despite the name, wide-column databases are ROW-oriented — they store all values from a row together. Examples: Bigtable, Accumulo, HBase.
7.3 Column compression
A product_sk column with billions of rows but only 100,000 distinct values:
[69, 69, 69, 69, 74, 31, 31, 68, 69, 69, 74, 31, 31, 31, 68, 69]
BITMAP ENCODING — one bitmap per distinct value, one bit per row:
product_sk=31: 0 0 0 0 0 1 1 0 0 0 0 1 1 1 0 0
product_sk=68: 0 0 0 0 0 0 0 1 0 0 0 0 0 0 1 0
product_sk=69: 1 1 1 1 0 0 0 0 1 1 0 0 0 0 0 1
product_sk=74: 0 0 0 0 1 0 0 0 0 0 1 0 0 0 0 0Bitmaps are usually sparse — mostly zeros — so they are run-length encoded: product_sk=31 becomes 5 zeros, 2 ones, 4 zeros, 3 ones, 2 zeros → 5,2,4,3,2.
Why this is so effective: the number of distinct values in a column is often small compared to the number of rows (a retailer has billions of sales but only 100,000 distinct products).
Roaring bitmaps switch between the plain-bitmap and run-length representations, using whichever is more compact.
Why bitmaps are perfect for warehouse queries:
WHERE product_sk IN (31, 68, 69) — load three bitmaps and bitwise OR them. Extremely efficient.
WHERE product_sk = 30 AND store_sk = 3 — load two bitmaps and bitwise AND them.
This works because the columns contain rows in the same order: the k-th bit of one column's bitmap is the same row as the k-th bit of another's.
Bitmaps also answer graph queries — e.g. find all users of a social network followed by user X who also follow user Y.
7.4 Sort order in column storage
Rows can be stored in insertion order (then inserting = appending to each column). But you can impose an order and use it as an indexing mechanism, as with SSTables.
⚠️ Sorting each column INDEPENDENTLY would make no sense — you'd no longer know which items belong to the same row. The data must be sorted an entire row at a time, even though it's stored by column.
Choosing sort keys:
- First sort key: pick from knowledge of common queries. If queries often target date ranges (last month), make
date_keyfirst — the query then scans only last month's rows. - Second sort key determines order among rows tying on the first. If
date_keyis first,product_sksecond groups all sales of the same product on the same day together, helping queries that group or filter by product within a date range.
Sorting also boosts compression — but unevenly:
| Column | Values | Compression |
|---|---|---|
Sort key 1 — date_key | 140 140 140 140 141 141 141 142 142 142 … | Long runs → run-length encodes to a few kilobytes, even for billions of rows |
Sort key 2 — product_sk | 31 31 69 74 31 69 69 31 68 74 … | More jumbled, so shorter runs |
| Sort key 3 and beyond | Essentially random order | Probably will not compress well |
The compression effect is strongest on the first sort key. Still, having the first few columns sorted is a win overall.
7.5 Writing to column-oriented storage
Writes in a warehouse tend to be bulk imports, often via ETL.
Why individual row inserts are terrible here: writing a row in the middle of a sorted table means rewriting all the compressed columns from the insertion position onward. But a bulk write of many rows at once amortizes that cost.
The solution is the same log-structured approach as LSM:
Done by Snowflake, Vertica, Apache Pinot, Apache Druid, and many others.
7.6 Query execution: compilation vs vectorization
A complex analytical SQL query becomes a query plan of stages called operators, possibly distributed across machines for parallel execution. The planner optimizes which operators, in what order, and where each runs.
The problem: for queries scanning millions of rows, it's not just disk bytes — it's CPU time. The simplest operator is an interpreter: while iterating over each row, it consults a data structure representing the query to find out which comparisons/calculations to perform on which columns. Too slow for analytics.
Two solutions, both used in practice:
| Approach | Mechanism |
|---|---|
| Query compilation | The engine generates code for executing the query — iterate rows, look at the columns of interest, perform comparisons, copy values to an output buffer if conditions hold. Compile that generated code to machine code (often via LLVM) and run it on column-encoded data in memory. Similar to JIT compilation in the JVM. |
| Vectorized processing | The query is interpreted, not compiled, but made fast by processing many values from a column in a batch rather than row by row. A fixed set of predefined operators is built into the database; pass arguments, get back a batch of results. |
Vectorization example:
Both exploit the same four CPU characteristics (this list is the real payoff of the section):
- Preferring sequential memory access over random access, to reduce cache misses
- Doing most work in tight inner loops — few instructions, no function calls — to keep the instruction pipeline busy and avoid branch mispredictions
- Parallelism: multiple threads and SIMD instructions
- Operating directly on COMPRESSED data without decoding it into a separate in-memory representation, saving memory allocation and copying
7.7 Materialized views and data cubes
| Virtual view | Materialized view | |
|---|---|---|
| What it is | A shortcut for writing queries | An actual copy of the query results, written to disk |
| On read | SQL engine expands it into the underlying query on the fly and processes the expanded query | Read the stored copy |
| On underlying data change | nothing to do | Must be updated — more work on writes |
Some databases update them automatically; Materialize is a system specializing in materialized view maintenance.
Materialized aggregates / data cubes (OLAP cubes). Warehouse queries often use COUNT, SUM, AVG, MIN, MAX. If many queries use the same aggregates, crunching raw data each time is wasteful. A data cube creates a grid of aggregates grouped by different dimensions.
| date ↓ / product → | 31 | 68 | 69 | 74 | Total |
|---|---|---|---|---|---|
| 140 | 149.6 | 31.0 | 238.9 | 1.20 | 420.7 |
| 141 | 132.9 | 253.9 | 38.3 | 482.1 | 907.2 |
| 142 | 191.4 | 818.7 | 213.5 | 767.8 | 1991.4 |
| Total | 473.9 | 1103.6 | 490.7 | 1251.1 | 3319.3 |
Each cell is SUM(net_price) of all facts with that date–product combination; the right column is sales by date regardless of product, and the bottom row sales by product regardless of date. Real fact tables have more dimensions — date, product, store, promotion, customer — giving a five-dimensional hypercube, summarizable along each axis.
| Data cube | |
|---|---|
| ✔ Advantage | Certain queries become very fast because they're effectively precomputed. Total sales per store yesterday = read a total along one dimension — no need to scan millions of rows |
| ✗ Disadvantage | No flexibility. You cannot compute "what proportion of sales came from items costing more than $100," because price isn't one of the dimensions |
Most data warehouses therefore keep as much raw data as possible and use aggregates like data cubes only as a performance boost for certain queries.