A columnar database stores each column of a table separately instead of storing each row as a contiguous record. That single change makes analytical queries, which scan many rows but touch few columns, one to two orders of magnitude cheaper, and it makes point updates and single-row lookups awkward. Engines such as ClickHouse, DuckDB, Snowflake, BigQuery, Redshift and Vertica all rest on this idea, and so do the Parquet and ORC file formats that lakehouse engines read.

File formats are covered on this site already, in the Parquet article and the ORC article. This article is about the engine: how data is laid out on disk, how it is encoded, how queries avoid reading most of it, how execution works on columns, and why the write path looks the way it does. It ends with a worked example and the operational failures you will actually meet.

Advertisement

Why columns: the IO arithmetic

Take a table of one billion events with 60 columns averaging 8 bytes each, roughly 480 GB stored row by row. A dashboard asks for daily revenue by country for one month. It needs three columns: timestamp, country and amount. A row store must read every page containing a matching row, and every page carries all 60 columns, so it reads close to the full row width for every row it touches. A column store reads only the three columns, 3/60 of the data before anything else happens, about 24 GB for a full scan.

Then compression multiplies the effect. Values within one column share a type and usually a narrow distribution: country has perhaps 200 distinct values, timestamps rise almost monotonically, amounts cluster. Encoded, those three columns may shrink by another factor of five to twenty. Finally, pruning skips blocks that cannot match the month or the country, so the engine may read a few hundred megabytes to answer a question over half a terabyte. Each factor is modest; together they are why analytical systems are columnar. The contrast with transaction processing is laid out in OLTP versus OLAP.

Anatomy: tables, parts, column chunks and granules

A column store: storage layout and the read pathInsertsbatched rowsNew partsorted by ORDER BYBackground mergesmall parts to largePart on disktsuser_idcountryamounteach column: granules of 8192 rowsplus sparse index and min/max stats1 Prunesparse index, zone maps2 Read needed columnsonly granules that survive3 Decode in batchesdictionary, RLE, delta4 Vectorised filterselection vector5 Late materialisefetch other columnsquery
Left: inserts become immutable sorted parts that background merges combine. Middle: inside a part each column is its own file, divided into granules with a sparse index and min/max statistics. Right: a query prunes, reads only the needed columns and granules, decodes in batches, filters with vectors and fetches remaining columns only for surviving rows.

Every column store has the same three levels, under different names. A table is split horizontally into large immutable units: parts in ClickHouse, row groups in Parquet and DuckDB, micro-partitions in Snowflake. Inside a unit each column is stored as its own contiguous chunk. Each chunk is divided into small blocks, granules or pages, which are the smallest thing the engine reads and the unit that carries statistics.

The horizontal unit exists so that all columns of a row can be found again: row 5,000,123 is at the same offset in every column chunk of that part. That positional alignment is what lets the engine filter on one column and then fetch the matching rows from another without storing row ids.

Advertisement

Encodings: making columns small and fast

Column stores apply lightweight encodings first and general-purpose compression, typically LZ4 or ZSTD, on top. The lightweight ones matter most because the engine can often compute on them directly.

  • Dictionary encoding replaces each string with a small integer code into a per-chunk or per-column dictionary. Country becomes one byte. Equality filters compare codes, and group-by can hash codes instead of strings.
  • Run-length encoding stores a value and a repeat count. On a column sorted by country, a billion rows collapse to a couple of hundred runs, and a filter or count can be evaluated per run.
  • Delta encoding stores differences between neighbours. Sorted timestamps become small deltas that fit in a byte or two.
  • Frame of reference and bit-packing subtract a block minimum and pack the remainders in the fewest bits that hold the block's range, so values between 1,000,000 and 1,000,500 need 9 bits each.

Choice is data-dependent: dictionaries fail on unique identifiers, runs need sort order, deltas need monotonic data. Good engines choose per chunk. The toy below shows dictionary encoding, frame-of-reference, zone maps and late materialisation in about thirty lines of Python with NumPy.

import numpy as np

GRANULE = 8192

