ClickHouse 1.1.54388: Materialized Views вместо ночных ETL-задач на 500M строк событий безопасности
Внедрили ClickHouse для аналитики событий безопасности у клиента. Materialized Views заменили ночной ETL - агрегации теперь строятся инкрементально на лету.
ClickHouse активно развивается в Яндексе как OLAP-движок для больших данных; версия 1.1.54388 стабилизировала Materialized Views
Несколько месяцев назад к нам пришёл клиент с типичной проблемой: журналы событий безопасности - firewall, IDS, аутентификация, endpoint-агенты - собираются исправно, несколько сотен миллионов строк накоплено, но аналитику по ним можно получить только к утру следующего дня. Ночная ETL-задача пересчитывала агрегаты за предыдущие сутки и грузила результат в PostgreSQL. Дежурный аналитик видел вчерашнее.
Мы взялись за проект по DWH и аналитике, и первый же вопрос был: на чём это делать. PostgreSQL с партиционированием - вариант, мы его пробовали на других проектах. Но здесь объём другой: данные только растут, события льются непрерывно, и аналитикам нужны срезы за произвольный период с фильтрами по десяткам атрибутов. ClickHouse стал очевидным кандидатом - Яндекс открыл его в 2016-м, к 2018-му он дорос до вменяемого состояния и уже работал в нескольких известных нам проектах за пределами Яндекса.
Почему ClickHouse, а не привычные инструменты
Честно - перед тем как ставить ClickHouse в продакшн клиенту, мы его несколько недель крутили на стенде. Несколько наблюдений из этого периода.
Скорость на колоночном хранении. Схема событий широкая: 40+ полей. На аналитическом запросе типа "сколько уникальных источников атак по типу события за неделю" PostgreSQL читает всю строку, ClickHouse - только нужные колонки. На сотнях миллионов строк разница не в процентах, а в разах.
MergeTree как основной движок. Таблицы на MergeTree физически сортируются по первичному ключу (ORDER BY), что при правильном выборе даёт естественный фильтр на чтение. Для событий безопасности ключ очевиден: (event_date, event_type, source_ip). Запрос с фильтром по дате и типу события читает маленький кусок файлов.
Materialized Views как инкрементальная агрегация. Вот это оказалось главным. Materialized View в ClickHouse - это не кэш снимка как в PostgreSQL (где REFRESH нужен вручную). Это цепочка: данные вставляются в source-таблицу, триггер на вставке прогоняет их через SELECT из MV, результат пишется в отдельную целевую таблицу. Агрегат обновляется атомарно на каждом батче вставки. Ночной ETL-задаче нечего делать.
Как устроена схема
Исходные данные - JSON-события от нескольких коллекторов, нормализованные в Kafka топик. Оттуда ClickHouse тянет через табличную функцию Kafka (появилась в 1.1.54337, работает через движок Kafka для таблиц). Схема получается такая:
CREATE TABLE security_events_raw (
event_time DateTime,
event_date Date DEFAULT toDate(event_time),
event_type LowCardinality(String),
source_ip String,
dest_ip String,
severity UInt8,
details String
) ENGINE = MergeTree()
PARTITION BY toYYYYMM(event_date)
ORDER BY (event_date, event_type, source_ip);
Поверх неё - три Materialized Views под разные разрезы аналитики. Пример для суточных агрегатов по типам событий и источникам:
CREATE MATERIALIZED VIEW security_by_day
ENGINE = SummingMergeTree()
PARTITION BY toYYYYMM(event_date)
ORDER BY (event_date, event_type, source_ip)
AS SELECT
event_date,
event_type,
source_ip,
count() AS event_count,
sum(severity) AS severity_sum,
max(severity) AS severity_max
FROM security_events_raw
GROUP BY event_date, event_type, source_ip;
SummingMergeTree - это движок, который при фоновом слиянии (merge) автоматически суммирует числовые колонки по одинаковым ключам. После вставки данных возможны промежуточные дубли в пределах одного куска до слияния, поэтому запросы к MV пишем с явным GROUP BY и sum() - не доверяем что слияние уже произошло.
Что получилось на реальных данных
Данные заливали постепенно: сначала историческую выгрузку из старой PostgreSQL-таблицы (это как раз те самые несколько сотен миллионов строк), потом переключили живой поток из Kafka. Вставка исторических данных шла батчами по 500k строк, MV обновлялись параллельно с вставкой - процесс занял несколько часов.
Несколько наблюдений по результату:
Аналитический запрос за месяц. На raw-таблице с несколькими сотнями миллионов строк - секунды, не минуты. На MV с агрегатами за месяц - меньше секунды. Для сравнения: тот же запрос на старой PostgreSQL-таблице с частичными индексами уходил в десятки секунд.
Ночной ETL отменён. Аналитик видит данные с задержкой в единицы минут - это время на доставку через Kafka. Не следующее утро.
Неожиданный нюанс с LowCardinality. Тип LowCardinality(String) для event_type - это словарное кодирование, работает как ENUM но без фиксированного набора значений. На запросах с GROUP BY event_type даёт заметный прирост скорости, потому что сравниваются целые числа, а не строки. Но в 1.1.54388 у него есть ограничение: LowCardinality в Materialized View иногда вёл себя непредсказуемо при слиянии кусков. Пришлось в MV использовать обычный String и только в raw-таблице держать LowCardinality.
Репликация. Развернули два узла с ReplicatedMergeTree через ZooKeeper. Репликация работает, но ZooKeeper здесь - отдельная зависимость, которую надо эксплуатировать. На одном узле ClickHouse отлично работает без ZK - если реплики не нужны, не усложняйте.
Что настораживает
ClickHouse - молодой инструмент. Яндекс его активно разрабатывает, коммиты идут практически ежедневно, версии выходят часто. Это и хорошо (быстро фиксят), и немного тревожно - release notes надо читать, обновления на продакшн надо тестировать. Прямо сейчас у нас на стенде лежит следующая версия, проверяем.
Документации на русском достаточно - Яндекс позаботился. Документация на английском существует, но местами отстаёт от реального поведения. clickhouse-users в Telegram - живой, вопросы отвечают быстро включая людей из Яндекса.
Мониторинг из коробки: системные таблицы system.query_log, system.merges, system.replication_queue дают достаточно информации для Prometheus-экспортёра. Grafana-дашборд нашли готовый, немного допилили под себя.
Промежуточный итог
Ночные ETL-задачи на этом проекте больше не нужны - MV берут их работу на себя инкрементально. Аналитик видит данные в реальном времени. Запросы по историческим данным работают за секунды там, где раньше были минуты или ночной пересчёт.
Проект ещё не закрыт: впереди миграция ещё нескольких источников событий и настройка ролевого доступа для аналитиков. Но базовая механика работает именно так, как мы рассчитывали.