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