ClickHouse Kafka Engine: события из брокера напрямую в таблицы, ETL на выброс
Настроили Kafka Engine в ClickHouse: топики читаются нативным консьюмером, данные попадают в ReplicatedMergeTree через MaterializedView - без Python-скриптов и без ночных батчей.
ClickHouse Kafka Engine - нативный консьюмер Kafka-топиков без внешнего ETL-процесса
В феврале мы описывали схему с Materialized Views в ClickHouse - события из Kafka забирал Python-консьюмер, трансформировал и писал батчами в таблицы. Схема рабочая, но с прослойкой: отдельный процесс, деплой, мониторинг, иногда отставание при перегрузке. Kafka Engine убирает эту прослойку полностью. ClickHouse сам становится консьюмером.
Контекст задачи
Новый проект: события с мобильного приложения - сессии, действия, покупки. Продюсер пишет в Kafka, нужна аналитика почти в реальном времени плюс исторические срезы. Объём - несколько миллионов событий в день, пики во время акций.
В прошлых проектах мы ставили Python-консьюмер: читал из Kafka, нормализовал JSON, батчами вставлял через HTTP API ClickHouse. 200 строк кода с обработкой ошибок, оффсет-менеджментом вручную и несколькими режимами ретрая. Работало, но каждый раз переписывалось под конкретную схему. На этот раз попробовали обойтись без него.
Как устроен Kafka Engine
Kafka Engine в ClickHouse создаёт виртуальную таблицу, которая при чтении отдаёт сообщения из топика. Сама по себе она ничего не хранит - это просто курсор в брокер. Смысл появляется в паре с MaterializedView: MV читает из Kafka-таблицы и пишет в целевую таблицу хранения. ClickHouse сам управляет консьюмер-группой и оффсетами.
Схема из трёх DDL-выражений:
-- 1. Kafka-таблица: транзитная точка входа
CREATE TABLE app_events_queue (
event_time DateTime,
user_id UInt64,
session_id String,
event_type LowCardinality(String),
screen LowCardinality(String),
item_id UInt64,
revenue_kopek UInt64
) ENGINE = Kafka
SETTINGS
kafka_broker_list = 'kafka-01:9092,kafka-02:9092',
kafka_topic_list = 'app-events',
kafka_group_name = 'clickhouse-app-events',
kafka_format = 'JSONEachRow',
kafka_num_consumers = 4;
-- 2. Целевая таблица хранения
CREATE TABLE app_events (
event_date Date,
event_time DateTime,
user_id UInt64,
session_id String,
event_type LowCardinality(String),
screen LowCardinality(String),
item_id UInt64,
revenue_kopek UInt64
) ENGINE = ReplicatedMergeTree('/clickhouse/tables/app_events', '{replica}')
PARTITION BY toYYYYMM(event_date)
ORDER BY (event_date, event_type, user_id);
-- 3. MaterializedView - связка между ними
CREATE MATERIALIZED VIEW app_events_mv TO app_events AS
SELECT
toDate(event_time) AS event_date,
event_time,
user_id,
session_id,
event_type,
screen,
item_id,
revenue_kopek
FROM app_events_queue;
Вот и всё. Python-консьюмер заменён десятью строками DDL.
Что происходит внутри
ClickHouse периодически опрашивает Kafka-топик, вычитывает батч сообщений, прогоняет через MaterializedView и вставляет результат в целевую таблицу. Интервал между опросами регулируется через stream_flush_interval_ms - у нас 5 секунд. Оффсеты ClickHouse хранит в Kafka (через стандартную консьюмер-группу), то есть после рестарта сервера продолжает с последнего подтверждённого сообщения.
Количество параллельных консьюмеров задаётся параметром kafka_num_consumers на каждой ноде ClickHouse. При двух нодах кластера и kafka_num_consumers = 4 получается 8 параллельных консьюмеров - по числу партиций топика. Больше консьюмеров чем партиций смысла нет, лишние будут простаивать.
Что убрали из инфраструктуры
В этом проекте Python-консьюмера нет вообще. Из пайплайна исчезло:
- Отдельный сервис с собственным деплоем, конфигом, healthcheck-ом и alerting-ом.
- Логика оффсет-менеджмента - ClickHouse берёт это на себя через стандартный consumer group API.
- Батч-логика с ретраями - при ошибке вставки ClickHouse повторяет попытку автоматически, оффсет не сдвигается.
- Отдельный мониторинг консьюмера - лаг читается из Kafka стандартными инструментами по группе
clickhouse-app-events.
Что обнаружили на практике
Схема сообщений в топике должна совпадать с колонками Kafka-таблицы без трансформаций. Если продюсер присылает поле ts а в таблице оно называется event_time - не подойдёт. У нас пришлось договориться с командой приложения о конкретном формате JSON, который ложится напрямую в схему. Это хорошая дисциплина, но требует координации.
Nullable-поля лучше не использовать в Kafka-таблице. Если в JSON нет поля - ClickHouse подставит дефолт для типа (0, пустая строка). Мы намеренно спроектировали схему без Nullable, вынеся дефолты на уровень продюсера. Причины те же что в феврале: GROUP BY по Nullable работает медленнее.
MaterializedView вставляет атомарно по батчу. Если MV падает на трансформации - батч не вставляется, оффсет не сдвигается, данные не теряются. Это приятно. Минус в том, что одно сломанное сообщение в батче стопорит весь батч. Пришлось добавить валидацию на стороне продюсера, чтобы не пускать мусор в топик.
Мониторинг лага - через Kafka, не через ClickHouse. Сам ClickHouse не показывает насколько он отстаёт от хвоста топика. Стандартный kafka-consumer-groups.sh --describe --group clickhouse-app-events даёт всю нужную картину - интегрировали в Grafana через kafka_exporter.
Когда Kafka Engine не подходит
Если нужна сложная трансформация данных на лету - enrichment из внешних таблиц, условная маршрутизация в разные таблицы по типу события, сложная нормализация - Kafka Engine становится неудобным. MaterializedView позволяет базовые трансформации SQL, но не более того. Для таких сценариев прослойка в виде Kafka Streams или отдельного консьюмера имеет смысл.
Если схема сообщений в топике часто меняется без согласования с командой ClickHouse - тоже риск. Schema Registry с Avro мог бы помочь, но в текущей конфигурации Kafka Engine поддерживает JSONEachRow и несколько других форматов, без нативной интеграции со Schema Registry. Живём с этим.
Промежуточный итог
200 строк Python ушли. Осталось три CREATE-выражения и конфигурация брокеров. Пайплайн от события в приложении до строки в ClickHouse - меньше 10 секунд при штатной нагрузке. Отдельного процесса для консьюминга нет, оффсетами управляет сам ClickHouse.
Kafka Engine - это не серебряная пуля для всех сценариев, но для случая «плоские события из Kafka в колоночное хранилище» он убирает целый класс инфраструктурной работы. Проекты с такими задачами ведём в рамках DWH и аналитики.