← Back to the section

In PostgreSQL you design tables around the entities of your domain and then add indexes to serve queries. In ClickHouse it is the other way around: first you list the analytical questions, and from them you derive the ORDER BY, the engine, and the pre-aggregates. A mistake in these decisions cannot be cured by "one more index" — only by recreating the table.

the same insert goes into the table and, through the view, into the pre-aggregate INSERT into order_events eu · order_paid · 120 ru · order_paid · 50 eu · order_refund · 40 eu · order_paid · 80 materialized view WHERE event_type = 'order_paid' sumState(amount) the refund row fails WHERE — no state is written revenue_by_region_daily eu ru 120 80 50 read: sumMerge(revenue)eu → 200ru → 50

A materialized view in ClickHouse is not a snapshot but an insert trigger: every inserted block runs through the view query, and what lands in the target table is an intermediate aggregate state, not a finished number. The refund row fails the condition and never reaches the pre-aggregate. The sum appears only at read time — sumMerge combines all the states of a region.

Sort order — the main decision

When ClickHouse writes data, it sorts the rows within each part by the columns from ORDER BY and builds a sparse index. At query time it reads not rows one by one, but blocks of 8,192 rows (granules). If the leading columns of ORDER BY match the WHERE condition, the engine skips whole unneeded granules.

The rule: put columns with a small number of unique values that you always filter by at the start of ORDER BY; put more unique columns and time at the end.

CREATE TABLE order_events (
    event_time   DateTime,
    event_type   LowCardinality(String),
    region       LowCardinality(String),
    customer_id  UInt64,
    order_id     UUID,
    amount       Decimal(18, 2)
)
ENGINE = MergeTree
PARTITION BY toYYYYMM(event_time)
ORDER BY (region, event_type, event_time);

A "by region over a period" query uses the key prefix and reads the fewest granules. If you set ORDER BY (event_time, region, event_type), filtering by event type stops cutting off granules: within each time range all types are mixed together.

An extra bonus: identical values sitting next to each other after sorting compress better with codecs — a good key both speeds up reads and shrinks on-disk size.

Data types: where you save space and time

Every column in ClickHouse is stored as a separate file, and its size directly affects scan speed. So the type is not a formality.

  • LowCardinality(String) — dictionary encoding for columns with hundreds or thousands of unique values: statuses, event types, regions, currencies. Less space, faster GROUP BY. With millions of unique values the dictionary bloats and becomes harmful.
  • Decimal(18, 2) — for money. Float64 introduces rounding errors in arithmetic, just as in any other database.
  • DateTime vs DateTime64(3) — second precision versus millisecond precision. DateTime64 takes twice as much space; use it only when milliseconds truly matter.
  • UInt8/16/32/64 — unsigned integers instead of bigint for everything. The smaller the type, the faster the scan.
  • Enum8('created' = 1, 'paid' = 2, ...) — more compact and stricter than LowCardinality, but adding a new value means an ALTER. Suitable for closed sets that change rarely.
  • Nullable(T) — creates a separate mask file for each column and forbids use in an ORDER BY prefix. By default it is better to rely on defaults (0, ''); use Nullable only when "no value" is semantically different from zero or an empty string.
  • Array(T), Map(K, V) — legitimate denormalization: tags, arbitrary event attributes. Together with the functions arrayJoin, has, mapKeys they cover most "flexible schema" cases without separate tables.

Aggregate queries: idioms that PostgreSQL does not have

ClickHouse is built for aggregation. A few built-in functions make queries shorter and faster compared to plain SQL:

SELECT
    toStartOfMonth(event_time)            AS month,
    count()                               AS events,
    countIf(event_type = 'order_paid')    AS paid_orders,
    sumIf(amount, event_type = 'order_paid') AS revenue,
    uniq(customer_id)                     AS customers,
    quantile(0.95)(amount)                AS p95_check
