Материализованные представления и real-time агрегации в ClickHouse

Материализованные представления в ClickHouse как триггеры на вставку, предагрегация через AggregatingMergeTree и SummingMergeTree, чтение потока из Kafka через Kafka engine, TTL и ретеншн, проекции — как строить обновляемые в реальном времени агрегаты

Считать GROUP BY по сырым событиям на каждый запрос дэшборда — расточительно: одни и те же миллиарды строк пересчитываются снова и снова. ClickHouse предлагает считать один раз — в момент вставки. Материализованное представление здесь не «закешированный SELECT», как в PostgreSQL, а триггер: пришла пачка данных — представление тут же досчитало агрегат в целевую таблицу. Так строятся дэшборды, которые обновляются в реальном времени и стоят копейки на чтении.

Это четвёртая статья серии «ClickHouse и аналитические БД» (предыдущие — «MergeTree из Go и Java» про parts, гранулы и вставку батчами, и сравнение драйверов с бенчмарком throughput). Здесь предполагается, что вставка уже устроена правильно — батчами, без построчного антипаттерна, — и разбирается, что происходит с данными дальше, внутри сервера, на каждую такую вставку.

Все числа ниже — из живого стенда clickhouse/materialized-views в публичном репозитории digital-cookbook: 1 000 000 строк основного датасета плюс отдельный однонодовый Kafka-кластер (KRaft) с продюсером на Go (franz-go), гоняющим 50 000 синтетических событий через ENGINE=Kafka в ClickHouse.

Материализованные представления и real-time агрегация: Kafka→ClickHouse, AggregatingMergeTree, проекции

В статье

MV как триггер на INSERT, а не снапшот

В PostgreSQL материализованное представление — это снапшот результата запроса на момент последнего REFRESH MATERIALIZED VIEW: данные статичны, пока их явно не пересчитать заново, целиком. В ClickHouse MATERIALIZED VIEW устроено принципиально иначе — это не хранилище данных, а триггер, подвешенный на исходную таблицу. Каждый INSERT в источник запускает SELECT из тела MV над только что вставленным блоком строк (а не над всей таблицей) и пишет результат в целевую таблицу:

CREATE MATERIALIZED VIEW demo.mv_events_to_agg_mv TO demo.mv_events_agg AS
SELECT
    country,
    toStartOfHour(event_time) AS hour,
    sumState(revenue) AS revenue_state,
    uniqState(user_id) AS users_state,
    countState() AS count_state
FROM demo.mv_events
GROUP BY country, hour

Синтаксис TO target_table явно называет таблицу-приёмник (альтернатива — неявная скрытая таблица .inner, если TO не указать); дальше в этой статье используется только явный вариант — с ним проще понимать, что где физически хранится. Ключевое следствие того, что MV — триггер, а не снапшот: он видит только строки, вставленные в источник после своего создания. Существующие на момент CREATE MATERIALIZED VIEW данные в целевую таблицу не попадают — это не побочный эффект, а прямое следствие модели «SELECT над вставленным блоком»: блока, соответствующего уже лежащим строкам, никогда не было.

sequenceDiagram participant App as Приложение participant Src as demo.mv_events (источник) participant MV as mv_events_to_agg_mv (триггер) participant Agg as demo.mv_events_agg (AggregatingMergeTree) App->>Src: INSERT 1000 строк (ДО создания MV) Note over Src: 1000 строк лежат в источнике App->>MV: CREATE MATERIALIZED VIEW ... TO Agg AS SELECT ... Note over Agg: count() = 0 — MV не видит старые строки App->>Src: INSERT новая пачка (ПОСЛЕ создания MV) Src-->>MV: блок вставленных строк MV->>MV: SELECT ... GROUP BY (только над этим блоком) MV->>Agg: INSERT partial-state строк Note over Agg: агрегат пополнился — только новые данные

