Files
clickstream-ch-kafka-supers…/docs/course/lessons/01_kafka_to_clickhouse.md
T
ddadmin d76c03655b docs(course): переведены уроки на стартовую историю
- Зачем:
  - учебный путь должен идти через генерацию и штатный пайплайн, а не через архивный сид.
- Что:
  - обновлены уроки 00-06 и стандарт урока под startup-history/backfill.
  - объяснено, что data/*.jsonl остаются кладовкой значений генератора.
  - тест-план переведён на новый штатный запуск и HITL-приёмку.
- Проверка:
  - rg -n \"make clean/up/ddl/data/transform|make data|kafka_load|LIMIT=|2022-11-28|26 из 50\" docs/course docs/TEST_PLAN.md.
  - git diff --cached --check.
2026-07-04 22:22:36 +03:00

276 lines
16 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# Урок 1. Заземление Kafka → ClickHouse (слой STG)
> Формат: **практика** — будешь сам запускать команды и менять код, не только читать.
> Пререквизит: пройден урок 0 (словарь Kafka — топик, партиция, offset, consumer-группа —
> уже знаком и виден в Kafka UI).
> Эталонный путь: [`sql/ddl/stg/10_stg.sql`](../../../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` в этом пути не источник аналитики; пока это только
кладовка значений для источника данных стенда.
```bash
make generated-history-analytics
make up
```
Теперь смотрим, что доехало до ClickHouse. Открой SQL-консоль:
`http://localhost:9123/play` (или Kafka UI на `http://localhost:8082`, чтобы тем же
взглядом, что в уроке 0, увидеть те же сообщения «с другого конца»).
```sql
-- Сколько сырых событий браузера приземлилось
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`)
Здесь живёт вся настройка чтения топика:
```sql
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 их нет:
```sql
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 — добавляем колонку в таблицу-приёмник:**
```sql
ALTER TABLE stg.browser_raw ADD COLUMN kafka_msg_ts DateTime;
```
**Шаг 2 — пересоздаём MV, чтобы оно заполняло новую колонку** (поменять SELECT можно и
«на месте» через `ALTER TABLE ... MODIFY QUERY`, но для наглядности пересоздадим целиком):
```sql
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`):
```sql
TRUNCATE TABLE stg.browser_raw;
```
Затем возвращаем стенд в чистое состояние и заново создаём стартовую историю, чтобы
Kafka-движок прочитал сообщения уже с новой схемой:
```bash
make generated-history-analytics
make up
```
**Смотрим результат:**
```sql
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 отвечаешь ты, а не движок.
**Верни как было** (откат — тоже две операции, и порядок важен):
```sql
DROP VIEW stg.mv_kafka_browser_to_stg; -- сначала MV, что ссылается на колонку
ALTER TABLE stg.browser_raw DROP COLUMN kafka_msg_ts;
```
```bash
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: когда доходит до типизации и
проверок качества, важнее видеть и держать под контролем каждый пересчёт, чем экономить
доли секунды на задержке.