ClickHouse Kafka Engine в продакшне: от событий в топике до аналитики за секунды
Запустили production-пайплайн: события из Kafka напрямую в ClickHouse через Kafka Table Engine и Materialized View. Задержка аналитики упала с часов до секунд.
ClickHouse 19.x развивает движок Kafka и Materialized View: потоковая инжиниция данных из Kafka в аналитическую СУБД стала простой и производительной.
В июле мы писали про замену SSIS-джобы на связку Kafka + ClickHouse у одного из клиентов. Пайплайн встал, работает, инцидентов почти не было. Параллельно с тем проектом у нас дозревал второй - похожий по структуре, но со своими особенностями. Вот о нём и расскажем, потому что несколько технических деталей оказались полезны и нам, и, возможно, будут полезны тем, кто идёт тем же путём.
Откуда задача
Клиент - e-commerce с нормальной по меркам этого сегмента нагрузкой: события о просмотрах, добавлениях в корзину, заказах. Всё это шло в Kafka с прошлого года, но аналитика строилась через batch-выгрузки: раз в несколько часов скрипт читал из топиков, писал в CSV, CSV грузился в PostgreSQL, поверх PostgreSQL сидел Metabase.
Аналитики хотели видеть данные быстрее. Не «в реальном времени» в маркетинговом смысле, а просто - что происходит прямо сейчас, а не три часа назад. Задача для DWH и аналитики: убрать batch, поставить потоковый пайплайн.
Почему Kafka Engine, а не ещё один консьюмер
Технически можно было написать консьюмер на Python или Go, который читает из Kafka и вставляет в ClickHouse через HTTP-интерфейс или нативный клиент. Многие так и делают. Но в ClickHouse 19.x Kafka Engine вырос до состояния, когда своим консьюмером можно не заморачиваться - СУБД становится консьюмером сама.
Выглядит это просто: создаёшь таблицу с движком Kafka, указываешь брокеры, топик и формат - всё, ClickHouse сам читает сообщения и делает это как нормальный Kafka-консьюмер с group_id и offset-менеджментом.
CREATE TABLE events_queue (
event_time DateTime,
user_id UInt64,
session_id String,
event_type LowCardinality(String),
page_url String,
item_id Nullable(UInt64)
) ENGINE = Kafka
SETTINGS
kafka_broker_list = 'kafka-1:9092,kafka-2:9092',
kafka_topic_list = 'site-events',
kafka_group_name = 'clickhouse-analytics',
kafka_format = 'JSONEachRow',
kafka_num_consumers = 2;
Поверх - материализованное представление, которое при каждой порции данных перекладывает их в основную таблицу на MergeTree:
CREATE MATERIALIZED VIEW events_mv TO events AS
SELECT
event_time,
user_id,
session_id,
event_type,
page_url,
item_id
FROM events_queue
WHERE event_type != '';
Условие WHERE event_type != '' тут не для красоты - оказалось, что в топике изредка прилетают битые сообщения с пустым типом, и их лучше фильтровать сразу в представлении, чем потом объяснять аналитикам откуда берутся строки-призраки.
Что трудоёмко в Kafka Engine
Есть нюансы, которые не сразу очевидны.
Отладка потока данных. Kafka Engine-таблица не хранит данные - она читает из топика и тут же отдаёт в представление. Сделать SELECT * FROM events_queue можно, но результат непредсказуем: получишь пачку сообщений, которые в момент запроса шли по потоку. Привыкаешь быстро, но поначалу сбивает с толку.
Схема и несовпадение типов. Если в топике прилетело сообщение с полем, которое не соответствует типу в таблице, ClickHouse по умолчанию пишет дефолтное значение и идёт дальше. Никакой ошибки, никакого предупреждения в момент записи - только тихо искажённые данные. Нашли это при сверке агрегатов: один из источников отправлял item_id как строку вместо числа для определённой категории товаров. Решили через явный каст в представлении и через мониторинг на стороне продюсеров.
Перебалансировка и лаг. При рестарте ClickHouse-нода выходит из консьюмер-группы, происходит ребаланс, несколько секунд данные не читаются. Для нас это приемлемо - важнее было не потерять данные, а небольшой лаг в аналитике никого не пугает. Но нужно понимать, что сообщения не теряются - они просто накапливаются в топике и будут прочитаны после рестарта.
Как выглядит сейчас
Пайплайн запущен три недели назад и работает. Задержка от события до видимости в Metabase - секунды, изредка десятки секунд при росте нагрузки. Это против трёх-четырёх часов при batch-подходе.
Мониторинг consumer_lag настроили через Prometheus + kafka_exporter: если ClickHouse начинает отставать от топика, алерт приходит в Slack раньше, чем кто-то успевает заметить задержку в дашборде. Порог поставили консервативно - пока ни одного реального алерта.
PostgreSQL оставили как есть - туда по-прежнему идут транзакционные данные, которые нужны бэкенду. ClickHouse занимается только аналитической нагрузкой, которую PostgreSQL тащил с видимым скрипом на больших выборках.
Что дальше смотрим
Следующий шаг - агрегатные материализованные представления с AggregatingMergeTree для тех дашбордов, которые считают одно и то же по всей истории при каждом запросе. Это уже не про Kafka Engine, а про то, как не пересчитывать агрегаты с нуля - разбираемся отдельно.
Kafka Engine в ClickHouse 19.x производит впечатление рабочего инструмента, а не лабораторного. Главное - знать, где он молчит о проблемах, и закрывать эти точки мониторингом и явными правилами трансформации.