sequenceDiagram
    participant App as Приложение
    participant Src as demo.mv_events (источник)
    participant MV as mv_events_to_agg_mv (триггер)
    participant Agg as demo.mv_events_agg (AggregatingMergeTree)

    App->>Src: INSERT 1000 строк (ДО создания MV)
    Note over Src: 1000 строк лежат в источнике
    App->>MV: CREATE MATERIALIZED VIEW ... TO Agg AS SELECT ...
    Note over Agg: count() = 0 — MV не видит старые строки
    App->>Src: INSERT новая пачка (ПОСЛЕ создания MV)
    Src-->>MV: блок вставленных строк
    MV->>MV: SELECT ... GROUP BY (только над этим блоком)
    MV->>Agg: INSERT partial-state строк
    Note over Agg: агрегат пополнился — только новые данные
MV как триггер на INSERT: срабатывает только на новые вставки, существующие в источнике данные не подхватывает

Стенд проверяет это утверждение не декларативно, а живым экспериментом: 1000 строк вставляются в demo.mv_events до создания mv_events_to_agg_mv, затем MV создаётся — и сразу после создания demo.mv_events_agg пуста, count() = 0, несмотря на 1000 строк, уже лежащих в источнике. Только после этого таблица очищается (TRUNCATE) и загружается заново — уже под работающим триггером — для честного сравнения агрегата с пересчётом дальше в статье. Практическое следствие: если MV создаётся поверх таблицы, в которой уже есть данные, требующие агрегации, нужен отдельный ручной бэкфилл — обычно INSERT INTO target SELECT ... FROM source GROUP BY ... тем же запросом, что в теле MV, выполненный один раз вручную, до или сразу после создания MV.

AggregatingMergeTree: -State и -Merge

Наивная идея — писать в целевую таблицу уже готовые агрегаты (sum(revenue), uniq(user_id)) — не работает, потому что каждый вызов MV видит только один блок вставки, а не всю историю по ключу (country, hour). Если просто писать sum(revenue) из каждого блока, для одного и того же часа в целевой таблице появится несколько строк с частичными суммами, которые надо будет ещё раз складывать при чтении — то есть проблема не решена, а просто отложена.

AggregatingMergeTree решает это иначе: колонки хранят не готовое значение, а промежуточное состояние агрегатной функции — тип AggregateFunction(sum, Decimal(10,2)), AggregateFunction(uniq, UInt64) и так далее. Записывается состояние функцией с суффиксом -State (sumState, uniqState, countState), а при слиянии двух parts движок умеет объединять два состояния в одно тем же способом, каким сама агрегатная функция объединяла бы промежуточные результаты — для sum это просто сложение чисел, для uniq — слияние HyperLogLog-структур, а не пересчёт с нуля:

CREATE TABLE demo.mv_events_agg
(
    country       LowCardinality(String),
    hour          DateTime,
    revenue_state AggregateFunction(sum, Decimal(10,2)),
    users_state   AggregateFunction(uniq, UInt64),
    count_state   AggregateFunction(count)
)
ENGINE = AggregatingMergeTree
ORDER BY (country, hour)

Читать состояние напрямую бессмысленно — это внутреннее бинарное представление, не число. Чтобы получить обычное значение, состояние сводят функцией с суффиксом -Merge (sumMerge, uniqMerge, countMerge), сгруппировав по тому же ключу, что и ORDER BY целевой таблицы — так объединяются все ещё не слитые фоном строки-состояния, относящиеся к одному часу и стране:

SELECT country, hour,
    sumMerge(revenue_state) AS revenue,
    uniqMerge(users_state)  AS users,
    countMerge(count_state) AS cnt
FROM demo.mv_events_agg
GROUP BY country, hour

Стенд сверяет это с прямым пересчётом теми же функциями, но по сырым данным (sum(revenue), uniqExact(user_id), count() из demo.mv_events), причём делает это до какого-либо явного OPTIMIZE FINAL — то есть пока в mv_events_agg ещё лежит множество несведённых частичных состояний (по одной строке на каждый вставленный батч на каждый ключ). Результат: обе стороны дают ровно 42495 групп (country, hour); countMerge() суммарно равен count() — 1 000 000 против 1 000 000; sumMerge(revenue) совпадает с sum(revenue) до копейки — 13445645.77 в обоих случаях. Единственная функция, которая теоретически могла разойтись, — uniq (это приближённый кардинальный оценщик на HyperLogLog, а не точный подсчёт), и на живом прогоне относительная погрешность uniqMerge(user_id) против uniqExact(user_id) составила 0,0000% — сходится точно, хотя формально алгоритм это не гарантирует.

