В этом руководстве показано, как поддерживать предварительно агрегированные rollup-данные из высоконагруженной таблицы событий с помощью materialized views. Вы создадите три объекта: сырую таблицу, таблицу rollup и materialized view, которая автоматически записывает данные в rollup.
Когда использовать этот шаблон
Используйте этот шаблон, когда:
- У вас есть поток событий, в который данные только добавляются (клики, просмотры страниц, IoT, журналы).
- Большинство запросов — это агрегации по временным диапазонам (по минутам/часам/дням).
- Вам нужно стабильное чтение данных менее чем за секунду без повторного сканирования всех необработанных строк.
Создайте таблицу сырых событий
CREATE TABLE events_raw
(
event_time DateTime,
user_id UInt64,
country LowCardinality(String),
event_type LowCardinality(String),
value Float64
)
ENGINE = MergeTree
PARTITION BY toYYYYMM(event_time)
ORDER BY (event_time, user_id)
TTL event_time + INTERVAL 90 DAY DELETEПримечания
PARTITION BY toYYYYMM(event_time)разбивает данные на небольшие партиции, которые легко удалять.ORDER BY (event_time, user_id)оптимизирует запросы с ограничением по времени и дополнительной фильтрацией.LowCardinality(String)экономит память для категориальных измерений.TTLудаляет необработанные данные через 90 дней (настройте в соответствии с вашими требованиями к сроку хранения).
Спроектируйте rollup-таблицу (агрегированную таблицу)
Мы будем выполнять предварительную агрегацию с часовой детализацией. Выбирайте уровень детализации под наиболее типичный временной диапазон анализа.
CREATE TABLE events_rollup_1h
(
bucket_start DateTime, -- start of the hour
country LowCardinality(String),
event_type LowCardinality(String),
users_uniq AggregateFunction(uniqExact, UInt64),
value_sum AggregateFunction(sum, Float64),
value_avg AggregateFunction(avg, Float64),
events_count AggregateFunction(count)
)
ENGINE = AggregatingMergeTree
PARTITION BY toYYYYMM(bucket_start)
ORDER BY (bucket_start, country, event_type)Мы храним состояния агрегатных функций (например, AggregateFunction(sum, ...)), которые компактно представляют частичные агрегаты и впоследствии могут быть объединены или финализированы.
Создайте materialized view, заполняющее rollup
Это materialized view автоматически срабатывает при вставке в events_raw и записывает состояния агрегатных функций в rollup.
CREATE MATERIALIZED VIEW mv_events_rollup_1h
TO events_rollup_1h
AS
SELECT
toStartOfHour(event_time) AS bucket_start,
country,
event_type,
uniqExactState(user_id) AS users_uniq,
sumState(value) AS value_sum,
avgState(value) AS value_avg,
countState() AS events_count
FROM events_raw
GROUP BY bucket_start, country, event_type;Вставьте немного тестовых данных
Вставьте немного тестовых данных:
INSERT INTO events_raw VALUES
(now() - INTERVAL 4 SECOND, 101, 'US', 'view', 1),
(now() - INTERVAL 3 SECOND, 101, 'US', 'click', 1),
(now() - INTERVAL 2 SECOND, 202, 'DE', 'view', 1),
(now() - INTERVAL 1 SECOND, 101, 'US', 'view', 1);Выполнение запросов к rollup
Вы можете либо объединять состояния при чтении, либо финализировать их:
SELECT
bucket_start,
country,
event_type,
uniqExactMerge(users_uniq) AS users,
sumMerge(value_sum) AS value_sum,
avgMerge(value_avg) AS value_avg,
countMerge(events_count) AS events
FROM events_rollup_1h
WHERE bucket_start >= now() - INTERVAL 1 DAY
GROUP BY ALL
ORDER BY bucket_start, country, event_type;SELECT
bucket_start,
country,
event_type,
uniqExactMerge(users_uniq) AS users,
sumMerge(value_sum) AS value_sum,
avgMerge(value_avg) AS value_avg,
countMerge(events_count) AS events
FROM events_rollup_1h
WHERE bucket_start >= now() - INTERVAL 1 DAY
GROUP BY ALL
ORDER BY bucket_start, country, event_type
SETTINGS final = 1; -- или используйте SELECT ... FINALФильтруйте по полям первичного ключа для наилучшей производительности
Чтобы увидеть, как индекс отсекает данные, можно использовать команду EXPLAIN:
EXPLAIN indexes=1
SELECT *
FROM events_rollup_1h
WHERE bucket_start BETWEEN now() - INTERVAL 3 DAY AND now()
AND country = 'US'; ┌─explain────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────┐
1. │ Expression ((Project names + Projection)) │
2. │ Expression │
3. │ ReadFromMergeTree (default.events_rollup_1h) │
4. │ Indexes: │
5. │ MinMax │
6. │ Keys: │
7. │ bucket_start │
8. │ Condition: and((bucket_start in (-Inf, 1758550242]), (bucket_start in [1758291042, +Inf))) │
9. │ Parts: 1/1 │
10. │ Granules: 1/1 │
11. │ Partition │
12. │ Keys: │
13. │ toYYYYMM(bucket_start) │
14. │ Condition: and((toYYYYMM(bucket_start) in (-Inf, 202509]), (toYYYYMM(bucket_start) in [202509, +Inf))) │
15. │ Parts: 1/1 │
16. │ Granules: 1/1 │
17. │ PrimaryKey │
18. │ Keys: │
19. │ bucket_start │
20. │ country │
21. │ Condition: and((country in ['US', 'US']), and((bucket_start in (-Inf, 1758550242]), (bucket_start in [1758291042, +Inf)))) │
22. │ Parts: 1/1 │
23. │ Granules: 1/1 │
└────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────┘Показанный выше план выполнения запроса демонстрирует использование трех типов индексов:
индекса MinMax, индекса партиции и индекса первичного ключа.
Каждый индекс использует поля, входящие в наш первичный ключ: (bucket_start, country, event_type).
Чтобы обеспечить максимальную эффективность фильтрации, убедитесь, что ваши запросы используют поля первичного ключа для отсечения данных.
Распространённые варианты
- Разная гранулярность: добавьте ежедневную агрегацию:
CREATE TABLE events_rollup_1d
(
bucket_start Date,
country LowCardinality(String),
event_type LowCardinality(String),
users_uniq AggregateFunction(uniqExact, UInt64),
value_sum AggregateFunction(sum, Float64),
value_avg AggregateFunction(avg, Float64),
events_count AggregateFunction(count)
)
ENGINE = AggregatingMergeTree
PARTITION BY toYYYYMM(bucket_start)
ORDER BY (bucket_start, country, event_type);Затем — второе materialized view:
CREATE MATERIALIZED VIEW mv_events_rollup_1d
TO events_rollup_1d
AS
SELECT
toDate(event_time) AS bucket_start,
country,
event_type,
uniqExactState(user_id),
sumState(value),
avgState(value),
countState()
FROM events_raw
GROUP BY ALL;- Сжатие: применяйте кодеки к большим столбцам (например,
Codec(ZSTD(3))) в raw-таблице. - Контроль затрат: основную нагрузку по хранению возложите на raw-таблицу, а roll-up'ы храните долго.
- Дозагрузка: при загрузке исторических данных вставляйте их в
events_raw, а materialized view автоматически построит roll-up'ы. Для существующих строк используйтеPOPULATEпри создании materialized view, если это уместно, либоINSERT SELECT.
Очистка и сроки хранения
- Увеличьте TTL для сырых данных (например, до 30/90 дней), но храните roll-up-данные дольше (например, 1 год).
- Вы также можете использовать TTL для перемещения старых частей в более дешёвое хранилище, если включено распределение по уровням хранения.
Устранение неполадок
- materialized view не обновляется? Проверьте, что вставка идет в events_raw (а не в таблицу rollup) и что целевая таблица materialized view указана правильно (
TO events_rollup_1h). - Медленные запросы? Убедитесь, что они обращаются к rollup (выполните запрос к таблице rollup напрямую) и что временные фильтры соответствуют гранулярности rollup.
- Несоответствия при дозагрузке? Используйте
SYSTEM FLUSH LOGSи проверьтеsystem.query_log/system.parts, чтобы подтвердить вставки и слияния.