FROM order_events
WHERE event_time >= '2026-01-01'
GROUP BY month
ORDER BY month;
  • countIf / sumIf / avgIf — conditional aggregates instead of CASE WHEN inside an aggregate. They read better and run faster.
  • uniq vs uniqExact — uniq counts approximately, with an error of about one percent, and is orders of magnitude cheaper in memory and CPU. Inside it is an adaptive sampling algorithm, not HyperLogLog — classic HLL comes from uniqCombined and uniqHLL12. uniqExact gives an exact answer, but at a high cost. For dashboards uniq is almost always enough.
  • quantile / quantileExact — the same pair for percentiles.
  • argMax(col, ts) — returns the value of col from the row with the maximum ts. Handy for the "last order status" without window functions and self-joins.

Materialized views: pre-aggregation on the fly

In PostgreSQL a materialized view is a snapshot of a query result that is refreshed on command. In ClickHouse it is different: a materialized view works as an insert trigger. Every INSERT into the source table is run through the view's query and appended to the target table right at write time.

The standard pairing is with AggregatingMergeTree:

CREATE TABLE revenue_by_region_daily (
    day      Date,
    region   LowCardinality(String),
    revenue  AggregateFunction(sum, Decimal(18, 2)),
    orders   AggregateFunction(uniq, UUID)
)
ENGINE = AggregatingMergeTree
ORDER BY (region, day);

CREATE MATERIALIZED VIEW revenue_by_region_daily_mv
TO revenue_by_region_daily AS
SELECT
    toDate(event_time)        AS day,
    region,
    sumState(amount)          AS revenue,
    uniqState(order_id)       AS orders
FROM order_events
WHERE event_type = 'order_paid'
GROUP BY day, region;

Reading is done through -Merge functions:

SELECT day, region, sumMerge(revenue) AS revenue, uniqMerge(orders) AS orders
FROM revenue_by_region_daily
GROUP BY day, region;

Why sumState/uniqState and not just sum/uniq? They store the intermediate state of the aggregate rather than the finished number. This means a per-day pre-aggregate correctly rolls up into months and years. For uniq this is essential: you cannot add two finished counts — you need the original structures. The dashboard reads a table of thousands of rows instead of raw billions.

The difference shows up in the average order value. Daily averages cannot be added and divided: a day with two orders would weigh as much as a day with a thousand. The state keeps a "sum and count" pair, and the answer from it is always right:

live example

import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;

public class StateMergeDemo {
    record Event(String day, String type, long amount) {}

    record State(long sum, long count) {
        State plus(State other) {
            return new State(sum + other.sum(), count + other.count());
        }
    }

    public static void main(String[] args) {
        List<Event> events = List.of(
                new Event("2026-01-05", "order_paid", 120),
                new Event("2026-01-05", "order_paid", 50),
                new Event("2026-01-05", "order_refund", 40),
                new Event("2026-01-06", "order_paid", 80));

        Map<String, State> daily = new LinkedHashMap<>();
        for (Event e : events) {
            if (!e.type().equals("order_paid")) {
                continue;
            }
            daily.merge(e.day(), new State(e.amount(), 1), State::plus);
        }
        daily.forEach((day, s) ->
                System.out.println("state " + day + ": sum=" + s.sum() + " count=" + s.count()));

        double avgOfAvg = daily.values().stream()
                .mapToDouble(s -> (double) s.sum() / s.count())
                .average().orElse(0);
        State merged = daily.values().stream().reduce(new State(0, 0), State::plus);

        System.out.println("average of daily averages: " + round2(avgOfAvg));
        System.out.println("merged states:             " + round2((double) merged.sum() / merged.count()));
    }

    static double round2(double value) {
        return Math.round(value * 100) / 100.0;
    }
}
Run

Running examples is part of paid access. There the same code runs inside the article: editor, run and check next to the paragraph. Three free days →

82.5 against 83.33 — the price of a finished number instead of a state.