class Column:
    def __init__(self, values):
        values = np.asarray(values)
        if values.dtype.kind in "OU":                       # strings: dictionary encode
            self.dict, codes = np.unique(values, return_inverse=True)
            self.data = codes.astype(np.uint16 if len(self.dict) < 65536 else np.uint32)
        else:                                               # numbers: frame of reference
            self.dict = None
            self.base = values.min()
            self.data = (values - self.base).astype(np.min_scalar_type(values.max() - self.base))
        g = [self.data[i:i + GRANULE] for i in range(0, len(self.data), GRANULE)]
        self.zmin = np.array([x.min() for x in g])          # zone map, in encoded space
        self.zmax = np.array([x.max() for x in g])

    def encode_const(self, v):
        if self.dict is not None:
            i = np.searchsorted(self.dict, v)
            return i if i < len(self.dict) and self.dict[i] == v else None
        return v - self.base

def scan_eq(col, value):
    code = col.encode_const(value)                          # compare codes, not strings
    if code is None:
        return np.empty(0, dtype=np.int64)                  # value absent: skip everything
    hits = []
    for g in np.nonzero((col.zmin <= code) & (col.zmax >= code))[0]:
        chunk = col.data[g * GRANULE:(g + 1) * GRANULE]
        hits.append(np.nonzero(chunk == code)[0] + g * GRANULE)   # vectorised compare
    return np.concatenate(hits) if hits else np.empty(0, dtype=np.int64)

def total_amount(country_col, amount_col, country):
    rows = scan_eq(country_col, country)                    # positions only
    return int((amount_col.data[rows].astype(np.int64) + amount_col.base).sum())  # late materialise

Notice that the filter never decodes a string: it looks up the literal once, then compares small integers, block by block, and skips any block whose min and max exclude the code. That is the whole trick in miniature.

Pruning: not reading data at all

Zone maps, also called min/max statistics, store the smallest and largest value of each column per granule or row group. A predicate like ts >= '2026-09-01' skips every block whose maximum is earlier. Parquet, ORC and DuckDB keep them automatically; in ClickHouse, partitions carry min/max for the partition key, and per-granule min/max on other columns is an opt-in minmax skip index. Either way they only help when values are clustered: on random data every block spans the full range and nothing is skipped.

That is why the sort key matters so much. ClickHouse sorts each part by the table's ORDER BY key and keeps a sparse primary index: one entry per granule, by default every 8,192 rows, holding the key values at the granule's first row. The index is small enough to sit in memory even for billions of rows, and a binary search over it finds the granule range for a key prefix. It is not a B-tree and does not find single rows; it narrows a scan to the granules that might contain them. Partitioning, such as by month, adds a coarser layer that drops whole parts. Engines add optional skip structures too: Bloom filters for high-cardinality equality lookups, set indexes, and bitmaps, which the bitmap index article explains in depth.

Execution: vectors, not rows

A classic row-at-a-time executor calls a next() function per row per operator, paying interpretation overhead on every value. Column engines process vectors of a few thousand values per call. A filter runs a tight loop over an array of codes and produces a selection vector of matching positions; the compiler turns that loop into SIMD instructions, and the data stays in CPU cache between operators. The Hive vectorisation article shows the same idea inside a SQL-on-Hadoop engine.

Late materialisation means carrying positions through the plan and fetching other columns only for rows that survive. If the country filter keeps 2% of rows, the engine decodes the amount column only at those positions, often skipping whole granules. Early materialisation, building full rows first, throws much of the columnar advantage away and is a common reason a poorly planned query on a column store is slower than expected.

Aggregation benefits most. Summing a column is a streaming loop; grouping by a dictionary-encoded column can use the codes as direct array indexes when the dictionary is small, avoiding a hash table entirely.

The write path: parts, merges and deletes

Columnar layout is hostile to single-row writes. Inserting one row means touching every column file, and updating a value in the middle of a compressed, encoded block means rewriting the block. So column stores turn writes into batches. The C-Store research design used a small write-optimised store merged into a read-optimised store; ClickHouse writes each insert as a new immutable sorted part and merges parts in the background, much like the compaction in an LSM tree. DuckDB and warehouse engines keep a write buffer or delta layer that queries read alongside the main columns.

This has direct consequences. Every insert creates a part, and each part costs files, metadata and merge work. Thousands of tiny inserts per second outrun the merger, the part count grows, queries open more files, and ClickHouse eventually rejects inserts with a too-many-parts error. The fix is to batch on the client, to use the server's asynchronous insert buffering, or to put a queue in front that writes large blocks, ideally many thousands of rows at a time.

