ADG Оставить заявку
Блог Данные и аналитика 5 мин чтения

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 и аналитики.

Контакт

Нужна такая же инженерная работа?

Опишите задачу и контекст. Ответим в течение рабочего дня, при необходимости подпишем NDA.