Databases / An explorable explanation

How a column
store works

Using ClickHouse as an example4 interactive modelsExplore the database landscape ↗

You're building an analytics service for four online shops, A, B, C, and D. Whenever someone opens a page, the shop sends your service a page-view record. The service collects these records so it can produce reports for each shop.

IllustrationShops send page views
Four online shops send page-view records to one analytics service. Each shop has its own path into the shared collection of records. Online shops Page-view records
Visits to any of the four shops produce records in the same analytics service. Each record includes the shop it came from, so the service can report on the shops separately.

For example, a visitor in Sweden opens the blue-mug page in shop B at 09:41:12. The record contains the shop, time, country, event type, visitor ID, and page address. We call the event's extra details its payload; for this page view, that's the page address.

A visitor opens a product page in shop B
Shop
B
Time
09:41:12
Country
Sweden
Event
Page view
User
visitor_42
Payload
/products/blue-mug

One report counts page views by country for each shop. To produce it, we need the shop and country from every record. The other fields don't affect these counts: a view from Sweden adds one to the total whether it happened at 09:41 or 09:42, and whichever product the visitor looked at.

Let's give the service 16 page views, including that visit to the blue-mug page. All of the records in this example are page views, so we don't need to filter by event type. Here's the report:

SQL example

Count the 16 page views

SELECT shop, country,
       count(*) AS views
FROM page_views
GROUP BY shop, country
ORDER BY shop, country;
The result, in either storage layout
shopcountryviews
ASweden2
AUK2
BSweden2
BUK2
CGermany2
CSweden2
DGermany2
DUK2

Each result row is one shop-and-country pair. The two views from Sweden in shop B include our blue-mug visit. We keep the other fields for other reports, but they don't enter this calculation.

So the query needs two fields from every record. Does the database have to read the other four as well?

01 / LAYOUT

Reading two fields out of six

If we store each record as a row, its six fields sit next to each other. The database reads chunks of the file called pages. A page containing the shop and country will also contain fields we don't need, and reading that page brings them all into memory.

A column store rearranges the values: all the shop values together, all the timestamps together, all the countries together, and so on. Now the query can read the shop and country columns without reading the others.

The saving comes from the fields we leave behind. Counting page views by shop and country requires those two columns, however many other fields each record contains. If a report needs all six fields, storing them in columns won't let us skip any of them. For a request that fetches one complete customer record, a row store with a suitable index may be a better fit.

Interactive experiment01

Choose what your query reads

Here are our 16 page views in two storage layouts. Each field takes up one cell, and each page holds eight cells. Shop and country are selected for the report; the marked cells belong to the blue-mug visit.

Run scan, then choose Store as columns and run it again. Compare the pages read. Try selecting every field too.

Controls

Result Updates when you use the controls

Marked cells: the blue-mug visitShop: B · Time: 09:41:12 · Country: Sweden · Event: Page view · User: visitor_42 · Payload: /products/blue-mug

For the shop-and-country report, the row layout spreads the requested values across all 12 pages. In the column layout, the 16 shop values fit into two pages and the 16 country values into another two. Both layouts produce the same counts; one reads 12 pages and the other reads four.

Text description of this view

What this model leaves out

16 records, six equal-sized fields, eight cells per page. A scan reads a whole page if it contains a requested field. No indexes, compression, caching, or per-row filters. These are page counts in this model, not measured database timings.

Columns also put similar values together, which gives compression something to work with. Our 16 page views contain just three country names. Instead of writing “Sweden” into every record that needs it, we can store the name once in a dictionary and put a short code in those records. For example, if 0 means Sweden, the blue-mug visit can store 0 as its country. Looking up that code gives us the original name.

Illustration · Dictionary encoding

Store each country name once

The country field from the same 16 page views, in the same order.

BeforeCountry names
AfterCodes
01 UK
1
02 Sweden ●
0
03 Germany
2
04 UK
1
05 Sweden
0
06 UK
1
07 Sweden
0
08 Germany
2
09 UK
1
10 Sweden
0
11 Germany
2
12 UK
1
13 Sweden
0
14 UK
1
15 Sweden
0
16 Germany
2

16 country names

16 codes

Shared dictionary

Keep this alongside the codes.

0
Sweden
1
UK
2
Germany

3 names, stored once each

Follow the blue-mug visit

Record 02 came from Sweden. Its position stays the same; its country is now stored as code 0.

0look up
Sweden

The dictionary turns that code back into the original name.

Every record still has a country. The encoded column holds a short code for each record, plus a dictionary of the three names. Box widths are schematic, not measured byte sizes.

The dictionary takes space too. Reusing its entries across many records can save space; if every record has a different country name, there's little repetition to remove and the extra codes may make things larger. This is one encoding idea, rather than a description of every ClickHouse column's on-disk format.

ClickHouse can also process a batch of column values together, instead of repeating the same bookkeeping for one record at a time. Compression and batch processing are further reasons column stores suit large analytical queries. These savings come on top of avoiding unused fields.