Two important limitations. First, the MV fires only on new inserts — data inserted before the view was created will not make it into the aggregate. Backfilling is done with a separate INSERT ... SELECT. Second, an error in a cascade of views can break the INSERT into the source table. Cascades deeper than one or two levels quickly become undebuggable.

ReplacingMergeTree: how to store "current state"

ClickHouse is an append-only system. To update a row, you insert a new version. When merging parts, ReplacingMergeTree keeps only the row with the maximum value of the specified field:

CREATE TABLE orders_latest (
    order_id    UUID,
    status      LowCardinality(String),
    amount      Decimal(18, 2),
    updated_at  DateTime64(3)
)
ENGINE = ReplacingMergeTree(updated_at)
ORDER BY order_id;

Every change to an order is a new insert with the same order_id. When merged, the row with the maximum updated_at remains. The problem is that the merge happens "someday", and until then duplicates are visible in queries. An honest read is one of two options:

-- Short option: merges versions right at read time
SELECT * FROM orders_latest FINAL WHERE status = 'paid';

-- Alternative: predictable cost on large tables
SELECT order_id, argMax(status, updated_at) AS status, argMax(amount, updated_at) AS amount
FROM orders_latest
GROUP BY order_id;

FINAL is convenient, but expensive on large tables. The argMax option is more verbose, but more predictable under load. In both cases this is a deliberate price for "updatability" in an append-only world.

JOIN and dictionaries

ClickHouse does have JOIN, but by default the right-hand table is loaded entirely into memory — the join_algorithm setting switches to other algorithms, some able to spill to disk. Joining a facts table with a regions reference is fine; two billion-row event tables mean out-of-memory.

How people work around it:

  • Denormalization at write time — the main technique. The category name, region, tariff go straight into the event row at load time. Extra disk space is cheaper than expensive JOINs at read time.
  • Dictionaries — reference data that ClickHouse loads itself from PostgreSQL, a file, or an HTTP source and refreshes on a schedule. Access is through the function dictGet('regions_dict', 'name', region_id) with no JOIN at all.
  • If a JOIN is unavoidable — small table on the right, filters applied before the JOIN, not after.

Common mistakes

MistakeWhat happensThe right way
SELECT * on a wide tableAll columns are read — the point of columnar storage is lostList only the columns you need
Point lookup by order_id that is not in ORDER BYA full scan on every callPoint reads belong in PostgreSQL; in ClickHouse — a separate table with the right key
Inserting one row at a time from the applicationTOO_MANY_PARTS error, inserts stallBatch inserts / async insert (details)
Nullable on all columns "just in case"An extra file per column, everything slowerDefaults; Nullable only for a semantic need
PARTITION BY toDate(...) over years of dataThousands of partitions, degraded inserts and readsMonth (toYYYYMM) as the default
Frequent ALTER ... UPDATE/DELETEA mutation queue, disk loadDELETE FROM for deletes, ReplacingMergeTree with row versions for edits
uniqExact/quantileExact in every dashboardWasted memory and CPU for precision no one will noticeuniq/quantile
A five-level cascade of materialized viewsInsert failures are hard to diagnoseOne or two levels; the rest via scheduled recomputation

In short

  • Design the schema from queries: list the analytical questions first, then the ORDER BY — low-cardinality columns (region, type) first, time last.
  • Types matter: LowCardinality(String) saves space and speeds up GROUP BY; Nullable adds a file per column, so rely on defaults.
  • countIf/sumIf/argMax are the core idioms of aggregate queries.
  • A materialized view is an insert trigger, not a snapshot: it fires only on new data; older rows need a separate INSERT ... SELECT.
  • sumState/uniqState keep the intermediate state, so daily pre-aggregates roll up into months; sumMerge computes the number at read time.
  • ReplacingMergeTree defers deduplication to the merge — read honestly with FINAL or argMax; JOIN loads the right table into memory, so prefer denormalization and dictionaries.