Отдельная граница, которую стоит держать в уме: finalizeAggregation() без GROUP BY финализирует состояние каждой строки саму по себе, не объединяя строки с одинаковым ключом, которые ещё не схлопнулись фоновым merge. На стенде после OPTIMIZE TABLE demo.mv_events_agg FINAL число строк в таблице падает с 152708 (частичные состояния сразу после bulk-загрузки — по одной строке на каждый затронутый батчем ключ) ровно до 42495 — по одной строке на ключ, — и только тогда finalizeAggregation() без GROUP BY тоже даёт корректный результат, совпадающий с -Merge. До OPTIMIZE FINAL тот же вызов без GROUP BY просто вернул бы частичные, ещё не объединённые значения по каждой лежащей строке отдельно.

SummingMergeTree: суммирование без -State

Для более простого случая — когда нужны только суммы и счётчики, без uniq/quantile/других сложных агрегатов — есть SummingMergeTree. Он не требует -State/-Merge вообще: числовые колонки вне ORDER BY просто складываются при фоновом слиянии строк с одинаковым ключом сортировки, как обычные числа, а не как состояния агрегатных функций:

CREATE TABLE demo.mv_events_summing
(
    country LowCardinality(String),
    hour    DateTime,
    revenue Decimal(18,2),
    events  UInt64
)
ENGINE = SummingMergeTree
ORDER BY (country, hour)

Живая проверка на стенде — минималистичная и оттого нагляднее любого объёма: семь строк вставляются по отдельности, пять с ключом (RU, 2026-06-15 10:00:00) и две с ключом (US, 2026-06-15 11:00:00). До OPTIMIZE TABLE ... FINAL SELECT по ключу RU возвращает все пять строк как есть — слияние ещё не произошло, SummingMergeTree ничего не пересчитывает на лету при чтении, только при merge. После OPTIMIZE FINAL на каждый ключ остаётся ровно одна строка: RU = 136.00 (10.50 + 20.25 + 5.00 + 100.00 + 0.25), US = 499.99 (300.00 + 199.99) — суммы совпали побайтово с ожидаемыми. Разница с AggregatingMergeTree — не в результате, а в цене: SummingMergeTree не умеет ничего, кроме сложения чисел по ключу, зато не требует ни отдельных -State-типов колонок, ни -Merge при чтении — обычный SELECT sum(revenue) ... GROUP BY на уже слитых данных работает и без специальных функций.

Проекции: альтернативный физический порядок

MV и AggregatingMergeTree/SummingMergeTree решают задачу «пересчитать агрегат один раз, а не на каждый запрос». Но иногда проблема не в повторном счёте, а в том, что запрос фильтрует по колонке, для которой физический порядок хранения (ORDER BY основной таблицы) не даёт пропуска гранул — ровно тот случай из статьи про MergeTree, где точечный фильтр по user_id (второй колонке ORDER BY (event_time, user_id)) не даёт выигрыша, потому что физически строки отсортированы по event_time, а не по user_id.

Проекция (projection) — это второй физический порядок данных внутри той же таблицы: ClickHouse хранит дополнительную, автоматически поддерживаемую копию данных, отсортированную иначе, и сам выбирает при планировании запроса, какую версию — основную или проекцию — дешевле прочитать:

ALTER TABLE demo.mv_events ADD PROJECTION user_id_proj (SELECT * ORDER BY user_id);
ALTER TABLE demo.mv_events MATERIALIZE PROJECTION user_id_proj;

ADD PROJECTION только объявляет проекцию, MATERIALIZE PROJECTION — отдельная, потенциально долгая операция (идёт как мутация, system.mutations), которая физически строит вторую копию данных в новом порядке для всех уже существующих parts. В отличие от MV, здесь не нужна ни отдельная целевая таблица, ни ручной бэкфилл — движок сам поддерживает проекцию синхронно с основной таблицей на любой последующей вставке.