02 / ORDER

Looking at one shop

We've counted page views for all the shops. When the owner of shop B opens their dashboard, though, they only need B's counts. We could save more work by skipping records from the other shops. Storing the values in columns doesn't help us find B: its values are still mixed in with A, C, and D.

SQL example

The same report, just for shop B

SELECT shop, country,
       count(*) AS views
FROM page_views
WHERE shop = 'B'
GROUP BY shop, country
ORDER BY shop, country;
Result: 4 page views
shopcountryviews
BSweden2
BUK2

Sorting by shop puts B's records next to each other. This moves whole records: the blue-mug visit's country, time, and payload must move with its shop value. We still store the fields in separate columns, but a given position in every column must refer to the same visit. Sorting each column independently would scramble the records.

We can divide that shared record order into groups called blocks. Each block covers corresponding positions across the columns. That gives us two ways to avoid reading data: leave out entire fields, or skip a range of records within the fields we need.

For each block, keep its smallest and largest shop key. A block whose keys run from C to D can't contain B. We can skip those positions in both the shop and country columns, while continuing to leave the other four columns unread.

When every block mixes records from all four shops, its range runs from A to D and can't rule out B. Sorting brings B's records together, so blocks containing only the other shops can be skipped. The report still counts two visits from Sweden and two from the UK; the sort order changes the work needed to find them.

Interactive experiment02

Find one shop’s records

This model groups our 16 records into four blocks of four. Only shop and country are drawn, but all six fields follow the same record order. These blocks group records; they aren't the eight-cell disk pages in the previous experiment.

Choose Sort by shop and follow the marked blue-mug visit. Its country stays beside its shop. Compare the block ranges, then try a different shop.

Controls

Result Updates when you use the controls

Text description of this view

What this model leaves out

The same 16 records, in blocks of four. Only shop and country are drawn; all six fields follow the same record order. We test each block's minimum and maximum shop key. ClickHouse's sparse primary index uses a different algorithm. Blocks here describe groups of records, not the eight-cell disk pages in the first model. Real granules are larger and may cross key boundaries.

ClickHouse uses a sparse primary index for this job. It records the key at the start of each group of sorted rows, called a granule. For a key such as shop, the query can search these ordered entries to find candidate granules, then check their rows.

Illustration · Sparse index
What the sparse index contains
Our sorted shop keys, with four rows per granule for this drawing.
Index entry: AA A A AGranule 0
Check the boundary
Index entry: BB B B BGranule 1
Contains B
Index entry: CC C C CGranule 2
Past B: skip
Index entry: DD D D DGranule 3
Past B: skip

The index stores the first key, not every key shown underneath it. An entry of A followed by B doesn't prove the A granule contains only A: it could end with B. A conservative lookup includes that boundary granule and checks its rows too.

The minimum-and-maximum method keeps each block's last key as well as its first. A sparse index has less information about each granule's contents. The diagram shows candidate ranges, not every optimisation in ClickHouse's search. Once granules are selected, separate mark files locate their data in the requested columns. Real granules are much larger than four rows.

This is why the sorting key matters. Putting each shop's events together is useful for our dashboard. It won't help as much if the next query asks for one user across all shops: that user's events could be scattered through the table. We could keep another arrangement for that query, but we'd have to store it and update it too.

03 / WRITES

Adding events to sorted data

So far we've rearranged data that was already there. Visitors keep opening pages, so the shops keep sending new records, and each one needs to end up in the right place in the sort order. Rewriting a large sorted file whenever an event arrives would mean doing a lot of work for a small addition.

ClickHouse's MergeTree storage engines handle this by sorting incoming batches separately. A batch is written into one or more parts. Each part contains sorted rows stored in columns, along with an index. Once written, the part is immutable: it won't be edited in place.

Illustration · Inside a partOne part, several columns
One open tray contains three colored columns with six aligned positions. A separate small index card points into the tray. Index Data part
The tray represents one part. Three columns are shown here; values at the same position belong to the same page-view record. The index helps a query locate ranges within the part. This drawing leaves out compression, marks, and the arrangement of files on disk.

A query can read matching rows from several parts, so new events don't have to wait for one big sorted file. But lots of small parts give each query more places to look. In the background, ClickHouse merges parts: it reads their contents and writes a new, larger sorted part to replace them.

Illustration · Merging sorted parts
One sorted partA · C
Another sorted partB · D
Read the next smallest key ↓
New sorted partA · B · C · D
These are shop keys from four new page views. Their other fields travel with them. The new part replaces both inputs; no records have been added or removed.

For a small example, imagine four page views arriving as four separate batches, each creating one part. Combining them into one part requires reading and rewriting those records. If the same four views arrive in one batch, we can write one part immediately and avoid those merges. Both cases store the same records, but the separate batches cause extra work.

Interactive experiment03

Insert and merge batches

Each batch in this model creates one sorted part. Merging combines the two smallest parts, and the counter records how many events have been rewritten.