Deletes and updates are handled logically. Engines mark rows deleted in a bitmap or hidden mask column that scans respect, and physical removal happens when the part is rewritten or merged. ClickHouse's lightweight deletes work this way; its older mutations rewrite whole parts and are expensive. Plan for bulk deletes by dropping partitions whenever the data allows, for example deleting a whole month at once.

Worked example: one query, traced

Consider this ClickHouse table and query.

CREATE TABLE events
(
    ts        DateTime,
    user_id   UInt64,
    country   LowCardinality(String),
    event     LowCardinality(String),
    amount    Decimal(12, 2),
    INDEX ts_mm ts TYPE minmax GRANULARITY 1
)
ENGINE = MergeTree
PARTITION BY toYYYYMM(ts)
ORDER BY (country, event, ts);

-- Reads country, ts and amount only; pruned by partition and by the sparse primary index.
SELECT toStartOfDay(ts) AS day, sum(amount)
FROM events
WHERE country = 'DE' AND ts >= '2026-09-01' AND ts < '2026-10-01'
GROUP BY day
ORDER BY day;

Trace the read. The partition key drops every part outside September 2026. Inside each remaining part the sparse index is searched on the key prefix (country): rows are sorted by country first, so all German rows sit in one contiguous granule range, and the engine reads only those granules. ts is the third key column, so its range filter narrows less through the index, but the ts_mm minmax skip index still skips granules. Of five columns only country, ts and amount are read, and LowCardinality stores country as dictionary codes. If Germany is 3% of traffic, the query reads roughly 3% of three columns for one month.

Now change the query to filter by user_id instead. The key does not start with user_id, so the sparse index cannot help, and user ids are spread across every granule, so a minmax skip index on it would skip almost nothing. The same table now scans the full month. Options are a Bloom-filter skip index on user_id, a projection or materialised view sorted by user_id, or accepting the scan if the query is rare. That is the central design decision in a column store: choose the sort key for the filters your heaviest queries use, in order of how selective and how common they are.

Failure modes

SymptomCauseFix
Inserts rejected, too many partsMany tiny inserts outrun mergesBatch client-side or use async inserts
Query scans everything despite a filterFilter column not in the sort key prefixReorder key, add skip index or projection
Zone maps never skipData unclustered on that columnSort or partition by it; check statistics
Point lookups slowSparse index finds granules, not rowsUse a row store or key-value store for lookups
SELECT * very slowEvery column read and decodedSelect only needed columns
Update-heavy table degradesMutations rewrite whole partsModel as append plus latest-version or deletes by partition
Memory spikes on joinsLarge hash tables built from right sidePut the smaller table on the build side; pre-aggregate
High-cardinality dictionary bloatsDictionary on near-unique stringsPlain encoding plus ZSTD for such columns

Trade-offs and choosing an engine

NeedGood fitWhy
Many point reads and updatesRow store (PostgreSQL, MySQL)Row-contiguous pages and B-tree indexes
Analytics inside one process or notebookDuckDBEmbedded, vectorised, reads Parquet directly
High-ingest event analytics, self-managedClickHouseMergeTree parts, sparse index, fast scans
Managed warehouse, elastic computeSnowflake, BigQuery, RedshiftSeparated storage and compute
Open files shared by many enginesParquet or ORC with a table formatEngine-neutral columnar storage
Mixed transactional and analyticalHTAP engines or CDC into a column storeKeeps each workload on its best layout

What to do next

  1. List your five heaviest analytical queries and the columns each one filters on; that list is your sort key candidate.
  2. Load a month of real data into DuckDB or ClickHouse with two different sort keys and compare bytes read per query.
  3. Check your ingest path: measure rows per insert and part count over a day, and batch if parts grow.
  4. Inspect per-column compressed and uncompressed sizes, and change encodings for the worst columns.
  5. Rewrite any SELECT * in dashboards to name columns.
  6. Decide how deletes will work before you need them: by partition, lightweight delete, or a versioned row model.
Key takeaway: A columnar engine wins by reading less: only needed columns, encoded compactly, with pruning skipping blocks that cannot match, and vectorised operators working on encoded values with late materialisation. The price is paid on the write side, where rows arrive as immutable sorted parts that must be merged, and deletes are logical until rewrite. Choose the sort key for your heaviest filters, batch your inserts, and select only the columns you need.