Стенд измеряет эффект проекции двойным способом — оба независимо доказывают, что запрос действительно пошёл через неё, а не через основную таблицу. Первый способ — EXPLAIN indexes = 1 на фильтре WHERE user_id = ?: план явно упоминает имя проекции user_id_proj. Второй, более жёсткий, способ — настройка force_optimize_projection = 1: если бы ClickHouse не смог использовать ни одну проекцию для этого запроса, сервер вернул бы ошибку вместо результата; успешное выполнение само по себе доказывает использование проекции, без разбора текста EXPLAIN. На живом прогоне точечный фильтр по user_id читает read_rows = 1000000 без проекции (optimize_use_projections = 0, форсированный полный скан) и read_rows = 122880 с проекцией (optimize_use_projections = 1, дефолт) — падение в 8,14 раза (1000000 / 122880), при том что count() результата — одни и те же 3 строки в обоих случаях, независимо от того, каким путём сервер их нашёл.

Цена проекции — не бесплатна: каждая добавленная проекция удваивает (или увеличивает кратно числу проекций) объём данных на диске под эту таблицу и объём работы фонового merge, потому что обе копии — основная и проекция — поддерживаются синхронно на каждой вставке. Проекция оправдана, когда нужен второй устойчивый паттерн доступа к той же таблице (например, точечный поиск по user_id поверх таблицы, физически отсортированной по времени), а не как способ ускорить произвольный редкий запрос.

TTL на партициях

Ретеншн в ClickHouse задаётся декларативно, через TTL в DDL — интервал, по истечении которого строки удаляются (или, в более сложных схемах, перемещаются на другой диск, тема статьи про S3-тиринг дальше в серии) автоматически, без внешнего cron-скрипта с DELETE:

CREATE TABLE demo.mv_events_ttl_demo
(
    event_time  DateTime,
    user_id     UInt64,
    event_type  LowCardinality(String),
    url         String,
    duration_ms UInt32,
    country     LowCardinality(String),
    revenue     Decimal(10,2)
)
ENGINE = MergeTree
ORDER BY (event_time, user_id)
PARTITION BY toDate(event_time)
TTL event_time + INTERVAL 3 DAY DELETE

PARTITION BY toDate(event_time) здесь — не случайный выбор гранулярности партиции, а прямое условие для дешёвого TTL: когда TTL-правило проверяется во время фонового merge и оказывается, что все строки партиции просрочены, ClickHouse может отбросить партицию целиком, а не переписывать каждый part построчно (тот же механизм дешёвого DROP PARTITION, что и в статье про MergeTree, только запускается автоматически по правилу, а не явной командой).

На стенде 500 строк с event_time на 5 суток в прошлом (заведомо просрочены относительно TTL в 3 дня) и 500 свежих строк вставляются в таблицу выше. Здесь стенд честно фиксирует наблюдение, которое стоит знать заранее: TTL — механизм ленивый, он применяется как часть фонового merge, а не мгновенно в момент истечения интервала. На простаивающем dev-сервере без другой нагрузки фоновый планировщик слияний иногда успевает подобрать свежевставленные мелкие parts и применить TTL за доли секунды — быстрее, чем следующий SELECT клиента. Единственный детерминированный момент, когда результат гарантирован независимо от того, успел фон отработать сам или нет, — явный OPTIMIZE TABLE ... FINAL. После него count() таблицы равен ровно 500 — только свежие строки, просроченные удалены; партиция с датой пятидневной давности возвращает 0 строк. При этом сам part просроченной партиции может ещё какое-то время оставаться записью в system.parts (active = 1, rows = 0) до отдельного, более позднего фонового цикла очистки пустых parts — «строки удалены» и «part физически исчез из метаданных» не одно и то же событие, и различие между ними иногда важно для мониторинга.

Живой пайплайн: Kafka → ClickHouse → AggregatingMergeTree

Всё разобранное выше — MV-триггер, AggregatingMergeTree, проекции — собирается в одну цепочку для типового сценария «поток событий из шины должен превращаться в живой, всегда актуальный агрегат». ENGINE=Kafka в ClickHouse — это не хранилище, а «вьюпорт» на топик: таблица с этим движком сама по себе не накапливает данные на диске и не запускает консюмер, пока к ней не подключён хотя бы один MV, читающий из неё, — именно подключение MV запускает фоновый поток чтения топика:

