@@ -1,20 +1,26 @@
# Урок 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, знакомый тебе по курсовой.
Первое, что делаем с потоком событий — складываем его в таблицу ** как есть** , ничего в нё м
не меняя. А рядом, в соседних колонках, пишем метаданные доставки: из какого топика и партиции
пришло сообщение, с каким offset'ом и в какое время . Это слой **STG ** — тот самый staging,
знакомый тебе по курсовой: первая «посадочная площадка», куда поток приземляется в сыром виде,
до любой обработки.
Зачем хранить сырой JSON строкой и не парсить его сразу:
@@ -24,10 +30,15 @@
- и, главное, чтобы приём не падал из-за одного кривого поля. Разбор JSON и проверки
качества — это уже следующий слой (урок 2), а STG принимает всё подряд.
`kafka_offset` здесь не просто метаданные: пара «партиция + offset» однозначно указывает
на конкретное сообщение в топике — п о ней всегда понятно, та же это запись или другая.
Сама таблица повторы при этом не отсеивает (`MergeTree` ничего не дедуплицирует) — если
понадобится, дубли убирают уже на следующих слоях .
Отдельно про `kafka_offset` — это не просто справочная метка. Помнишь из урока 0: пара
«партиция + offset» однозначно указывает на конкретное сообщение в топике. П о ней всегда
видно, та же это запись или другая, — пригодится, когда дальше начнём сверять данные между
слоями .
При этом **повторы STG не отсеивает ** . Движок этих таблиц — `MergeTree` , и он не
дедуплицирует, то есть не убирает строки-дубли: что пришло, то и легло, даже если две записи
окажутся одинаковыми. Если дубли потом помешают — их убирают уже на следующих слоях, а STG
держит всё подряд.
> **В проде иначе. ** На потоке в десятки тысяч сообщений в секунду читателей будет
> несколько, и Kafka сама делит работу между ними. В этом уроке — один читатель
@@ -73,6 +84,8 @@ LIMIT 5;
## 3. Загляни внутрь (`sql/ddl/stg/10_stg.sql` )
### Три кирпича слоя
Весь слой STG собран из **трёх кирпичей ** , и каждый топик повторяет одну и ту же тройку:
| Кирпич | Объект | Движок | Что делает |
@@ -81,13 +94,21 @@ LIMIT 5;
| 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` .
Связка работает так. Т аблица с `ENGINE = Kafka` (кирпич 2) — это не хранилище, а «кран» к
топику: через неё ClickHouse читает сообщения, но **сами данные она не копит** . Забирает их
третий кирпич — **Materialized View ** (MV) .
Стоит задержаться на двух местах файла.
И тут стоит остановиться на самом слове. Обычное представление (view) — это сохранённый
запрос: данные оно считает только тогда, когда его спросишь. «Materialized» (материализованное)
значит другое: оно срабатывает **само ** на каждую новую порцию из источника и сразу
складывает результат в постоянную таблицу. Получается цепочка: сообщение появилось в топике →
MV тут же подхватило его и положило в `MergeTree` -таблицу `*_raw` (кирпич 1), где оно и лежит.
**Таблица-источник Kafka ( `kafka_*_raw` )** — здесь живёт вся настройка чтения:
Дальше за держимся на двух местах файла.
### Таблица-источник Kafka (`kafka_*_raw` )
Здесь живёт вся настройка чтения топика:
```sql
ENGINE = Kafka
@@ -100,12 +121,19 @@ SETTINGS
kafka_handle_error_mode = 'stream'; -- кривое сообщение не рвёт чтение топика
```
`JSONAsString` — почему мы и можем класть `raw` строкой: ClickHouse не пытается разобрать
JSON на этом этапе. `kafka_handle_error_mode = 'stream'` — ровно то «STG принимает всё»
из секции 1: битое сообщение не уронит консьюмера.
Две настройки тут — самые важные для всего урока:
**Materialized View** — здесь сообщение превращается в строку таблицы. Метаданные берутся
из виртуальных колонок Kafka-движка (`_topic` , `_partition` , `_offset` , `_timestamp_ms` ):
- `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
@@ -118,10 +146,10 @@ SELECT
FROM stg.kafka_browser_raw;
```
Тут без фокусов: каждая виртуальная колонка ложится в свою. `_timestamp_ms` — это уже
готовый `DateTime64(3)` (время сообщения с точностью до миллисекунд), поэтому идёт в
`kafka_ts` как есть, без преобразований . В секции 4 ты положишь рядом ещё одно время из
Kafka и увидишь, чем они отличаются.
Тут без фокусов: каждая виртуальная колонка ложится в свою. Одно место стоит запомнить —
`_timestamp_ms` : это уже готовый `DateTime64(3)` (время сообщения с точностью до миллисекунд),
поэтому оно идёт в `kafka_ts` как есть, без всякого преобразования . В секции 4 ты положишь
рядом ещё одно время из Kafka и увидишь, чем они отличаются.
---