diff --git a/docs/course/lessons/00_kafka_intro.md b/docs/course/lessons/00_kafka_intro.md index b7fc161..67ebc4d 100644 --- a/docs/course/lessons/00_kafka_intro.md +++ b/docs/course/lessons/00_kafka_intro.md @@ -76,16 +76,20 @@ consumer читает в своём темпе, не трогая producer'а. - **Offset** — порядковый номер сообщения в партиции (0, 1, 2, …); - **Timestamp** — когда сообщение легло в Kafka (время *доставки*, не время самого события); -- **Value** — тело: JSON события целиком, например: +- **Value** — тело: JSON события целиком. Например, у сообщения с offset `0`: ```json -{"event_id": "8cca1c7d-...", "event_timestamp": "2026-01-01 00:01:00.000000", - "event_type": "pageview", "browser_name": "Chrome", "browser_language": "sat_IN"} +{"event_id": "b3336fa4-184d-44e5-9878-98c8f8058496", + "event_timestamp": "2026-01-01 00:00:00.000000", "event_type": "pageview", + "click_id": "76a146f2-adc3-4065-8140-9fc42523b554", + "browser_name": "Opera", + "browser_user_agent": "Opera/9.44.(Windows NT 5.0; iw-IL) Presto/2.9.183 Version/12.00", + "browser_language": "ha_NG"} ``` - Загляни внутрь Value: у события есть своё `event_timestamp` (когда оно случилось, модельное время стенда), и оно отличается от Kafka-Timestamp (когда оно попало в топик). +Kafka-Timestamp зависит от времени импорта, поэтому у тебя будет своё значение. Два разных времени у одной записи — запомни этот момент, в уроке 1 он всплывёт уже на стороне ClickHouse. diff --git a/docs/course/lessons/02_stg_to_ods.md b/docs/course/lessons/02_stg_to_ods.md index b4e5f1f..6cbe71a 100644 --- a/docs/course/lessons/02_stg_to_ods.md +++ b/docs/course/lessons/02_stg_to_ods.md @@ -93,25 +93,23 @@ ODS мы пересобираем целиком, одной задачей Airf ``` Статистика ODS: - ┌─table──────────────────────┬─rows─┐ - │ ods.browser_event │ ... │ - │ ods.location_event │ ... │ - │ ods.device_by_click │ ... │ - │ ods.geo_by_click │ ... │ - │ ods.browser_event_errors │ 0 │ - │ ods.location_event_errors │ 0 │ - │ ods.device_by_click_errors │ 0 │ - │ ods.geo_by_click_errors │ 0 │ - └────────────────────────────┴──────┘ + ┌─table──────────────────────┬───rows─┐ + │ ods.browser_event │ 280437 │ + │ ods.location_event │ 280437 │ + │ ods.device_by_click │ 26083 │ + │ ods.geo_by_click │ 26083 │ + │ ods.browser_event_errors │ 0 │ + │ ods.location_event_errors │ 0 │ + │ ods.device_by_click_errors │ 0 │ + │ ods.geo_by_click_errors │ 0 │ + └────────────────────────────┴────────┘ ``` - Прочитаем эту табличку — в ней три вещи, которые стоит заметить. **Все четыре `*_errors` — по нулям.** Значит, стартовая история чистая: ни одна запись не дала ошибки разбора, столбец `parse_errors` у всех пустой. Это нормально — данные стенда аккуратные. Ошибки мы увидим в секции 4, когда сами их устроим. - **`browser` и `location` идут в одном зерне события.** Сколько событий пришло, столько строк и ожидаем увидеть после типизации, если ключи валидны. @@ -134,9 +132,10 @@ SELECT count() AS stg_rows, FROM stg.geo_raw; ``` -`distinct_clicks` должен быть меньше или равен `stg_rows` и совпадать с числом строк в -`ods.geo_by_click`. Значит, это схлопнутые повторы по `click_id`, а не пропавшие данные. -Ничего не потерялось молча. +На эталонном мире запрос возвращает `stg_rows = 280437` и +`distinct_clicks = 26083`. Второе число совпадает с числом строк в +`ods.geo_by_click`. Значит, это схлопнутые повторы по `click_id`, а не пропавшие +данные. Ничего не потерялось молча. --- @@ -245,11 +244,10 @@ ORDER BY (click_id) Урок про типы — так давай **намеренно ошибёмся типом** и посмотрим, что будет. Это самый поучительный момент урока. -Возьмём координату `geo_latitude` — широту. Это дробное число, например `50.82709`. Достаём мы +Возьмём координату `geo_latitude` — широту. Это дробное число, например `-7.60361`. Достаём мы её через `toFloat64OrNull` — «привести к дробному числу». Заменим тип на целочисленный — -`toInt64OrNull`, «привести к целому». Для строки `"50.82709"` целого числа не получится +`toInt64OrNull`, «привести к целому». Для строки `"-7.60361"` целого числа не получится (там точка, дробная часть), и функция вернёт `NULL`. То есть широта просто исчезнет. - Из секции 3 помним: разбор продублирован, поэтому правок будет **две** — в обоих `INSERT` блока `GEO EVENTS`. Открой `sql/ods/20_stg_to_ods.sql`, найди оба вхождения и в каждом замени @@ -271,10 +269,13 @@ make transform И смотрим на ту же «Статистику ODS». Таблица ошибок гео, которая была пустой, теперь полная: ``` - │ ods.geo_by_click │ ... │ - │ ods.geo_by_click_errors │ ... │ ← было 0 + │ ods.geo_by_click │ ≈ 26083 │ + │ ods.geo_by_click_errors │ 280437 │ ← было 0 ``` +В таблице ошибок число точное. В `ods.geo_by_click` после фонового схлопывания +останется `26083` строки, но сразу после прогона число может быть больше. + А в самой основной таблице широта пропала — но не молча, рядом стоит метка: ```sql @@ -285,11 +286,10 @@ LIMIT 4; ``` ┌─click_id─────┬─geo_latitude─┬─geo_longitude─┬─parse_errors─────────┐ -│ 58cdfc1e-... │ ᴺᵁᴸᴸ │ -0.2 │ ['bad_geo_latitude'] │ -│ 9ffd819b-... │ ᴺᵁᴸᴸ │ 85.37752 │ ['bad_geo_latitude'] │ +│ cee12466-... │ ᴺᵁᴸᴸ │ -8.07257 │ ['bad_geo_latitude'] │ +│ a8e39850-... │ ᴺᵁᴸᴸ │ 37.92792 │ ['bad_geo_latitude'] │ └──────────────┴──────────────┴───────────────┴──────────────────────┘ ``` - Вот теперь видно всё разом — и DQ-split, и «двойной учёт» из секции 3 вживую. Строки с валидным `click_id` остались в основной таблице с пометкой `bad_geo_latitude`. И те же записи попали в @@ -304,7 +304,6 @@ LIMIT 4; > это молча срезало миллисекунды, и время по всему стенду уехало в `1970-01-21`. Ничто на это > не указывало — поймали только прогоном на стенде. Мораль урока: тип выбирают осознанно, даже > когда функция «не падает». - **Верни как было.** Откати обе правки — верни `toFloat64OrNull` в оба места. Если запутался, проще одной командой откатить весь файл к версии из репозитория: @@ -317,7 +316,6 @@ make transform После этого `geo_by_click_errors` снова `0`, широта на месте. А если стенд совсем «поплыл», пройди [канонический сброс](../README.md#подготовка-и-канонический-сброс). - --- @@ -329,7 +327,6 @@ make transform | почему `device`/`geo` меньше событий | запрос `uniqExact(click_id)` по `stg.geo_raw` | число различных `click_id` совпадает с `ods.geo_by_click` | | правка из секции 4 | блок «Статистика ODS» | `ods.geo_by_click_errors` прыгнул с `0` на ненулевое число | | правка из секции 4 | `SELECT geo_latitude, parse_errors FROM ods.geo_by_click` | широта `NULL`, в `parse_errors` — `bad_geo_latitude` | - --- @@ -340,7 +337,6 @@ make transform - скрин блока «Статистика ODS», где после правки `ods.geo_by_click_errors` ушёл с `0` на ненулевое число; - либо выборка из `ods.geo_by_click` с пустой широтой и меткой `bad_geo_latitude` рядом. - И проверь себя на словах — примерно эти вопросы всплывут на еженедельном созвоне: diff --git a/docs/course/lessons/03_ods_to_dds.md b/docs/course/lessons/03_ods_to_dds.md index f965f85..c48aeae 100644 --- a/docs/course/lessons/03_ods_to_dds.md +++ b/docs/course/lessons/03_ods_to_dds.md @@ -74,10 +74,10 @@ ``` Статистика DDS: - ┌─table─────┬─rows─┐ - │ dds.click │ ... │ - │ dds.event │ ... │ - └───────────┴──────┘ + ┌─table─────┬───rows─┐ + │ dds.click │ 26083 │ + │ dds.event │ 280437 │ + └───────────┴────────┘ ``` Прочитаем эти две строки. @@ -95,10 +95,9 @@ ``` ┌─check_date─┬─layer─┬─table_name──────────┬─check_name────┬─check_value─┐ - │ 2026-06-05 │ dds │ event_without_click │ orphan_events │ 0 │ + │ 2026-07-23 │ dds │ event_without_click │ orphan_events │ 0 │ └────────────┴───────┴─────────────────────┴───────────────┴─────────────┘ ``` - В колонке `check_date` стоит `today()` из кода витрины, так что у тебя там будет сегодняшняя дата — не пугайся, если она не совпадёт с примером. @@ -106,7 +105,6 @@ `orphan_events = 0` — ни одной сироты. Каждое событие нашло свой клик в `dds.click`. На чистой стартовой истории так и должно быть: данные аккуратные, ничего не потерялось. В секции 4 мы сироту устроим сами — и эта строка оживёт. - Проверь нолик сам, не верь на слово. Открой SQL-консоль `http://localhost:9123/play` (пользователь `default`, пароль `123456`) и посчитай сирот напрямую: @@ -122,7 +120,6 @@ WHERE click_id IS NOT NULL `NOT IN (SELECT ...)` читается прямо по словам: «click_id события **не входит** в список всех click_id из `dds.click`». То есть событие ссылается на клик, которого в карточках кликов нет. Сейчас таких ноль — запомни этот запрос, в секции 4 он покажет другое число. - --- @@ -308,7 +305,6 @@ make transform После этого `orphan_events` снова `0`, придуманное событие исчезло. А если стенд совсем «поплыл», пройди [канонический сброс](../README.md#подготовка-и-канонический-сброс). - --- @@ -321,7 +317,6 @@ make transform | почему `click` меньше `event` | запрос `count()` по `dds.click` и `dds.event` | много событий ссылаются на меньшее число кликов — норма | | правка из секции 4 (вставили сироту) | запрос `count()` сирот в play-консоли | `0 → 1` | | та же сирота через `dm.v_events_enriched` | `SELECT device_type, geo_country ...` | поля клика пустые (`NULL`) — это `LEFT JOIN` | - --- @@ -332,7 +327,6 @@ make transform - скрин запроса со счётчиком сирот: было `0`, после вставки стало `1`; - либо выборка из `dm.v_events_enriched` по событию-сироте, где `device_type` и `geo_country` пусты, — `LEFT JOIN` сохранил событие без клика. - И проверь себя на словах — примерно эти вопросы всплывут на еженедельном созвоне: diff --git a/docs/course/lessons/04_airflow_orchestration.md b/docs/course/lessons/04_airflow_orchestration.md index 704e5da..40b4f5c 100644 --- a/docs/course/lessons/04_airflow_orchestration.md +++ b/docs/course/lessons/04_airflow_orchestration.md @@ -81,8 +81,9 @@ Airflow. Главная единица Airflow — **DAG** (Directed Acyclic Gra заново наполнит их из ODS. Для учебного стенда это удобный чистый прогон: результат повторяемый, старые эксперименты не мешают. -Когда DAG завершится, открой его граф. На чистой стартовой истории все задачи должны быть -зелёными. Найди внутри группы `transform` две задачи подряд: +Когда DAG завершится, открой его граф. На чистой стартовой истории выбранные задачи +должны быть зелёными, а невыбранная ветка `skip_truncate` — в состоянии `skipped`. +Найди внутри группы `transform` две задачи подряд: - `check_dds_integrity` — SQL-задача, которая считает сирот; - `assert_dds_integrity` — Python-задача, которая решает, можно ли идти дальше. @@ -99,9 +100,9 @@ WHERE layer = 'dds' AND check_name = 'orphan_events'; ``` -Ожидаем `check_value = 0`. Это тот же смысл, что в уроке 3, только теперь число появилось внутри -управляемого прогона Airflow. - +На проверенном стенде `check_date = 2026-07-23`, `check_value = 0`. Дата берётся +из `today()`, поэтому у тебя будет своя. Это тот же смысл, что в уроке 3, только +теперь число появилось внутри управляемого прогона Airflow. --- @@ -314,7 +315,6 @@ WHERE click_id IS NOT NULL Снова должно быть `0`. Если стенд после экспериментов совсем запутался, пройди [канонический сброс](../README.md#подготовка-и-канонический-сброс). - --- @@ -322,12 +322,11 @@ WHERE click_id IS NOT NULL | Действие | Где смотреть | Что ожидать | |----------|--------------|-------------| -| `etl_pipeline` с `{"full_refresh": true}` | Airflow graph | все задачи зелёные | +| `etl_pipeline` с `{"full_refresh": true}` | Airflow graph | выбранная ветка зелёная, `skip_truncate` в `skipped` | | чистый прогон | `dm.dq_summary`, строка `orphan_events` | `0` | | вставка события-сироты | прямой SQL-счётчик сирот | `0 → 1` | | `etl_pipeline` с `{"full_refresh": false}` после вставки | task `transform.assert_dds_integrity` | task красная, DAG failed | | откат через `{"full_refresh": true}` | прямой SQL-счётчик сирот | снова `0` | - --- diff --git a/docs/course/lessons/07_lab_next_day.md b/docs/course/lessons/07_lab_next_day.md index 4022e34..4489dbe 100644 --- a/docs/course/lessons/07_lab_next_day.md +++ b/docs/course/lessons/07_lab_next_day.md @@ -82,19 +82,18 @@ docker compose run --rm --no-deps generator \ events visits users min_event_timestamp max_event_timestamp model_t0 model_t_end profile ``` -Сверь результат с таблицей. Значения будут вписаны после контрольного прогона -эталонного мира: +Сверь результат с контрольными значениями: | Поле | Ожидаем после дня 4 | |------|----------------------| -| `events` | `<события-день-4>` | -| `visits` | `<визиты-день-4>` | -| `users` | `<пользователи-день-4>` | -| `min_event_timestamp` | `<минимальное-время>` | -| `max_event_timestamp` | `<максимальное-время>` | -| `model_t0` | `<левая-граница>` | -| `model_t_end` | `<правая-граница-дня-4>` | -| `profile` | `<профиль>` | +| `events` | `374092` | +| `visits` | `34801` | +| `users` | `5388` | +| `min_event_timestamp` | `2026-01-01 00:00:00.000000` | +| `max_event_timestamp` | `2026-01-04 23:59:58.581616` | +| `model_t0` | `2026-01-01T00:00:00+00:00` | +| `model_t_end` | `2026-01-05T00:00:00+00:00` | +| `profile` | `daily-wave` | Здесь числа накопительные: `events`, `visits` и `users` относятся ко всему миру от `model_t0` до `model_t_end`, а не только к четвёртому дню. @@ -128,7 +127,9 @@ ORDER BY visit_start; ``` Запрос должен найти хотя бы один визит, чьи события лежат по разные стороны -полуночи. Подойдёт любой стык внутри мира: 1→2, 2→3 или 3→4. +полуночи. На контрольном мире он находит `35`, `22` и `26` визитов на стыках +1→2, 2→3 и 3→4. Для проверки нового дня достаточно увидеть `26` визитов на +стыке 3→4. > **Как слова связаны со схемой.** В manifest написано «визит», а в SQL такой > визит обозначен `click_id`. Отдельной таблицы визитов нет: контекст лежит в @@ -233,8 +234,8 @@ Airflow проигрывать все пропущенные интервалы 2. `make generated-history-chain-check` заканчивается без ошибки; 3. в `Events over Time` появился день 5. -Точные числа после дня 5 не фиксируем: это результат учебной правки. - +После дня 5 контрольный manifest показывает `468025` событий, `43482` визита, +`6713` пользователей и `model_t_end = 2026-01-06T00:00:00+00:00`. ### Верни как было @@ -249,9 +250,9 @@ Airflow проигрывать все пропущенные интервалы | Действие | Где смотреть | Что ожидать | |----------|--------------|-------------| | запустить `world_next_day` с пустой формой | Airflow graph | все пять задач зелёные, `etl_pipeline` вызван с `full_refresh=true` | -| прочитать manifest после дня 4 | `kafka-manifest-summary` | числа совпадают с таблицей секции 2 | +| прочитать manifest после дня 4 | `kafka-manifest-summary` | числа совпадают с таблицей секции 2 | | проверить накопленную цепочку | `make generated-history-chain-check` | команда завершается без ошибки | -| найти переходящие визиты | запрос в `dds.event` | хотя бы один `click_id` пересекает полночь | +| найти переходящие визиты | запрос в `dds.event` | на стыке 3→4 найдено `26` визитов | | снять паузу с DAG | Airflow runs | на ближайшей получасовой границе приезжает один запланированный день | | вернуть паузу и пройти сброс | Airflow и Superset | расписание выключено, снова видны три эталонных дня | diff --git a/docs/course/lessons/08_lab_continue.md b/docs/course/lessons/08_lab_continue.md index 2e2aa96..68a9553 100644 --- a/docs/course/lessons/08_lab_continue.md +++ b/docs/course/lessons/08_lab_continue.md @@ -55,7 +55,7 @@ Materialized View сразу складывает сообщения в STG. Д SELECT 'STG' AS layer, max(parseDateTime64BestEffortOrNull( - JSONExtractString(raw, 'event_timestamp'), 6, 'UTC' + JSONExtractString(raw, 'event_timestamp'), 6 )) AS max_event_ts FROM stg.browser_raw @@ -76,7 +76,8 @@ FROM dm.v_events_enriched; ``` На эталонном мире правые границы слоёв должны совпадать. - +После импорта трёх эталонных дней все четыре слоя заканчиваются на +`2026-01-03 23:59:59.066766`. ### Запускаем живое продолжение @@ -217,7 +218,9 @@ make generator-continue Запусти `etl_pipeline` с пустой формой. После зелёного прогона выполни: ```sql -WITH parseDateTime64BestEffort('', 6, 'UTC') AS boundary +WITH parseDateTime64BestEffort( + substring('', 1, 19), 6 +) AS boundary SELECT user_domain_id, countIf(visit_start < boundary) AS visits_before, @@ -233,7 +236,6 @@ GROUP BY user_domain_id HAVING min(visit_start) < boundary AND max(visit_start) >= boundary LIMIT 10; ``` - Подставь вместо `` границу из манифеста. Успех — хотя бы одна строка. Иначе подожди несколько минут, перезапусти `etl_pipeline` и повтори @@ -241,6 +243,10 @@ LIMIT 10; пользователь с разными визитами пересекает границу замороженного мира. Мир расходится: итоги превышают эталонные и различаются у менти. Это правильно. +Например, для границы `2026-01-06T00:00:00+00:00` найден пользователь +`5938d296-14cf-49cb-801e-82370518cf59`: `13` визитов до границы и `2` после +неё. Твои числа и идентификатор будут другими. + ### Верни как было Останови живой поток: @@ -262,7 +268,7 @@ make generator-down | Действие | Где смотреть | Что ожидать | |----------|--------------|-------------| | выполнить `make generator-continue` | Prometheus Targets | `generator` переходит в `UP` | -| измерить свежесть до и после ETL | запрос из секции 2 и часы | записаны модельное отставание и настенное ожидание | +| измерить свежесть до и после ETL | запрос из секции 2 и часы | например: до ETL — `3:45:01.519535`, после — `0:45:55.977582` модельного времени; настенное ожидание — около `29` секунд | | выполнить `make generator-down` | Grafana и Kafka | события перестают поступать, отставание читателей стекает к нулю | | снова выполнить `make generator-continue` | Grafana, Kafka и STG | поток продолжается, offset-ы и правая граница снова растут | | найти вернувшегося пользователя | запрос по границе `model_t_end` | есть пользователь с визитами до и после границы |