SysDesignPrep.com
Study guide 102 of 183

OLTP versus OLAP: analytics, columnar storage and data warehouses

Why analytics needs a different store from your application database, row versus columnar storage, data warehouses and lakes, star schemas, pre-aggregation and rollups, and real-time analytics with ClickHouse or Druid.

Reading is half of it. See this used in a real interview: walk through Design an Ad Click Aggregator →

Sooner or later every design gets a dashboard: clicks per hour, revenue by country, the top 100 videos this week. Answering those questions from the application database is a classic mistake that slows the product for everyone. Analytics workloads are shaped differently, and the systems built for them store data differently. This guide covers when and how to split them.

Two kinds of workload

OLTP (transactions)OLAP (analytics)
Typical queryfetch or update one ordersum revenue by country for last quarter
Rows toucheda few, by keymillions to billions
Columns touchedmost of the rowa few of many
Writesmany small, concurrentlarge batches or streams, mostly appends
Latency targetmillisecondssub-second to minutes
ExamplesPostgres, MySQL, DynamoDBClickHouse, BigQuery, Snowflake, Redshift, Druid

Running an OLAP query on an OLTP database scans huge amounts of data, evicts the hot working set from memory, and competes for I/O with user requests. Even on a read replica it is slow, because row storage is the wrong layout.

Row versus columnar storage

A row store keeps each row’s fields together on disk: perfect for "give me order 123", wasteful for "sum the amount column over a billion rows", because it reads every other column too.

A column store keeps each column’s values together:

  • A query reads only the columns it uses: 3 columns of 50 means about 6 % of the data.
  • Values in one column are similar, so they compress extremely well (run-length, dictionary, delta encoding), often 5 to 20 times.
  • Execution is vectorised: the CPU processes batches of values from one column at a time.

The cost: updating or fetching a single row means touching many column files, so column stores favour appends and bulk loads over row-level updates.

Warehouses, lakes and lakehouses

  • A data warehouse (BigQuery, Snowflake, Redshift) is a managed columnar SQL database for analytics, loaded by batch or streaming pipelines.
  • A data lake is cheap object storage (S3) holding raw and processed files, usually in columnar formats such as Parquet, queried by engines like Spark, Trino or Athena.
  • A lakehouse adds table formats (Iceberg, Delta Lake) on top of the lake so files behave like tables with schemas, updates and time travel.

A common shape: application databases and event streams feed raw data into the lake; batch jobs clean and model it; the warehouse (or lakehouse) serves analysts and dashboards. See batch and stream processing.

Getting data out of the application

  • Change data capture from the OLTP database into Kafka and then into the warehouse, which keeps the warehouse within minutes of production without load on the primary. See event sourcing, CQRS and CDC.
  • Event streams from the application (clicks, views, searches) written straight to Kafka and then to the lake. These are often the bulk of analytics data and never belonged in the OLTP database at all.
  • Nightly exports for small or legacy systems.

Modelling for analytics

Analytics schemas are denormalised on purpose. A star schema has a large fact table (one row per event: an order line, a click) with foreign keys to small dimension tables (customer, product, date, country). Queries join the fact table to a few dimensions and aggregate. Partition fact tables by date so a query for last week reads only last week’s files.

Pre-aggregation and rollups

Even columnar scans of billions of rows take time and money. When the same aggregates are asked constantly (a dashboard refreshed every minute), compute them ahead:

  • Rollup tables: clicks per ad per minute, per hour, per day. A dashboard over a month reads 30 daily rows per ad instead of millions of raw clicks.
  • Materialised views that the engine keeps updated as data arrives (ClickHouse and Druid do this well).
  • Approximate aggregates for distinct counts and percentiles (HyperLogLog, t-digest), which can be merged across rollups. See probabilistic data structures.

Keep the raw events (cheaply, in the lake) so rollups can be recomputed when definitions change or bugs are found.

Real-time analytics

When the dashboard must be seconds fresh (ad spend, live sales, operational metrics), a real-time OLAP store such as ClickHouse, Apache Druid or Apache Pinot ingests directly from Kafka and answers aggregate queries in under a second over recent data. Stream processors (Flink) can pre-aggregate into windows before writing. Design an Ad Click Aggregator is built around exactly this.

Numbers to reason with

  • Columnar compression: 5–20× smaller than row storage for typical event data.
  • Scanning: a modern engine scans hundreds of millions to billions of values per second per node for simple aggregates.
  • Rollups: per-minute rollups shrink raw event volume by the number of events per key per minute, often 100× or more.

In the interview

When the requirements include reporting or dashboards, say plainly: "Analytics runs on a separate columnar store fed by CDC and event streams, never on the primary database." Then say what freshness it needs (batch overnight or real-time), whether you pre-aggregate, and how raw data is kept for recomputation.

Checklist

  • Analytics separated from the OLTP database.
  • Columnar store chosen for the freshness needed: warehouse (minutes to hours) or real-time OLAP (seconds).
  • How data gets there: CDC, event streams, or exports.
  • Fact and dimension tables, partitioned by date.
  • Rollups or materialised views for hot dashboards.
  • Raw events retained for recomputation.

Open in your browser to sign in

Google does not allow sign-in inside this app's built-in browser. Open this page in Safari and sign in there. The link opens this same page.

Tap the ⋯ or share button at the top or bottom of the screen, then Open in browser. Or copy the link and paste it into Safari.