ClickHouse 19.x: заменили ETL-конвейер на Kafka + MaterializedView и забыли про ночные задачи
Перешли на ClickHouse 19.x со стабильным MergeTree и materialized views: аналитика по 200 млн строк за секунды без отдельного OLAP-хранилища.
ClickHouse 19.x - стабильные релизы с улучшенным MergeTree и materialized views вышли в начале 2019 года
Прошлой осенью мы описывали первый опыт с ClickHouse Materialized Views на проекте по событиям безопасности. С тех пор вышло несколько релизов серии 19.x, и картина стала яснее: команда ClickHouse серьёзно поработала над стабильностью MergeTree-движков и несколькими раздражающими граблями с materialized views. В феврале запустили второй проект - на этот раз поведенческая аналитика для e-commerce. Разбираем что изменилось и что нет.
Откуда задача
Клиент накапливал события пользовательского поведения - клики, просмотры, воронки, транзакции. Несколько источников, разная структура, всё сваливается в Kafka. К моменту разговора с нами таблица в PostgreSQL перешагнула за 200 миллионов строк и начала показывать характер: аналитический запрос «воронка за прошлую неделю по категориям» уходил в несколько минут даже с индексами. На дашборд нельзя было смотреть без кофе - иначе не хватало терпения.
ETL-пайплайн выглядел классически: ночью cron-задача пережёвывала сырые данные из Kafka через набор Python-скриптов, складывала агрегаты в PostgreSQL, утром аналитик видел вчерашнее. Воронки в реальном времени - мечта.
Вариантов было несколько: отдельный Redshift, BigQuery, Druid. Но клиент не хотел ни сложности Druid, ни уходить в облако с данными о поведении пользователей. ClickHouse лежал на стенде с осени, и после нескольких месяцев наблюдений за проектом в продакшне выглядел как разумный выбор. Тем более что свежие 19.x-релизы закрыли несколько проблем, которые нас нервировали раньше.
Что поменялось в 19.x
Честно: главное что дала серия 19.x - это предсказуемость. Не взрывной прирост новых фич, а то что существующие механики стали работать так, как написано в документации.
Materialized Views и LowCardinality. Осенью мы обходили баг: LowCardinality-колонки в Materialized View вели себя непредсказуемо при слиянии кусков. В 19.x это починили. Теперь LowCardinality работает и в исходной таблице, и в MV - словарное кодирование работает насквозь, и GROUP BY по строковым колонкам с небольшой кардинальностью стал заметно быстрее.
MergeTree и partition pruning. Запросы с фильтром по партиции стали надёжнее исключать ненужные куски на этапе планирования. Звучит скучно, но для запроса «события за последние 7 дней» это разница между тем, читается ли месяц данных или только нужная неделя.
Kafka engine. Небольшие, но приятные улучшения в управлении консьюмерами: стал надёжнее вести себя при перебалансировке, перестал иногда терять оффсеты при рестарте сервера. На предыдущих версиях у нас был кейс когда после рестарта ClickHouse начинал читать Kafka с начала партиции - сейчас такого не наблюдали.
Как устроена схема
Структура похожа на то, что делали осенью, но под другую задачу. Три слоя:
Первый - таблица с Kafka-движком, чисто транзитная, без хранения:
CREATE TABLE events_queue (
event_time DateTime,
user_id UInt64,
session_id String,
event_type LowCardinality(String),
category LowCardinality(String),
item_id Nullable(UInt64),
amount Nullable(Float64)
) ENGINE = Kafka
SETTINGS kafka_broker_list = 'kafka:9092',
kafka_topic_list = 'user-events',
kafka_group_name = 'clickhouse-consumer',
kafka_format = 'JSONEachRow';
Второй - целевая raw-таблица на ReplicatedMergeTree, куда всё пишется через Materialized View:
CREATE TABLE events (
event_date Date,
event_time DateTime,
user_id UInt64,
session_id String,
event_type LowCardinality(String),
category LowCardinality(String),
item_id UInt64,
amount Float64
) ENGINE = ReplicatedMergeTree('/clickhouse/tables/events', '{replica}')
PARTITION BY toYYYYMM(event_date)
ORDER BY (event_date, event_type, category, user_id);
Третий - агрегирующие MV под конкретные аналитические запросы. Для воронок по категориям:
CREATE MATERIALIZED VIEW events_daily_by_category
ENGINE = SummingMergeTree()
PARTITION BY toYYYYMM(event_date)
ORDER BY (event_date, category, event_type)
AS SELECT
event_date,
category,
event_type,
count() AS event_count,
uniqExact(user_id) AS unique_users,
uniqExact(session_id) AS unique_sessions,
sumIf(amount, amount > 0) AS revenue
FROM events
GROUP BY event_date, category, event_type;
Запрос к MV за месяц - меньше секунды. К raw-таблице за тот же период с аналогичной агрегацией - несколько секунд. Для дашборда обоих достаточно.
Что отменили в ETL-пайплайне
Cron-задачу убрали полностью. Python-скрипты, которые жили на отдельной машине и гоняли данные ночью, - тоже. Kafka читается напрямую движком ClickHouse, агрегаты обновляются инкрементально при каждом батче вставки. Задержка между событием и его появлением в аналитике - единицы минут, это время доставки через Kafka плюс интервал между батчами вставки (у нас настроен на 5 секунд через stream_flush_interval_ms).
Отдельного OLAP-хранилища нет: ClickHouse сам и хранит сырые события, и строит агрегаты. Это уменьшает количество движущихся частей и снимает вопрос синхронизации между хранилищами.
Где всё ещё не идеально
SummingMergeTree и промежуточные куски. После вставки до фонового слияния кусков агрегат в MV может быть неточным - несколько кусков с одним и тем же ключом ещё не слиты. Запросы к MV надо писать с явным GROUP BY и агрегатными функциями, не доверяя что слияние уже произошло. Это не баг, это задокументированное поведение - но разработчики, привыкшие к PostgreSQL, на это регулярно натыкаются.
Nullable-колонки и производительность. В исходной схеме событий у клиента было много Nullable-полей - не все события несут все атрибуты. ClickHouse Nullable реализован через отдельный битовый массив, и GROUP BY по Nullable-колонке работает медленнее чем по обычной. Перенесли nullable-логику на уровень приложения-продюсера: в Kafka пишем дефолтные значения (0 для числовых, пустую строку для строк), Nullable убрали из схемы. Стало быстрее.
ZooKeeper как обязательная зависимость для репликации. На двух узлах без ZooKeeper - нельзя. Это отдельный сервис, за которым надо следить. Если ZooKeeper лёг - репликация встала, вставки идут, но на один узел. Держим ZooKeeper на трёх нодах отдельно от ClickHouse.
Промежуточный итог
Ночной ETL исчез. Аналитики видят данные с задержкой в минуты, а не к следующему утру. Запросы по 200 миллионам строк уходят в секунды на raw-таблице и в доли секунды на MV. Выделенного OLAP-кластера нет.
Серия 19.x стала достаточно стабильной, чтобы тащить в продакшн клиентов без ощущения что стоишь на льду. Осенью это ощущение было. Сейчас - меньше. По работам с DWH и аналитикой - это сейчас основной инструментальный стек для аналитики событий.