Learn Labs
4. Storage and Retrieval

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:

LayerResponsibilityExamples
Data catalogWhich 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 formatWhich 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 formatHow 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 engineParse 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-orientedColumn-oriented
r1: date|prod|store|promo|cust|qty|price|cost|… 100+ …
r2: … 100+ columns …
r3: … 100+ columns …
date_key: [140,140,140,141,141,…]
product_sk: [69, 69, 74, 31, 31,…]
store_sk: [4, 5, 5, 8, 8,…]
quantity: [1, 3, 1, 5, 2,…]
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 0

Bitmaps 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_key first — the query then scans only last month's rows.
  • Second sort key determines order among rows tying on the first. If date_key is first, product_sk second 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:

ColumnValuesCompression
Sort key 1 — date_key140 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_sk31 31 69 74 31 69 69 31 68 74 …More jumbled, so shorter runs
Sort key 3 and beyondEssentially random orderProbably 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:

merge and write new files in bulkrecent writeswritesrow-oriented sorted in-memory storelike a memtableimmutable column-encoded filesobject storage suits this wellquery enginecombines bothQueries examine both the column files and the recent in-memory writes. The engine hides this: to an analyst,inserts, updates and deletes are immediately reflected.
Figure 4.7.6The 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:

ApproachMechanism
Query compilationThe 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 processingThe 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:

product_sk columnequality operator= 'bananas'bitmap Astore_sk columnequality operator= store 42bitmap Bbitwise ANDwhat SIMD hardware is forresultbananas in store 42
Figure 4.7.7Vectorization example

Both exploit the same four CPU characteristics (this list is the real payoff of the section):

  1. Preferring sequential memory access over random access, to reduce cache misses
  2. Doing most work in tight inner loops — few instructions, no function calls — to keep the instruction pipeline busy and avoid branch mispredictions
  3. Parallelism: multiple threads and SIMD instructions
  4. 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 viewMaterialized view
What it isA shortcut for writing queriesAn actual copy of the query results, written to disk
On readSQL engine expands it into the underlying query on the fly and processes the expanded queryRead the stored copy
On underlying data changenothing to doMust 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 →31686974Total
140149.631.0238.91.20420.7
141132.9253.938.3482.1907.2
142191.4818.7213.5767.81991.4
Total473.91103.6490.71251.13319.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
✔ AdvantageCertain 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
✗ DisadvantageNo 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.


On this page