Files
clickstream-ch-kafka-supers…/docs/course/lessons/01_kafka_to_clickhouse.md
T
ddadmin 77e4fb6119 docs(course): согласованы уроки со startup-history
- Зачем:
  - курс должен проходить на чистом стенде без архивного сида и скрытых шагов.
- Что:
  - уроки 00, 01 и 05 согласованы с генераторными топиками, no-live default и consumer lag.
  - упражнение с kafka_msg_ts переведено на повторную заливку без сброса схемы.
  - учебный путь в операционной документации ведёт через startup-history.
- Проверка:
  - make generated-history-analytics; make up; doc rg checks; git diff --check.
2026-07-05 22:01:00 +03:00

17 KiB
Raw Blame History

Урок 1. Заземление Kafka → ClickHouse (слой STG)

Формат: практика — будешь сам запускать команды и менять код, не только читать. Пререквизит: пройден урок 0 (словарь Kafka — топик, партиция, offset, consumer-группа — уже знаком и виден в Kafka UI). Эталонный путь: sql/ddl/stg/10_stg.sql.

Поток данных одной строкой: Kafka → kafka_*_raw (ENGINE=Kafka) → MV → *_raw (MergeTree)

О чём урок простыми словами: смотрим, как сообщение из Kafka-топика само, без нашего участия, превращается в строку таблицы ClickHouse — и почему на этом первом слое мы кладём JSON целиком, ничего в нём не разбирая.


1. Зачем и где в проде

Первое, что делаем с потоком событий — складываем его в таблицу как есть, ничего в нём не меняя. А рядом, в соседних колонках, пишем метаданные доставки: из какого топика и партиции пришло сообщение, с каким offset'ом и в какое время. Это слой STG — тот самый staging, знакомый тебе по курсовой: первая «посадочная площадка», куда поток приземляется в сыром виде, до любой обработки.

Зачем хранить сырой JSON строкой и не парсить его сразу:

  • чтобы отладке было на что опереться: когда дальше по пайплайну что-то «не сходится», всегда есть исходное сообщение для сверки;
  • чтобы любой следующий слой можно было пересобрать из STG, не вычитывая Kafka заново;
  • и, главное, чтобы приём не падал из-за одного кривого поля. Разбор JSON и проверки качества — это уже следующий слой (урок 2), а STG принимает всё подряд.

Отдельно про kafka_offset — это не просто справочная метка. Помнишь из урока 0: пара «партиция + offset» однозначно указывает на конкретное сообщение в топике. По ней всегда видно, та же это запись или другая, — пригодится, когда дальше начнём сверять данные между слоями.

При этом повторы STG не отсеивает. Движок этих таблиц — MergeTree, и он не дедуплицирует, то есть не убирает строки-дубли: что пришло, то и легло, даже если две записи окажутся одинаковыми. Если дубли потом помешают — их убирают уже на следующих слоях, а STG держит всё подряд.

В проде иначе. На потоке в десятки тысяч сообщений в секунду читателей будет несколько, и Kafka сама делит работу между ними. В этом уроке — один читатель (kafka_num_consumers = 1) и маленький срез данных: нам важно понять, как поток вообще попадает в базу, а не выжимать скорость.


2. Руки: убедись, что данные текут

