Описание
Комбинатор MergeState
можно применять к функции avg,
чтобы объединить частичные состояния агрегатной функции типа AverageFunction(avg, T) и
вернуть новое промежуточное состояние агрегации.
Пример использования
Комбинатор MergeState особенно полезен в сценариях многоуровневой агрегации,
когда нужно объединять предварительно агрегированные состояния и сохранять их
в виде состояний (а не финализировать) для дальнейшей обработки. Для
наглядности рассмотрим пример, в котором мы преобразуем отдельные метрики
производительности серверов в иерархические агрегации на нескольких уровнях:
уровень сервера → уровень региона → уровень датацентра.
Сначала создадим таблицу для хранения исходных данных:
CREATE TABLE raw_server_metrics
(
timestamp DateTime DEFAULT now(),
server_id UInt32,
region String,
datacenter String,
response_time_ms UInt32
)
ENGINE = MergeTree()
ORDER BY (region, server_id, timestamp);Мы создадим целевую таблицу для агрегации на уровне сервера и определим Incremental materialized view, выполняющую роль триггера вставки для этой таблицы:
CREATE TABLE server_performance
(
server_id UInt32,
region String,
datacenter String,
avg_response_time AggregateFunction(avg, UInt32)
)
ENGINE = AggregatingMergeTree()
ORDER BY (region, server_id);
CREATE MATERIALIZED VIEW server_performance_mv
TO server_performance
AS SELECT
server_id,
region,
datacenter,
avgState(response_time_ms) AS avg_response_time
FROM raw_server_metrics
GROUP BY server_id, region, datacenter;То же самое сделаем на региональном уровне и на уровне датацентра:
CREATE TABLE region_performance
(
region String,
datacenter String,
avg_response_time AggregateFunction(avg, UInt32)
)
ENGINE = AggregatingMergeTree()
ORDER BY (datacenter, region);
CREATE MATERIALIZED VIEW region_performance_mv
TO region_performance
AS SELECT
region,
datacenter,
avgMergeState(avg_response_time) AS avg_response_time
FROM server_performance
GROUP BY region, datacenter;
-- таблица и materialized view уровня датацентра
CREATE TABLE datacenter_performance
(
datacenter String,
avg_response_time AggregateFunction(avg, UInt32)
)
ENGINE = AggregatingMergeTree()
ORDER BY datacenter;
CREATE MATERIALIZED VIEW datacenter_performance_mv
TO datacenter_performance
AS SELECT
datacenter,
avgMergeState(avg_response_time) AS avg_response_time
FROM region_performance
GROUP BY datacenter;Затем вставим образец исходных данных в исходную таблицу:
INSERT INTO raw_server_metrics (timestamp, server_id, region, datacenter, response_time_ms) VALUES
(now(), 101, 'us-east', 'dc1', 120),
(now(), 101, 'us-east', 'dc1', 130),
(now(), 102, 'us-east', 'dc1', 115),
(now(), 201, 'us-west', 'dc1', 95),
(now(), 202, 'us-west', 'dc1', 105),
(now(), 301, 'eu-central', 'dc2', 145),
(now(), 302, 'eu-central', 'dc2', 155);Мы напишем три запроса для каждого уровня:
SELECT
server_id,
region,
avgMerge(avg_response_time) AS avg_response_ms
FROM server_performance
GROUP BY server_id, region
ORDER BY region, server_id;┌─server_id─┬─region─────┬─avg_response_ms─┐
│ 301 │ eu-central │ 145 │
│ 302 │ eu-central │ 155 │
│ 101 │ us-east │ 125 │
│ 102 │ us-east │ 115 │
│ 201 │ us-west │ 95 │
│ 202 │ us-west │ 105 │
└───────────┴────────────┴─────────────────┘SELECT
region,
datacenter,
avgMerge(avg_response_time) AS avg_response_ms
FROM region_performance
GROUP BY region, datacenter
ORDER BY datacenter, region;┌─region─────┬─datacenter─┬────avg_response_ms─┐
│ us-east │ dc1 │ 121.66666666666667 │
│ us-west │ dc1 │ 100 │
│ eu-central │ dc2 │ 150 │
└────────────┴────────────┴────────────────────┘SELECT
datacenter,
avgMerge(avg_response_time) AS avg_response_ms
FROM datacenter_performance
GROUP BY datacenter
ORDER BY datacenter;┌─datacenter─┬─avg_response_ms─┐
│ dc1 │ 113 │
│ dc2 │ 150 │
└────────────┴─────────────────┘Можно вставить дополнительные данные:
INSERT INTO raw_server_metrics (timestamp, server_id, region, datacenter, response_time_ms) VALUES
(now(), 101, 'us-east', 'dc1', 140),
(now(), 201, 'us-west', 'dc1', 85),
(now(), 301, 'eu-central', 'dc2', 135);Давайте снова проверим производительность на уровне датацентра. Обратите внимание, как вся цепочка агрегации автоматически обновилась:
SELECT
datacenter,
avgMerge(avg_response_time) AS avg_response_ms
FROM datacenter_performance
GROUP BY datacenter
ORDER BY datacenter;┌─datacenter─┬────avg_response_ms─┐
│ dc1 │ 112.85714285714286 │
│ dc2 │ 145 │
└────────────┴────────────────────┘