CREATE TABLE demo.mv_events_queue
(
    event_time  DateTime,
    user_id     UInt64,
    event_type  LowCardinality(String),
    url         String,
    duration_ms UInt32,
    country     LowCardinality(String),
    revenue     Decimal(10,2)
)
ENGINE = Kafka
SETTINGS
    kafka_broker_list = 'kafka:9092',
    kafka_topic_list = 'mv-events-stream',
    kafka_group_name = 'ch-mv-events-consumer',
    kafka_format = 'JSONEachRow',
    kafka_num_consumers = 1,
    kafka_flush_interval_ms = 3000,
    kafka_skip_broken_messages = 0

Полная цепочка стенда — пять объектов, два MV-триггера подряд: demo.mv_events_queue (ENGINE=Kafka, читает топик) → MV → demo.mv_events_kafka_raw (обычный MergeTree, реальное хранение сырых событий на диске) → MV → demo.mv_events_kafka_agg (AggregatingMergeTree, живой предагрегат — та же схема -State, что и в основной части статьи). Первый MV просто копирует и материализует поток из «вьюпорта» в постоянную таблицу, второй агрегирует уже из неё — Kafka-таблица сама никогда не участвует в GROUP BY, только в передаче потока дальше:

flowchart LR KAFKA["Kafka topic
mv-events-stream"] QUEUE["mv_events_queue
ENGINE=Kafka
(вьюпорт, без хранения)"] MV1["MV mv_events_queue_mv
(триггер на приход сообщения)"] RAW["mv_events_kafka_raw
MergeTree
(реальное хранение)"] MV2["MV mv_events_kafka_agg_mv
(триггер на INSERT в raw)"] AGG["mv_events_kafka_agg
AggregatingMergeTree
(-State по country, hour)"] KAFKA -->|consumer group,
kafka_flush_interval_ms| QUEUE QUEUE --> MV1 MV1 --> RAW RAW --> MV2 MV2 --> AGG style QUEUE fill:#f4d9c6,stroke:#c67a4a style AGG fill:#c9e4c5,stroke:#5b8a5e

flowchart LR
  KAFKA["Kafka topic
mv-events-stream"] QUEUE["mv_events_queue
ENGINE=Kafka
(вьюпорт, без хранения)"] MV1["MV mv_events_queue_mv
(триггер на приход сообщения)"] RAW["mv_events_kafka_raw
MergeTree
(реальное хранение)"] MV2["MV mv_events_kafka_agg_mv
(триггер на INSERT в raw)"] AGG["mv_events_kafka_agg
AggregatingMergeTree
(-State по country, hour)"] KAFKA -->|consumer group,
kafka_flush_interval_ms| QUEUE QUEUE --> MV1 MV1 --> RAW RAW --> MV2 MV2 --> AGG style QUEUE fill:#f4d9c6,stroke:#c67a4a style AGG fill:#c9e4c5,stroke:#5b8a5e
Пайплайн Kafka → ClickHouse: два MV-триггера подряд, ENGINE=Kafka не хранит данные

Продюсер стенда написан на Go поверх franz-go — асинхронная отправка (Produce с callback на каждую запись) и один финальный Flush, а не синхронный ProduceSync на каждое сообщение, который был бы слишком медленным для десятков тысяч записей подряд:

var success int64
for i := 0; i < n; i++ {
    msg := syntheticKafkaEvent(rng)
    rec := &kgo.Record{Topic: topic, Value: []byte(msg)}
    cl.Produce(ctx, rec, func(_ *kgo.Record, produceErr error) {
        if produceErr != nil {
            if firstErr == nil {
                firstErr = produceErr
            }
            return
        }
        atomic.AddInt64(&success, 1)
    })
}
if flushErr := cl.Flush(ctx); flushErr != nil {
    return int(atomic.LoadInt64(&success)), time.Since(start), fmt.Errorf("flush: %w", flushErr)
}

На живом прогоне продюсер отправил 50000 из 50000 сообщений за 107,993158 мс (462992 msg/s) — сама отправка в Kafka быстрая, продюсер не является узким местом пайплайна. Дальше стенд опрашивает ClickHouse (интервал 500 мс, таймаут 120 с) до тех пор, пока count() в demo.mv_events_kafka_raw не сравняется с числом отправленных сообщений: все 50000 из 50000 строк стали видны в ClickHouse, но только через 46,752125741 с (94 опроса).

Это честная находка, которую стоит разобрать по механизму, а не списывать на «Kafka медленная» или тем более на throughput: продюсер закончил работу за 108 мс, а видимость данных в ClickHouse отстала почти на 47 секунд. Разрыв объясняется двумя накладывающимися факторами на стороне консюмера, а не сетью и не throughput. Первый — холодный старт консюмера ENGINE=Kafka: подключение к брокеру, вступление в consumer group, получение назначенных партиций занимает заметное время до того, как консюмер вообще начинает читать сообщения. Второй — kafka_flush_interval_ms = 3000 (само по себе уже укороченное с дефолтных 7500 мс для этого стенда): консюмер копит прочитанные сообщения и сбрасывает их в целевую MergeTree-таблицу пачками по этому интервалу, а не построчно по мере чтения. На пути «сообщение отправлено → видно в ClickHouse» большую часть времени съедает не передача данных, а инициализация и батчинг на приёмной стороне — тот же принцип, что уже разбирался для async inserts в статье про MergeTree: сервер сознательно откладывает видимость ради меньшего числа операций записи на диск.

После того как все строки видны, стенд повторяет ту же сверку -Merge против прямого пересчёта, что и в основной части статьи, но теперь на данных, реально прошедших через Kafka, а не загруженных батчем из CSV: агрегат через -Merge даёт 30 групп (country, hour), прямой пересчёт по сырым — тоже 30; countMerge() суммарно равен count() — 50000 против 50000; sumMerge(revenue) совпадает с пересчётом до копейки — 757510,25 в обоих случаях; погрешность uniqMerge(user_id) против uniqExact(user_id) — снова 0,0000%. Механика -State/-Merge, проверенная выше на batch-загрузке, работает идентично и на данных, пришедших через потоковый источник, — с точки зрения AggregatingMergeTree неважно, откуда взялся INSERT, вызвавший MV-триггер.

Что выбрать под задачу

  • Простой счётчик или сумма без сложных агрегатных функций → SummingMergeTree + MV: минимум кода, обычные числовые колонки.
  • Уникальные пользователи, перцентили, другие приближённые или составные агрегаты → AggregatingMergeTree + -State/-Merge: нужны специальные типы колонок и функции чтения, но работает с любой агрегатной функцией ClickHouse.
  • Ускорить существующий паттерн запросов к таблице, не заводя отдельную целевую таблицу и не думая о бэкфилле, → проекция: дороже по диску и merge, но синхронна с основной таблицей автоматически.
  • Ретеншн по времени без внешнего cron/DELETETTL на партициях, с оговоркой на ленивость применения — не полагаться на конкретный момент без явного OPTIMIZE FINAL, если момент важен для теста или мониторинга.
  • Поток событий из шины (Kafka и аналоги) должен превращаться в живой агрегат → ENGINE=Kafka + MV-триггер + MergeTree/AggregatingMergeTree на приёмной стороне; концептуальный разбор самой Kafka как системы — в серии про Kafka.

Во всех случаях общий принцип один: ClickHouse платит за агрегацию один раз, при вставке, а не при каждом чтении — но это работает только в одну сторону, вперёд по времени. MV не ретроактивен, TTL ленив, проекции синхронны только с момента материализации — если нужен пересчёт по уже лежащим данным, это всегда отдельное, явное действие, а не побочный эффект создания объекта. Дальше в серии — то, что происходит с этими же parts под нагрузкой на масштабе кластера: шардирование, репликация и Keeper.

Версии в прогоне: ClickHouse 26.6.1.1193, clickhouse-go v2.47.0, Kafka apache/kafka:4.3.1 (single-node KRaft), franz-go v1.21.5.

Обсуждение в Telegram

Присоединиться →

Комментарии