Este tutorial mostra como manter rollups pré-agregados de uma tabela de eventos de alto volume usando visões materializadas. Você criará três objetos: uma tabela de dados brutos, uma tabela de rollup e a visão materializada que grava no rollup automaticamente.
Quando usar este padrão
Use este padrão quando:
- Você tem um fluxo de eventos somente de acréscimo (cliques, visualizações de página, IoT, logs).
- A maioria das consultas são agregações em intervalos de tempo (por minuto/hora/dia).
- Você quer leituras consistentes em menos de um segundo sem precisar varrer novamente todas as linhas brutas.
Criar a tabela de eventos brutos
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 DELETEObservações
PARTITION BY toYYYYMM(event_time)mantém as partições pequenas e fáceis de remover.ORDER BY (event_time, user_id)dá suporte a consultas com intervalo de tempo definido + filtro secundário.LowCardinality(String)economiza memória para dimensões categóricas.TTLremove os dados brutos após 90 dias (ajuste conforme seus requisitos de retenção).
Defina a tabela de rollup (agregada)
Vamos fazer a pré-agregação com granularidade horária. Escolha a granularidade de acordo com a janela de análise mais comum.
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)Armazenamos estados de agregação (por exemplo, AggregateFunction(sum, ...)), que representam, de forma compacta, agregações parciais e podem ser combinados ou finalizados posteriormente.
Crie uma visão materializada que alimenta o rollup
Esta visão materializada é executada automaticamente nas inserções em events_raw e grava estados de agregação no 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;Insira alguns dados de exemplo
Insira alguns dados de exemplo:
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);Consultando o rollup
Você pode mesclar os estados no momento da leitura ou finalizá-los:
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; -- ou use SELECT ... FINALFiltre pelos campos da chave primária para obter o melhor desempenho
Você pode usar o comando EXPLAIN para ver como o índice é usado para descartar dados:
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 │
└────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────┘O plano de execução da consulta acima mostra três tipos de índices sendo usados:
um índice MinMax, um índice de partição e um índice de chave primária.
Cada índice usa campos especificados em nossa chave primária: (bucket_start, country, event_type).
Para obter o melhor desempenho de filtragem, você deve garantir que suas consultas usem campos da chave primária para descartar dados.
Variações comuns
- Diferentes granularidades: adicione um rollup diário:
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);Em seguida, uma segunda visão materializada:
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;- Compressão: aplique codecs em colunas grandes (exemplo:
Codec(ZSTD(3))) na tabela bruta. - Controle de custos: concentre a retenção mais pesada na tabela bruta e mantenha agregações de longa duração.
- Backfilling: ao carregar dados históricos, faça a inserção em
events_rawe deixe a visão materializada criar as agregações automaticamente. Para linhas existentes, usePOPULATEao criar a visão materializada, se fizer sentido, ouINSERT SELECT.
Limpeza e retenção
- Aumente o TTL dos dados brutos (por exemplo, 30/90 dias), mas mantenha os roll-ups por mais tempo (por exemplo, 1 ano).
- Você também pode usar o TTL para mover partes antigas para um armazenamento mais barato, caso o armazenamento em camadas esteja habilitado.
Solução de problemas
- A visão materializada não está sendo atualizada? Verifique se as inserções vão para events_raw (não para a tabela de rollup) e se o destino da visão materializada está correto (
TO events_rollup_1h). - Consultas lentas? Confirme se elas estão atingindo o rollup (consulte a tabela de rollup diretamente) e se os filtros de tempo estão alinhados com a granularidade do rollup.
- Inconsistências no backfill? Use
SYSTEM FLUSH LOGSe verifiquesystem.query_log/system.partspara confirmar as inserções e os merges.