From d76c03655bea9029ef7006fc75a5de92e2bc74c9 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sat, 4 Jul 2026 22:22:36 +0300 Subject: [PATCH] =?UTF-8?q?docs(course):=20=D0=BF=D0=B5=D1=80=D0=B5=D0=B2?= =?UTF-8?q?=D0=B5=D0=B4=D0=B5=D0=BD=D1=8B=20=D1=83=D1=80=D0=BE=D0=BA=D0=B8?= =?UTF-8?q?=20=D0=BD=D0=B0=20=D1=81=D1=82=D0=B0=D1=80=D1=82=D0=BE=D0=B2?= =?UTF-8?q?=D1=83=D1=8E=20=D0=B8=D1=81=D1=82=D0=BE=D1=80=D0=B8=D1=8E?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - учебный путь должен идти через генерацию и штатный пайплайн, а не через архивный сид. - Что: - обновлены уроки 00-06 и стандарт урока под startup-history/backfill. - объяснено, что data/*.jsonl остаются кладовкой значений генератора. - тест-план переведён на новый штатный запуск и HITL-приёмку. - Проверка: - rg -n \"make clean/up/ddl/data/transform|make data|kafka_load|LIMIT=|2022-11-28|26 из 50\" docs/course docs/TEST_PLAN.md. - git diff --cached --check. --- docs/TEST_PLAN.md | 62 +++++-------- docs/course/LESSON_STANDARD.md | 8 +- docs/course/README.md | 17 ++-- docs/course/lessons/00_kafka_intro.md | 34 +++++--- docs/course/lessons/01_kafka_to_clickhouse.md | 23 ++--- docs/course/lessons/02_stg_to_ods.md | 87 +++++++++---------- docs/course/lessons/03_ods_to_dds.md | 60 ++++++------- .../lessons/04_airflow_orchestration.md | 23 ++--- docs/course/lessons/05_monitoring.md | 11 ++- docs/course/lessons/06_superset_bi.md | 68 ++++++--------- 10 files changed, 192 insertions(+), 201 deletions(-) diff --git a/docs/TEST_PLAN.md b/docs/TEST_PLAN.md index 3a9ea5d..37f3e9e 100644 --- a/docs/TEST_PLAN.md +++ b/docs/TEST_PLAN.md @@ -8,15 +8,17 @@ Проверить, что стек `Kafka + ClickHouse + Airflow + Superset + Prometheus/Grafana`: - стабильно поднимается; -- загружает данные по пути `kafka_load -> STG -> etl_pipeline -> ODS/DDS/DM`; +- загружает данные по пути `startup-history/backfill -> Kafka -> STG -> ODS/DDS/DM`; - корректно обрабатывает «грязные» записи (ошибки фиксируются в ODS, пайплайн не падает); - отдает метрики и дашборды мониторинга. ## Общие принципы -- По умолчанию используем малый срез (`limit=50`) для быстрых и повторяемых проверок. -- Полный прогон (`limit=0`) выполняем отдельно как long-run сценарий. -- Основной путь запуска — через Airflow DAG. +- По умолчанию используем быстрый профиль стартовой истории (`PROFILE=ci`). +- Полный прогон выполняем отдельно через `PROFILE=daily-wave`. +- Основной ручной путь запуска — через Airflow DAG `generator_control`. +- Консольный чистый прогон `make generated-history-analytics` остаётся коротким + повторяемым сценарием для smoke и CI. - Критерий успеха: не только `Success` DAG, но и проверки данных/ошибок/мониторинга. --- @@ -28,10 +30,10 @@ ### A.1 Подготовка окружения ```bash -# Полная очистка стенда -make clean +# Чистый быстрый прогон: backfill -> Kafka -> STG -> ODS -> DDS -> DM -> Superset +make generated-history-analytics -# Запуск инфраструктуры +# Поднять остальные UI-сервисы после чистого прогона make up # Проверка контейнеров @@ -42,16 +44,7 @@ docker compose ps - `airflow-init` в `Exited (0)`; - остальные сервисы в `Up` (включая `superset`, `prometheus`, `grafana`, `kafka-exporter`, `statsd-exporter`). -### A.2 DDL и минимальная загрузка данных - -```bash -# Инициализация схемы -docker compose exec -T airflow-webserver airflow dags trigger ddl_init - -# Быстрый ingest: по 50 строк на поток -docker compose exec -T airflow-webserver airflow dags trigger kafka_load \ - --conf '{"limit": 50, "reset_topics": true}' -``` +### A.2 Проверка стартовой истории в STG Проверки: @@ -72,12 +65,7 @@ SELECT 'geo_raw', count() FROM stg.geo_raw Ожидаем: во всех 4 таблицах `cnt > 0`. -### A.3 ETL и проверки слоев - -```bash -docker compose exec -T airflow-webserver airflow dags trigger etl_pipeline \ - --conf '{"full_refresh": true}' -``` +### A.3 Проверки слоев Проверки: @@ -123,7 +111,7 @@ curl -s -u admin:admin "http://localhost:3000/api/dashboards/uid/airflow-overvie Критерий успеха smoke: - сервисы подняты; -- 3 DAG (`ddl_init`, `kafka_load`, `etl_pipeline`) успешны; +- стартовая история доведена до DM; - данные проходят до DM; - мониторинг и дашборды доступны. @@ -133,29 +121,22 @@ curl -s -u admin:admin "http://localhost:3000/api/dashboards/uid/airflow-overvie Ожидаемая длительность: ~30-60 минут. -### B.1 Полная загрузка и полный ETL +### B.1 Полная стартовая история и ETL ```bash -# Чистый старт -make clean && make up - -# DDL -docker compose exec -T airflow-webserver airflow dags trigger ddl_init - -# Полный ingest (все строки) -docker compose exec -T airflow-webserver airflow dags trigger kafka_load \ - --conf '{"limit": 0, "reset_topics": true}' - -# Полный ETL -docker compose exec -T airflow-webserver airflow dags trigger etl_pipeline \ - --conf '{"full_refresh": true}' +# История на 2 суток с видимой суточной волной +PROFILE=daily-wave make generated-history-analytics +make up ``` ### B.2 Проверка объемов и DQ ```bash -# Фактические размеры входа -wc -l data/*.jsonl +# Диапазон модельного времени в витрине +docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query " +SELECT min(event_ts) AS min_event_ts, max(event_ts) AS max_event_ts, count() AS events +FROM dm.v_events_enriched +" # Сводка по слоям docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query " @@ -185,6 +166,7 @@ docker compose exec -T clickhouse clickhouse-client --user=default --password=12 Критерии успеха full: - STG заполнен по всем 4 потокам; - ODS/DDS/DM заполнены; +- диапазон `event_ts` соответствует созданной стартовой истории; - `dm.dq_summary` не пуста и содержит метрики всех слоев. ### B.3 Тест восстановления stop/start diff --git a/docs/course/LESSON_STANDARD.md b/docs/course/LESSON_STANDARD.md index 04f0971..1f46624 100644 --- a/docs/course/LESSON_STANDARD.md +++ b/docs/course/LESSON_STANDARD.md @@ -21,7 +21,7 @@ запись и увидеть `+1` в таблице ошибок `parse_errors`; сменить `kafka_group_name` и увидеть, как топик читается заново. Менти меняет — видит эффект — объясняет. **Верни как было** — каждая правка завершается явным шагом отката к чистому - состоянию (откатить изменение либо `make clean/up/ddl/data/transform`), чтобы + состоянию (откатить изменение либо `make generated-history-analytics && make up`), чтобы самостоятельный менти не застрял со сломанным стендом без ментора. (Урок 0 — без этого шага, только наблюдение.) 5. **Проверь себя** — самопроверка (раздел 4). @@ -112,7 +112,7 @@ - штатные быстрые проверки (smoke) из `docs/TEST_PLAN.md`; - встроенные проверки в `etl_pipeline_dag.py` (DAG падает на пустой витрине или нарушении целостности); -- `make clean/up/ddl/data/transform` для сброса и повтора. +- `make generated-history-analytics && make up` для чистого сброса и повтора штатного пути. В каждом уроке — маленькая табличка самопроверки в формате **действие → где смотреть → что ожидать**, а для управляемой правки — какой @@ -125,8 +125,8 @@ - **Проверяй на стенде, а не «на глаз».** Любой категоричный claim про числа или целостность подтверждай командой на стенде. Так нашёлся баг с `kafka_ts` в уроке 1 (значение типа `DateTime64` молча резалось через `toInt64`, и время по всему стенду уехало в `1970`) и - подтвердилось, что «26 из 50» в уроке 2 — это схлопывание повторов по `click_id` (проверено - `uniqExact`), а не потеря данных. + подтвердилось, что меньший размер таблиц контекста в уроке 2 — это схлопывание повторов + по `click_id` (проверяется через `uniqExact`), а не потеря данных. - **Имена таблиц и колонок в тексте сверяй с реальным DDL.** То, что написано в уроке, должно совпадать с `sql/ddl/...`, иначе менти запутается, когда выполнит запрос на стенде. - **Один паттерн на урок; будущие темы только анонсируй.** Не раскрывай то, что относится к diff --git a/docs/course/README.md b/docs/course/README.md index bfaa579..ff1b11c 100644 --- a/docs/course/README.md +++ b/docs/course/README.md @@ -25,16 +25,21 @@ слова с живым стендом. - **Железо:** стек тяжёлый — Kafka, ClickHouse, Airflow, Superset, Prometheus и Grafana поднимаются одновременно. Нужна машина, которая это потянет. -- **Подними стенд и залей малый срез** (из корня репозитория) — этого хватит, чтобы - начать, и прогон быстрый: +- **Подними стенд и создай стартовую историю** (из корня репозитория) — этого хватит, + чтобы начать, и прогон быстрый: ```bash - make up # поднять контейнеры - make ddl # создать схему в ClickHouse (базы, таблицы, VIEW) - LIMIT=50 make data # залить по 50 строк на топик в Kafka + make generated-history-analytics + make up ``` - Дальше каждый урок в секции «Руки» сам напоминает, что перезапустить и с каким срезом. + Эта команда проводит штатный путь стенда: готовый источник данных создаёт стартовую + историю, события попадают в Kafka, затем в STG, ODS, DDS, DM и Superset. Файлы + `data/*.jsonl` пока остаются только кладовкой готовых значений для генератора + (браузеры, страны, устройства, UTM), а не источником аналитического контура. `make up` + после неё поднимает остальные UI-сервисы курса: Kafka UI, Airflow, Prometheus и Grafana. + + Дальше каждый урок в секции «Руки» сам напоминает, что перезапустить. Точные шаги, параметры и troubleshooting — в [`docs/OPERATIONS.md`](../OPERATIONS.md). ## Уроки diff --git a/docs/course/lessons/00_kafka_intro.md b/docs/course/lessons/00_kafka_intro.md index b24d436..3b450a6 100644 --- a/docs/course/lessons/00_kafka_intro.md +++ b/docs/course/lessons/00_kafka_intro.md @@ -7,7 +7,7 @@ > Эталонный путь: сам Kafka UI на `http://localhost:8082`. > > Поток данных одной строкой: -> `make data → топики Kafka (партиции, offset'ы) → consumer-группа ClickHouse вычитывает` +> `готовый источник данных стенда → топики Kafka (партиции, offset'ы) → consumer-группа ClickHouse вычитывает` > > О чём урок простыми словами: ходим по Kafka UI и разглядываем поток — где лежат события, > кто их читает и как Kafka помнит, до какого места уже дочитано. @@ -22,11 +22,15 @@ consumer читает в своём темпе, не трогая producer'а. Если читатель отстал или прилёг — события не теряются, они лежат в топике и ждут. -На нашем стенде роли уже расставлены: `make data` играет роль трекера и заливает события -в топики, а ClickHouse — это consumer, который их вычитывает. В этом уроке мы не запускаем -пайплайн и ничего не меняем — мы открываем Kafka UI и **узнаём в живом кластере** те самые -понятия из видео. Это разминка: в уроке 1 ты уже руками увидишь, как эти же сообщения -доезжают до таблиц ClickHouse. +На нашем стенде роли уже расставлены: готовый источник данных стенда играет роль трекера +и пишет события в топики, а ClickHouse — это consumer, который их вычитывает. Для курса нам +важно не устройство источника, а сам путь данных: Kafka → STG → ODS → DDS → DM → Superset. +Файлы `data/*.jsonl` здесь не источник аналитики; пока это только кладовка готовых значений +для источника данных стенда. + +В этом уроке мы не запускаем пайплайн и ничего не меняем — мы открываем Kafka UI и +**узнаём в живом кластере** те самые понятия из видео. Это разминка: в уроке 1 ты уже руками +увидишь, как эти же сообщения доезжают до таблиц ClickHouse. > **В проде так же, только крупнее.** Здесь один брокер и крошечный срез данных. В бою > брокеров несколько, топик разбит на много партиций, читателей в группе — тоже несколько, @@ -37,8 +41,14 @@ consumer читает в своём темпе, не трогая producer'а. ## 2. Наблюдай: открой Kafka UI -Стенд уже должен быть поднят (`make up`) и в топиках должны лежать события -(`LIMIT=50 make data` из урока 1 — или любой прошлый прогон). Открой Kafka UI: +Стенд уже должен быть поднят, а в топиках должны лежать события после стартовой истории: + +```bash +make generated-history-analytics +make up +``` + +Открой Kafka UI: `http://localhost:8082`. Ходи по нему свободно — это режим чтения, сломать тут ничего нельзя. Пройди по трём экранам и просто посмотри. @@ -54,20 +64,20 @@ consumer читает в своём темпе, не трогая producer'а. топиков) — его Kafka использует сама, мы его не трогаем. У каждого нашего топика в колонке с партициями стоит **1**: топик маленький, делить не на что. -**Сообщения в топике.** Открой `browser_events` → вкладку *Messages*. Это и есть события, -которые залил `make data`. У каждого сообщения видно: +**Сообщения в топике.** Открой `browser_events` → вкладку *Messages*. Это и есть события +стенда. У каждого сообщения видно: - **Offset** — порядковый номер сообщения в партиции (0, 1, 2, …); - **Timestamp** — когда сообщение легло в Kafka (время *доставки*, не время самого события); - **Value** — тело: JSON события целиком, например: ```json -{"event_id": "8cca1c7d-...", "event_timestamp": "2022-11-28 20:51:05.627882", +{"event_id": "8cca1c7d-...", "event_timestamp": "2026-01-01 00:01:00.000000", "event_type": "pageview", "browser_name": "Chrome", "browser_language": "sat_IN"} ``` Загляни внутрь Value: у события есть своё `event_timestamp` (когда оно случилось, -здесь — 2022 год), и оно отличается от 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 2c12e1a..4d31d81 100644 --- a/docs/course/lessons/01_kafka_to_clickhouse.md +++ b/docs/course/lessons/01_kafka_to_clickhouse.md @@ -49,13 +49,14 @@ ## 2. Руки: убедись, что данные текут -Поднимаем стенд, создаём схему и заливаем **малый срез** (50 строк на топик — этого -хватает, чтобы всё увидеть, и прогон быстрый): +Поднимаем стенд и создаём стартовую историю. Это штатный путь курса: готовый источник +данных стенда пишет события в Kafka, ClickHouse читает их в STG, затем batch строит +ODS, DDS и DM. Файлы `data/*.jsonl` в этом пути не источник аналитики; пока это только +кладовка значений для источника данных стенда. ```bash -make up # поднять инфраструктуру (Kafka, ClickHouse, ...) -make ddl # создать базы и таблицы в ClickHouse (в т.ч. слой STG) -LIMIT=50 make data # залить по 50 строк каждого файла в Kafka-топики +make generated-history-analytics +make up ``` Теперь смотрим, что доехало до ClickHouse. Открой SQL-консоль: @@ -196,11 +197,12 @@ FROM stg.kafka_browser_raw; TRUNCATE TABLE stg.browser_raw; ``` -Затем перезаливаем — с пересозданием топиков, чтобы Kafka-движок перечитал сообщения -с начала (без этого он считает их уже прочитанными и ничего нового не подхватит): +Затем возвращаем стенд в чистое состояние и заново создаём стартовую историю, чтобы +Kafka-движок прочитал сообщения уже с новой схемой: ```bash -LIMIT=50 RESET_TOPICS=1 make data +make generated-history-analytics +make up ``` **Смотрим результат:** @@ -233,7 +235,8 @@ ALTER TABLE stg.browser_raw DROP COLUMN kafka_msg_ts; make ddl # пересоздаёт эталонный MV из 10_stg.sql — схема снова как в репозитории ``` -Если запутался в состоянии — всегда есть полный сброс: `make clean && make up && make ddl && LIMIT=50 make data`. +Если запутался в состоянии — всегда есть полный чистый прогон: +`make generated-history-analytics && make up`. --- @@ -241,7 +244,7 @@ make ddl # пересоздаёт эталонный MV из 10_stg.sql — | Действие | Где смотреть | Что ожидать | |----------|--------------|-------------| -| `LIMIT=50 make data` | `SELECT count() FROM stg.browser_raw` | счётчик > 0 и растёт после загрузки | +| `make generated-history-analytics && make up` | `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` | колонка заполнена временем сообщения | diff --git a/docs/course/lessons/02_stg_to_ods.md b/docs/course/lessons/02_stg_to_ods.md index 8265017..6c75c39 100644 --- a/docs/course/lessons/02_stg_to_ods.md +++ b/docs/course/lessons/02_stg_to_ods.md @@ -84,27 +84,27 @@ ODS мы пересобираем целиком, одной задачей Airf ## 2. Руки: смотрим базовый прогон -Поднимаем стенд, создаём схему, заливаем **малый срез** (50 строк на топик) и запускаем -трансформацию: +Поднимаем стенд и создаём стартовую историю. Это штатный путь курса: готовый источник +данных стенда пишет события в Kafka, ClickHouse читает их в STG, затем batch строит +ODS, DDS и DM. Файлы `data/*.jsonl` пока остаются только кладовкой значений для этого +источника, а не источником аналитического контура. ```bash -make up # поднять инфраструктуру -make ddl # создать базы и таблицы (в т.ч. слой ODS) -LIMIT=50 make data # залить по 50 строк каждого файла в Kafka → STG -make transform # батч STG → ODS → DDS → DM (нас интересует первый шаг) +make generated-history-analytics +make up ``` -`make transform` прогоняет всю цепочку слоёв сразу, но прямо в консоли печатает то, что нам -нужно сейчас, — блок **«Статистика ODS»**. Это просто счётчики строк по всем восьми таблицам -слоя (четыре основных и четыре с ошибками): +Команда прогоняет всю цепочку слоёв и прямо в консоли печатает то, что нам нужно сейчас, — +блок **«Статистика ODS»**. Это просто счётчики строк по всем восьми таблицам слоя (четыре +основных и четыре с ошибками). Пример формы вывода: ``` Статистика ODS: ┌─table──────────────────────┬─rows─┐ - │ ods.browser_event │ 50 │ - │ ods.location_event │ 50 │ - │ ods.device_by_click │ 26 │ - │ ods.geo_by_click │ 26 │ + │ 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 │ @@ -114,35 +114,34 @@ make transform # батч STG → ODS → DDS → DM (нас инте Прочитаем эту табличку — в ней три вещи, которые стоит заметить. -**Все четыре `*_errors` — по нулям.** Значит, наш срез чистый: ни одна запись не дала ошибки -разбора, столбец `parse_errors` у всех пустой. Это нормально — данные в демо аккуратные. +**Все четыре `*_errors` — по нулям.** Значит, стартовая история чистая: ни одна запись не дала +ошибки разбора, столбец `parse_errors` у всех пустой. Это нормально — данные стенда аккуратные. Ошибки мы увидим в секции 4, когда сами их устроим. -**`browser` и `location` дали 50 из 50.** Сколько событий пришло — столько и легло, один к -одному. +**`browser` и `location` идут в одном зерне события.** Сколько событий пришло, столько строк +и ожидаем увидеть после типизации, если ключи валидны. -**А `device` и `geo` — только 26 из 50.** Вот это уже интересно. Половина куда-то делась? Нет. -И это важно понять, иначе дальше будет казаться, что данные текут. +**А `device` и `geo` обычно меньше, чем событий.** Вот это уже интересно. Часть строк +куда-то делась? Нет. И это важно понять, иначе дальше будет казаться, что данные текут. Дело в том, что эти две таблицы хранят не события, а **контекст клика**: с какого устройства -был клик и из какой точки на карте. Ключ у них — `click_id`. А в срезе на 50 событий разных -кликов всего 26: на один клик приходится несколько событий, и `click_id` у них повторяется. -Движок таблицы (про него — в секции 3) схлопывает повторы по ключу, оставляя по одной строке -на клик. Отсюда и 26. +был клик и из какой точки на карте. Ключ у них — `click_id`. Разных кликов меньше, чем событий: +на один клик приходится несколько событий, и `click_id` у них повторяется. Движок таблицы +(про него — в секции 3) схлопывает повторы по ключу, оставляя по одной строке на клик. Проверь это сам, а не верь на слово. Открой SQL-консоль `http://localhost:9123/play` -(пользователь `default`, пароль `123456`) и посчитай, сколько в срезе *различных* `click_id`: +(пользователь `default`, пароль `123456`) и посчитай, сколько в STG *различных* `click_id`: ```sql --- Всего строк в STG — 50, но различных click_id среди них — ровно 26 +-- Строк событий больше, чем различных click_id SELECT count() AS stg_rows, uniqExact(toUUIDOrNull(JSONExtractString(raw, 'click_id'))) AS distinct_clicks FROM stg.geo_raw; ``` -Получишь `stg_rows = 50`, `distinct_clicks = 26` — ровно столько, сколько строк в -`ods.geo_by_click`. Значит, 26 — это схлопнутые повторы, а не пропавшие данные. Ничего не -потерялось молча. +`distinct_clicks` должен быть меньше или равен `stg_rows` и совпадать с числом строк в +`ods.geo_by_click`. Значит, это схлопнутые повторы по `click_id`, а не пропавшие данные. +Ничего не потерялось молча. --- @@ -174,7 +173,7 @@ toFloat64OrNull(JSONExtractString(raw, 'geo_latitude')) AS Весь смысл — в суффиксе `OrNull`. Если значение **не** приводится к нужному типу (вместо UUID пришёл мусор), функция не падает с ошибкой, а просто возвращает `NULL`. Это ровно то правило стенда, что и в STG — «грязная запись не валит пайплайн», — только теперь на уровне типов. -Один кривой `event_id` станет `NULL` и будет помечен, а остальные 49 строк спокойно доедут. +Один кривой `event_id` станет `NULL` и будет помечен, а остальные строки спокойно доедут. > Кстати, про `AS`: эти строки живут в блоке `WITH` в начале запроса. `WITH` — это просто > способ заранее посчитать значение и дать ему имя, чтобы ниже по запросу ссылаться на него @@ -230,7 +229,7 @@ arrayFilter(x -> x != '', [ > если поменять разбор только в одном из двух мест, они разойдутся. В секции 4 мы как раз этим > воспользуемся — и увидим, чем грозит такой рассинхрон. -### Движок: откуда взялись 26 строк +### Движок: почему строк контекста меньше И последнее место — строчка про движок основных таблиц: @@ -241,8 +240,8 @@ ORDER BY (click_id) `ReplacingMergeTree` — это таблица, которая схлопывает строки с одинаковым ключом (ключ берётся из `ORDER BY`), оставляя самую свежую по `src_ingest_ts` — времени загрузки в ODS. Вот она, -причина «26 из 50» из секции 2: у `device` и `geo` много строк с одинаковым `click_id`, и -движок оставляет по одной на клик. +причина разницы из секции 2: у `device` и `geo` много строк с одинаковым `click_id`, и движок +оставляет по одной на клик. --- @@ -276,8 +275,8 @@ make transform И смотрим на ту же «Статистику ODS». Таблица ошибок гео, которая была пустой, теперь полная: ``` - │ ods.geo_by_click │ 26 │ - │ ods.geo_by_click_errors │ 50 │ ← было 0 + │ ods.geo_by_click │ ... │ + │ ods.geo_by_click_errors │ ... │ ← было 0 ``` А в самой основной таблице широта пропала — но не молча, рядом стоит метка: @@ -295,11 +294,10 @@ LIMIT 4; └──────────────┴──────────────┴───────────────┴──────────────────────┘ ``` -Вот теперь видно всё разом — и DQ-split, и «двойной учёт» из секции 3 вживую. 26 строк -остались в основной таблице (ключ `click_id` цел) с пометкой `bad_geo_latitude`. И те же -записи попали в число 50 строк `geo_by_click_errors`. Долгота на месте, а широты больше нет: -один неверный тип — и целое поле потеряно по всему слою. Заметили это `parse_errors` и таблица -ошибок — для того DQ-split и нужен. +Вот теперь видно всё разом — и DQ-split, и «двойной учёт» из секции 3 вживую. Строки с валидным +`click_id` остались в основной таблице с пометкой `bad_geo_latitude`. И те же записи попали в +`geo_by_click_errors`. Долгота на месте, а широты больше нет: один неверный тип — и целое поле +потеряно по всему слою. Заметили это `parse_errors` и таблица ошибок — для того DQ-split и нужен. > **Бывает и хуже — тихо, совсем без метки.** Здесь нас спас суффикс `OrNull`: неверный тип > дал `NULL`, а `NULL` мы умеем замечать (на него и сработал `parse_errors`). По-настоящему @@ -319,7 +317,7 @@ make transform ``` После этого `geo_by_click_errors` снова `0`, широта на месте. А если стенд совсем «поплыл» — -всегда есть полный сброс: `make clean && make up && make ddl && LIMIT=50 make data && make transform`. +всегда есть полный чистый прогон: `make generated-history-analytics && make up`. --- @@ -327,9 +325,9 @@ make transform | Действие | Где смотреть | Что ожидать | |----------|--------------|-------------| -| `make transform` (базовый прогон) | блок «Статистика ODS» | `browser`/`location` = 50, `device`/`geo` = 26, все `*_errors` = 0 | -| почему 26, а не 50 | запрос `uniqExact(click_id)` по `stg.geo_raw` | 26 различных `click_id` — это схлопывание повторов, а не потеря | -| правка из секции 4 | блок «Статистика ODS» | `ods.geo_by_click_errors` прыгнул `0 → 50` | +| базовый прогон | блок «Статистика ODS» | основные таблицы не пустые, все `*_errors` = 0 | +| почему `device`/`geo` меньше событий | запрос `uniqExact(click_id)` по `stg.geo_raw` | число различных `click_id` совпадает с `ods.geo_by_click` | +| правка из секции 4 | блок «Статистика ODS» | `ods.geo_by_click_errors` прыгнул с `0` на ненулевое число | | правка из секции 4 | `SELECT geo_latitude, parse_errors FROM ods.geo_by_click` | широта `NULL`, в `parse_errors` — `bad_geo_latitude` | --- @@ -338,7 +336,8 @@ make transform После урока у тебя на руках — видимый результат (одно на выбор): -- скрин блока «Статистика ODS», где после правки `ods.geo_by_click_errors` ушёл с `0` на `50`; +- скрин блока «Статистика ODS», где после правки `ods.geo_by_click_errors` ушёл с `0` + на ненулевое число; - либо выборка из `ods.geo_by_click` с пустой широтой и меткой `bad_geo_latitude` рядом. И проверь себя на словах — примерно эти вопросы всплывут на еженедельном созвоне: diff --git a/docs/course/lessons/03_ods_to_dds.md b/docs/course/lessons/03_ods_to_dds.md index 86ff12b..495c10e 100644 --- a/docs/course/lessons/03_ods_to_dds.md +++ b/docs/course/lessons/03_ods_to_dds.md @@ -65,35 +65,35 @@ ## 2. Руки: смотрим базовый прогон -Поднимаем стенд, создаём схему, заливаем **малый срез** (50 строк на топик) и запускаем -трансформацию: +Поднимаем стенд и создаём стартовую историю. Это штатный путь курса: готовый источник +данных стенда пишет события в Kafka, ClickHouse читает их в STG, затем batch строит +ODS, DDS и DM. Файлы `data/*.jsonl` пока остаются только кладовкой значений для этого +источника, а не источником аналитического контура. ```bash -make up # поднять инфраструктуру -make ddl # создать базы и таблицы (в т.ч. слой DDS) -LIMIT=50 make data # залить по 50 строк каждого файла в Kafka → STG -make transform # батч STG → ODS → DDS → DM +make generated-history-analytics +make up ``` -`make transform` прогоняет всю цепочку слоёв и по дороге печатает в консоль блок **«Статистика -DDS»** — счётчики строк по двум нашим сущностям: +Команда прогоняет всю цепочку слоёв и по дороге печатает в консоль блок **«Статистика DDS»** — +счётчики строк по двум нашим сущностям: ``` Статистика DDS: ┌─table─────┬─rows─┐ - │ dds.click │ 26 │ - │ dds.event │ 50 │ + │ dds.click │ ... │ + │ dds.event │ ... │ └───────────┴──────┘ ``` Прочитаем эти две строки. -**`dds.event` — 50.** Сколько событий пришло, столько карточек и собралось: одно событие — одна -строка. Ровно как `ods.browser_event` из прошлого урока. +**`dds.event` — события.** Сколько событий пришло с валидным ключом, столько карточек и +собралось: одно событие — одна строка. Ровно как `ods.browser_event` из прошлого урока. -**`dds.click` — 26, а не 50.** И это та же история, что мы уже разбирали в уроке 2. Карточка -клика — одна на клик, а в срезе на 50 событий разных кликов всего 26 (на один клик приходится -несколько событий). Поэтому 50 событий ссылаются на 26 кликов — это нормально, так и должно быть. +**`dds.click` обычно меньше, чем `dds.event`.** И это та же история, что мы уже разбирали +в уроке 2. Карточка клика — одна на клик, а событий на один клик может быть несколько. +Поэтому много событий ссылаются на меньшее число кликов — это нормально, так и должно быть. Теперь — главный счётчик урока. Он печатается чуть ниже, в блоке **«Сводка по качеству данных»** (это таблица `dm.dq_summary`, куда стенд складывает метрики по всем слоям). Найди в ней строку @@ -108,9 +108,9 @@ DDS»** — счётчики строк по двум нашим сущност В колонке `check_date` стоит `today()` из кода витрины, так что у тебя там будет сегодняшняя дата — не пугайся, если она не совпадёт с примером. -`orphan_events = 0` — ни одной сироты. Каждое из 50 событий нашло свой клик в `dds.click`. На -чистом демо-срезе так и должно быть: данные аккуратные, ничего не потерялось. В секции 4 мы -сироту устроим сами — и эта строка оживёт. +`orphan_events = 0` — ни одной сироты. Каждое событие нашло свой клик в `dds.click`. На +чистой стартовой истории так и должно быть: данные аккуратные, ничего не потерялось. В секции +4 мы сироту устроим сами — и эта строка оживёт. Проверь нолик сам, не верь на слово. Открой SQL-консоль `http://localhost:9123/play` (пользователь `default`, пароль `123456`) и посчитай сирот напрямую: @@ -158,9 +158,9 @@ SELECT click_id FROM ods.geo_by_click ... строим карточки: так не потеряется клик, который есть, например, в `geo`, но почему-то не доехал в `device`. -> На нашем срезе `device` и `geo` содержат один и тот же набор из 26 кликов, так что универсум -> тоже 26. Но код написан так, чтобы пережить случай, когда наборы **разойдутся**, — и это -> правильно: в проде они расходятся постоянно. +> На чистой стартовой истории `device` и `geo` должны содержать один и тот же набор кликов, +> так что универсум совпадает с обоими источниками. Но код написан так, чтобы пережить случай, +> когда наборы **разойдутся**, — и это правильно: в проде они расходятся постоянно. ### `argMax`: одна строка на клик, самая свежая @@ -221,7 +221,7 @@ LEFT JOIN ( ...снапшот geo... ) AS g ON g.click_id = c.click_id ### Сироты: событие без клика -Мы собрали `dds.click` (26 карточек кликов) и `dds.event` (50 карточек событий). Внутри каждого +Мы собрали `dds.click` (карточки кликов) и `dds.event` (карточки событий). Внутри каждого события лежит `click_id` — ссылка на клик. И вот тут возникает вопрос целостности из секции 1: **а на каждую ли ссылку есть карточка клика?** @@ -237,16 +237,16 @@ WHERE click_id IS NOT NULL Заметь разницу с предыдущим пунктом. Пустое гео — это когда у **клика** не подтянулся свой контекст (внутренний пропуск в карточке, но сам клик есть). А сирота — это когда у **события** нет вообще никакого клика (порвана связь между сущностями). Это разные дырки: первую видно по -пустым полям внутри карточки, вторую — отдельным счётчиком. На чистом срезе сирот ноль — сейчас -мы это изменим. +пустым полям внутри карточки, вторую — отдельным счётчиком. На чистой стартовой истории сирот +ноль — сейчас мы это изменим. --- ## 4. Управляемая правка: заведём сироту -Сирота на чистом срезе не появится сама — данные слишком аккуратные. Поэтому **создадим её -руками**: добавим в `dds.event` одно событие, которое ссылается на клик, которого в `dds.click` -нет. И посмотрим, как оживёт счётчик сирот и как себя поведёт `LEFT JOIN`. +Сирота на чистой стартовой истории не появится сама — данные слишком аккуратные. Поэтому +**создадим её руками**: добавим в `dds.event` одно событие, которое ссылается на клик, которого +в `dds.click` нет. И посмотрим, как оживёт счётчик сирот и как себя поведёт `LEFT JOIN`. Открой SQL-консоль `http://localhost:9123/play` и вставь придуманное событие: @@ -309,7 +309,7 @@ make transform ``` После этого `orphan_events` снова `0`, придуманное событие исчезло. А если стенд совсем «поплыл» — -полный сброс: `make clean && make up && make ddl && LIMIT=50 make data && make transform`. +полный чистый прогон: `make generated-history-analytics && make up`. --- @@ -317,9 +317,9 @@ make transform | Действие | Где смотреть | Что ожидать | |----------|--------------|-------------| -| `make transform` (базовый прогон) | блок «Статистика DDS» | `dds.click` = 26, `dds.event` = 50 | +| базовый прогон | блок «Статистика DDS» | `dds.click` и `dds.event` не пустые | | `make transform` (базовый прогон) | блок «Сводка по качеству», строка `orphan_events` | `0` | -| почему `click` = 26, а `event` = 50 | запрос `count()` по `dds.click` и `dds.event` | 50 событий ссылаются на 26 кликов — норма | +| почему `click` меньше `event` | запрос `count()` по `dds.click` и `dds.event` | много событий ссылаются на меньшее число кликов — норма | | правка из секции 4 (вставили сироту) | запрос `count()` сирот в play-консоли | `0 → 1` | | та же сирота через `dm.v_events_enriched` | `SELECT device_type, geo_country ...` | поля клика пустые (`NULL`) — это `LEFT JOIN` | diff --git a/docs/course/lessons/04_airflow_orchestration.md b/docs/course/lessons/04_airflow_orchestration.md index 2b3ec5c..a62d3e8 100644 --- a/docs/course/lessons/04_airflow_orchestration.md +++ b/docs/course/lessons/04_airflow_orchestration.md @@ -56,12 +56,14 @@ Airflow. Главная единица Airflow — **DAG** (Directed Acyclic Gra ## 2. Руки: запускаем DAG и смотрим зелёный прогон -Подними стенд, создай схему и залей малый срез: +Подними стенд и создай стартовую историю. Это штатный путь курса: готовый источник +данных стенда пишет события в Kafka, ClickHouse читает их в STG, затем batch строит +ODS, DDS и DM. Файлы `data/*.jsonl` пока остаются только кладовкой значений для этого +источника, а не источником аналитического контура. ```bash +make generated-history-analytics make up -make ddl -LIMIT=50 make data ``` Открой Airflow: `http://localhost:8080` (логин `admin`, пароль `admin`). Найди DAG @@ -75,13 +77,13 @@ LIMIT=50 make data заново наполнит их из ODS. Для учебного стенда это удобный чистый прогон: результат повторяемый, старые эксперименты не мешают. -Когда DAG завершится, открой его граф. На чистом срезе все задачи должны быть зелёными. Найди -внутри группы `transform` две задачи подряд: +Когда DAG завершится, открой его граф. На чистой стартовой истории все задачи должны быть +зелёными. Найди внутри группы `transform` две задачи подряд: - `check_dds_integrity` — SQL-задача, которая считает сирот; - `assert_dds_integrity` — Python-задача, которая решает, можно ли идти дальше. -На чистом срезе `assert_dds_integrity` зелёная: сирот нет, пайплайн прошёл в DM. +На чистой стартовой истории `assert_dds_integrity` зелёная: сирот нет, пайплайн прошёл в DM. Проверь то же число в ClickHouse play-консоли `http://localhost:9123/play`: @@ -305,14 +307,15 @@ WHERE click_id IS NOT NULL AND click_id NOT IN (SELECT click_id FROM dds.click); ``` -Снова должно быть `0`. Если стенд после экспериментов совсем запутался, полный сброс остаётся -тем же: +Снова должно быть `0`. Если стенд после экспериментов совсем запутался, сделай штатный +чистый прогон: ```bash -make clean && make up && make ddl && LIMIT=50 make data +make generated-history-analytics +make up ``` -После этого запусти `etl_pipeline` с `{"full_refresh": true}`. +После этого при необходимости запусти `etl_pipeline` с `{"full_refresh": true}`. --- diff --git a/docs/course/lessons/05_monitoring.md b/docs/course/lessons/05_monitoring.md index db96b03..d693cc2 100644 --- a/docs/course/lessons/05_monitoring.md +++ b/docs/course/lessons/05_monitoring.md @@ -57,15 +57,18 @@ ## 2. Руки: открываем дашборды и targets -Подними стенд и прогони маленький срез, если он ещё не поднят: +Подними стенд и создай стартовую историю, если он ещё не поднят. Это штатный путь курса: +готовый источник данных стенда пишет события в Kafka, ClickHouse читает их в STG, затем +batch строит ODS, DDS и DM. Файлы `data/*.jsonl` пока остаются только кладовкой значений +для этого источника, а не источником аналитического контура. ```bash +make generated-history-analytics make up -make ddl -LIMIT=50 make data ``` -После этого запусти `etl_pipeline` в Airflow с конфигом: +Если хочешь отдельно посмотреть Airflow-run после готовой истории, запусти `etl_pipeline` +в Airflow с конфигом: ```json {"full_refresh": true} diff --git a/docs/course/lessons/06_superset_bi.md b/docs/course/lessons/06_superset_bi.md index ca1f53b..c0604df 100644 --- a/docs/course/lessons/06_superset_bi.md +++ b/docs/course/lessons/06_superset_bi.md @@ -63,42 +63,28 @@ DM-витрина — это SQL-объект в ClickHouse. Она задаёт ## 2. Руки: запускаем Superset и смотрим дашборд -Подними стенд и прогони полный демо-датасет: +Подними стенд и создай стартовую историю. Это штатный путь курса: готовый источник +данных стенда пишет события в Kafka, ClickHouse читает их в STG, затем batch строит +ODS, DDS, DM и обновляет Superset. Файлы `data/*.jsonl` пока остаются только кладовкой +значений для этого источника, а не источником аналитического контура. ```bash +make generated-history-analytics make up -make ddl -make data -make transform ``` -`make transform` прогоняет цепочку STG → ODS → DDS → DM вне Airflow. Для этого урока так -быстрее: нам нужен готовый DM-слой, а не разбор DAG. Отладочный срез через -`LIMIT=50 make data` можно использовать для быстрых экспериментов, но эталонный dashboard -и числа урока рассчитаны на полном наборе данных. - -Теперь инициализируй Superset: - -```bash -make superset-init -``` - -Эта команда создаёт или обновляет: +Команда уже прогоняет цепочку STG → ODS → DDS → DM и создаёт metadata Superset: - подключение `clickhouse_dwh`; - 6 datasets поверх `dm.*`. - -После этого создай или обнови charts и сам dashboard: - -```bash -make superset-dashboard -``` - -Эта команда создаёт или обновляет: - - 10 charts; - dashboard `E-commerce Analytics Dashboard`. +Для этого урока нам нужен готовый DM-слой и уже созданный dashboard, а не разбор DAG. +Если Superset metadata нужно пересоздать отдельно, используй `make superset-init`. +Если нужно обновить только charts и dashboard после правки в `superset/create_dashboard.py`, +используй `make superset-dashboard`. + Открой Superset: `http://localhost:8088` (логин `admin`, пароль `admin`). Если Superset предлагает сменить пароль после первого входа, для учебного стенда можно нажать @@ -133,33 +119,33 @@ http://localhost:8088/superset/dashboard/ecommerce-analytics/ > **Зерно (grain), он же уровень гранулярности — это что считается одной строкой > таблицы.** У события зерно > «одно событие = одна строка» (ключ `event_id`), у визита — «один визит = одна -> строка» (ключ `click_id`). Это разные зёрна: событий 1000, а визитов 99, потому -> что в одном визите много событий. Складывать строки разного зерна в одно число -> бессмысленно — это всё равно что сложить «штуки яблок» и «корзины яблок». +> строка» (ключ `click_id`). Это разные зёрна: событий обычно больше, чем визитов, +> потому что в одном визите много событий. Складывать строки разного зерна в одно +> число бессмысленно — это всё равно что сложить «штуки яблок» и «корзины яблок». Поэтому чарт держит **одно зерно — event**: берёт по одной канонической таблице событий на слой (`browser_raw → browser_event → event → v_events_enriched`), а не сумму по слою. Если просуммировать все таблицы слоя, в один столбец попадут таблицы разного -зерна (события 1000 + визиты 99 + пустые error-таблицы) и получится «воронка потерь», -которой на самом деле нет. +зерна: события, визиты и технические таблицы ошибок. Получится «воронка потерь», которой +на самом деле нет. -Шаг **1050 → 1000** на первом переходе — это не потеря данных, а дедупликация -at-least-once потока по `event_id` в ODS (`ReplacingMergeTree`): в STG приехало 1050 строк, -но уникальных `event_id` среди них — 1000 (часть событий Kafka доставила повторно). Дальше -число стабильно. Настоящие проблемы качества (ошибки парсинга, осиротевшие события) на -чистых демо-данных равны нулю и лежат в `dm.dq_summary` отдельными `check_name` — их -разбирали уроки 3–4. +Если на первом переходе число уменьшается, это не обязательно потеря данных. В ODS работает +дедупликация at-least-once потока по `event_id` (`ReplacingMergeTree`): в STG могут приехать +повторы, а дальше остаётся каноническое число уникальных событий. Настоящие проблемы качества +(ошибки парсинга, осиротевшие события) на чистой стартовой истории равны нулю и лежат в +`dm.dq_summary` отдельными `check_name` — их разбирали уроки 3–4. ### Фильтр даты -В демо-данных события датированы `2022-11-28`. В текущей конфигурации dashboard фильтр даты -открывается как `No filter`. Если у тебя осталась старая metadata Superset и native filter -**Date Range** стоит в значении `Last week`, часть графиков может быть пустой, хотя данные есть. +В стартовой истории события идут по модельному времени стенда. В текущей конфигурации +dashboard фильтр даты открывается как `No filter`. Если у тебя осталась старая metadata +Superset и native filter **Date Range** стоит в значении `Last week`, часть графиков может +быть пустой, хотя данные есть. Для этого урока поставь в фильтре даты одно из двух: - `No filter`; -- или ручной диапазон вокруг `2022-11-28`. +- или ручной диапазон вокруг дат, которые видны в `event_timestamp`. После этого нажми **Apply filters**. Теперь смотри на dashboard как аналитик: какие графики отвечают на бизнес-вопросы, а какие только показывают техническое устройство конвейера. @@ -415,7 +401,7 @@ metadata Superset. | `make superset-init` | Superset → **Settings → Database Connections** | есть подключение `clickhouse_dwh` | | открыть **Datasets** | Superset UI | есть datasets `v_events_enriched`, `v_top_pages_daily`, `dq_summary` | | открыть dashboard | Superset UI | видны KPI, маркетинг, география и прохождение строк по слоям | -| поставить **Date Range → No filter** | dashboard filters | графики не скрываются из-за даты `2022-11-28` | +| поставить **Date Range → No filter** | dashboard filters | графики не скрываются из-за даты модельного времени | | поменять `row_limit` у `Page Funnel` на `3` и запустить `make superset-dashboard` | chart `Page Funnel` | не больше трёх страниц | | вернуть `row_limit` на `20` и запустить `make superset-dashboard` | chart `Page Funnel` | ограничение снова до 20 страниц |