Set Events per batch to 1 and insert four batches. Merge them until one part remains. Then reset and insert a single batch of four events.

Controls

Result Updates when you use the controls

Text description of this view

What this model leaves out

Each inserted batch creates exactly one part. A merge combines the two smallest parts; rewritten rows count both inputs. Real engines use more complex merge scheduling, partition boundaries, buffers, and compressed byte sizes.

Of course, collecting a larger batch takes time. If the dashboard needs new events quickly, you may not want to wait. ClickHouse's asynchronous inserts can collect batches on the server. Their acknowledgement settings control whether the sender waits for that batch to be written before receiving success.

What if an event was wrong? With ReplacingMergeTree, you can insert a corrected version and let a later merge reconcile the versions. Until that merge happens, both versions may still be present. A query that needs the latest version must account for that. And if you've already used the old event to update a daily total, correcting the event doesn't by itself fix the total.

04 / FAILURE

A write can succeed before it reaches the second copy

We've followed a page view into a sorted part. But what if the machine storing that part fails? This is a deployment question shared by many databases. To finish our ClickHouse example, we'll look at a setup that keeps a second copy.

A new page view arrives at the receiving replica, which stores it and sends a copy to the second replica. Sending takes time. Until it arrives, only the receiving replica has this record. We're adding a new event here, so there are no old and corrected versions to reconcile.

When should we tell the analytics service its write succeeded? We could reply as soon as the receiving replica stores it, or wait for the second copy.

If the receiving replica loses its storage before the second copy arrives, the page view is lost under either policy. Waiting doesn't copy the data any faster. But an immediate reply would have told the service its write succeeded, even though the only copy was then lost. Waiting avoids giving that success response before a second copy exists.

Once the second replica has the page view, losing the receiving replica leaves a copy available. Waiting therefore delays success until the write can survive that particular failure.

Interactive experiment04

Wait for a second copy?

A new page view has just been stored on the receiving replica. Advance time to send it to the second replica, or lose the receiving replica's storage before the copy arrives.

Choose Lose receiving replica immediately. Then enable Wait for second replica, which starts a fresh write, and lose it again. Compare what the writer heard. Finally, reset and let the copy arrive before causing the failure.

Controls

Result Updates when you use the controls

Text description of this view

What this model leaves out

One new page view, two data replicas, and logical time steps. The record is stored on the receiving replica at time zero; the second replica receives it after the chosen delay. Failure means permanently losing the receiving replica's storage. No version replacement, elections, repair, network partitions, or disk-flush simulation. The checkbox only controls when the writer receives success. This model does not implement consensus.

Self-managed ReplicatedMergeTree copies data asynchronously. A write can succeed before another replica has it. Insert-quorum settings let you require more replicas to confirm receipt first. That still doesn't mean every replica you might read from has caught up.

ClickHouse Cloud uses a different arrangement. With SharedMergeTree, compute nodes read parts from shared object storage, and metadata coordinates which parts are available. The storage service keeps the durable data, so adding a compute node doesn't mean giving it a separate full copy. Readers still need to learn that a new part exists.

Illustration · Compute and storageWhere do the copies live?
Left: each of two compute blocks connects to its own data tray, and the trays replicate between them. Right: two compute blocks connect to one shared data tray. Compute Compute Replica copies Shared data

Separate storageEach replica holds a copy.

Shared storageCompute reads the same store.

On the left, each replica keeps its own copy of the data. On the right, both compute nodes use the same storage service. That service has its own redundancy; the single tray doesn't mean a single disk. Metadata coordination and local caches aren't shown. The experiment above models the arrangement on the left.
05 / THE CHOICE

When would you use this?

For our dashboard, this arrangement makes sense. Most events are added once and read many times. Queries count and group values from a few columns, and a shop filter can rule out much of the data. I'd consider a column store for that workload.

Now consider the checkout in an online shop. Two people try to buy the last item. The database needs to update inventory and create an order without selling that item twice. Reading fewer columns won't establish that rule; we need a transaction that keeps those changes together. I'd start by looking at a transactional relational database for that part of the application.

You can recognise some of the same ideas in other analytical databases. DuckDB runs inside your process, but also benefits from column layouts and avoiding unnecessary reads. Analytical warehouses may use similar techniques while arranging storage and compute differently. The part handling and replication behaviour we've looked at here are specific to the ClickHouse engines we've named.

If you're comparing one of those systems with ClickHouse, start with an actual query from your application. Work out which columns it needs and whether the sorting order lets it skip records. Then check how new data and corrections reach that query. Those details will tell you more than the label "column store" alone.

← Return to the database landscape

Sources & model notes

This is a draft of the new article format. The experiments use small, invented datasets to make the behaviour visible. They count pages, records, and copies; they don't predict how long a query will take on your database. Each experiment has a note explaining its assumptions.

Teaching references: Sam Who's controlled experiments, Bartosz Ciechanowski's progressive models, and Julia Evans's concrete scenes. Original illustrations generated with imagegen; labels and data in the working models are rendered by code.