Поднимаем стенд и создаём стартовую историю. Это штатный путь курса: готовый источник данных стенда пишет события в Kafka, ClickHouse читает их в STG, затем batch строит ODS, DDS и DM. Файлы data/*.jsonl в этом пути не источник аналитики; пока это только кладовка значений для источника данных стенда.

make generated-history-analytics
make up

Теперь смотрим, что доехало до ClickHouse. Открой SQL-консоль: http://localhost:9123/play (или Kafka UI на http://localhost:8082, чтобы тем же взглядом, что в уроке 0, увидеть те же сообщения «с другого конца»).

-- Сколько сырых событий браузера приземлилось
SELECT count() FROM stg.browser_raw;

-- Как выглядит приземлённая строка: метаданные доставки + сырой JSON
SELECT kafka_topic, kafka_partition, kafka_offset, kafka_ts, raw
FROM stg.browser_raw
ORDER BY kafka_offset
LIMIT 5;

Что важно заметить:

  • raw — это целый JSON строкой, никто его пока не разбирал;
  • поля kafka_topic, kafka_partition, kafka_offset, kafka_ts заполнены: видно, откуда именно приехала запись;
  • offset'ы идут по возрастанию и не повторяются.

3. Загляни внутрь (sql/ddl/stg/10_stg.sql)

Три кирпича слоя

Весь слой STG собран из трёх кирпичей, и каждый топик повторяет одну и ту же тройку:

Кирпич Объект Движок Что делает
1 stg.browser_raw MergeTree хранит сырьё + метаданные, постоянно
2 stg.kafka_browser_raw ENGINE = Kafka читает топик, ничего не хранит
3 stg.mv_kafka_browser_to_stg MATERIALIZED VIEW перекладывает из (2) в (1) на лету

Связка работает так. Таблица с ENGINE = Kafka (кирпич 2) — это не хранилище, а «кран» к топику: через неё ClickHouse читает сообщения, но сами данные она не копит. Забирает их третий кирпич — Materialized View (MV).

И тут стоит остановиться на самом слове. Обычное представление (view) — это сохранённый запрос: данные оно считает только тогда, когда его спросишь. «Materialized» (материализованное) значит другое: оно срабатывает само на каждую новую порцию из источника и сразу складывает результат в постоянную таблицу. Получается цепочка: сообщение появилось в топике → MV тут же подхватило его и положило в MergeTree-таблицу *_raw (кирпич 1), где оно и лежит.

Дальше задержимся на двух местах файла.

Таблица-источник Kafka (kafka_*_raw)

Здесь живёт вся настройка чтения топика:

ENGINE = Kafka
SETTINGS
    kafka_broker_list = 'kafka:29092',      -- адрес брокера внутри Docker-сети
    kafka_topic_list  = 'browser_events',   -- какой топик читаем
    kafka_group_name  = 'ch_stg_browser',   -- consumer-группа: по ней Kafka помнит offset'ы
    kafka_format      = 'JSONAsString',      -- берём сообщение целиком, как строку
    kafka_num_consumers = 1,
    kafka_handle_error_mode = 'stream';      -- кривое сообщение не рвёт чтение топика

Две настройки тут — самые важные для всего урока:

  • kafka_format = 'JSONAsString' — вот почему мы и можем класть raw одной строкой: ClickHouse берёт тело сообщения как текст и не пытается разобрать JSON на этом этапе;
  • kafka_handle_error_mode = 'stream' — это ровно то «STG принимает всё» из секции 1: одно битое сообщение не уронит консьюмера, чтение топика продолжится.

Materialized View: как сообщение становится строкой

Здесь сообщение из топика превращается в строку таблицы. Откуда MV берёт метаданные доставки? Из виртуальных колонок Kafka-движка — это служебные поля (_topic, _partition, _offset, _timestamp_ms), которые движок подставляет к каждому сообщению сам, хотя в теле JSON их нет:

SELECT
    now64(3) AS ingest_ts,
    _topic AS kafka_topic,
    _partition AS kafka_partition,
    _offset AS kafka_offset,
    _timestamp_ms AS kafka_ts,
    raw
FROM stg.kafka_browser_raw;

Тут без фокусов: каждая виртуальная колонка ложится в свою. Одно место стоит запомнить — _timestamp_ms: это уже готовый DateTime64(3) (время сообщения с точностью до миллисекунд), поэтому оно идёт в kafka_ts как есть, без всякого преобразования. В секции 4 ты положишь рядом ещё одно время из Kafka и увидишь, чем они отличаются.


4. Управляемая правка: протащи ещё одно поле метаданных

Задача: добавить в stg.browser_raw ещё одну колонку — то же время сообщения, но из другой виртуальной колонки Kafka: _timestamp (тип DateTime, точность до секунды) — и увидеть её заполненной.

Соль урока — в порядке из двух шагов: MV не «дотянет» новое поле само, его сначала нужно добавить в таблицу-приёмник. Схема приёмника и MV меняются вместе.

Шаг 1 — добавляем колонку в таблицу-приёмник:

ALTER TABLE stg.browser_raw ADD COLUMN kafka_msg_ts DateTime;

Шаг 2 — пересоздаём MV, чтобы оно заполняло новую колонку (поменять SELECT можно и «на месте» через ALTER TABLE ... MODIFY QUERY, но для наглядности пересоздадим целиком):

DROP VIEW stg.mv_kafka_browser_to_stg;

CREATE MATERIALIZED VIEW stg.mv_kafka_browser_to_stg
TO stg.browser_raw
AS
SELECT
    now64(3) AS ingest_ts,
    _topic AS kafka_topic,
    _partition AS kafka_partition,
    _offset AS kafka_offset,
    _timestamp_ms AS kafka_ts,
    _timestamp AS kafka_msg_ts,    -- ещё одно время из Kafka, но с точностью только до секунды
    raw
FROM stg.kafka_browser_raw;

Перезаливаем срез, чтобы новая колонка заполнилась. Сначала чистим только таблицу приёмника — иначе рядом останутся строки из секции 2, вставленные ещё до ADD COLUMN, и в них kafka_msg_ts будет пустой (1970-01-01):

TRUNCATE TABLE stg.browser_raw;

Затем повторно создаём стартовую историю без очистки volumes. Это важно: CLEAN_START=0 сохраняет твою новую колонку и пересозданное MV, но добавляет свежие сообщения в Kafka, чтобы ClickHouse прочитал их уже с новой схемой.

CLEAN_START=0 make generated-history-analytics

Смотрим результат:

SELECT kafka_ts, kafka_msg_ts
FROM stg.browser_raw
ORDER BY kafka_offset
LIMIT 5;

Видно два времени одного и того же сообщения: kafka_ts (из _timestamp_ms) — с миллисекундами, kafka_msg_ts (из _timestamp) — округлённое до секунды. Момент тот же, точность разная: Kafka отдаёт эту метку через две виртуальные колонки сразу.

Порядок важен, а ClickHouse молчит. Пересоздай MV с _timestamp, не добавив колонку kafka_msg_ts в таблицу, — ошибки не будет. ClickHouse сопоставляет колонки MV с таблицей по имени и лишнюю просто отбрасывает: данные текут, а твоё новое поле молча теряется. Обратный случай — колонку добавил, а в MV не указал — тоже не упадёт: поле заполнится дефолтом. Вывод: за синхронность схемы и MV отвечаешь ты, а не движок.

Верни как было (откат — тоже две операции, и порядок важен):

DROP VIEW stg.mv_kafka_browser_to_stg;          -- сначала MV, что ссылается на колонку
ALTER TABLE stg.browser_raw DROP COLUMN kafka_msg_ts;
make ddl     # пересоздаёт эталонное MV из 10_stg.sql, схема снова как в репозитории

Если запутался в состоянии — всегда есть полный чистый прогон: make generated-history-analytics && make up.


5. Проверь себя

Действие Где смотреть Что ожидать
make generated-history-analytics && make up SELECT count() FROM stg.browser_raw счётчик > 0
глянуть строку SELECT raw FROM stg.browser_raw LIMIT 1 валидный JSON целиком, неразобранный
глянуть offset'ы SELECT kafka_offset FROM stg.browser_raw ORDER BY kafka_offset идут по возрастанию, без дублей
правка из секции 4 SELECT kafka_msg_ts FROM stg.browser_raw LIMIT 5 колонка заполнена временем сообщения

6. Что должно получиться

После урока у тебя на руках — видимый результат (одно на выбор):

  • строка stg.browser_raw с заполненными kafka_topic, kafka_partition, kafka_offset, kafka_ts и raw;
  • либо результат твоей правки — две колонки времени, kafka_ts и kafka_msg_ts, рядом.

И проверь себя на словах: сможешь объяснить, чем таблица с ENGINE = Kafka отличается от MergeTree и зачем между ними нужен Materialized View? Короткая зацепка для ответа — что случится с данными, если MV убрать? Примерно такие вопросы по теме урока всплывут на еженедельном созвоне.


Мост к уроку 2

Мы приземлили поток через Materialized View — почти «в реальном времени». Возникает честный вопрос: раз MV так удобно перекладывает данные, почему дальше, в ODS, мы не продолжаем на MV, а уходим в батч? Ответ — в уроке 2: когда доходит до типизации и проверок качества, важнее видеть и держать под контролем каждый пересчёт, чем экономить доли секунды на задержке.