diff --git a/docs/course/LEARNING_PLAN.md b/docs/course/LEARNING_PLAN.md index ce2dcdb..db33bd6 100644 --- a/docs/course/LEARNING_PLAN.md +++ b/docs/course/LEARNING_PLAN.md @@ -59,7 +59,7 @@ Kafka из роадмапа — оно даёт словарь терминов. | Урок | Путь | Файлы | Вердикт | |------|------|-------|---------| -| 1 | STG (Kafka→CH) | `sql/ddl/stg/10_stg.sql` | **годно как есть** (+1 строка «зачем») | +| 1 | STG (Kafka→CH) | `sql/ddl/stg/10_stg.sql` | **точечно править** (баг конвертации `kafka_ts`, найден на стенде) | | 2 | STG→ODS | `sql/ods/20_stg_to_ods.sql` + DDL `sql/ddl/ods/20_ods.sql` | **точечно править** | | 3 | ODS→DDS (+ демоут DM) | `sql/dds/30_ods_to_dds.sql` + DDL `sql/ddl/dds/30_dds.sql`; DM `sql/dm/40_dds_to_dm.sql` + `sql/ddl/dm/40_dm.sql` | **точечно править** | | 4 | Airflow DAG | `airflow/dags/etl_pipeline_dag.py` | **точечно править** | @@ -73,12 +73,19 @@ Kafka из роадмапа — оно даёт словарь терминов. Правки точечные, делаем не отдельным забегом, а в составе соответствующего урока. **Урок 1 — STG (`sql/ddl/stg/10_stg.sql`):** -- [ ] Добавить одну строку «зачем» к `fromUnixTimestamp64Milli(toInt64(_timestamp_ms))` - в MV — это единственный неочевидный трюк файла (конвертация Kafka-таймстампа). +- [x] **Баг конвертации `kafka_ts` (найден прогоном на стенде 2026-06-03).** + Было `fromUnixTimestamp64Milli(toInt64(_timestamp_ms))` во всех 4 MV: `_timestamp_ms` + это `DateTime64(3)`, `toInt64()` срезает его до **секунд**, и `fromUnixTimestamp64Milli` + читает секунды как миллисекунды → `kafka_ts` = `1970-01-21` по всему стенду. + Исправлено на `_timestamp_ms AS kafka_ts` (прямое присвоение, мс сохраняются). + Прежний аудит ошибочно пометил путь «годно как есть» — потому что его не прогоняли. - Колоночные комментарии-пересказы (`-- Партиция`, `-- Имя топика`) безвредны — не трогаем. - В финале урока — мост к уроку 2: вопрос «почему дальше не MV?». **Урок 2 — STG→ODS, типизация + DQ-split (`sql/ods/20_stg_to_ods.sql` + DDL `sql/ddl/ods/20_ods.sql`):** +- [ ] **Кандидат на war-story по типам:** баг `kafka_ts` из урока 1 (`toInt64(DateTime64)` + молча срезал миллисекунды → 1970). Это идеальная иллюстрация темы урока — «тихая + потеря данных на неверном типе». Решить при написании, давать ли как пример. - [ ] Добавить однострочный «поток данных» в шапку `20_stg_to_ods.sql` (`stg.*_raw → ods.* + ods.*_errors`) — п.3 чеклиста. - [ ] Добавить строку «зачем»: основная таблица = валидный ключ (но может иметь diff --git a/docs/course/LESSON_STANDARD.md b/docs/course/LESSON_STANDARD.md index 1df84de..a4d97a6 100644 --- a/docs/course/LESSON_STANDARD.md +++ b/docs/course/LESSON_STANDARD.md @@ -24,10 +24,12 @@ самостоятельный менти не застрял со сломанным стендом без ментора. (Урок 0 — без этого шага, только наблюдение.) 5. **Проверь себя** — самопроверка (раздел 3). -6. **Принеси на сессию** — **конкретный артефакт** плюс вопросы: скрин видимого - результата правки (например, красный DAG или `+1` в таблице ошибок), один абзац - «своими словами» про паттерн урока или запрос, который пришлось написать. Артефакт - делает самопроверку проверяемой, а сверку — предметной (критерий успеха `PRD` §5). +6. **Что должно получиться** — **конкретный видимый результат**, который менти проверяет + сам: скрин эффекта правки (например, красный DAG или `+1` в таблице ошибок) плюс один + абзац «своими словами» про паттерн урока или запрос, который пришлось написать. Это и + есть самопроверка (раздел 3); те же вопросы «своими словами» менти разбирает с ментором + на еженедельном созвоне. (Урок 0 — без правки, поэтому результат здесь — что менти + увидел при наблюдении.) ## 2. Стандарт качества эталонного кода diff --git a/docs/course/PRD.md b/docs/course/PRD.md index 34ef49c..f661f19 100644 --- a/docs/course/PRD.md +++ b/docs/course/PRD.md @@ -6,6 +6,9 @@ > поверхность потребления. Обязательных уроков стало 0–5, опциональный Superset — урок 6. > Затронуты §4 (скоуп) и §3/§5 (ожидаемый такт — ~день на урок). Это изменение рамки, > а не план реализации. +> Поправка 2026-06-03 (терминология): режим сопровождения — еженедельный **созвон**, а не +> «сессия»; менти проходит материал сам и ничего «не приносит», а на созвоне ментор +> разбирает затыки и проверяет глубину понимания. Затронуты §1, §3, §5. > Назначение документа: зафиксировать для будущих сессий, что это за курс, зачем > он, что входит в скоуп работ, а что нет. Это договорная **рамка**, а не план > реализации и не стандарт уроков (см. раздел «Связанные документы»). @@ -30,8 +33,8 @@ data/*.jsonl → Kafka → ClickHouse (Kafka engine + MV → STG) → Airflow ET прошедших базу (SQL, моделирование, Python, Git, Docker, Airflow). Цель курса — превратить стенд в **самостоятельный учебный материал**, по которому -продвинутый менти проходит ключевые паттерны сам, а ментор подключается на обычной -еженедельной сессии (что получилось / что нет / вопросы / план на неделю). Долгий +продвинутый менти проходит ключевые паттерны сам, а ментор подключается на обычном +еженедельном созвоне (что получилось / что нет / вопросы / план на неделю). Долгий разбор кода вживую форматом не предусмотрен — поэтому материал обязан быть самодостаточным. @@ -61,7 +64,7 @@ Kafka → ClickHouse → BI). - **Режим:** самостоятельный, асинхронный. Менти клонирует репозиторий, поднимает стенд у себя (`make up`) и идёт по урокам из `docs/course/` рядом с кодом. Уроки короткие и односоставные — ожидаемый срок прохождения одного **около дня**. -- **Роль ментора:** еженедельная сессия-сверка (покрывает несколько уроков), без +- **Роль ментора:** еженедельный созвон-сверка (покрывает несколько уроков), без построчного разбора кода. - **Железо:** стек тяжёлый (Kafka + ClickHouse + Airflow + Superset + Prometheus + Grafana одновременно). Считаем наличие подходящего железа данностью; стек не режем @@ -103,7 +106,7 @@ Kafka → ClickHouse → BI). около дня на урок: уроки короткие и односоставные). - Может своими словами объяснить паттерн урока и привязать его к обычной кликстрим-аналитике. -- На сессии приносит осмысленные вопросы по сути, а не «застрял на запуске». +- На созвоне задаёт осмысленные вопросы по сути, а не «застрял на запуске». - Эталонный код проходит «тест одного прохода» (см. `LESSON_STANDARD.md`). ## 6. Связанные документы и порядок работ diff --git a/docs/course/README.md b/docs/course/README.md index 2dde4a1..012c506 100644 --- a/docs/course/README.md +++ b/docs/course/README.md @@ -1,8 +1,15 @@ # Курс «Кликстрим на ClickHouse» (со звёздочкой) -Продвинутый курс для менти, уже прошедших базовую программу: ключевые паттерны -инженерии данных на стенде Kafka + ClickHouse + Airflow + мониторинг. Самостоятельный -учебный материал; ментор подключается на еженедельной сессии-сверке. +Продвинутый курс для менти, уже прошедших базовую программу: основные паттерны +инженерии данных на стенде Kafka + ClickHouse + Airflow + мониторинг. Менти проходит +материал самостоятельно, а на еженедельном созвоне с ментором разбирает затыки и отвечает +на вопросы по теме — чтобы проверить глубину понимания. + +Сам курс — это блок «со звёздочкой» открытого роадмапа Data Engineer +([de.dementev.space](https://de.dementev.space/), он же +[на GitHub](https://github.com/dementev-dev/de-roadmap)). Роадмап открыт и устроен так, +что идти по нему можно в одиночку; с ментором — ощутимо короче дорога: добавляются +персональный план, разбор домашек и код-ревью. ## Что где искать @@ -11,7 +18,7 @@ | [`PRD.md`](./PRD.md) | Рамка: зачем курс, цели, аудитория, скоуп, критерии успеха | Чтобы понять «что и зачем». Замороженный документ | | [`LEARNING_PLAN.md`](./LEARNING_PLAN.md) | План обучения: карта уроков, маршрут, аудит эталонных путей | Чтобы понять «в каком порядке и из чего» | | [`LESSON_STANDARD.md`](./LESSON_STANDARD.md) | Стандарт уроков: шаблон урока, качество кода, самопроверка | Рабочий чеклист при написании каждого урока | -| `lessons/` *(будет)* | Сами уроки, по одному файлу | Прохождение курса менти | +| [`lessons/`](./lessons/) | Сами уроки, по одному файлу (есть: урок 1) | Прохождение курса менти | ## Порядок чтения diff --git a/docs/course/lessons/01_kafka_to_clickhouse.md b/docs/course/lessons/01_kafka_to_clickhouse.md new file mode 100644 index 0000000..742ce0e --- /dev/null +++ b/docs/course/lessons/01_kafka_to_clickhouse.md @@ -0,0 +1,244 @@ +# Урок 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: когда доходит до типизации и +проверок качества, важнее видеть и держать под контролем каждый пересчёт, чем экономить +доли секунды на задержке. diff --git a/sql/ddl/stg/10_stg.sql b/sql/ddl/stg/10_stg.sql index cc8a821..1f03d60 100644 --- a/sql/ddl/stg/10_stg.sql +++ b/sql/ddl/stg/10_stg.sql @@ -128,7 +128,7 @@ SELECT _topic AS kafka_topic, _partition AS kafka_partition, _offset AS kafka_offset, - fromUnixTimestamp64Milli(toInt64(_timestamp_ms)) AS kafka_ts, + _timestamp_ms AS kafka_ts, raw FROM stg.kafka_browser_raw; @@ -140,7 +140,7 @@ SELECT _topic AS kafka_topic, _partition AS kafka_partition, _offset AS kafka_offset, - fromUnixTimestamp64Milli(toInt64(_timestamp_ms)) AS kafka_ts, + _timestamp_ms AS kafka_ts, raw FROM stg.kafka_location_raw; @@ -152,7 +152,7 @@ SELECT _topic AS kafka_topic, _partition AS kafka_partition, _offset AS kafka_offset, - fromUnixTimestamp64Milli(toInt64(_timestamp_ms)) AS kafka_ts, + _timestamp_ms AS kafka_ts, raw FROM stg.kafka_device_raw; @@ -164,6 +164,6 @@ SELECT _topic AS kafka_topic, _partition AS kafka_partition, _offset AS kafka_offset, - fromUnixTimestamp64Milli(toInt64(_timestamp_ms)) AS kafka_ts, + _timestamp_ms AS kafka_ts, raw FROM stg.kafka_geo_raw;