Обработка обновлений в режиме реального времени в ClickHouse
Данные типа mutable, как правило, нежелательны в базах данных OLAP, и ClickHouse не исключение. Как и некоторые другие OLAP-продукты, ClickHouse изначально даже и не поддерживал обновления данных. Чуть позже в систему добавили опцию обновления, но это, как и многие другие вещи, было сделано «по-кликхаусовски».
Даже сейчас обновления ClickHouse асинхронны, что затрудняет их использование в интерактивных приложениях. Тем не менее вопрос внесения изменений в данные по-прежнему актуален. Можно ли сделать это в рамках ClickHouse? Конечно, да.
Краткая история обновлений в ClickHouse
Еще в 2016 году команда ClickHouse опубликовала статью под названием Как обновить данные в ClickHouse. В то время ClickHouse еще не поддерживал функцию модификации данных. Для их обновлений можно было использовать только специальные структуры вставки, а данные приходилось сбрасывать по разделам…
Под давлением требований GDPR команда ClickHouse добавила операторы UPDATE и DELETE. Updates и Deletes в ClickHouse до сих пор является одной из самых популярных в блоге Altinity. Эти асинхронные, неатомарные обновления реализуются операторами ALTER TABLE UPDATE и потенциально могут перемешивать большое количество данных, что действительно полезно для нечастых обновлений, когда немедленный результат не нужен. «Обычные» SQL-обновления по-прежнему отсутствуют, хотя они ежегодно фигурируют в дорожной карте системы. Если нужно провести обновления данных в режиме реального времени, необходимо использовать другой подход. Давайте рассмотрим конкретный пример и сравним различные способы реализации этой операции в ClickHouse.
Пример
Рассмотрим систему, генерирующую различные виды оповещений. Для их просмотра пользователи или алгоритмы машинного обучения время от времени запрашивают базу данных. Операции подтверждения должны модифицировать запись об оповещении в базе данных. После подтверждения оповещения должны исчезнуть из поля зрения пользователей. Это совсем не похоже на операцию OLTP, которая в буквальном смысле чужда ClickHouse.
Поскольку мы не можем использовать обновления, нам придется вставить измененную запись. Когда две записи находятся в базе данных, нам нужен эффективный способ получить последнюю из них. Для этого мы попробуем три разных подхода:
- ReplacingMergeTree
- Агрегатные функции
- AggregatingMergeTree
ReplacingMergeTree
Для начала создадим таблицу, в которой будут храниться оповещения:
CREATE TABLE alerts( tenant_id UInt32, alert_id String, timestamp DateTime Codec(Delta, LZ4), alert_data String, acked UInt8 DEFAULT 0, ack_time DateTime DEFAULT toDateTime(0), ack_user LowCardinality(String) DEFAULT '' ) ENGINE = ReplacingMergeTree(ack_time) PARTITION BY tuple() ORDER BY (tenant_id, timestamp, alert_id);
Для простоты все столбцы, содержащие оповещения, упакованы в общий столбец alert_data. Но Вы же понимаете, что alert может содержать десятки или даже сотни столбцов. Кроме того, в нашем примере alert_id - это случайная строка.
Обратите внимание на движок ReplacingMergeTree. ReplacingMergeTee - это специальный движок таблицы, который заменяет данные по первичному ключу (ORDER BY) - более новая версия строки с тем же значением ключа заменит старую. «Новизна» определяется колонкой, в нашем случае – это «ack_time». Замена происходит во время фоновой операции слияния. Она происходит не сразу, и нет никакой гарантии, что она вообще произойдет, поэтому возникает риск несогласованности результатов запроса. Однако в ClickHouse реализован специальный синтаксис для работы с такими таблицами, и мы будем использовать его в запросах ниже.
Прежде чем выполнять запросы, давайте заполним таблицу данными. Мы генерируем 10 миллионов оповещений для 1000 пользователей:
INSERT INTO alerts(tenant_id, alert_id, timestamp, alert_data)
SELECT
toUInt32(rand(1)%1000+1) AS tenant_id,
randomPrintableASCII(64) as alert_id,
toDateTime('2020-01-01 00:00:00') + rand(2)%(3600*24*30) as timestamp,
randomPrintableASCII(1024) as alert_data
FROM numbers(10000000);
Далее подтвердим 99% оповещений, предоставив новые значения для столбцов 'acked', 'ack_user' и 'ack_time'. Вместо обновления мы просто вставим новую строку:.
INSERT INTO alerts (tenant_id, alert_id, timestamp, alert_data, acked, ack_user, ack_time)
SELECT tenant_id, alert_id, timestamp, alert_data,
1 as acked,
concat('user', toString(rand()%1000)) as ack_user, now() as ack_time
FROM alerts WHERE cityHash64(alert_id) % 99 != 0;
Если мы обратимся к этой таблице прямо сейчас, то увидим что-то вроде этого:
SELECT count() FROM alerts ┌──count()─┐ │ 19898060 │ └──────────┘
Таким образом, таблице содержит как подтвержденные, так и не подтвержденные строки. Поэтому замены пока не происходит. Чтобы увидеть «конечную версию» данных, мы должны добавить ключевое слово FINAL.
SELECT count() FROM alerts FINAL ┌──count()─┐ │ 10000000 │ └──────────┘ 1 rows in set. Elapsed: 3.693 sec. Processed 19.90 million rows, 1.71 GB (5.39 million rows/s., 463.39 MB/s.)
Теперь подсчет идет правильно, но посмотрите на время выполнения запроса!!! При использовании FINAL ClickHouse приходится сканировать все строки и объединять их по первичному ключу. Это дает правильный ответ, но со значительными накладными расходами. Давайте посмотрим, сможем ли мы добиться большего, отфильтровав только те строки, которые не были подтверждены.
SELECT count() FROM alerts FINAL WHERE NOT acked ┌─count()─┐ │ 101940 │ └─────────┘ 1 rows in set. Elapsed: 3.570 sec. Processed 19.07 million rows, 1.64 GB (5.34 million rows/s., 459.38 MB/s.)
Время выполнения запроса и объем обрабатываемых данных одинаковы, хотя количество данных значительно уменьшилось. Фильтрация не помогает ускорить выполнение запроса. При увеличении размера таблицы затраты могут быть еще более значительными.
NB: Для удобства чтения все запросы и время их выполнения представлены так, как будто они выполняются в 'clickhouse-client'. На самом деле мы пробовали обрабатывать запросы несколько раз, чтобы убедиться в согласованности результатов и подтвердить их с помощью утилиты 'clickhouse-benchmark'.
Хорошо, запрос всей таблицы не так уж и полезен. Можем ли мы все же использовать ReplacingMergeTree в наше случае? Давайте выберем случайный идентификатор пользователя и выберем все записи, которые еще не были подтверждены - представьте, что у нас есть дашборд, в который в любой момент может заглянуть каждый пользователь. Я - фанат Рэй Брэдбери, поэтому я выбрал 451. Поскольку 'alert_data' - это случайные даты, рассчитаем контрольную сумму и будем использовать ее для подтверждения того, что результаты одинаковы при использовании нескольких подходов:
SELECT count(), sum(cityHash64(*)) AS data FROM alerts FINAL WHERE (tenant_id = 451) AND (NOT acked) ┌─count()─┬─────────────────data─┐ │ 90 │ 18441617166277032220 │ └─────────┴──────────────────────┘ 1 rows in set. Elapsed: 0.278 sec. Processed 106.50 thousand rows, 119.52 MB (383.45 thousand rows/s., 430.33 MB/s.)
Это было очень быстро! За 278 мс мы смогли запросить все неподтвержденные данные. Почему на этот раз все получилось так быстро? Разница заключается в условии фильтрации. 'tenant_id' является частью первичного ключа, поэтому ClickHouse может фильтровать данные перед FINAL. В этом случае ReplacingMergeTree эффективен.
Попробуем также использовать пользовательский фильтр и запросить количество оповещений, подтвержденных каждым конкретным пользователем. Кардинальность столбца та же - у нас 1000 пользователей и мы можем попробовать user451.
SELECT count() FROM alerts FINAL WHERE (ack_user = 'user451') AND acked ┌─count()─┐ │ 9725 │ └─────────┘ 1 rows in set. Elapsed: 4.778 sec. Processed 19.04 million rows, 1.69 GB (3.98 million rows/s., 353.21 MB/s.)
Теперь все очень медленно, потому что мы не использовали индекс. ПоэтомуClickHouse просканировал все 19,04 миллиона строк. Обратите внимание на то, что мы не можем добавить 'ack_user' в индекс, так как это нарушит семантику ReplacingMergeTree. Однако мы можем проделать трюк с PREWHERE:
SELECT count() FROM alerts FINAL PREWHERE (ack_user = 'user451') AND acked ┌─count()─┐ │ 9725 │ └─────────┘ 1 rows in set. Elapsed: 0.639 sec. Processed 19.04 million rows, 942.40 MB (29.80 million rows/s., 1.48 GB/s.)
PREWHERE - это специальная подсказка для ClickHouse, позволяющая применить фильтр по-другому. Обычно ClickHouse в состоянии автоматически перемещать условия в PREWHERE, поэтому пользователю ни о чем не стоит беспокоиться. В этот раз все прошло по-другому, поэтому хорошо, что мы все проверили!
Агрегатные функции
ClickHouse известен тем, что поддерживает большое количество агрегатных функций. В последних версиях их насчитывается более 100. (!) В сочетании с 9 комбинаторами агрегатных функций (https://clickhouse.tech/docs/en/query_language/agg_functions/combinators/) это дает опытному пользователю «зеленый свет». В нашем случае мы ничего не хотим усложнять, поэтому будем использовать только 3 функции: 'argMax', 'max' и 'any'.
Тот же запрос для 451-го пользователя можно выполнить с помощью агрегатной функции 'argMax'.
SELECT count(), sum(cityHash64(*)) data FROM (
SELECT tenant_id, alert_id, timestamp,
argMax(alert_data, ack_time) alert_data,
argMax(acked, ack_time) acked,
max(ack_time) ack_time_,
argMax(ack_user, ack_time) ack_user
FROM alerts
GROUP BY tenant_id, alert_id, timestamp
)
WHERE tenant_id=451 AND NOT acked;
┌─count()─┬─────────────────data─┐
│ 90 │ 18441617166277032220 │
└─────────┴──────────────────────┘
1 rows in set. Elapsed: 0.059 sec. Processed 73.73 thousand rows, 82.74 MB (1.25 million rows/s., 1.40 GB/s.)
Тот же результат, то же количество строк, но производительность в 4 раза выше! Это и есть эффективность агрегации ClickHouse. Основным недостатком является то, что запрос становится немного сложнее. Но мы можем сделать упростить его.
Отметим, что при подтверждении оповещения мы обновляем только 3 столбца:
- acked: 0 => 1
- ack_time: 0 => now()
- ack_user: '' => 'user1'
Во всех трех случаях значение столбца увеличивается! Поэтому вместо громоздкого 'argMax' мы можем использовать 'max'. Поскольку нам не нужно изменять 'alert_data', так как необходимости в агрегации по этому столбцу нет. В ClickHouse есть хорошая агрегатная функция 'any'. Она выбирает любое значение без лишних накладных расходов:
SELECT count(), sum(cityHash64(*)) data FROM (
SELECT tenant_id, alert_id, timestamp,
any(alert_data) alert_data,
max(acked) acked,
max(ack_time) ack_time,
max(ack_user) ack_user
FROM alerts
GROUP BY tenant_id, alert_id, timestamp
)
WHERE tenant_id=451 AND NOT acked;
┌─count()─┬─────────────────data─┐
│ 90 │ 18441617166277032220 │
└─────────┴──────────────────────┘
1 rows in set. Elapsed: 0.055 sec. Processed 73.73 thousand rows, 82.74 MB (1.34 million rows/s., 1.50 GB/s.)
Итак, запрос становится проще и обрабатывается немного быстрее! Причина заключается в том, что с функцией 'any' ClickHouse не нужно вычислять 'max' для столбца 'alert_data'!
AggregatingMergeTree
AggregatingMergeTree - одна из самых мощных функций ClickHouse. В сочетании с материализованными представлениями она позволяет агрегировать данные в режиме реального времени. В предыдущем подходе мы использовали агрегатные функции. Сможем ли мы сделать это еще более эффективно с помощью AggregatingMergeTree?
Мы обновляем строку только один раз, поэтому для группы нужно агрегировать только две строки. В этом случае AggregatingMergeTree - не самый лучший вариант. Однако мы можем пойти на хитрость. Мы знаем, что сначала оповещения всегда вставляются как неподтвержденные, и только затем их подтверждают. Когда пользователь подтверждает оповещение, нужно изменить только 3 столбца. Можно ли сэкономить место на диске и повысить производительность, если не дублировать данные для остальных столбцов?
Давайте создадим таблицу, в которой будет реализовано агрегирование с помощью агрегатной функции 'max'. Вместо 'max' мы могли бы использовать 'any', но тогда столбцы должны быть NULL - 'any' будет выбирать не NULL.
DROP TABLE alerts_amt_max; CREATE TABLE alerts_amt_max ( tenant_id UInt32, alert_id String, timestamp DateTime Codec(Delta, LZ4), alert_data SimpleAggregateFunction(max, String), acked SimpleAggregateFunction(max, UInt8), ack_time SimpleAggregateFunction(max, DateTime), ack_user SimpleAggregateFunction(max, LowCardinality(String)) ) Engine = AggregatingMergeTree() ORDER BY (tenant_id, timestamp, alert_id);
Поскольку исходные данные были случайными, заполним новую таблицу, используя существующие данные из 'alerts'. Для этого делаем 2 вставки: одна для неподтвержденных оповещений, другая - для подтвержденных:
INSERT INTO alerts_amt_max SELECT * FROM alerts WHERE NOT acked; INSERT INTO alerts_amt_max SELECT tenant_id, alert_id, timestamp, '' as alert_data, acked, ack_time, ack_user FROM alerts WHERE acked;
Обратите внимание на то, что для подтвержденных событий мы вставляем пустую строку вместо 'alert_data'. Агрегатная функция заполнит этот пробел. В реальном приложении мы можем просто пропустить все столбцы, которые не меняются, и дать им значения по умолчанию.
Когда у нас есть данные, давайте в первую очередь проверим их размер:
SELECT
table,
sum(rows) AS r,
sum(data_compressed_bytes) AS c,
sum(data_uncompressed_bytes) AS uc,
uc / c AS ratio
FROM system.parts
WHERE active AND (database = 'last_state')
GROUP BY table
┌─table──────────┬────────r─┬───────────c─┬──────────uc─┬──────────────ratio─┐
│ alerts │ 19039439 │ 20926009562 │ 21049307710 │ 1.0058921003373666 │
│ alerts_amt_max │ 19039439 │ 10723636061 │ 10902048178 │ 1.0166372782501314 │
└────────────────┴──────────┴─────────────┴─────────────┴────────────────────┘
Благодаря случайным строкам сжатия практически нет.
Теперь попробуем выполнить запрос к таблице aggregate:
SELECT count(), sum(cityHash64(*)) data FROM (
SELECT tenant_id, alert_id, timestamp,
max(alert_data) alert_data,
max(acked) acked,
max(ack_time) ack_time,
max(ack_user) ack_user
FROM alerts_amt_max
GROUP BY tenant_id, alert_id, timestamp
)
WHERE tenant_id=451 AND NOT acked;
┌─count()─┬─────────────────data─┐
│ 90 │ 18441617166277032220 │
└─────────┴──────────────────────┘
1 rows in set. Elapsed: 0.036 sec. Processed 73.73 thousand rows, 40.75 MB (2.04 million rows/s., 1.13 GB/s.)
Благодаря AggregatingMergeTree мы обрабатываем меньше данных (40 МБ, а 82 МБ ранее), поэтому процесс обработки запроса выполняется в разы эффективнее.
Материализация обновления
ClickHouse сделает все возможное, чтобы объединить данные в фоновом режиме, удалив дубликаты строк и выполнив агрегирование. Однако иногда имеет смысл принудительно завершить слияние, например, для того, чтобы освободить место на диске. Это можно сделать с помощью оператора OPTIMIZE FINAL. OPTIMIZE - это блокирующая и достаточно дорогая операция, поэтому ее нельзя выполнять слишком часто. Давайте посмотрим, повлияет ли она на производительность запроса.
OPTIMIZE TABLE alerts FINAL Ok. 0 rows in set. Elapsed: 105.675 sec. OPTIMIZE TABLE alerts_amt_max FINAL Ok. 0 rows in set. Elapsed: 70.121 sec.
После OPTIMIZE FINAL обе таблицы содержат одинаковое количество строк и полностью идентичные данные.
┌─table──────────┬────────r─┬───────────c─┬──────────uc─┬────────────ratio─┐ │ alerts │ 10000000 │ 10616223201 │ 10859490300 │ 1.02291465565429 │ │ alerts_amt_max │ 10000000 │ 10616223201 │ 10859490300 │ 1.02291465565429 │ └────────────────┴──────────┴─────────────┴─────────────┴──────────────────┘
Разница в производительности между разными подходами становится менее очевидной:
|
После вставок |
После OPTIMIZE FINAL |
|
|
ReplacingMergeTree FINAL |
0.278 |
0.037 |
|
argMax |
0.059 |
0.034 |
|
any/max |
0.055 |
0.029 |
|
AggregatingMergeTree |
0.036 |
0.026 |
Заключение
ClickHouse предлагает богатый набор инструментов для обработки обновлений в режиме реального времени, такие как ReplacingMergeTree, CollapsingMergeTree (не рассматривается в этой статье), AggregatingMergeTree и агрегатные функции. Все эти подходы имеют три общие черты:
- Данные обновляются путем вставки новой версии. Вставки в ClickHouse выполняются очень быстро.
- Существуют эффективные способы эмулировать семантику обновления, аналогичную семантике OLTP-баз данных.
- Фактическое обновление данных происходит не сразу.
Выбор того или иного подхода зависит от конкретного случая использования приложения. ReplacingMergeTree является простым и наиболее удобным для пользователя методом, но может использоваться только для таблиц малого и среднего размера или в том случае, когда данные запрашиваются по первичному ключу. Агрегатные функции обеспечивают большую производительность, но требует переписывания запросов. И наконец, AggregatingMergeTree позволяет экономить место в памяти за счет сохранения только измененных столбцов. Это хорошие и достаточно эффективные инструменты, которые должны быть в арсенале профессионального пользователя ClickHouse.




