Files
ddadminandClaude Opus 4.8 8f03d8b347 docs(course): маркеры «сверить-на-стенде» заменены живыми числами (#23)
Зачем: тикет #23 — каждое цитируемое число курса сверено с живым
стендом; до этого в уроках стояли заглушки.

Что:
- уроки 00–04: 18 маркеров заполнены числами свежего импорта
  (280 437 событий, нули в таблицах ошибок, счётчики ODS/DDS);
- лаба 07: таблица manifest после дня 4 (374 092 / 34 801 / 5 388),
  переходящие визиты по стыкам (35/22/26), числа после дня 5;
- лаба 08: каноническая граница трёх дней, пример замера свежести
  (лаг 3:45 модельного времени до догона, ETL ~29 с) и вернувшегося
  пользователя; две живые поправки разбора времени: убран
  принудительный UTC в разборе STG и суффикс +00:00 в сравнении
  границы (ловились только на живом стенде).

Проверка: детерминизм подтверждён двумя независимыми циклами
сброс→импорт→инкремент (числа manifest и checksum_sha256 дней 4 и 5
совпали бит в бит); chain-check зелёный; сценарий лабы 08 прогнан
вживую, включая стоп/продолжение и красный full_refresh=false из
урока 4; grep «сверить-на-стенде» пуст; ссылки и якоря целы;
make test (219+31) и make lint зелёные; /ai-text-lint по лабам — без
существенных находок.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-23 18:15:46 +03:00

298 lines
17 KiB
Markdown
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# Лаба 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` — стоять на паузе.
---
## Мост к следующему шагу
Здесь новый день приехал как законченная и повторяемая порция. В следующей
лабе граница исчезнет: события пойдут непрерывно, а разные слои начнут
отставать друг от друга.