diff --git a/.scratch/handoffs/20260723-1248-course-labs-redesign.md b/.scratch/handoffs/20260723-1248-course-labs-redesign.md new file mode 100644 index 0000000..b4382ad --- /dev/null +++ b/.scratch/handoffs/20260723-1248-course-labs-redesign.md @@ -0,0 +1,15 @@ +# Handoff: редизайн лаб курса (issue #7) — ЗАВЕРШЁН + +Обновлено: 2026-07-23, конец прогона. Работа выполнена целиком. + +- PR: https://github.com/dementev-dev/clickstream-ch-kafka-superset-demo/pull/25 + (closes #7; ветка `feature/course-labs-redesign`, 5 коммитов от спеки до + сверки на стенде). Дочерние тикеты #21/#22/#23 закрыты, чекбоксы в #7 + проставлены. Чистка легаси-тестов — отдельный issue #24. +- Осталось человеку: ревью и merge PR; после merge ветку можно удалить. +- Стенд после сверки НЕ канонический: мир дорос до дня 5 и жил в continue; + генератор остановлен, `world_next_day` на паузе. Перед занятиями менти — + канонический сброс по README курса. +- Урок прогона: пауза DAG во время работающего прогона замораживает его + задачи; `unpause` сразу после `make up` может быть перекрыт + `is_paused_upon_creation` (гонка с первым парсингом). diff --git a/CONTEXT.md b/CONTEXT.md index 1c25e86..2b436fb 100644 --- a/CONTEXT.md +++ b/CONTEXT.md @@ -63,8 +63,10 @@ user_domain_id (пользователь, постоянный) ### Возвращающийся пользователь (returning user) Пользователь, открывающий **более одного** визита (`click_id`) во времени, с -межсессионными паузами. Именно возвраты дают расхождение `users < sessions` — -то, чего нет на сиде (`users == sessions`) и что отличает поток от статики. +межсессионными паузами. Эталонный мир уже содержит возвраты: это замороженный +LIVE-мир с 4 056 пользователями и 26 083 визитами, поэтому уже при импорте +выполняется `users < sessions`. Режим `continue` продолжает этот мир: итоги +растут дальше и зависят от длительности запуска у каждого менти. ### Три значения слова «сид» @@ -116,7 +118,7 @@ user_domain_id (пользователь, постоянный) **Модельное время стенда** отвязано от настенных часов: генератор крутит внутренние часы, а драйвер задаёт скорость — ×1 (как реальное время), ×K -(ускоренно, учебная «ручка» урока 7) или `K → ∞` (мгновенная заливка прошлого). +(ускоренно, [учебная ручка в лабе 08](./docs/course/lessons/08_lab_continue.md#modelnoe-vremya)) или `K → ∞` (мгновенная заливка прошлого). При ускорении «сейчас» стенда уходит вперёд настенного времени — это свойство, не баг (ключ аналитики — `event_timestamp`). Решение и режимы — [ADR-0005](./docs/adr/0005-generator-model-clock.md). diff --git a/docs/course/LEARNING_PLAN.md b/docs/course/LEARNING_PLAN.md index 2d582aa..fbaba3c 100644 --- a/docs/course/LEARNING_PLAN.md +++ b/docs/course/LEARNING_PLAN.md @@ -1,7 +1,9 @@ # План обучения: курс «Кликстрим на ClickHouse» > Дата: 2026-06-03 (аудит путей выполнен; середина расщеплена — -> один паттерн на урок, всего 7 уроков, см. §1–2). +> один паттерн на урок, см. §1–2). +> Поправка 2026-07-23: весь маршрут стал обязательным. После урока 6 добавлены +> лабы 07–08 про следующий день и живое продолжение; обе лабы написаны. > Назначение: высокоуровневый маршрут менти по курсу — карта уроков, порядок, > результаты аудита эталонных путей. Рамка курса (зачем/что/скоуп) — в `PRD.md`; > как устроен отдельный урок — в `LESSON_STANDARD.md`. @@ -13,9 +15,11 @@ ## 1. Маршрут -Порядок линейный — уроки 1→4 повторяют сам пайплайн (STG → ODS → DDS → оркестрация), -а урок 0 — разминка перед ним. Мониторинг (5) и Superset (6) более самостоятельны и -работают как надстройки поверх готовых данных. +Порядок линейный, все уроки и лабы обязательны. Урок 0 — разминка, а уроки 1→4 +повторяют сам пайплайн (STG → ODS → DDS → оркестрация). Мониторинг (5) и +Superset (6) работают как надстройки поверх готовых данных. Лабы 07–08 затем +растят импортированный мир: сначала детерминированным следующим днём, потом +живым продолжением. Принцип нарезки — **один прод-паттерн на урок** (`LESSON_STANDARD` §2). Поэтому середина пайплайна разнесена: типизация+DQ (ODS) и сборка сущностей (DDS) — разные @@ -36,11 +40,13 @@ Kafka из роадмапа — оно даёт словарь терминов. | 3 | ODS → DDS: сборка сущностей | `sql/dds/30_ods_to_dds.sql` + DDL `sql/ddl/dds/30_dds.sql` (argMax, UNION, сироты) | обязательный | руки | | 4 | Оркестрация в Airflow | `airflow/dags/etl_pipeline_dag.py` (зависимости, гейты, остановка при нарушениях) | обязательный | руки | | 5 | Мониторинг | Prometheus + Grafana + экспортёры | обязательный | наблюдение | -| 6 | BI-витрина | Superset поверх ClickHouse | опциональный | руки | +| 6 | BI-витрина | Superset поверх ClickHouse | обязательный | руки | +| 7 | Лаба: следующий день | `airflow/dags/world_next_day_dag.py` | написана, обязательная | руки | +| 8 | Лаба: живое продолжение | `make generator-continue` и `etl_pipeline` | написана, обязательная | руки | -Урок 0 — обязательная разминка (без правок кода, только наблюдение); уроки 1–4 — -с управляемыми правками; урок 5 ближе к наблюдению (глубину уточняем, см. -«Открытые вопросы» в `PRD.md`). +Урок 0 — обязательная разминка без правок кода. В уроках 1–6 и лабах 07–08 +менти запускает стенд или делает управляемую правку; каждый такой эксперимент +заканчивается явным возвратом к чистому состоянию. **Где «MV vs батч»:** контраст из цели №2 PRD — это мост уроков 1→2. Урок 1 (STG) заканчивается вопросом «мы приземлили поток через MV — почему дальше не MV?»; урок 2 @@ -64,6 +70,7 @@ Kafka из роадмапа — оно даёт словарь терминов. | 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` | **точечно править** | | 5 | Мониторинг | `configs/*`, экспортёры | **годно как есть** (для режима наблюдения) | +| 7 | Следующий день | `airflow/dags/world_next_day_dag.py` | **точечно править** | > Карты путей раздела 2 включают и DDL целевых таблиц (`sql/ddl/{ods,dds}/*`), а не > только батч-трансформации: без формы целевых таблиц слой читается неполно. @@ -140,8 +147,9 @@ Kafka из роадмапа — оно даёт словарь терминов. ## 4. Что дальше -Все уроки 0–6 написаны (`lessons/`), правки кода §3.1 применены в составе уроков — -чекбоксы выше сверены с реальным кодом 2026-06-06. Базовый контент курса собран. +Все уроки 0–6 и лабы 07–08 написаны (`lessons/`). Правки кода §3.1 применены +в составе уроков; чекбоксы выше сверены с реальным кодом 2026-06-06. Маршрут +курса собран полностью. Исходный план фазы написания (оставлен как контекст): пишем уроки по одному, по шаблону из `LESSON_STANDARD.md`, начиная с урока 1 (Kafka→CH); урок 0 (разминка на diff --git a/docs/course/LESSON_STANDARD.md b/docs/course/LESSON_STANDARD.md index 1f46624..c08c898 100644 --- a/docs/course/LESSON_STANDARD.md +++ b/docs/course/LESSON_STANDARD.md @@ -21,7 +21,8 @@ запись и увидеть `+1` в таблице ошибок `parse_errors`; сменить `kafka_group_name` и увидеть, как топик читается заново. Менти меняет — видит эффект — объясняет. **Верни как было** — каждая правка завершается явным шагом отката к чистому - состоянию (откатить изменение либо `make generated-history-analytics && make up`), чтобы + состоянию (откатить изменение либо пройти + [канонический сброс](./README.md#подготовка-и-канонический-сброс)), чтобы самостоятельный менти не застрял со сломанным стендом без ментора. (Урок 0 — без этого шага, только наблюдение.) 5. **Проверь себя** — самопроверка (раздел 4). @@ -112,7 +113,8 @@ - штатные быстрые проверки (smoke) из `docs/TEST_PLAN.md`; - встроенные проверки в `etl_pipeline_dag.py` (DAG падает на пустой витрине или нарушении целостности); -- `make generated-history-analytics && make up` для чистого сброса и повтора штатного пути. +- [канонический сброс](./README.md#подготовка-и-канонический-сброс) для чистого + возврата и повтора штатного пути. В каждом уроке — маленькая табличка самопроверки в формате **действие → где смотреть → что ожидать**, а для управляемой правки — какой diff --git a/docs/course/PRD.md b/docs/course/PRD.md index 3ce0043..b020218 100644 --- a/docs/course/PRD.md +++ b/docs/course/PRD.md @@ -3,7 +3,7 @@ > Статус: прообраз PRD. Дата: 2026-06-01. > Поправка 2026-06-03 (разморозка по делу): середина пайплайна расщеплена — ODS и DDS > теперь разные уроки (принцип «один паттерн на урок»), витрины DM демотированы в -> поверхность потребления. Обязательных уроков стало 0–5, опциональный Superset — урок 6. +> поверхность потребления. Superset вынесен в урок 6. > Затронуты §4 (скоуп) и §3/§5 (ожидаемый такт — ~день на урок). Это изменение рамки, > а не план реализации. > Поправка 2026-06-03 (терминология): режим сопровождения — еженедельный **созвон**, а не @@ -16,6 +16,10 @@ > стенда больше не является. Схема потока в §1 обновлена под генератор (источник данных > сменился, см. ADR-0006). В §7 закрыта развилка про урок о генераторе и добавлена > развилка про кластерную конфигурацию. +> Поправка 2026-07-23 (маршрут курса): зафиксированы три режима работы с миром — +> импорт эталонной базы, пакетная дозаливка следующего дня и живое продолжение. +> Урок 6 стал обязательным, после него добавлены обязательные лабы 07–08. +> Затронуты §1–4 и §7. > Назначение документа: зафиксировать для будущих сессий, что это за курс, зачем > он, что входит в скоуп работ, а что нет. Это договорная **рамка**, а не план > реализации и не стандарт уроков (см. раздел «Связанные документы»). @@ -30,7 +34,7 @@ Репозиторий — рабочий сквозной стенд кликстрим-DWH: ``` -generator (backfill/live) → Kafka → ClickHouse (Kafka engine + MV → STG) → Airflow ETL (STG→ODS→DDS→DM) → Superset +generator (импорт / живой поток) → Kafka → ClickHouse (Kafka engine + MV → STG) → Airflow ETL (STG→ODS→DDS→DM) → Superset ↘ Prometheus / Grafana (мониторинг) ``` @@ -85,7 +89,9 @@ generator (backfill/live) → Kafka → ClickHouse (Kafka engine + MV → STG) 3. Читать и объяснять **оркестрацию в Airflow**: DAG, зависимости задач, проверки качества данных, остановку пайплайна при нарушениях. 4. Понимать, **как устроен мониторинг** пайплайна (метрики, экспортёры, дашборды). -5. (Опционально) Подключать **BI-витрину** поверх ClickHouse (Superset). +5. Подключать **BI-витрину** поверх ClickHouse (Superset). +6. Различать пакетную дозаливку дня и живой поток, понимать границы времени + и свежесть данных. Сквозная цель — не «посмотреть, как работает», а **уметь пересказать паттерн своими словами и привязать его к обычной кликстрим-аналитике** (трекер событий → @@ -96,8 +102,8 @@ Kafka → ClickHouse → BI). - **Аудитория:** продвинутые менти, прошедшие базовую программу. Пишем обобщённо, но затачиваем под реальный первый прогон, а не под гипотетических будущих менти. - **Режим:** самостоятельный, асинхронный. Менти клонирует репозиторий, готовит - стенд штатным путём (`make generated-history-analytics && make up`) и идёт по - урокам из `docs/course/` рядом с кодом. + стенд по [канонической инструкции](./README.md#подготовка-и-канонический-сброс) + и идёт по урокам из `docs/course/` рядом с кодом. Уроки короткие и односоставные — ожидаемый срок прохождения одного **около дня**. - **Роль ментора:** еженедельный созвон-сверка (покрывает несколько уроков), без построчного разбора кода. @@ -108,10 +114,12 @@ Kafka → ClickHouse → BI). ## 4. Скоуп ### Входит -- **Уроки 0–5 (обязательные):** вводный урок по Kafka, заземление Kafka→CH, STG→ODS +- **Уроки 0–6 (обязательные):** вводный урок по Kafka, заземление Kafka→CH, STG→ODS (типизация + DQ), ODS→DDS (сборка сущностей), Airflow, мониторинг. Принцип нарезки — один прод-паттерн на урок; контраст «где Materialized View, а где батч» проходит - мостом уроков 1→2. + мостом уроков 1→2. Урок 6 закрывает BI-слой в Superset. +- **Лабы 07–08 (обязательные):** пакетная дозаливка следующего дня, затем живое + продолжение потока. - Витрины **DM — не отдельный урок**: их показываем в деле там, где их потребляют (мониторинг и BI). См. `LEARNING_PLAN.md` §1–2. - **Аудит и точечная полировка эталонных путей** этих уроков до учебного качества @@ -121,10 +129,6 @@ Kafka → ClickHouse → BI). Подробная карта уроков (файлы стенда, статус, режим, вердикты аудита) — в плане обучения `LEARNING_PLAN.md`. -### Опционально -- **Урок 6: Superset (BI-витрина).** Делаем, если останется ресурс; обязательные - уроки он не блокирует. - ### Не входит - Переписывание всего стенда: полируем только эталонные пути обязательных уроков, остальной код стенда остаётся под капотом. @@ -171,8 +175,7 @@ Kafka → ClickHouse → BI). (зеркало урока 4) — урок 5 даёт «сломал-увидел», а не чистое наблюдение. См. `LEARNING_PLAN.md` §3.1. - ~~Нужна ли BI-витрина (Superset) уже в первой версии.~~ **Решено:** урок 6 написан и - синхронизирован с реальным дашбордом; остаётся опциональным (обязательные уроки не - блокирует). + синхронизирован с реальным дашбордом. С 2026-07-23 он входит в обязательный маршрут. Закрыто позже: diff --git a/docs/course/README.md b/docs/course/README.md index 5b0608e..55c14a1 100644 --- a/docs/course/README.md +++ b/docs/course/README.md @@ -27,30 +27,36 @@ поднимаются одновременно. Нужна машина, которая это потянет. - **Инструменты:** `Docker` с `docker compose`, `make`, `bash`, `curl`, `git` и `uv`. `uv` нужен для локальных Python-проверок и команд разработки. -- **Подними стенд и создай стартовую историю** (из корня репозитория) — этого хватит, - чтобы начать, и прогон быстрый: - ```bash - make generated-history-analytics - make up - ``` +## Подготовка и канонический сброс - Эта команда проводит штатный путь стенда: готовый источник данных создаёт стартовую - историю, события попадают в Kafka, затем в STG, ODS, DDS, DM и Superset. Это тот же - путь, что ручной вариант из README (Airflow UI и операция `backfill`), но одной - командой — выбери один из двух, оба ведут к одинаковому стенду. Файлы - `data/*.jsonl` пока остаются только кладовкой готовых значений для генератора - (браузеры, страны, устройства, UTM), а не источником аналитического контура. `make up` - после неё поднимает остальные UI-сервисы курса: Kafka UI, Airflow, Prometheus и Grafana. +Первый запуск и возврат к чистому эталонному миру идут одним путём. При первом запуске +пропусти `make clean`; для полного сброса выполни все три шага: - Дальше каждый урок в секции «Руки» сам напоминает, что перезапустить. - Точные шаги, параметры и troubleshooting — в [`docs/OPERATIONS.md`](../OPERATIONS.md). +1. Выполни `make clean`. Команда удалит данные стенда и метаданные Superset: сохранённые + в нём настройки и дашборды тоже придётся создать заново. +2. Выполни `make up` и открой Airflow на `http://localhost:8080` (`admin/admin`). + На свежем стенде DAG-и стоят на паузе. Подготовь и запусти их по порядку: + - `ddl_init` — сними паузу и запусти с пустой формой; + - `etl_pipeline` — только сними паузу: его вызовет следующий DAG; + - `world_init` — сними паузу и запусти с пустой формой после успешного + `ddl_init`. +3. Когда `world_init` завершится успешно и витрины DM будут готовы, выполни + `make superset-init`. + +Так события из эталонного мира попадут в Kafka, затем в STG, ODS, DDS и DM, а Superset +получит готовые наборы данных и дашборд. Файлы `data/*.jsonl` остаются только кладовкой +готовых значений для генератора (браузеры, страны, устройства, UTM), а не источником +аналитического контура. + +Точные параметры и разбор ошибок — в [`docs/OPERATIONS.md`](../OPERATIONS.md). ## Уроки Проходи по порядку — уроки 0→4 повторяют сам пайплайн (Kafka → STG → ODS → DDS → -оркестрация), а 5 и 6 надстраиваются поверх готовых данных. В колонке «режим»: -**наблюдение** — только смотрим, **руки** — запускаешь и меняешь сам. +оркестрация), а 5 и 6 надстраиваются поверх готовых данных. Лабы 7 и 8 затем +растят эталонный мир двумя способами. В колонке «режим»: **наблюдение** — только +смотрим, **руки** — запускаешь и меняешь сам. | # | Урок | Режим | О чём | |---|------|-------|-------| @@ -60,7 +66,9 @@ | 3 | [ODS → DDS: сборка сущностей](./lessons/03_ods_to_dds.md) | руки | Собираем `click` и `event` из кусочков (argMax, JOIN) и встречаем «сирот» | | 4 | [Оркестрация в Airflow](./lessons/04_airflow_orchestration.md) | руки | Всю цепочку — в один DAG с зависимостями и честным гейтом целостности | | 5 | [Мониторинг: Prometheus и Grafana](./lessons/05_monitoring.md) | наблюдение + мини-правка | Смотрим систему со стороны; гасим сервис — видим, как краснеет алерт | -| 6 | [BI-витрина в Superset](./lessons/06_superset_bi.md) | руки (опционально) | Дашборд поверх ClickHouse: KPI, динамика, воронка | +| 6 | [BI-витрина в Superset](./lessons/06_superset_bi.md) | руки | Дашборд поверх ClickHouse: KPI, динамика, воронка | +| 7 | [Лаба: следующий день](./lessons/07_lab_next_day.md) | руки | Пакетный инкремент дня и границы времени | +| 8 | [Лаба: живое продолжение](./lessons/08_lab_continue.md) | руки | Живой поток и свежесть данных | ## Как проходить @@ -76,22 +84,18 @@ ## Проверка чистого маршрута -Перед проверкой уроков 0, 1 и 5 подними стенд с нуля: - -```bash -make generated-history-analytics -make up -``` +Перед проверкой уроков 0, 1 и 5 пройди +[канонический сброс](#подготовка-и-канонический-сброс). Что должен подтвердить человек: - урок 0: в Kafka UI видны четыре топика событий и понятны служебные топики генератора; -- урок 1: после `CLEAN_START=0 make generated-history-analytics` учебная колонка - `kafka_msg_ts` не пропадает и заполняется; +- урок 1: после запуска `world_next_day` учебная колонка `kafka_msg_ts` не пропадает + и заполняется; - урок 5: Prometheus targets `clickhouse`, `kafka`, `airflow` находятся в `UP`, - а `Kafka No Messages Produced` трактуется с учётом того, запущен live-генератор - или только стартовая история. + а `Kafka No Messages Produced` трактуется с учётом того, включён живой поток + или стенд работает на импортированной базе. ## Границы diff --git a/docs/course/lessons/00_kafka_intro.md b/docs/course/lessons/00_kafka_intro.md index 53544f6..67ebc4d 100644 --- a/docs/course/lessons/00_kafka_intro.md +++ b/docs/course/lessons/00_kafka_intro.md @@ -41,12 +41,9 @@ consumer читает в своём темпе, не трогая producer'а. ## 2. Наблюдай: открой Kafka UI -Стенд уже должен быть поднят, а в топиках должны лежать события после стартовой истории: - -```bash -make generated-history-analytics -make up -``` +Подготовь стенд по +[канонической инструкции курса](../README.md#подготовка-и-канонический-сброс). +После успешного `world_init` в топиках уже лежат события эталонного мира. Открой Kafka UI: `http://localhost:8082`. Ходи по нему свободно — это режим чтения, сломать тут ничего нельзя. @@ -79,15 +76,20 @@ make up - **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/01_kafka_to_clickhouse.md b/docs/course/lessons/01_kafka_to_clickhouse.md index e345d91..12bc4d0 100644 --- a/docs/course/lessons/01_kafka_to_clickhouse.md +++ b/docs/course/lessons/01_kafka_to_clickhouse.md @@ -49,15 +49,9 @@ ## 2. Руки: убедись, что данные текут -Поднимаем стенд и создаём стартовую историю. Это штатный путь курса: готовый источник -данных стенда пишет события в Kafka, ClickHouse читает их в STG, затем batch строит -ODS, DDS и DM. Файлы `data/*.jsonl` в этом пути не источник аналитики; пока это только -кладовка значений для источника данных стенда. - -```bash -make generated-history-analytics -make up -``` +Подготовь стенд по +[канонической инструкции курса](../README.md#подготовка-и-канонический-сброс). +После успешного `world_init` эталонный мир уже прошёл путь Kafka → STG → ODS → DDS → DM. Теперь смотрим, что доехало до ClickHouse. Открой SQL-консоль: `http://localhost:9123/play` (или Kafka UI на `http://localhost:8082`, чтобы тем же @@ -189,27 +183,23 @@ SELECT FROM stg.kafka_browser_raw; ``` -**Перезаливаем срез, чтобы новая колонка заполнилась.** Сначала чистим только таблицу -приёмника — иначе рядом останутся строки из секции 2, вставленные ещё *до* -`ADD COLUMN`, и в них `kafka_msg_ts` будет пустой (`1970-01-01`): +**Дозаливаем следующий день, чтобы новая колонка заполнилась.** Не очищай +`stg.browser_raw`: `world_next_day` растит уже импортированный мир, а его финальная +проверка сверяет весь накопленный результат. В старых строках, вставленных до +`ADD COLUMN`, `kafka_msg_ts` останется пустым (`1970-01-01`) — ниже мы их отфильтруем. -```sql -TRUNCATE TABLE stg.browser_raw; -``` - -Затем повторно создаём стартовую историю **без очистки volumes**. Это важно: -`CLEAN_START=0` сохраняет твою новую колонку и пересозданное MV, но добавляет свежие -сообщения в Kafka, чтобы ClickHouse прочитал их уже с новой схемой. - -```bash -CLEAN_START=0 make generated-history-analytics -``` +Открой Airflow и кнопкой **Trigger DAG** запусти `world_next_day` с пустой формой. +Он дозальёт в Kafka следующий день, а ClickHouse прочитает новые сообщения уже с +изменённой схемой. Такой прогон занимает несколько минут — дождись состояния +`success`. Здесь мы только нажимаем готовую кнопку; подробно `world_next_day` +разберём в лабе 07. **Смотрим результат:** ```sql SELECT kafka_ts, kafka_msg_ts FROM stg.browser_raw +WHERE kafka_msg_ts > toDateTime(0) ORDER BY kafka_offset LIMIT 5; ``` @@ -224,19 +214,8 @@ LIMIT 5; > теряется. Обратный случай — колонку добавил, а в 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 generated-history-analytics && make up`. +**Верни как было.** После дозаливки мир уже вырос, поэтому верни схему и данные +[каноническим сбросом](../README.md#подготовка-и-канонический-сброс). --- @@ -244,10 +223,10 @@ make ddl # пересоздаёт эталонное MV из 10_stg.sql, с | Действие | Где смотреть | Что ожидать | |----------|--------------|-------------| -| `make generated-history-analytics && make up` | `SELECT count() FROM stg.browser_raw` | счётчик > 0 | +| каноническая подготовка | `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` | колонка заполнена временем сообщения | +| `world_next_day` после правки из секции 4 | `SELECT kafka_msg_ts FROM stg.browser_raw WHERE kafka_msg_ts > toDateTime(0) LIMIT 5` | новые строки заполнены временем сообщения | --- diff --git a/docs/course/lessons/02_stg_to_ods.md b/docs/course/lessons/02_stg_to_ods.md index 6c75c39..6cbe71a 100644 --- a/docs/course/lessons/02_stg_to_ods.md +++ b/docs/course/lessons/02_stg_to_ods.md @@ -84,32 +84,25 @@ ODS мы пересобираем целиком, одной задачей Airf ## 2. Руки: смотрим базовый прогон -Поднимаем стенд и создаём стартовую историю. Это штатный путь курса: готовый источник -данных стенда пишет события в Kafka, ClickHouse читает их в STG, затем batch строит -ODS, DDS и DM. Файлы `data/*.jsonl` пока остаются только кладовкой значений для этого -источника, а не источником аналитического контура. +Подготовь стенд по +[канонической инструкции курса](../README.md#подготовка-и-канонический-сброс). +После успешного `world_init` эталонный мир уже прошёл путь Kafka → STG → ODS → DDS → DM. -```bash -make generated-history-analytics -make up -``` - -Команда прогоняет всю цепочку слоёв и прямо в консоли печатает то, что нам нужно сейчас, — -блок **«Статистика ODS»**. Это просто счётчики строк по всем восьми таблицам слоя (четыре -основных и четыре с ошибками). Пример формы вывода: +Ниже — форма блока **«Статистика ODS»**, который печатает `make transform`. Это просто +счётчики строк по всем восьми таблицам слоя (четыре основных и четыре с ошибками): ``` Статистика 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 │ + └────────────────────────────┴────────┘ ``` Прочитаем эту табличку — в ней три вещи, которые стоит заметить. @@ -139,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`, а не пропавшие +данные. Ничего не потерялось молча. --- @@ -250,9 +244,9 @@ ORDER BY (click_id) Урок про типы — так давай **намеренно ошибёмся типом** и посмотрим, что будет. Это самый поучительный момент урока. -Возьмём координату `geo_latitude` — широту. Это дробное число, например `50.82709`. Достаём мы +Возьмём координату `geo_latitude` — широту. Это дробное число, например `-7.60361`. Достаём мы её через `toFloat64OrNull` — «привести к дробному числу». Заменим тип на целочисленный — -`toInt64OrNull`, «привести к целому». Для строки `"50.82709"` целого числа не получится +`toInt64OrNull`, «привести к целому». Для строки `"-7.60361"` целого числа не получится (там точка, дробная часть), и функция вернёт `NULL`. То есть широта просто исчезнет. Из секции 3 помним: разбор продублирован, поэтому правок будет **две** — в обоих `INSERT` @@ -275,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 @@ -289,8 +286,8 @@ 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'] │ └──────────────┴──────────────┴───────────────┴──────────────────────┘ ``` @@ -316,8 +313,9 @@ git checkout -- sql/ods/20_stg_to_ods.sql make transform ``` -После этого `geo_by_click_errors` снова `0`, широта на месте. А если стенд совсем «поплыл» — -всегда есть полный чистый прогон: `make generated-history-analytics && make up`. +После этого `geo_by_click_errors` снова `0`, широта на месте. А если стенд совсем +«поплыл», пройди +[канонический сброс](../README.md#подготовка-и-канонический-сброс). --- diff --git a/docs/course/lessons/03_ods_to_dds.md b/docs/course/lessons/03_ods_to_dds.md index 495c10e..c48aeae 100644 --- a/docs/course/lessons/03_ods_to_dds.md +++ b/docs/course/lessons/03_ods_to_dds.md @@ -65,25 +65,19 @@ ## 2. Руки: смотрим базовый прогон -Поднимаем стенд и создаём стартовую историю. Это штатный путь курса: готовый источник -данных стенда пишет события в Kafka, ClickHouse читает их в STG, затем batch строит -ODS, DDS и DM. Файлы `data/*.jsonl` пока остаются только кладовкой значений для этого -источника, а не источником аналитического контура. +Подготовь стенд по +[канонической инструкции курса](../README.md#подготовка-и-канонический-сброс). +После успешного `world_init` эталонный мир уже прошёл путь Kafka → STG → ODS → DDS → DM. -```bash -make generated-history-analytics -make up -``` - -Команда прогоняет всю цепочку слоёв и по дороге печатает в консоль блок **«Статистика DDS»** — -счётчики строк по двум нашим сущностям: +Ниже — форма блока **«Статистика DDS»**, который печатает `make transform`: счётчики +строк по двум нашим сущностям. ``` Статистика DDS: - ┌─table─────┬─rows─┐ - │ dds.click │ ... │ - │ dds.event │ ... │ - └───────────┴──────┘ + ┌─table─────┬───rows─┐ + │ dds.click │ 26083 │ + │ dds.event │ 280437 │ + └───────────┴────────┘ ``` Прочитаем эти две строки. @@ -101,7 +95,7 @@ make up ``` ┌─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 │ └────────────┴───────┴─────────────────────┴───────────────┴─────────────┘ ``` @@ -308,8 +302,9 @@ DDS — `make transform` чистит `dds.event` (`TRUNCATE`) и наполня make transform ``` -После этого `orphan_events` снова `0`, придуманное событие исчезло. А если стенд совсем «поплыл» — -полный чистый прогон: `make generated-history-analytics && make up`. +После этого `orphan_events` снова `0`, придуманное событие исчезло. А если стенд +совсем «поплыл», пройди +[канонический сброс](../README.md#подготовка-и-канонический-сброс). --- diff --git a/docs/course/lessons/04_airflow_orchestration.md b/docs/course/lessons/04_airflow_orchestration.md index a62d3e8..40b4f5c 100644 --- a/docs/course/lessons/04_airflow_orchestration.md +++ b/docs/course/lessons/04_airflow_orchestration.md @@ -56,15 +56,19 @@ Airflow. Главная единица Airflow — **DAG** (Directed Acyclic Gra ## 2. Руки: запускаем DAG и смотрим зелёный прогон -Подними стенд и создай стартовую историю. Это штатный путь курса: готовый источник -данных стенда пишет события в Kafka, ClickHouse читает их в STG, затем batch строит -ODS, DDS и DM. Файлы `data/*.jsonl` пока остаются только кладовкой значений для этого -источника, а не источником аналитического контура. +Подготовь стенд по +[канонической инструкции курса](../README.md#подготовка-и-канонический-сброс). +После успешного `world_init` эталонный мир уже прошёл путь Kafka → STG → ODS → DDS → DM. -```bash -make generated-history-analytics -make up -``` +Перед разбором `etl_pipeline` вспомни пульт курса как лесенку: + +1. `ddl_init` создаёт схему ClickHouse; +2. `world_init` импортирует эталонный мир и запускает его обработку; +3. `world_next_day` дозаливает следующий день и снова запускает обработку. + +Первые две ступени ты уже прошёл при подготовке. Третью пока только запомни: подробно +её разберём в лабе 07. Внутри двух последних ступеней работает тот самый +`etl_pipeline`, который мы сейчас откроем отдельно. Открой Airflow: `http://localhost:8080` (логин `admin`, пароль `admin`). Найди DAG `etl_pipeline` и запусти его через **Trigger DAG with config**: @@ -77,8 +81,9 @@ make up заново наполнит их из ODS. Для учебного стенда это удобный чистый прогон: результат повторяемый, старые эксперименты не мешают. -Когда DAG завершится, открой его граф. На чистой стартовой истории все задачи должны быть -зелёными. Найди внутри группы `transform` две задачи подряд: +Когда DAG завершится, открой его граф. На чистой стартовой истории выбранные задачи +должны быть зелёными, а невыбранная ветка `skip_truncate` — в состоянии `skipped`. +Найди внутри группы `transform` две задачи подряд: - `check_dds_integrity` — SQL-задача, которая считает сирот; - `assert_dds_integrity` — Python-задача, которая решает, можно ли идти дальше. @@ -95,8 +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. --- @@ -307,15 +313,8 @@ WHERE click_id IS NOT NULL AND click_id NOT IN (SELECT click_id FROM dds.click); ``` -Снова должно быть `0`. Если стенд после экспериментов совсем запутался, сделай штатный -чистый прогон: - -```bash -make generated-history-analytics -make up -``` - -После этого при необходимости запусти `etl_pipeline` с `{"full_refresh": true}`. +Снова должно быть `0`. Если стенд после экспериментов совсем запутался, пройди +[канонический сброс](../README.md#подготовка-и-канонический-сброс). --- @@ -323,7 +322,7 @@ make up | Действие | Где смотреть | Что ожидать | |----------|--------------|-------------| -| `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 | diff --git a/docs/course/lessons/05_monitoring.md b/docs/course/lessons/05_monitoring.md index 0555a1c..01bd03f 100644 --- a/docs/course/lessons/05_monitoring.md +++ b/docs/course/lessons/05_monitoring.md @@ -57,24 +57,10 @@ ## 2. Руки: открываем дашборды и targets -Подними стенд и создай стартовую историю, если он ещё не поднят. Это штатный путь курса: -готовый источник данных стенда пишет события в Kafka, ClickHouse читает их в STG, затем -batch строит ODS, DDS и DM. Файлы `data/*.jsonl` пока остаются только кладовкой значений -для этого источника, а не источником аналитического контура. - -```bash -make generated-history-analytics -make up -``` - -Учти: путь `generated-history-analytics` собирает витрины напрямую, без запуска -`etl_pipeline`, поэтому панели про задачи Airflow в `Airflow Overview` останутся -пустыми, пока ты хотя бы раз не запустишь `etl_pipeline` сам (это делалось в -уроке 4). Запустить его можно в Airflow с конфигом: - -```json -{"full_refresh": true} -``` +Подготовь стенд по +[канонической инструкции курса](../README.md#подготовка-и-канонический-сброс). +После успешного `world_init` эталонный мир уже прошёл путь Kafka → STG → ODS → DDS → DM, +а в Airflow есть завершённый прогон `etl_pipeline`. Нам нужны не идеальные объёмы, а живой стенд, в котором есть Kafka-топики, строки в ClickHouse и хотя бы один прогон Airflow. @@ -92,10 +78,10 @@ make up | `airflow` | `statsd-exporter:9102` | `statsd-exporter` отдаёт метрики Airflow в формате Prometheus | | `generator` | `generator:9109` | live-генератор отдаёт свои метрики, только когда явно запущен | -У `clickhouse`, `kafka` и `airflow` состояние должно быть `UP`. `generator` на штатном -backfill-only стенде может быть `DOWN`, потому что `make up` не запускает live-генератор. -Это нормально для курса до явного `make generator-continue`. Если один из трёх основных -target `DOWN`, Grafana дальше будет показывать `No data` или старые значения. +У `clickhouse`, `kafka` и `airflow` состояние должно быть `UP`. `generator` на базе +импортированного мира может быть `DOWN`: живой поток включается отдельно через +`make generator-continue`. Если один из трёх основных target `DOWN`, Grafana дальше +будет показывать `No data` или старые значения. То же можно проверить из терминала: @@ -125,7 +111,7 @@ curl -s http://localhost:9090/api/v1/targets | grep -o '"health":"[^"]*"' - `ClickHouse Overview`; - `Kafka Overview`; - `Airflow Overview`; -- `Generator Overview` — про live-генератор; на backfill-only стенде он пуст, +- `Generator Overview` — про живой поток; на базе импортированного мира он пуст, как и target `generator` выше, и в этом уроке не понадобится. Открой каждый и смотри не на красоту графиков, а на смысл: какой слой стенда он показывает и @@ -272,11 +258,11 @@ Grafana. | Airflow Alerts | `High Task Failure Rate` | растёт rate failed tasks | | Airflow Alerts | `High DAG Parse Time` | DAG-файлы долго парсятся | -Не все эти правила обязаны быть тихими в учебном стенде. По умолчанию `make up` и -`make generated-history-analytics` **не запускают live-генератор**, поэтому после готовой -стартовой истории новые сообщения перестают приходить. Из-за этого `Kafka No Messages Produced` -может перейти в `Alerting` на полностью здоровом backfill-only стенде. Если хочешь проверить -это правило в спокойном состоянии, явно включи live: +Не все эти правила обязаны быть тихими в учебном стенде. После импорта эталонной базы +живой поток не запущен, поэтому новые сообщения не приходят. Из-за этого +`Kafka No Messages Produced` может перейти в `Alerting` на полностью здоровом стенде: +алерт честно говорит, что потока сейчас нет, а не что импорт сломан. Если хочешь +проверить правило в спокойном состоянии, явно включи живой поток: ```bash make generator-continue @@ -288,6 +274,10 @@ make generator-continue make generator-down ``` +Остановка генератора не удаляет уже приехавшие данные и сохранённое состояние. +После этого необязательного опыта верни эталонный мир по +[канонической инструкции курса](../README.md#подготовка-и-канонический-сброс). + --- ## 4. Управляемая правка: остановим scheduler и увидим алерт @@ -395,7 +385,6 @@ make recover-monitoring ## Мост к следующему шагу -Теперь стенд закрывает полный учебный маршрут: Kafka принимает поток, ClickHouse раскладывает -слои, Airflow управляет порядком, а Prometheus и Grafana показывают состояние системы. Дальше -этот же стенд можно использовать не как разовый набор уроков, а как тренажёр: менять данные, -ломать отдельные места, смотреть, где появляется сигнал, и объяснять по метрикам, что произошло. +Теперь ты видишь состояние пайплайна со стороны: где идут данные, где растёт +отставание и где сработал алерт. В уроке 6 у готовых витрин появится +потребитель — дашборд в Superset. diff --git a/docs/course/lessons/06_superset_bi.md b/docs/course/lessons/06_superset_bi.md index 1addf96..8411d6e 100644 --- a/docs/course/lessons/06_superset_bi.md +++ b/docs/course/lessons/06_superset_bi.md @@ -63,17 +63,10 @@ DM-витрина — это SQL-объект в ClickHouse. Она задаёт ## 2. Руки: запускаем Superset и смотрим дашборд -Подними стенд и создай стартовую историю. Это штатный путь курса: готовый источник -данных стенда пишет события в Kafka, ClickHouse читает их в STG, затем batch строит -ODS, DDS, DM и обновляет Superset. Файлы `data/*.jsonl` пока остаются только кладовкой -значений для этого источника, а не источником аналитического контура. - -```bash -make generated-history-analytics -make up -``` - -Команда уже прогоняет цепочку STG → ODS → DDS → DM и создаёт metadata Superset: +Подготовь стенд по +[канонической инструкции курса](../README.md#подготовка-и-канонический-сброс). +После успешного `world_init` эталонный мир уже прошёл путь Kafka → STG → ODS → DDS → DM, +а `make superset-init` создал метаданные Superset: - подключение `clickhouse_dwh`; - 6 datasets поверх `dm.*`. @@ -432,9 +425,8 @@ metadata Superset. --- -## Мост после курса +## Мост к следующему шагу -Теперь у тебя есть сквозная цепочка: Kafka → ClickHouse STG → ODS → DDS → DM → мониторинг → -Superset. Следующий честный вопрос уже не про этот стенд, а про продакшен: какие витрины стоит -материализовать, какие права дать BI-пользователям и как не превратить dashboard в единственный -источник правды вместо версионированного SQL в репозитории. +Теперь у тебя есть сквозная цепочка от Kafka до дашборда в Superset. В лабе 07 +ты добавишь в этот мир следующий модельный день и проверишь границу двух +суточных порций. diff --git a/docs/course/lessons/07_lab_next_day.md b/docs/course/lessons/07_lab_next_day.md new file mode 100644 index 0000000..4489dbe --- /dev/null +++ b/docs/course/lessons/07_lab_next_day.md @@ -0,0 +1,297 @@ +# Лаба 07. Следующий день и границы времени + +> Формат: **практика** — добавишь к эталонному миру ещё один модельный день, +> проверишь стык дней и ненадолго включишь расписание. +> Пререквизит: пройден урок 6; ты умеешь запускать DAG в Airflow, читать граф +> задач и смотреть дневную динамику в Superset. +> Эталонный путь: +> [`airflow/dags/world_next_day_dag.py`](../../../airflow/dags/world_next_day_dag.py). +> +> Поток данных одной строкой: +> `эталонный мир (3 дня) → world_next_day → Kafka → STG → etl_pipeline → ODS → DDS → DM (4 дня)` +> +> О чём лаба простыми словами: ночью в хранилище часто приезжает очередной день +> данных. Ты добавишь такой день и проверишь место, где две суточные порции +> соприкасаются. + +--- + +## 1. Зачем и где в проде + +Пакетный инкремент — это новая порция данных за ограниченный период. Например, +ночью в хранилище приезжают события за вчера. Старые дни остаются на месте, а +пайплайн добавляет и проверяет новый. + +Граница порции требует отдельного внимания. Визит может начаться перед +полуночью, а закончиться после неё. Если обрабатывать каждый день как +изолированный файл, половины такого визита окажутся в разных запусках. Отсюда +появляются проверки стыков, запоздавшие данные и пересчёт хвоста прошлого дня. + +В этой лабе `world_next_day` добавляет ровно 24 модельных часа. Сам генератор +накапливает контрольные числа, а затем DAG запускает полный пересчёт +`etl_pipeline`. На маленьком стенде такой путь понятен и повторяем. + +> **В проде иначе.** Суточную задачу обычно привязывают к интервалу данных: +> повторный запуск того же интервала должен дать тот же результат. Здесь каждый +> ручной запуск сдвигает правую границу мира ещё на сутки. Это удобно для +> тренажёра, но не делает DAG идемпотентным. + +--- + +## 2. Руки: добавляем день 4 + +Подготовь стенд по +[канонической инструкции курса](../README.md#подготовка-и-канонический-сброс). +После успешного `world_init` в мире есть три модельных дня. + +### Запускаем `world_next_day` + +Открой Airflow: `http://localhost:8080` (логин `admin`, пароль `admin`). +Найди DAG `world_next_day` и запусти его через **Trigger DAG** с пустой формой. +Параметры ему не нужны. + +В графе должны позеленеть задачи: + +```text +check_etl_not_paused + → precheck_next_day + → run_next_day + → trigger_etl + → check_after_etl +``` + +`run_next_day` дописывает день 4 в Kafka. Затем `trigger_etl` запускает +`etl_pipeline` с полным пересчётом DDS и DM. Последняя задача сверяет ClickHouse +с manifest — паспортом мира, где лежат границы времени и контрольные числа. + +Если `check_etl_not_paused` красная, сними паузу с `etl_pipeline` и повтори +запуск. `world_next_day` не может вызвать зависимый DAG, пока тот стоит на паузе. + +### Сверяем manifest + +После зелёного прогона выполни: + +```bash +docker compose run --rm --no-deps generator \ + python -m clickstream_generator.startup_history_artifact kafka-manifest-summary +``` + +Команда печатает одну строку с полями в таком порядке: + +```text +events visits users min_event_timestamp max_event_timestamp model_t0 model_t_end profile +``` + +Сверь результат с контрольными значениями: + +| Поле | Ожидаем после дня 4 | +|------|----------------------| +| `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`, а не только к четвёртому дню. + +Теперь проверь все завершённые порции и их стыки: + +```bash +make generated-history-chain-check +``` + +Проверка должна закончиться без ошибки. Она подтверждает, что каждая порция +не пуста, контрольные числа STG совпадают с manifest, после правой границы нет +хвоста, а переходящие визиты не потеряли связанные записи. + +### Находим переходящие визиты + +Открой SQL-консоль `http://localhost:9123/play` (пользователь `default`, пароль +`123456`) и выполни: + +```sql +SELECT + click_id, + min(event_ts) AS visit_start, + max(event_ts) AS visit_end +FROM dds.event +WHERE event_ts IS NOT NULL + AND click_id IS NOT NULL +GROUP BY click_id +HAVING toDate(min(event_ts)) <> toDate(max(event_ts)) +ORDER BY visit_start; +``` + +Запрос должен найти хотя бы один визит, чьи события лежат по разные стороны +полуночи. На контрольном мире он находит `35`, `22` и `26` визитов на стыках +1→2, 2→3 и 3→4. Для проверки нового дня достаточно увидеть `26` визитов на +стыке 3→4. + +> **Как слова связаны со схемой.** В manifest написано «визит», а в SQL такой +> визит обозначен `click_id`. Отдельной таблицы визитов нет: контекст лежит в +> `dds.click`, а события визита — в `dds.event`. + +Почему запрос идёт прямо в `dds.event`? Представление +`dm.v_session_overview` группирует строки по `event_date` и `click_id`. +Переходящий визит там уже разрезан на две строки: одна до полуночи, другая +после. Чтобы увидеть визит целиком, сначала собираем все его события по +`click_id`, а потом сравниваем даты самого раннего и самого позднего события. + +### Смотрим день-к-дню + +Открой в Superset дашборд `E-commerce Analytics Dashboard`. Поставь фильтр +**Date Range → No filter** и найди график `Events over Time`. После обновления +на нём должны появиться непустые пятиминутные интервалы дня 4 после дня 3. +Высоту отдельной точки не сравнивай с накопительными числами manifest: +manifest уже сверили отдельно в начале секции. + +--- + +## 3. Загляни внутрь + +Открой +[`airflow/dags/world_next_day_dag.py`](../../../airflow/dags/world_next_day_dag.py). +Это короткая оркестрация вокруг трёх действий: проверить точку продолжения, +дописать сутки и пересчитать аналитику. + +### Почему сначала идут проверки + +`check_etl_not_paused` убеждается, что Airflow сможет запустить +`etl_pipeline`. `precheck_next_day` проверяет, что: + +- живой генератор сейчас не работает; +- состояние генератора согласовано с правой границей manifest; +- накопительные счётчики не повреждены. + +Только после этого `run_next_day` добавляет новую порцию. Так два режима, +`next-day` и `continue`, не пишут в один мир одновременно. + +### Почему ETL запускается с `full_refresh` + +`TriggerDagRunOperator` передаёт в `etl_pipeline`: + +```python +conf={"full_refresh": True} +``` + +Значит, DDS и DM пересобираются из всей накопленной ODS, а не только из нового +дня. Это простой путь для учебного стенда, но его цена растёт вместе с миром. + +Открой в Grafana дашборд `Airflow Overview` и панель `Task Duration (avg)`. +Сравни длительности задач двух запусков `world_next_day`: ручного дня 4 и +запланированного дня 5 из следующей секции. + +Генерация одной суточной порции должна занимать примерно одинаковое время: +счётчики manifest обновляются только новыми данными. Это исправлено в +[issue #5](https://github.com/dementev-dev/clickstream-ch-kafka-superset-demo/issues/5). +Полный ETL перечитывает растущий мир, поэтому со временем дорожает. Это +осознанный долг из +[issue #8](https://github.com/dementev-dev/clickstream-ch-kafka-superset-demo/issues/8), +а не поломка этой лабы. + +### Что делает расписание + +У DAG задано расписание `*/30 * * * *`: запуск каждые 30 минут. При создании +DAG стоит на паузе, поэтому мир сам не растёт. `catchup=False` запрещает +Airflow проигрывать все пропущенные интервалы с 2024 года. После снятия паузы +будет ждать только ближайший новый интервал. + +--- + +## 4. Управляемая правка: включаем один запуск по расписанию + +Сейчас день 4 был добавлен вручную. На короткое время разрешим Airflow добавить +день 5 автоматически. + +### Шаг 1. Зафиксируй правую границу + +Ещё раз выполни команду `kafka-manifest-summary` из секции 2 и запиши +`model_t_end`. После запланированного запуска она должна сдвинуться ровно на +24 модельных часа. + +### Шаг 2. Сними паузу + +В списке DAG-ов Airflow включи переключатель `world_next_day`. Не нажимай +**Trigger DAG**: на этот раз нужен именно запуск по расписанию. + +Ближайший запуск появится на следующей получасовой границе. Следи за списком +запусков и графом DAG. + +### Шаг 3. Останови расписание после одного успеха + +Как только один запланированный запуск станет зелёным, снова поставь +`world_next_day` на паузу. Сделай это до следующей получасовой границы, иначе +приедет ещё один день. + +Проверь три результата: + +1. `kafka-manifest-summary` показывает новую правую границу и новые + накопительные числа; +2. `make generated-history-chain-check` заканчивается без ошибки; +3. в `Events over Time` появился день 5. + +После дня 5 контрольный manifest показывает `468025` событий, `43482` визита, +`6713` пользователей и `model_t_end = 2026-01-06T00:00:00+00:00`. + +### Верни как было + +Сначала убедись, что `world_next_day` снова стоит на паузе. Но одной паузы +недостаточно: мир уже вырос до пяти дней. Верни чистый трёхдневный эталон по +[канонической инструкции курса](../README.md#подготовка-и-канонический-сброс). + +--- + +## 5. Проверь себя + +| Действие | Где смотреть | Что ожидать | +|----------|--------------|-------------| +| запустить `world_next_day` с пустой формой | Airflow graph | все пять задач зелёные, `etl_pipeline` вызван с `full_refresh=true` | +| прочитать manifest после дня 4 | `kafka-manifest-summary` | числа совпадают с таблицей секции 2 | +| проверить накопленную цепочку | `make generated-history-chain-check` | команда завершается без ошибки | +| найти переходящие визиты | запрос в `dds.event` | на стыке 3→4 найдено `26` визитов | +| снять паузу с DAG | Airflow runs | на ближайшей получасовой границе приезжает один запланированный день | +| вернуть паузу и пройти сброс | Airflow и Superset | расписание выключено, снова видны три эталонных дня | + +Ответь своими словами: + +- почему `dm.v_session_overview` не подходит для поиска визита через полночь; +- зачем проверять стык двух суточных порций; +- почему `catchup=False` не означает «у DAG нет расписания»; +- что случится, если дважды вручную запустить `world_next_day`; +- почему промышленная задача «день D» должна повторно обрабатывать тот же день, + а учебный DAG при повторе добавляет день D+1. + +Последние два вопроса — проверка идемпотентности. Идемпотентная операция при +повторе с тем же входом оставляет тот же результат. `world_next_day` устроен +иначе: его входом служит текущая правая граница мира, и после каждого успеха +эта граница меняется. + +--- + +## 6. Что должно получиться + +После лабы сохрани: + +- скрин зелёного ручного запуска `world_next_day`, который добавил день 4; +- скрин зелёного запланированного запуска `world_next_day`, который добавил + день 5; +- вывод `make generated-history-chain-check` без ошибки; +- результат SQL-запроса с переходящим через полночь `click_id`; +- скрин дневного графика с днём 5 после дня 4; +- короткое объяснение своими словами: чем пакетный инкремент отличается от + полного пересчёта и зачем отдельно проверять границу порций. + +Перед переходом дальше стенд должен быть возвращён к трёхдневному эталонному +миру, а `world_next_day` — стоять на паузе. + +--- + +## Мост к следующему шагу + +Здесь новый день приехал как законченная и повторяемая порция. В следующей +лабе граница исчезнет: события пойдут непрерывно, а разные слои начнут +отставать друг от друга. diff --git a/docs/course/lessons/08_lab_continue.md b/docs/course/lessons/08_lab_continue.md new file mode 100644 index 0000000..68a9553 --- /dev/null +++ b/docs/course/lessons/08_lab_continue.md @@ -0,0 +1,308 @@ +# Лаба 08. Живой поток и свежесть данных + +> Формат: **практика** — запустишь живой поток, остановишь его и продолжишь с +> сохранённого состояния. +> Пререквизит: пройдена лаба 07; стенд возвращён к эталонному миру, а +> `world_next_day` стоит на паузе. +> Эталонные пути: +> [`Makefile`](../../../Makefile) и +> [`airflow/dags/etl_pipeline_dag.py`](../../../airflow/dags/etl_pipeline_dag.py). +> +> Поток данных одной строкой: +> `generator-continue → Kafka → STG (сразу) → etl_pipeline (по запуску) → ODS → DDS → DM` +> +> О чём лаба простыми словами: события будут приезжать непрерывно. Ты увидишь, +> что Kafka и STG уже обновились, а аналитические слои ещё ждут пакетного +> запуска. + +--- + +## 1. Зачем и где в проде + +У данных есть **свежесть** — задержка между событием в источнике и моментом, +когда его увидит потребитель. Для одного отчёта допустимы сутки, для другого +нужны минуты. Обещание «данные готовы не позже чем через N минут» влияет на +всю цепочку. + +В этом стенде приём потоковый: ClickHouse постоянно читает Kafka, а +Materialized View сразу складывает сообщения в STG. Дальнейшая обработка +пакетная: ODS, DDS и DM обновляет только запуск `etl_pipeline`. + +Поэтому **потоковый приём не равен потоковой обработке**. Новые строки уже +могут лежать в STG, пока витрина для аналитика показывает старую правую +границу. + +> **В проде иначе.** Требование к свежести задают вместе с потребителем данных. +> Под него выбирают расписание, размер порции, проверки и допустимое +> отставание. Слово «поток» само по себе ничего не обещает о времени появления +> данных в отчёте. + +--- + +## 2. Руки: оживляем мониторинг и догоняем витрины + +Подготовь стенд по +[канонической инструкции курса](../README.md#подготовка-и-канонический-сброс). +После `world_init` живой генератор выключен, а эталонный мир заканчивается на +одной и той же модельной границе у всех менти. + +### Запоминаем начальную свежесть + +Открой SQL-консоль `http://localhost:9123/play` (пользователь `default`, пароль +`123456`) и выполни: + +```sql +SELECT + 'STG' AS layer, + max(parseDateTime64BestEffortOrNull( + JSONExtractString(raw, 'event_timestamp'), 6 + )) AS max_event_ts +FROM stg.browser_raw + +UNION ALL + +SELECT 'ODS', max(event_ts) +FROM ods.browser_event + +UNION ALL + +SELECT 'DDS', max(event_ts) +FROM dds.event + +UNION ALL + +SELECT 'DM', max(event_ts) +FROM dm.v_events_enriched; +``` + +На эталонном мире правые границы слоёв должны совпадать. +После импорта трёх эталонных дней все четыре слоя заканчиваются на +`2026-01-03 23:59:59.066766`. + +### Запускаем живое продолжение + +В терминале выполни: + +```bash +make generator-continue +``` + +Команда запускает генератор в режиме `continue`: он берёт сохранённое состояние +эталонного мира и продолжает его, не создавая новый мир с нуля. + +Вернись к экранам из урока 5: + +- в Prometheus target `generator` переходит в `UP`; +- в `Generator Overview` оживает `Total Events/min (all 4 topics)`; +- в `Kafka Overview` растут сообщения и конечные offset-ы партиций; +- в Kafka UI из урока 0 у новых сообщений растут offset-ы; +- `Kafka No Messages Produced` выходит из `Alerting`, когда Grafana увидит + устойчивый поток. + +У алерта есть окно `for: 10m`: в `Alerting` он переходит только после десяти +минут тишины. Когда поток вернётся, условие перестанет выполняться при +ближайшей оценке. График событий и offset-ы должны начать двигаться раньше. + +Подожди несколько минут и снова выполни запрос по слоям. Правая граница STG +уйдёт вперёд. ODS, DDS и DM останутся на границе последнего +`etl_pipeline`. Запиши время по обычным часам и границы STG и DM. Границу STG +считай контрольной. + +### Догоняем пакетные слои + +Открой Airflow и запусти `etl_pipeline` через **Trigger DAG** с пустой формой. +По умолчанию это полный пересчёт. При запуске включи секундомер. + +После зелёного прогона останови секундомер и повтори запрос. Запиши новые +границы STG и DM: DM должна достичь контрольной границы STG. Генератор всё ещё +работает, поэтому STG вскоре снова может оказаться чуть свежее. Это ожидаемое +расслоение, а не потеря строк. + +--- + +## 3. Загляни внутрь + +Открой [`Makefile`](../../../Makefile) и найди две команды этой лабы: + +```make +generator-continue: + COMPOSE_BIN="$(COMPOSE)" bash ./scripts/run_generator.sh continue + +generator-down: + $(COMPOSE) stop generator +``` + +`generator-down` останавливает только контейнер генератора. Kafka, ClickHouse и +читатели STG продолжают работать. Поэтому они успевают дочитать уже записанные +сообщения. + +`generator-continue` запускает генератор с сохранением прежнего состояния. В нём +остаются популяция пользователей, активные визиты, случайное состояние и +модельные часы. После рестарта продолжается тот же мир. + + +> **Модельное время.** Профиль `daily-wave` ускоряет внутренние часы стенда в +> 60 раз: модельные сутки проходят примерно за полчаса. Поэтому «сейчас» на +> дашборде может обгонять настенные часы. Ручка тренажёра называется +> `GEN_MODEL_TIME_SPEED`; в этой лабе её не меняем. + +### Где проходит граница между потоком и пакетом + +В уроке 1 ты видел Kafka engine и Materialized View. Пока сервисы работают, +этот путь сам переносит сообщения из Kafka в `stg.*_raw`. + +У `etl_pipeline` в +[`airflow/dags/etl_pipeline_dag.py`](../../../airflow/dags/etl_pipeline_dag.py) +стоит `schedule=None`. ODS, DDS и DM не догоняют поток сами: нужен ручной +запуск DAG. Получается две разные гарантии свежести: + +- Kafka → STG работает постоянно; +- STG → ODS → DDS → DM обновляется по запуску `etl_pipeline`. + +Если потребителю нужна витрина не старше пяти минут, ручной DAG такую +гарантию не даёт. Понадобится расписание или другая обработка, а затем +измерение фактического отставания. + +--- + +## 4. Управляемая правка: останавливаем и продолжаем поток + +Проверим, переживает ли мир остановку генератора. + +### Шаг 1. Зафиксируй состояние до остановки + +Пока генератор работает, запиши: + +- правую границу STG из запроса секции 2; +- значение `Total Events/min (all 4 topics)`; +- текущие конечные offset-ы в `Kafka Overview`. + +### Шаг 2. Останови генератор + +```bash +make generator-down +``` + +После остановки: + +- target `generator` переходит в `DOWN`; +- `Total Events/min (all 4 topics)` падает к нулю или перестаёт обновляться; +- конечные offset-ы перестают расти; +- читатели ClickHouse дочитывают уже записанные сообщения, и отставание + стекает к нулю; +- после окна ожидания `Kafka No Messages Produced` переходит в `Alerting`. + +Как и в уроке 5, Kafka UI или панель lag могут не показать прогресс групп +`ch_stg_*`. Тогда смотри `system.kafka_consumers` со стороны ClickHouse. +Главный признак здесь такой: новые offset-ы больше не появляются, а строки, +уже записанные до остановки, доходят в STG. + +### Шаг 3. Продолжи тот же мир + +```bash +make generator-continue +``` + +Проверь обратную картину: + +- target `generator` снова в `UP`; +- график событий и offset-ы снова растут; +- `Kafka No Messages Produced` возвращается в норму при ближайшей оценке; +- правая граница STG стала больше значения, записанного до остановки. + +Генератор взял прежнюю популяцию и прежнюю точку продолжения. Остановка не +превратила поток в новый независимый мир. + +### Шаг 4. Догони аналитику и увидь расхождение мира + +Запусти `etl_pipeline` с пустой формой. После зелёного прогона выполни: + +```sql +WITH parseDateTime64BestEffort( + substring('', 1, 19), 6 +) AS boundary +SELECT + user_domain_id, + countIf(visit_start < boundary) AS visits_before, + countIf(visit_start >= boundary) AS visits_after, + min(visit_start) AS first_visit_ts, max(visit_start) AS last_visit_ts +FROM ( + SELECT user_domain_id, click_id, min(event_ts) AS visit_start + FROM dm.v_events_enriched + WHERE user_domain_id IS NOT NULL + GROUP BY user_domain_id, click_id +) +GROUP BY user_domain_id +HAVING min(visit_start) < boundary AND max(visit_start) >= boundary +LIMIT 10; +``` + +Подставь вместо `` границу из манифеста. Успех — хотя бы одна +строка. Иначе подожди несколько минут, перезапусти `etl_pipeline` и повтори +запрос. Это симметрия лаб: в лабе 07 визит пересекал полночь, здесь +пользователь с разными визитами пересекает границу замороженного мира. +Мир расходится: итоги превышают эталонные и различаются у менти. Это правильно. + +Например, для границы `2026-01-06T00:00:00+00:00` найден пользователь +`5938d296-14cf-49cb-801e-82370518cf59`: `13` визитов до границы и `2` после +неё. Твои числа и идентификатор будут другими. + +### Верни как было + +Останови живой поток: + +```bash +make generator-down +``` + +После `continue` мир стал недетерминированным: его числа зависят от длительности +запуска. Верни стенд к эталону по +[канонической инструкции курса](../README.md#подготовка-и-канонический-сброс). +Не пересказывай шаги по памяти: эта ссылка остаётся единственным учебным +описанием полного сброса. + +--- + +## 5. Проверь себя + +| Действие | Где смотреть | Что ожидать | +|----------|--------------|-------------| +| выполнить `make generator-continue` | Prometheus Targets | `generator` переходит в `UP` | +| измерить свежесть до и после 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` | есть пользователь с визитами до и после границы | + +Ответь своими словами: + +- почему STG опережает DM, а `etl_pipeline` догоняет лишь очередной снимок; +- что сохраняет `continue` и почему итоги менти расходятся; +- какую свежесть можно честно обещать при ручном запуске и при расписании + каждые 30 минут; +- почему минутная свежесть лежит за пределами этого стенда. + +--- + +## 6. Что должно получиться + +После лабы сохрани: + +- скрин `generator` в `DOWN` после `make generator-down` и скрин `generator` + в `UP` после повторного `make generator-continue`; +- границы STG до остановки и после продолжения; +- строку пользователя с визитами до и после `model_t_end`; +- замеры до и после `etl_pipeline`: разницу `STG − DM` в модельном времени и + настенное ожидание до появления зафиксированной границы STG в DM; +- вывод своими словами: что можно обещать при ручном запуске, что — при + расписании каждые 30 минут и почему этот стенд не обещает минутную свежесть. + +В конце `generator` должен быть остановлен, а стенд — возвращён к эталонному +миру. + +--- + +## Мост после курса + +Замер отвечает на вопрос из секции 1: ручной запуск не ограничивает ожидание, +а расписание раз в 30 минут дало бы около 30 минут плюс время прогона. Для +минутной свежести нужна другая обработка. diff --git a/docs/specs/2026-07-23-course-labs-redesign.md b/docs/specs/2026-07-23-course-labs-redesign.md new file mode 100644 index 0000000..d778c21 --- /dev/null +++ b/docs/specs/2026-07-23-course-labs-redesign.md @@ -0,0 +1,219 @@ +# Редизайн лаб курса под три режима менти + +Статус: принято, в работе (ветка `feature/course-labs-redesign`). +Источник: issue #7; аудит курса против стенда и штурм педагогики лаб +с пользователем 2026-07-23. Опирается на +[спеку редизайна пути менти](./2026-07-19-mentee-path-redesign.md) +(реализована, issue #9 закрыт) и термины из `CONTEXT.md` +(три режима менти, эталонный мир, модельное время). + +## Проблема + +Курс в `docs/course/` (7 уроков + 4 метадокумента) целиком построен на +старом пути `make generated-history-analytics && make up`: ни один документ +не знает про `world_init`, `import`, эталонный мир и `make superset-init`. +После редизайна пути менти курс в текущем виде не заработает. + +При этом педагогическое ядро уроков цело: SQL, DAG'и, мониторинг и +Superset актуальны — ломаются только блоки подготовки стенда и обвязка. +А два новых режима роста мира (`next-day`, `continue`) вообще не имеют +уроков, хотя спека пути менти прямо называет их «разные педагогики и +разные лабы». + +## Цели + +- Менти проходит курс по новому пути: `make up` → Airflow UI + (`ddl_init` → `world_init` с пустой формой) → `make superset-init`. +- Оба режима роста мира получают по обязательному уроку-лабе с + собственным учебным паттерном. +- Числа в уроках совпадают у всех менти число-в-число (следствие + `import` эталонного мира) и сверены с живым стендом. +- Смешанных состояний не бывает: вся перестройка — один PR. + +## Целевая модель + +### Маршрут курса + +Маршрут остаётся линейным, все уроки — обязательные: + +- **Урок 6 (Superset) становится обязательным.** Его опциональность + родом из PRD («делаем, если останется ресурс») — экономия ресурса на + написание, а урок давно написан; причина истекла. Ожидания + работодателей к дата-инженеру теперь включают витринно-BI-слой. +- **Две новые лабы после урока 6, обе обязательные:** + `07_lab_next_day.md`, затем `08_lab_continue.md`. Порядок осознанный: + next-day детерминирована (числа сойдутся у всех — можно печатать + точные ожидания), continue — живая и недетерминированная, «второй шаг + свободы». +- Обе лабы строятся по шаблону `LESSON_STANDARD.md` (шапка + 6 секций). + +### Новая семантика сброса + +Для лаб «верни как было» — это не откат правки в git: мир вырос. +Канонический сброс — тот же путь, что и первый вход, и описывается он +в одном месте (README курса), уроки ссылаются: + +1. `make clean` (сносит и данные, и Superset — это надо назвать явно); +2. `make up` → Airflow UI: `ddl_init` → `world_init` с пустой формой; +3. `make superset-init`. + +Идём UI-путём менти, а не мейнтейнерским `make startup-history-import`. +Старая команда сброса `make generated-history-analytics && make up` +уходит из учебных текстов (остаётся мейнтейнеру и CI). + +### Лаба 07 — next-day + +Паттерн одной фразой: **пакетный инкремент дня и его границы** +(прод-аналог — «ночью приехал вчерашний день»). + +- **Руки:** триггер беспараметрного `world_next_day` → день 4 в + витринах; сверка чисел manifest (точные ожидаемые значения); + `make generated-history-chain-check` на стыке; день-к-дню на дашборде. +- **Центральное открытие — переходящие визиты:** менти SQL-запросом + находит визиты, чьи события лежат по обе стороны полуночи, и осознаёт, + почему суточная нарезка режет живые сессии — откуда берутся проверки + стыков, late data и пересчёт вчерашнего хвоста. Этого нет ни в одном + уроке 0–6. Важно для текста лабы: запрос идёт в `dds.event` напрямую + (`GROUP BY click_id` + `HAVING toDate(min(event_ts)) <> + toDate(max(event_ts))`) — витрина `dm.v_session_overview` не годится, + она группирует по `event_date` и режет переходящий визит на две + строки. Это новый приём, который лаба сама и учит (уроки давали + argMax/JOIN, но не агрегацию визита целиком). Годится любой стык дней: + стыки 1→2 и 2→3 внутри эталонного мира есть точно; живы ли визиты на + стыке 3→4 (заморозка мира могла закрыть открытые сессии) — проверяется + на стенде при реализации. Здесь же — короткая врезка про соответствие + терминов: «визит» из manifest = `click_id` в SQL (в схеме нет таблицы + «визитов», есть `dds.click`). +- **Управляемая правка:** включить расписание `world_next_day` в + Airflow UI → день 5 приезжает сам → выключить обратно. Спека пути + менти прямо проектировала paused-расписание под этот шаг лабы; заодно + менти трогает paused/`catchup`. +- **Цена роста — коротким наблюдением:** в длительностях задач Airflow + генерация дня ~постоянна (инкрементальные счётчики manifest, issue #5), + а `full_refresh` ETL растёт с миром — осознанный долг, ссылка на #8. +- **Самопроверка — крючок про идемпотентность:** «триггерни дважды — + что будет? почему прод-джобы за день устроены иначе?». +- Эталонный путь: `airflow/dags/world_next_day_dag.py`. + +### Лаба 08 — continue + +Паттерн одной фразой: **живой поток и свежесть данных**. + +- **Ядро — потоковый приём ≠ потоковая обработка:** Kafka и STG + пополняются сами (MV урока 1), витрины DM стоят до прогона + `etl_pipeline`; запустил — догнали и снова отстают. Свежесть данных + как явное понятие: «какую свежесть обещаем потребителю?». +- **Руки:** `make generator-continue` → мониторинг урока 5 оживает + (target генератора в UP, алерт `Kafka No Messages Produced` гаснет, + events/min шевелится, офсеты урока 0 растут) → наблюдение расслоения + свежести → догон через `etl_pipeline`. +- **Управляемая правка — останови поток и продолжи:** + `make generator-down` → алерт срабатывает, потребитель дочитывает лаг + до нуля; `make generator-continue` → мир продолжается с места + остановки (популяция и модельные часы пережили рестарт). Это суть + режима continue и рабочий паттерн устойчивости. По шаблону + LESSON_STANDARD лаба всё равно завершается явным шагом «верни как + было» — канонический сброс (см. выше): мир после continue + недетерминирован, и к следующему прохождению стенд возвращается к + эталону. +- **Финальное наблюдение — «мир расходится»:** после continue числа + менти перестают совпадать с эталонными, и это правильно; возвраты + пользователей вживую разводят `users < sessions` (на статике было + невозможно — см. `CONTEXT.md`). +- **Модельное время — только врезкой:** «стенд ускорен в 60 раз, чтобы + сутки потока уложились в ~полчаса; "сейчас" дашборда может обгонять + настенные часы», ручка `GEN_MODEL_TIME_SPEED` — одной строкой. + Это механизм тренажёра, не рабочий паттерн — секционного веса не даём. + Висячую ссылку `CONTEXT.md` про «урок 7» поправить на эту врезку. + +Симметрия курса: 07 — про границы времени в пакетном мире, 08 — про +свежесть в потоковом; обе лабы растят один и тот же импортированный мир +двумя способами. + +### Существующие уроки и метадокументы + +По вердиктам аудита: + +- **Все уроки 0–6:** блок подготовки в §2 заменить на новый путь + (import); сам блок вынести в одно каноническое место (README курса), + уроки ссылаются на него. +- **Все цифры уроков 0–6 недействительны.** Старые значения считались + на старом мире (другой профиль, другой объём); с переходом на + эталонный мир каждое цитируемое число каждого урока пересчитывается + заново. Это полноправный этап работы, а не финальная галочка — даже в + уроках, где правки текста минимальны. +- **Урок 1:** управляемая правка «добавь колонку и получи свежие + сообщения» переводится с перегенерации backfill на дозаливку через + `world_next_day` (детерминированно). Известная цена: прогон DAG'а + занимает минуты — урок называет её честно; сам DAG подаётся кнопкой- + анонсом без разбора («подробно — в лабе 07»), чтобы не красть у лабы + её материал. +- **Урок 4:** добавить лесенку `ddl_init` → `world_init` → + `world_next_day` как контекст оркестрации (менти уже прошёл её руками). +- **Урок 5:** словарь «backfill-only стенд» заменить на «база import / + живой поток»; сценарий алерта `Kafka No Messages Produced` переписать + в этих терминах. +- **Урок 6:** снять пометку «опционально»; цифры сверить с эталонным + миром. +- **README курса:** переписать блоки подготовки и «Проверка чистого + маршрута»; таблицу уроков дополнить лабами. +- **PRD:** не переписывать — одна датированная поправка по его же + конвенции (три режима менти, урок 6 обязателен, лабы 07–08). +- **LEARNING_PLAN:** переписать маршрут и статус, сохранить таблицу + аудита эталонных путей; дополнить её `world_next_day_dag.py`. +- **LESSON_STANDARD:** заменить команду сброса на import-путь. + +## Чего здесь не делаем + +- Не трогаем код стенда: DAG'и, SQL, генератор, инфраструктура — вне + скоупа; работа только с документами курса, README курса и `CONTEXT.md` + (одна ссылка). +- Не чиним рост стоимости `full_refresh` (issue #8) — в лабе 07 он + только показывается. +- Не пишем урок про модельное время — понижено до врезки в лабе 08. +- Не переносим приватные менторские материалы — граница из README курса + остаётся. +- Не переделываем эталонный код уроков 0–6 сверх замены обвязки: аудит + подтвердил, что ядро актуально. + +## Проверка + +- Чистый стенд: менти проходит весь курс 0–8 по текстам уроков, ни разу + не встретив `generated-history-analytics` и `backfill`. +- Каждая цитируемая цифра уроков и лаб сверена с живым стендом после + `import` эталонного мира (обязательный финальный шаг реализации). +- Лаба 07: после `world_next_day` числа manifest совпадают с + напечатанными в лабе; `make generated-history-chain-check` зелёный; + SQL-запрос из лабы находит хотя бы один переходящий визит (на любом + стыке дней). +- Детерминизм инкремента подтверждён до печати чисел: два независимых + прогона `world_next_day` от свежего import дают одинаковые числа + manifest и контрольную сумму. +- Лаба 08: сценарий стоп/продолжение проходит без потерь (лаг стекает к + нулю, после продолжения числа согласованы с manifest). +- Все внутренние ссылки курса живы; `make test` / `make lint` зелёные + (доки код не трогают — проверка от регрессий по касанию). + +## Решения и отклонённые варианты + +- **Лабы — секциями внутри уроков 4/5** — отклонено: уроки распухают, + лабы теряют самостоятельность; выбраны отдельные файлы. +- **Урок 6 оставить опциональным, финал лабы 07 — только SQL** — + отклонено: причина опциональности истекла, рынок ждёт BI-навыков; + урок 6 становится обязательным, лабы опираются на дашборд. +- **Модельное время как секция или отдельный урок** — отклонено: + механизм тренажёра, не рабочий паттерн; врезка. +- **Управляемая правка лабы 08 через ручку ×K** — отклонено по той же + причине; выбран стоп/продолжение потока. +- **Дробить работу на несколько PR** — отклонено: половинчатый курс + (часть уроков про import, часть про backfill) хуже любого из крайних + состояний; этапность — коммитами внутри одной ветки. +- **Цифры «ориентировочно», без сверки со стендом** — отклонено: + воспроизводимость число-в-число — главный козырь эталонного мира. + +## Влияние на документацию + +Вся работа и есть документация: `docs/course/**` (7 уроков + 2 лабы + +4 метадокумента), одна ссылка в `CONTEXT.md`. Корневой `README.md` и +`docs/OPERATIONS.md` уже описывают новый путь — не трогаем. Issue #7 +ссылается на эту спеку и ведёт чек-лист шагов. diff --git a/generator/tests/test_world_dags_contract.py b/generator/tests/test_world_dags_contract.py index 748cd82..7be7f3b 100644 --- a/generator/tests/test_world_dags_contract.py +++ b/generator/tests/test_world_dags_contract.py @@ -188,7 +188,45 @@ def test_course_readme_lists_uv_before_first_command(): text = (REPO_ROOT / "docs" / "course" / "README.md").read_text(encoding="utf-8") assert "uv" in text - assert text.index("uv") < text.index("make generated-history-analytics") + assert text.index("uv") < text.index("make clean") + + +def test_course_canonical_setup_unpauses_dags_before_world_init(): + """Канонический блок называет три DAG, снятие паузы и порядок etl→world_init.""" + # Прозаические формулировки не проверяем — их сторожит ревью, не pytest. + # Держим только структурные факты: якорь-заголовок существует, в блоке + # названы три DAG-а, упомянута пауза, а etl_pipeline идёт до world_init. + text = (REPO_ROOT / "docs" / "course" / "README.md").read_text(encoding="utf-8") + setup = text.split("## Подготовка и канонический сброс", maxsplit=1)[1].split( + "## Уроки", + maxsplit=1, + )[0] + + assert "пауз" in setup + for dag_id in ("ddl_init", "etl_pipeline", "world_init"): + assert dag_id in setup + assert setup.index("etl_pipeline") < setup.index("world_init") + + +def test_next_day_lab_uses_provisioned_superset_dashboard_name(): + """Лаба следующего дня ведёт на существующий дашборд Superset.""" + text = ( + REPO_ROOT / "docs" / "course" / "lessons" / "07_lab_next_day.md" + ).read_text(encoding="utf-8") + + assert "E-commerce Analytics Dashboard" in text + assert "Clickstream Analytics" not in text + + +def test_monitoring_lesson_links_to_canonical_reset(): + """Урок мониторинга ведёт обратно к каноническому сбросу курса.""" + # Проверяем живую перекрёстную ссылку: битый якорь — это дефект, а не + # переписанная проза. Точные фразы упражнения не пиним. + text = ( + REPO_ROOT / "docs" / "course" / "lessons" / "05_monitoring.md" + ).read_text(encoding="utf-8") + + assert "../README.md#подготовка-и-канонический-сброс" in text def test_startup_history_runbook_warns_about_daily_wave_idle_gap():