# Лаба 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` — стоять на паузе. --- ## Мост к следующему шагу Здесь новый день приехал как законченная и повторяемая порция. В следующей лабе граница исчезнет: события пойдут непрерывно, а разные слои начнут отставать друг от друга.