Files
clickstream-ch-kafka-supers…/docs/course/lessons/01_kafka_to_clickhouse.md
T
ddadminandClaude Opus 4.8 f4d3118a42 docs(course): добавлен урок 1 и выровнена рамка под созвоны
- Зачем:
  - нужен первый урок курса по эталонному пути STG, а рамка курса описывала сопровождение как «сессию-сверку», хотя по факту это самостоятельная работа + еженедельный созвон.
- Что:
  - добавлен docs/course/lessons/01_kafka_to_clickhouse.md (Kafka → ClickHouse, слой STG) по шаблону LESSON_STANDARD.
  - в sql/ddl/stg/10_stg.sql исправлен баг kafka_ts во всех 4 MV: toInt64(DateTime64) срезал миллисекунды, kafka_ts по всему стенду был 1970-01-21; теперь _timestamp_ms присваивается напрямую (downstream на kafka_ts не опирается).
  - урок 1 §3/§4 приведены к исправленному коду; врезка про рассинхрон MV↔таблица описывает реальное поведение (молчаливый сброс лишней колонки, не ошибка).
  - «сессия/сессия-сверка» → «созвон» в PRD (датированная поправка), LESSON_STANDARD §6 (секция «Что должно получиться») и README курса; добавлена ссылка на открытый DE-роадмап.
  - в LEARNING_PLAN исправлен вердикт аудита урока 1 (баг найден прогоном), war-story про toInt64(DateTime64) припаркована в урок 2.
- Проверка:
  - прогон на стенде: make up && make ddl && LIMIT=50 make data → kafka_ts = 2026-… с миллисекундами; правка §4 (ALTER + пересоздание MV + TRUNCATE + перезаливка) и откат отработали.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-03 23:03:24 +03:00

245 lines
14 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)`
---
## 1. Зачем и где в проде
Первое, что делаем с потоком событий — складываем его в таблицу как есть, а рядом пишем
метаданные доставки: из какого топика и партиции пришло сообщение, с каким offset'ом и
временем. Это слой STG — тот самый staging, знакомый тебе по курсовой.
Зачем хранить сырой JSON строкой и не парсить его сразу:
- чтобы отладке было на что опереться: когда дальше по пайплайну что-то «не сходится»,
всегда есть исходное сообщение для сверки;
- чтобы любой следующий слой можно было пересобрать из STG, не вычитывая Kafka заново;
- и, главное, чтобы приём не падал из-за одного кривого поля. Разбор JSON и проверки
качества — это уже следующий слой (урок 2), а STG принимает всё подряд.
`kafka_offset` здесь не просто метаданные: пара «партиция + offset» однозначно указывает
на конкретное сообщение в топике — по ней всегда понятно, та же это запись или другая.
Сама таблица повторы при этом не отсеивает (`MergeTree` ничего не дедуплицирует) — если
понадобится, дубли убирают уже на следующих слоях.
> **В проде иначе.** На потоке в десятки тысяч сообщений в секунду читателей будет
> несколько, и Kafka сама делит работу между ними. В этом уроке — один читатель
> (`kafka_num_consumers = 1`) и маленький срез данных: нам важно понять, как поток вообще
> попадает в базу, а не выжимать скорость.
---
## 2. Руки: убедись, что данные текут
Поднимаем стенд, создаём схему и заливаем **малый срез** (50 строк на топик — этого
хватает, чтобы всё увидеть, и прогон быстрый):
```bash
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, увидеть те же сообщения «с другого конца»).
```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` — это не хранилище, а «кран» к топику. Сама по себе
она данные не копит; данные забирает Materialized View и складывает их в обычную
`MergeTree`-таблицу. Сообщение появилось в топике → MV тут же положило его в `*_raw`.
Стоит задержаться на двух местах файла.
**Таблица-источник 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'; -- кривое сообщение не рвёт чтение топика
```
`JSONAsString` — почему мы и можем класть `raw` строкой: ClickHouse не пытается разобрать
JSON на этом этапе. `kafka_handle_error_mode = 'stream'` — ровно то «STG принимает всё»
из секции 1: битое сообщение не уронит консьюмера.
**Materialized View** — здесь сообщение превращается в строку таблицы. Метаданные берутся
из виртуальных колонок Kafka-движка (`_topic`, `_partition`, `_offset`, `_timestamp_ms`):
```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
LIMIT=50 RESET_TOPICS=1 make data
```
**Смотрим результат:**
```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 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: когда доходит до типизации и
проверок качества, важнее видеть и держать под контролем каждый пересчёт, чем экономить
доли секунды на задержке.