- Зачем:
- на уроке 2 решили писать разжёванным языком; уроки 0–1 и шапки курса
остались в сжатом регистре, а слово «черновик»/«Режим: руки» путало менти.
- Что:
- урок 1: расшифрованы staging, MergeTree-дедуп, Materialized View и
виртуальные колонки; секция «Загляни внутрь» разбита на ###-подзаголовки;
плотные абзацы разбиты на пункты; добавлен зачин «О чём урок простыми словами».
- урок 0: добавлен зачин «О чём урок простыми словами» (лёгкая полировка).
- шапки всех уроков: «Статус: черновик. Режим: руки/наблюдение» заменены на
понятное «Формат: практика/наблюдение — …».
- PRD/LEARNING_PLAN/LESSON_STANDARD: убрано слово «черновик» из статуса.
- Проверка:
- grep -rn "черновик" docs/course/ — пусто;
- прочитать урок 1 сверху вниз: термины раскрыты на первом употреблении.
16 KiB
Урок 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. Руки: убедись, что данные текут
Поднимаем стенд, создаём схему и заливаем малый срез (50 строк на топик — этого хватает, чтобы всё увидеть, и прогон быстрый):
make up # поднять инфраструктуру (Kafka, ClickHouse, ...)
make ddl # создать базы и таблицы в ClickHouse (в т.ч. слой STG)
LIMIT=50 make data # залить по 50 строк каждого файла в Kafka-топики
Теперь смотрим, что доехало до 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;
Затем перезаливаем — с пересозданием топиков, чтобы Kafka-движок перечитал сообщения с начала (без этого он считает их уже прочитанными и ничего нового не подхватит):
LIMIT=50 RESET_TOPICS=1 make data
Смотрим результат:
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 clean && make up && make ddl && LIMIT=50 make data.
5. Проверь себя
| Действие | Где смотреть | Что ожидать |
|---|---|---|
LIMIT=50 make data |
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: когда доходит до типизации и проверок качества, важнее видеть и держать под контролем каждый пересчёт, чем экономить доли секунды на задержке.