From 1dd79a3287c82b2f17f24e9b97e014dadaa8d42e Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Thu, 23 Jul 2026 13:36:30 +0300 Subject: [PATCH] =?UTF-8?q?docs(course):=20=D0=BB=D0=B0=D0=B1=D1=8B=2007?= =?UTF-8?q?=20(next-day)=20=D0=B8=2008=20(continue)=20+=20=D0=BC=D0=B5?= =?UTF-8?q?=D1=82=D0=B0=D0=B4=D0=BE=D0=BA=D1=83=D0=BC=D0=B5=D0=BD=D1=82?= =?UTF-8?q?=D1=8B=20(#22)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Зачем: два новых режима роста мира не имели уроков, а метадокументы курса не знали о новом маршруте. Что: лаба 07 «Следующий день и границы времени» — инкремент дня, сверка manifest, переходящие визиты запросом в dds.event, включение расписания на один запуск, цена роста full_refresh; лаба 08 «Живой поток и свежесть данных» — расслоение свежести слоёв, стоп/продолжение генератора, users < sessions, врезка про модельное время; LEARNING_PLAN и README курса согласованы с маршрутом 0–8; в CONTEXT.md починена ссылка «урок 7» (теперь ведёт на врезку лабы 08). Числа со стенда помечены маркером «сверить-на-стенде». Проверка: обе лабы по шаблону LESSON_STANDARD (6 секций, явный «верни как было»); внутренние ссылки разрешаются; git diff --check чистый. Co-Authored-By: Claude Opus 4.8 (1M context) --- CONTEXT.md | 2 +- docs/course/LEARNING_PLAN.md | 30 +-- docs/course/README.md | 9 +- docs/course/lessons/07_lab_next_day.md | 293 ++++++++++++++++++++++++ docs/course/lessons/08_lab_continue.md | 299 +++++++++++++++++++++++++ 5 files changed, 614 insertions(+), 19 deletions(-) create mode 100644 docs/course/lessons/07_lab_next_day.md create mode 100644 docs/course/lessons/08_lab_continue.md diff --git a/CONTEXT.md b/CONTEXT.md index 1c25e86..8c600bf 100644 --- a/CONTEXT.md +++ b/CONTEXT.md @@ -116,7 +116,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 fe6f413..fbaba3c 100644 --- a/docs/course/LEARNING_PLAN.md +++ b/docs/course/LEARNING_PLAN.md @@ -1,9 +1,9 @@ # План обучения: курс «Кликстрим на ClickHouse» > Дата: 2026-06-03 (аудит путей выполнен; середина расщеплена — -> один паттерн на урок, всего 7 уроков, см. §1–2). -> Поправка 2026-07-23: урок 6 стал обязательным; после него добавлены обязательные -> лабы 07–08 про следующий день и живое продолжение. +> один паттерн на урок, см. §1–2). +> Поправка 2026-07-23: весь маршрут стал обязательным. После урока 6 добавлены +> лабы 07–08 про следующий день и живое продолжение; обе лабы написаны. > Назначение: высокоуровневый маршрут менти по курсу — карта уроков, порядок, > результаты аудита эталонных путей. Рамка курса (зачем/что/скоуп) — в `PRD.md`; > как устроен отдельный урок — в `LESSON_STANDARD.md`. @@ -15,10 +15,11 @@ ## 1. Маршрут -Порядок линейный — уроки 1→4 повторяют сам пайплайн (STG → ODS → DDS → оркестрация), -а урок 0 — разминка перед ним. Мониторинг (5) и Superset (6) работают как надстройки -поверх готовых данных. Лабы 07–08 затем растят импортированный мир: сначала -детерминированным следующим днём, потом живым продолжением. +Порядок линейный, все уроки и лабы обязательны. Урок 0 — разминка, а уроки 1→4 +повторяют сам пайплайн (STG → ODS → DDS → оркестрация). Мониторинг (5) и +Superset (6) работают как надстройки поверх готовых данных. Лабы 07–08 затем +растят импортированный мир: сначала детерминированным следующим днём, потом +живым продолжением. Принцип нарезки — **один прод-паттерн на урок** (`LESSON_STANDARD` §2). Поэтому середина пайплайна разнесена: типизация+DQ (ODS) и сборка сущностей (DDS) — разные @@ -40,12 +41,12 @@ Kafka из роадмапа — оно даёт словарь терминов. | 4 | Оркестрация в Airflow | `airflow/dags/etl_pipeline_dag.py` (зависимости, гейты, остановка при нарушениях) | обязательный | руки | | 5 | Мониторинг | Prometheus + Grafana + экспортёры | обязательный | наблюдение | | 6 | BI-витрина | Superset поверх ClickHouse | обязательный | руки | -| 7 | Лаба: следующий день | `airflow/dags/world_next_day_dag.py` | в работе, обязательный | руки | -| 8 | Лаба: живое продолжение | `make generator-continue` и `etl_pipeline` | в работе, обязательный | руки | +| 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 @@ -146,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/README.md b/docs/course/README.md index 8cfbc51..18e4a3b 100644 --- a/docs/course/README.md +++ b/docs/course/README.md @@ -52,8 +52,9 @@ ## Уроки Проходи по порядку — уроки 0→4 повторяют сам пайплайн (Kafka → STG → ODS → DDS → -оркестрация), а 5 и 6 надстраиваются поверх готовых данных. В колонке «режим»: -**наблюдение** — только смотрим, **руки** — запускаешь и меняешь сам. +оркестрация), а 5 и 6 надстраиваются поверх готовых данных. Лабы 7 и 8 затем +растят эталонный мир двумя способами. В колонке «режим»: **наблюдение** — только +смотрим, **руки** — запускаешь и меняешь сам. | # | Урок | Режим | О чём | |---|------|-------|-------| @@ -64,8 +65,8 @@ | 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, динамика, воронка | -| 7 | Лаба: следующий день | руки, в работе | Пакетный инкремент дня и границы времени | -| 8 | Лаба: живое продолжение | руки, в работе | Живой поток и свежесть данных | +| 7 | [Лаба: следующий день](./lessons/07_lab_next_day.md) | руки | Пакетный инкремент дня и границы времени | +| 8 | [Лаба: живое продолжение](./lessons/08_lab_continue.md) | руки | Живой поток и свежесть данных | ## Как проходить 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..44e3264 --- /dev/null +++ b/docs/course/lessons/07_lab_next_day.md @@ -0,0 +1,293 @@ +# Лаба 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` | `<события-день-4>` | +| `visits` | `<визиты-день-4>` | +| `users` | `<пользователи-день-4>` | +| `min_event_timestamp` | `<минимальное-время>` | +| `max_event_timestamp` | `<максимальное-время>` | +| `model_t0` | `<левая-граница>` | +| `model_t_end` | `<правая-граница-дня-4>` | +| `profile` | `<профиль>` | + +Здесь числа накопительные: `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; +``` + +Запрос должен найти хотя бы один визит, чьи события лежат по разные стороны +полуночи. Подойдёт любой стык внутри мира: 1→2, 2→3 или 3→4. + +> **Как слова связаны со схемой.** В manifest написано «визит», а в SQL такой +> визит обозначен `click_id`. Отдельной таблицы визитов нет: контекст лежит в +> `dds.click`, а события визита — в `dds.event`. + +Почему запрос идёт прямо в `dds.event`? Представление +`dm.v_session_overview` группирует строки по `event_date` и `click_id`. +Переходящий визит там уже разрезан на две строки: одна до полуночи, другая +после. Чтобы увидеть визит целиком, сначала собираем все его события по +`click_id`, а потом сравниваем даты самого раннего и самого позднего события. + +### Смотрим день-к-дню + +Открой в Superset дашборд `Clickstream Analytics`. Поставь фильтр +**Date Range → No filter** и найди график `Events over Time`. После обновления +на нём должен появиться день 4. Точную высоту точки не угадывай: она должна +сойтись с данными ClickHouse и 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 не фиксируем: это результат учебной правки. + + +### Верни как было + +Сначала убедись, что `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` | хотя бы один `click_id` пересекает полночь | +| снять паузу с 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`; +- вывод `make generated-history-chain-check` без ошибки; +- результат SQL-запроса с переходящим через полночь `click_id`; +- скрин дневного графика с добавленным днём; +- короткое объяснение своими словами: чем пакетный инкремент отличается от + полного пересчёта и зачем отдельно проверять границу порций. + +Перед переходом дальше стенд должен быть возвращён к трёхдневному эталонному +миру, а `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..dd35157 --- /dev/null +++ b/docs/course/lessons/08_lab_continue.md @@ -0,0 +1,299 @@ +# Лаба 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, 'UTC' + )) 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; +``` + +На эталонном мире правые границы слоёв должны совпадать. + + +### Запускаем живое продолжение + +В терминале выполни: + +```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`. + +### Догоняем пакетные слои + +Открой Airflow и запусти `etl_pipeline` через **Trigger DAG** с пустой формой. +По умолчанию это полный пересчёт. + +После зелёного прогона повтори запрос. ODS, DDS и 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 +SELECT + uniqExact(user_domain_id) AS users, + uniqExact(click_id) AS sessions +FROM dm.v_events_enriched +WHERE user_domain_id IS NOT NULL; +``` + +В живом мире один пользователь может вернуться и открыть новый визит. Поэтому +после достаточного продолжения ожидаем `users < sessions`. На эталонной +статике было `users == sessions`. + +Точные числа здесь не печатаем: длительность живого запуска у каждого своя. +`users = <значение>`, `sessions = <значение>`. + + +Если равенство ещё сохранилось, дай генератору поработать несколько минут, +снова запусти `etl_pipeline` и повтори запрос. Важно увидеть сам возврат +пользователя, а не угадать конкретное число. + +### Верни как было + +Останови живой поток: + +```bash +make generator-down +``` + +После `continue` мир стал недетерминированным: его числа зависят от длительности +запуска. Верни стенд к эталону по +[канонической инструкции курса](../README.md#подготовка-и-канонический-сброс). +Не пересказывай шаги по памяти: эта ссылка остаётся единственным учебным +описанием полного сброса. + +--- + +## 5. Проверь себя + +| Действие | Где смотреть | Что ожидать | +|----------|--------------|-------------| +| выполнить `make generator-continue` | Prometheus Targets | `generator` переходит в `UP` | +| посмотреть живой генератор | `Generator Overview` | `Total Events/min (all 4 topics)` становится ненулевым | +| сравнить свежесть до ETL | SQL-запрос по слоям | STG свежее ODS, DDS и DM | +| запустить `etl_pipeline` | Airflow и SQL-запрос | пакетные слои догоняют снимок STG | +| выполнить `make generator-down` | Grafana и Kafka | события перестают поступать, отставание читателей стекает к нулю | +| снова выполнить `make generator-continue` | Grafana, Kafka и STG | поток продолжается, offset-ы и правая граница снова растут | +| пересчитать аналитику | запрос `users` и `sessions` | после возвратов пользователей выполняется `users < sessions` | +| остановить поток и пройти сброс | Grafana и SQL | `generator` выключен, слои снова совпадают с эталонным миром | + +Ответь своими словами: + +- что означает свежесть данных и кто задаёт требование к ней; +- почему новые строки уже есть в STG, но их ещё нет в DDS; +- почему запуск `etl_pipeline` догоняет поток лишь до очередного снимка; +- что сохраняет режим `continue`; +- почему после живого продолжения числа разных менти расходятся; +- откуда берётся `users < sessions`. + +--- + +## 6. Что должно получиться + +После лабы сохрани: + +- скрин Prometheus Targets с `generator` в `UP`; +- скрин живого `Total Events/min (all 4 topics)`; +- два результата запроса свежести: до и после `etl_pipeline`; +- наблюдение остановки и продолжения по offset-ам или правой границе STG; +- результат запроса с `users < sessions`; +- короткое объяснение своими словами: почему потоковый приём не гарантирует + потоковую обработку до витрины. + +В конце `generator` должен быть остановлен, а стенд — возвращён к эталонному +миру. + +--- + +## Мост после курса + +Теперь у тебя есть два способа растить один мир: повторяемая суточная порция и +живое продолжение. Следующий практический вопрос уже зависит от продукта: +какую свежесть обещать потребителю и какую цену платить за пересчёт.