diff --git a/.scratch/generator-model-time-startup-history/issues/06-generated-history-as-analytics-source.md b/.scratch/generator-model-time-startup-history/issues/06-generated-history-as-analytics-source.md index 3cc43f0..e7c5ecd 100644 --- a/.scratch/generator-model-time-startup-history/issues/06-generated-history-as-analytics-source.md +++ b/.scratch/generator-model-time-startup-history/issues/06-generated-history-as-analytics-source.md @@ -1,4 +1,4 @@ -Status: ready-for-agent +Status: ready-for-human # Стартовая история как источник аналитики @@ -24,26 +24,204 @@ STG→ODS→DDS→DM строится на ней, а Superset работает ## Acceptance criteria -- [ ] Штатная команда запуска чистого стенда создаёт или загружает стартовую +- [x] Штатная команда запуска чистого стенда создаёт или загружает стартовую историю и прогоняет её до DM-витрин. -- [ ] `kafka_load_dag` или заменяющий его путь больше не использует архивный +- [x] `kafka_load_dag` или заменяющий его путь больше не использует архивный `data/*.jsonl` как источник аналитики. -- [ ] ClickHouse-проверки из задачи 5 доступны как повторяемая команда для +- [x] ClickHouse-проверки из задачи 5 доступны как повторяемая команда для координатора или CI. -- [ ] Основные DM-витрины, на которых стоят дашборды, непустые и показывают +- [x] Основные DM-витрины, на которых стоят дашборды, непустые и показывают данные генерации. -- [ ] Superset-датасеты и дашборды технически открываются на данных генерации; +- [x] Superset-датасеты и дашборды технически открываются на данных генерации; ручная оценка формы графиков остаётся финальной HITL-приёмкой. -- [ ] `README.md`, `docs/OPERATIONS.md` и `generator/README.md` больше не +- [x] `README.md`, `docs/OPERATIONS.md` и `generator/README.md` больше не описывают архивный сид как основной источник аналитики. -- [ ] Если уроки или дашборды требуют нетривиальной переделки, создан follow-up +- [x] Если уроки или дашборды требуют нетривиальной переделки, создан follow-up issue вместо расширения этой задачи. -- [ ] Если учебные материалы ещё описывают архивный сид как основной источник, +- [x] Если учебные материалы ещё описывают архивный сид как основной источник, создан отдельный follow-up на миграцию уроков по ADR-0006. -- [ ] Описан повторный чистый прогон: какие данные очищаются и какие команды +- [x] Описан повторный чистый прогон: какие данные очищаются и какие команды выполняются, чтобы координатор мог надёжно перепроверить результат. ## Blocked by - `.scratch/generator-model-time-startup-history/issues/05-startup-history-backfill-to-clickhouse.md` - Review gate из `PRD.md`: распределения и два пути генерации после задачи 5 + +## Фактический прогон worker-а + +Дата проверки: 2026-06-14. + +Штатная команда: + +```bash +make generated-history-analytics +``` + +Команда выполняет чистый прогон: `docker compose down -v --remove-orphans`, +запуск ClickHouse и Kafka, DDL, `GEN_RUN_MODE=backfill`, batch +STG -> ODS -> DDS -> DM, инициализацию Superset и итоговую проверку +`make generated-history-check`. + +Быстрая команда issue 06 проверяет путь backfill -> DM -> Superset. Стык +backfill/live, отсутствие дублей на границе и однородность визитов через +`T_end` проверены в issue 05 и не входят в быстрый штатный check issue 06. + +Быстрый проверочный профиль по умолчанию: + +- `GEN_SEED=4242` +- `GEN_MODEL_T0=2026-01-01T00:00:00+00:00` +- `GEN_MODEL_T_END=2026-01-01T06:00:00+00:00` +- `GEN_MODEL_TIMEZONE=UTC` +- `GEN_MODEL_TIME_SPEED=1` +- `GEN_TICK_SECONDS=60` +- `GEN_LAMBDA_BASE_PER_MIN=60` +- `GEN_JITTER_PCT=0` +- `GEN_MIN_EVENTS_PER_TICK=1` +- `GEN_MAX_EVENTS_PER_TICK=1000` + +Суточный профиль остаётся доступен явно: + +```bash +GEN_MODEL_T_END=2026-01-02T00:00:00+00:00 make generated-history-analytics +``` + +Backfill завершился: + +- `events=16054` +- `visits=1516` +- `users=445` + +Batch STG -> ODS -> DDS -> DM: + +- STG, суммарно по 4 потокам: `64216` +- `ods.browser_event=16054` +- `ods.location_event=16054` +- `ods.device_by_click=1516` +- `ods.geo_by_click=1516` +- `dds.event=16054` +- `dds.click=1516` +- `dds.event_without_click=0` +- `dm.v_events_enriched` в `dm.dq_summary`: `16054` + +Повторная команда проверки: + +```bash +make generated-history-check +``` + +Результат ClickHouse: + +- `events=16054` +- `visits=1516` +- `users=445` +- `min_event_ts=2026-01-01 00:00:00.000000` +- `max_event_ts=2026-01-01 05:59:59.353407` +- `digest=B7E183DD1835E6593CC4B33C7F5B2817` + +Возвраты: + +- `users=445` +- `returning_users=351` +- `returning_share=0.7887640449438202` +- `max_visits_per_user=8` + +Форма длины визита: + +- `visits=1516` +- `short_visit_share=0.17678100263852242` +- `median_events_per_visit=8` +- `avg_events_per_visit=10.589709762532982` +- `capped_visit_share=0.06398416886543536` +- `median_duration_sec=204` +- `p95_duration_sec=835` +- `max_events_per_visit=30` + +Ordered funnel: + +- `home=1387` +- `products=849` +- `cart=511` +- `payment=332` +- `confirmation=182` +- `monotonic_ok=1` +- `confirmation_share=0.13121845710165825` + +Contains funnel: + +- `home=1387` +- `products=1254` +- `cart=880` +- `payment=568` +- `confirmation=324` +- `monotonic_ok=1` +- `confirmation_share=0.2335976928622927` + +Основные DM-витрины: + +- `dm.v_events_enriched=16054` +- `dm.v_daily_traffic=797` +- `dm.v_top_pages_daily=6` +- `dm.v_utm_effectiveness=33` +- `dm.v_session_overview=1516` +- `dm.dq_summary=20` + +Superset technical check: + +- `superset_datasets=6` +- `superset_dashboards=1` +- `superset_dashboard_charts=10` +- `dashboard_url=http://localhost:8088/superset/dashboard/ecommerce-analytics/` +- `login_api_status=200` +- `dashboard_api_status=200` +- `dashboard_title=🛒 E-commerce Analytics Dashboard` +- зарегистрированная Superset database `clickhouse_dwh` читает + `dm.v_events_enriched` через SQLAlchemy engine: + `superset_clickhouse_events=16054`. + +Проверки: + +- `bash -n scripts/run_generated_history_analytics.sh` — ok. +- `bash -n scripts/check_generated_analytics.sh` — ok. +- `make generated-history-analytics` — ok. +- `make generated-history-check` — ok. +- `make generator-test` — 134 passed. + +Документы запуска обновлены: + +- `README.md` +- `docs/OPERATIONS.md` +- `generator/README.md` +- `docs/REPO_MAP.md` +- `docs/SUPERSET_DASHBOARD.md` + +Follow-up: + +- `.scratch/generator-model-time-startup-history/issues/07-migrate-course-from-archive-seed.md` + — миграция учебных материалов с архивного сида на генерацию. + +## Риски и что не проверено + +- Ручная оценка формы dashboard глазами не выполнялась: это финальная HITL-приёмка. +- Суточный backfill не прогонялся до конца в этом срезе: он доступен через env, + но для координатора выбран быстрый 6-часовой профиль. +- `kafka_load_dag.py` физически оставлен как архивный ручной путь; штатные + документы больше не ведут через него как основной источник аналитики. +- `scripts/check_generated_analytics.sh` рассчитан на штатный UTC-профиль + (`GEN_MODEL_T0/T_END` с `+00:00`). Для произвольного timezone offset в этих + переменных нужна отдельная нормализация границ. +- Проверка `no_2022_rows` доказывает отсутствие архивного сида внутри быстрого + чистого модельного диапазона. Она не доказывает отсутствие старых строк на + грязном стенде вне этого диапазона; штатная команда закрывает это через + `docker compose down -v`. + +## Review gate + +- Саморевью worker-а нашло слабые доказательства по проверкам issue 05 и Superset: + добавлены возвраты, форма длины визита, ordered/contains funnel, Superset login + API, dashboard API и чтение ClickHouse через Superset database. +- Независимый reviewer gate: `gate pass`, блокирующих находок нет. +- Minor-находка reviewer-а по старой подсказке `make data` в `scripts/run_batch.sh` + исправлена до коммита. +- Проверки после исправлений: `make generated-history-check` — ok, + `bash -n scripts/*.sh` — ok. diff --git a/.scratch/generator-model-time-startup-history/issues/07-migrate-course-from-archive-seed.md b/.scratch/generator-model-time-startup-history/issues/07-migrate-course-from-archive-seed.md new file mode 100644 index 0000000..6846212 --- /dev/null +++ b/.scratch/generator-model-time-startup-history/issues/07-migrate-course-from-archive-seed.md @@ -0,0 +1,40 @@ +Status: needs-triage + +# Миграция учебных материалов с архивного сида на генерацию + +## Parent + +`.scratch/generator-model-time-startup-history/PRD.md` + +## Why + +После перевода штатного аналитического контура на стартовую историю генератора +часть учебных материалов всё ещё описывает `data/*.jsonl` как основной источник +данных стенда. Это нельзя править внутри issue 06: потребуется пройти уроки и +сохранить понятный учебный путь. + +## What to build + +Обновить курс и демо-материалы так, чтобы основной путь был: + +```text +startup-history/backfill -> Kafka -> STG -> ODS -> DDS -> DM -> Superset +``` + +Архивный `data/*.jsonl` оставить только как временную фактуру генератора. + +## Acceptance criteria + +- [ ] `docs/course/` больше не ведёт ученика через `make data` или `kafka_load` + как основной путь получения аналитических данных. +- [ ] Уроки явно объясняют, что `data/*.jsonl` пока остаётся кладовкой значений + для генератора, а не источником аналитического контура. +- [ ] Демо-шпаргалки и тест-план согласованы с новым штатным путём запуска. +- [ ] Если для уроков нужны новые скриншоты или ручная оценка dashboard, это + вынесено в HITL-приёмку. + +## Notes + +Найденные места для начала: `docs/course/PRD.md`, +`docs/course/lessons/06_superset_bi.md`, `docs/DEMO_CHEATSHEET_5MIN.md`, +`docs/DEMO_SCRIPT_10_15MIN.md`, `docs/TEST_PLAN.md`. diff --git a/Makefile b/Makefile index c703a95..02cbd5f 100644 --- a/Makefile +++ b/Makefile @@ -1,4 +1,5 @@ .PHONY: up down clean ddl data transform logs \ + generated-history-analytics generated-history-check \ reload-monitoring recover-monitoring \ superset-init superset-dashboard superset-ui superset-restart \ generator-up generator-down generator-logs generator-restart \ @@ -37,6 +38,14 @@ data: transform: bash ./scripts/run_batch.sh +# Чистый аналитический прогон: стартовая история генератора -> STG -> ODS -> DDS -> DM -> Superset +generated-history-analytics: + COMPOSE_BIN="$(COMPOSE)" bash ./scripts/run_generated_history_analytics.sh + +# Повторяемая проверка после прогона стартовой истории +generated-history-check: + COMPOSE_BIN="$(COMPOSE)" bash ./scripts/check_generated_analytics.sh + # Перезагрузка конфигурации мониторинга (после изменений в provisioning) reload-monitoring: @echo "=== Перезагрузка сервисов мониторинга ===" diff --git a/README.md b/README.md index 3a6e5a5..5a77ee7 100644 --- a/README.md +++ b/README.md @@ -9,11 +9,14 @@ витринами. Поток данных коротко: -- **bootstrap**: `data/*.jsonl → Airflow (kafka_load) → Kafka → ClickHouse (слой STG) → - Airflow (etl_pipeline: STG → ODS → DDS → DM) → Superset`. -- **steady-stream**: `generator-service → Kafka → ClickHouse (STG) → Airflow (etl_pipeline) +- **стартовая история**: `generator backfill → Kafka → ClickHouse (STG) → + batch STG → ODS → DDS → DM → Superset`. +- **живое продолжение**: `generator live → Kafka → ClickHouse (STG) → batch ETL → Superset`. +Файлы `data/*.jsonl` больше не основной источник аналитики. Пока они остаются +архивной кладовкой значений для генератора: браузеры, страны, устройства и UTM. + ## Куда дальше - **Хочешь учиться** — открой [курс «Кликстрим на ClickHouse»](./docs/course/README.md). @@ -25,34 +28,26 @@ ## Быстрый старт -Стенд управляется через Airflow — это основной рабочий способ. Отдельные shell-скрипты в -`scripts/` оставлены как запасной вариант для локальных прогонов (см. -[OPERATIONS](./docs/OPERATIONS.md)). +Штатный чистый запуск строит аналитику из стартовой истории генератора. Команда +очищает volumes ClickHouse и Kafka, создаёт стартовую историю, доводит её до DM и +проверяет Superset metadata. ```bash -# 1. Поднять весь стек -make up -docker compose ps # убедиться, что контейнеры запустились +make generated-history-analytics +docker compose ps ``` -Дальше — три шага в Airflow (веб-интерфейс `http://localhost:8080`, логин и пароль -`admin`/`admin`). Сними каждый DAG с паузы (кнопка Unpause) и запусти по очереди: - -1. `ddl_init` — создаёт базы, таблицы и представления в ClickHouse. -2. `kafka_load` — заливает события из `data/*.jsonl` в Kafka. -3. `etl_pipeline` — прогоняет цепочку STG → ODS → DDS → DM. - -Те же шаги можно запускать из командной строки — это удобно для скриптов: +По умолчанию это быстрый проверочный профиль на 6 часов модельного времени. +Суточную историю можно прогнать отдельно: ```bash -docker compose exec -T airflow-webserver airflow dags trigger ddl_init +GEN_MODEL_T_END=2026-01-02T00:00:00+00:00 make generated-history-analytics +``` -# Загрузить первые 100 строк каждого файла (limit=0 — загрузить всё) -docker compose exec -T airflow-webserver airflow dags trigger kafka_load \ - --conf '{"limit": 100, "reset_topics": true}' +Повторить только техническую проверку после уже выполненного прогона: -docker compose exec -T airflow-webserver airflow dags trigger etl_pipeline \ - --conf '{"full_refresh": true}' +```bash +make generated-history-check ``` Проверить, что данные дошли до витрин: @@ -78,17 +73,16 @@ docker compose exec -T clickhouse clickhouse-client --user=default --password=12 Готовый дашборд в Superset: `http://localhost:8088/superset/dashboard/ecommerce-analytics/` — он создаётся -автоматически через минуту-две после `make up`. Состав и настройка дашборда описаны в +во время `make generated-history-analytics`. Состав и настройка дашборда описаны в [SUPERSET_DASHBOARD](./docs/SUPERSET_DASHBOARD.md). ## Как устроен поток данных ```mermaid flowchart LR - subgraph AF["Airflow"] - D1["ddl_init"] - D2["kafka_load"] - D3["etl_pipeline"] + subgraph GEN["Generator"] + BF["backfill"] + LIVE["live"] end subgraph Kafka["Kafka"] @@ -102,11 +96,12 @@ flowchart LR DM["DM: витрины VIEW"] end - D2 -->|загрузка JSONL| Kafka -->|Kafka MV| STG + BF -->|стартовая история| Kafka + LIVE -->|продолжение| Kafka + Kafka -->|Kafka MV| STG STG -->|batch| ODS -->|batch| DDS -->|VIEW| DM - D1 -.->|DDL| CH - D3 -.->|batch| ODS & DDS + DDL["DDL"] -.-> CH ``` «Грязные» записи не роняют пайплайн: ошибки разбора складываются в `ods.*_errors` и в diff --git a/docs/OPERATIONS.md b/docs/OPERATIONS.md index dd7c1bf..ac0b235 100644 --- a/docs/OPERATIONS.md +++ b/docs/OPERATIONS.md @@ -9,8 +9,12 @@ - `make up` (или `docker compose up -d`) - `make down` (остановить и удалить контейнеры/сети проекта) - `make clean` (полная очистка: `down -v --remove-orphans`) +- `make generated-history-analytics` (штатный чистый прогон: стартовая история + генератора -> Kafka/STG -> ODS -> DDS -> DM -> Superset) +- `make generated-history-check` (повторяемая проверка ClickHouse и Superset + после прогона стартовой истории) - `make ddl` (применяет SQL из `sql/ddl/00_databases.sql` и `sql/ddl/*/*.sql` в ClickHouse) -- `make data` (пересоздаёт топики и заливает данные в Kafka; по умолчанию полный объём, срез — `LIMIT=50 make data`) +- `make data` (архивный путь: заливает `data/*.jsonl` в Kafka; не основной источник аналитики) - `make transform` (запускает batch-процесс ODS -> DDS -> DM) - `make superset-init` (повторная инициализация Superset: подключение к ClickHouse, датасеты, дашборд) - `docker compose ps` @@ -33,6 +37,10 @@ ## Airflow DAGs +Штатный аналитический путь больше не начинается с `kafka_load`: чистый стенд +получает данные из стартовой истории генератора. DAG-и ниже остаются для +ручных экспериментов, отладки и совместимости учебного стенда. + ### `ddl_init` - Запуск: ручной (`Trigger DAG`) @@ -41,6 +49,7 @@ ### `kafka_load` +- Архивный путь, не основной источник аналитики. - Запуск: ручной (`Trigger DAG with config`) - Параметры: - `limit` (`int`, default `0`) — количество строк (`0` = все) @@ -62,6 +71,8 @@ - `full_refresh` (`bool`, default `true`) — очистить DDS перед загрузкой - `wait_stg_timeout_sec` (`int`, default `600`, minimum `30`) — сколько секунд задача `wait_for_stg_data` ждёт появления данных в STG, прежде чем упасть по таймауту - Зависимость: требует наличия данных в STG (от `kafka_load` или `make data`) +- В штатном сценарии STG наполняет `make generated-history-analytics` через + backfill генератора. - Гейт целостности DDS: `check_dds_integrity` считает события без клика, а `assert_dds_integrity` роняет DAG при `orphan_events > 0`. Проверка идёт после `load_dds` и до `load_dm_summary`, чтобы DM не собирался поверх нарушенной связи @@ -122,6 +133,29 @@ GEN_STATE_RESET=true GEN_LAMBDA_BASE_PER_MIN=60 docker compose up -d generator ### Стартовая история через backfill +Штатная команда чистого прогона: + +```bash +make generated-history-analytics +``` + +Она выполняет полный сброс volumes, поднимает ClickHouse и Kafka, применяет DDL, +запускает `GEN_RUN_MODE=backfill`, прогоняет batch STG -> ODS -> DDS -> DM, +инициализирует Superset и запускает техническую проверку. Для координатора или CI +короткая повторная проверка после уже готового стенда: + +```bash +make generated-history-check +``` + +По умолчанию команда использует быстрый проверочный профиль: 6 часов модельного +времени (`GEN_MODEL_T_END=2026-01-01T06:00:00+00:00`). Суточную историю можно +прогнать отдельно, явно задав правую границу: + +```bash +GEN_MODEL_T_END=2026-01-02T00:00:00+00:00 make generated-history-analytics +``` + `GEN_RUN_MODE=backfill` быстро проматывает модельное прошлое от `GEN_MODEL_T0` до `GEN_MODEL_T_END` без сна. В Kafka попадают события только за полуоткрытый отрезок `[T0, T_end)`. В compact-topic `generator_state` сохраняется state на @@ -129,8 +163,9 @@ GEN_STATE_RESET=true GEN_LAMBDA_BASE_PER_MIN=60 docker compose up -d generator контрольными числами. При live-запуске с теми же настройками генератор видит, что state совпадает с manifest, и стартует ровно с `T_end` без настенной дельты. -Для чистого повтора проще всего пересоздать volumes. Это сбрасывает ClickHouse, -Kafka-топики данных и compact-topic state. +Для чистого повтора пересоздавайте volumes. Это сбрасывает ClickHouse, +Kafka-топики данных, state и manifest генератора. `make generated-history-analytics` +делает это по умолчанию (`CLEAN_START=1`). ```bash make clean @@ -141,7 +176,7 @@ GEN_RUN_MODE=backfill \ GEN_STATE_RESET=true \ GEN_SEED=4242 \ GEN_MODEL_T0=2026-01-01T00:00:00+00:00 \ -GEN_MODEL_T_END=2026-01-02T00:00:00+00:00 \ +GEN_MODEL_T_END=2026-01-01T06:00:00+00:00 \ GEN_MODEL_TIMEZONE=UTC \ GEN_MODEL_TIME_SPEED=1 \ GEN_TICK_SECONDS=60 \ @@ -154,6 +189,9 @@ sleep 10 bash scripts/run_batch.sh ``` +Ручной сценарий выше нужен для отладки. В обычной проверке используйте +`make generated-history-analytics`, чтобы не забыть Superset и итоговый check. + Manifest можно посмотреть так: ```bash @@ -175,7 +213,7 @@ docker compose run -d --name startup-history-live --no-deps \ -e GEN_STATE_RESET=false \ -e GEN_SEED=4242 \ -e GEN_MODEL_T0=2026-01-01T00:00:00+00:00 \ - -e GEN_MODEL_T_END=2026-01-02T00:00:00+00:00 \ + -e GEN_MODEL_T_END=2026-01-01T06:00:00+00:00 \ -e GEN_MODEL_TIMEZONE=UTC \ -e GEN_MODEL_TIME_SPEED=1 \ -e GEN_TICK_SECONDS=60 \ @@ -195,7 +233,7 @@ bash scripts/run_batch.sh ```sql WITH toDateTime64('2026-01-01 00:00:00', 6) AS t0, - toDateTime64('2026-01-02 00:00:00', 6) AS t_end + toDateTime64('2026-01-01 06:00:00', 6) AS t_end SELECT uniqExact(user_domain_id) AS users, uniqExact(click_id) AS visits, @@ -219,7 +257,7 @@ docker compose exec -T clickhouse clickhouse-client \ --query " WITH toDateTime64('2026-01-01 00:00:00', 6) AS t0, - toDateTime64('2026-01-02 00:00:00', 6) AS t_end + toDateTime64('2026-01-01 06:00:00', 6) AS t_end SELECT hex(sipHash128(groupArray(tuple( event_id, click_id, @@ -245,7 +283,7 @@ FROM ( ```sql WITH toDateTime64('2026-01-01 00:00:00', 6) AS t0, - toDateTime64('2026-01-02 00:00:00', 6) AS t_end, + toDateTime64('2026-01-01 06:00:00', 6) AS t_end, users AS ( SELECT user_domain_id, uniqExact(click_id) AS visits FROM dm.v_events_enriched @@ -266,7 +304,7 @@ FROM users; ```sql WITH toDateTime64('2026-01-01 00:00:00', 6) AS t0, - toDateTime64('2026-01-02 00:00:00', 6) AS t_end, + toDateTime64('2026-01-01 06:00:00', 6) AS t_end, 30 AS max_session_events, sessions AS ( SELECT @@ -296,7 +334,7 @@ FROM sessions; ```sql WITH toDateTime64('2026-01-01 00:00:00', 6) AS t0, - toDateTime64('2026-01-02 00:00:00', 6) AS t_end, + toDateTime64('2026-01-01 06:00:00', 6) AS t_end, sessions AS ( SELECT click_id, @@ -327,7 +365,7 @@ FROM sessions; ```sql WITH toDateTime64('2026-01-01 00:00:00', 6) AS t0, - toDateTime64('2026-01-02 00:00:00', 6) AS t_end, + toDateTime64('2026-01-01 06:00:00', 6) AS t_end, sessions AS ( SELECT click_id, @@ -356,8 +394,8 @@ FROM sessions; ```sql WITH toDateTime64('2026-01-01 00:00:00', 6) AS t0, - toDateTime64('2026-01-02 00:00:00', 6) AS t_end, - toDateTime64('2026-01-02 00:10:00', 6) AS t_live_end + toDateTime64('2026-01-01 06:00:00', 6) AS t_end, + toDateTime64('2026-01-01 06:10:00', 6) AS t_live_end SELECT count() AS events, uniqExact(event_id) AS unique_events, @@ -373,7 +411,7 @@ WHERE event_ts >= t0 AND event_ts < t_live_end; ```sql WITH - toDateTime64('2026-01-02 00:00:00', 6) AS t_end, + toDateTime64('2026-01-01 06:00:00', 6) AS t_end, crossing AS ( SELECT click_id, @@ -518,32 +556,19 @@ curl -s -u admin:admin -X POST http://localhost:3000/api/admin/provisioning/dash docker compose restart grafana ``` -## Рекомендуемый сценарий (фаза 2) +## Рекомендуемый сценарий ```bash -# 1. Запуск инфраструктуры -make up +# Полный чистый путь: генерация -> STG -> ODS -> DDS -> DM -> Superset +make generated-history-analytics -# 2. Инициализация схемы (один раз) -# Airflow UI -> DAGs -> ddl_init -> Trigger DAG - -# 3. Загрузка данных через Airflow -# Airflow UI -> DAGs -> kafka_load -> Trigger DAG with config -# Параметры по умолчанию: limit=0, reset_topics=true - -# 4. Запуск ETL -# Airflow UI -> DAGs -> etl_pipeline -> Trigger DAG with config -# {"full_refresh": true} - -# 5. Проверка результатов -docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query "SELECT count() FROM ods.browser_event" -docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query "SELECT count() FROM dds.event" -docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query "SELECT * FROM dm.dq_summary" +# Повторная техническая проверка без пересоздания данных +make generated-history-check ``` ## Быстрые проверки -- Kafka ingest: наличие данных в `stg.*` и типизированных строк в `ods.*`. +- Kafka ingest: наличие данных генератора в `stg.*` и типизированных строк в `ods.*`. - Airflow UI: `http://localhost:8080` показывает DAG `ddl_init`, `kafka_load`, `etl_pipeline`. - BI: витрина `dm.v_events_enriched` отвечает за разумное время при фильтре по дате. @@ -760,6 +785,7 @@ curl -s -X POST -u admin:admin http://localhost:3000/api/admin/provisioning/aler docker compose up -d clickhouse docker compose up -d --force-recreate superset-init superset ``` -- После `docker compose down -v` нужно повторно прогнать: `ddl_init` -> `kafka_load` -> `etl_pipeline`. -- После `make clean`/`down -v` Superset стартует, но витрины `dm.*` ещё пустые или отсутствуют до прогона ETL; после `ddl_init` -> `kafka_load` -> `etl_pipeline` выполнить `make superset-init`. -- Для демо по умолчанию использовать малый срез данных; полный прогон делать осознанно. +- После `docker compose down -v` нужно повторно прогнать `make generated-history-analytics`. +- После `make clean`/`down -v` Superset стартует, но витрины `dm.*` ещё пустые или + отсутствуют до прогона стартовой истории; используйте `make generated-history-analytics`. +- Архивную загрузку `make data` использовать только для ручных экспериментов. diff --git a/docs/REPO_MAP.md b/docs/REPO_MAP.md index 05bebb1..9e160fd 100644 --- a/docs/REPO_MAP.md +++ b/docs/REPO_MAP.md @@ -4,10 +4,10 @@ ## Исполняемые файлы -### Airflow (основной путь запуска) +### Airflow (ручной и учебный путь запуска) - `airflow/dags/ddl_init_dag.py` — инициализация схемы ClickHouse -- `airflow/dags/kafka_load_dag.py` — загрузка в Kafka из JSONL +- `airflow/dags/kafka_load_dag.py` — архивная загрузка в Kafka из JSONL; не основной источник аналитики - `airflow/dags/etl_pipeline_dag.py` — ETL процесс STG -> ODS -> DDS -> DM - `airflow/dags/utils/kafka_helpers.py` — helper-функции для Kafka - `airflow/dags/utils/sql_helpers.py` — чтение и подготовка SQL-файлов для DAG @@ -35,17 +35,20 @@ DDL (форма таблиц): - `superset/init_superset.py` — подключение к ClickHouse + создание датасетов - `superset/create_dashboard.py` — сборка дашборда с чартами -### Скрипты (запасной путь, не основной) +### Скрипты -Shell-скрипты `scripts/*` (и обёртки `make ddl`/`make data`/`make transform`) — это локальный fallback в обход Airflow. Основной путь запуска — DAG-и Airflow (см. выше). +Shell-скрипты `scripts/*` и Makefile-обёртки дают повторяемый локальный запуск. +Основной чистый путь аналитики — `make generated-history-analytics`. - `scripts/apply_clickhouse_ddl.sh` — применение DDL -- `scripts/load_kafka_data.sh` — загрузка в Kafka +- `scripts/load_kafka_data.sh` — архивная загрузка `data/*.jsonl` в Kafka - `scripts/run_batch.sh` — batch-процесс +- `scripts/run_generated_history_analytics.sh` — чистый прогон стартовой истории до DM и Superset +- `scripts/check_generated_analytics.sh` — проверка DM-витрин и Superset metadata на данных генерации ## Данные и конфиги -- `data/*.jsonl` — исходные данные (могут быть грязными) +- `data/*.jsonl` — архивная фактура для генератора; не основной источник аналитики - `configs/` — конфиги ClickHouse, Prometheus, Grafana - `configs/prometheus.yml` — конфигурация Prometheus (scrape targets для ClickHouse, Kafka, Airflow) - `configs/statsd_mapping.yml` — маппинг StatsD → Prometheus метрик для Airflow diff --git a/docs/SUPERSET_DASHBOARD.md b/docs/SUPERSET_DASHBOARD.md index 78b4541..edfed2e 100644 --- a/docs/SUPERSET_DASHBOARD.md +++ b/docs/SUPERSET_DASHBOARD.md @@ -9,20 +9,13 @@ ### 1. Запуск инфраструктуры ```bash -# Запуск всех сервисов -make up - -# Применение DDL в ClickHouse -make ddl - -# Загрузка данных в Kafka -make data - -# Запуск ETL-пайплайна (ODS → DDS → DM) -make transform +# Чистый прогон: стартовая история генератора -> DM -> Superset +make generated-history-analytics ``` -### 2. Инициализация Superset +Команда очищает volumes, генерирует стартовую историю, прогоняет batch +STG -> ODS -> DDS -> DM и создаёт metadata Superset. Если данные уже +подготовлены и нужно только пересобрать Superset: ```bash # Автоматическая инициализация (создание подключения и датасетов) @@ -32,7 +25,7 @@ make superset-init make superset-dashboard ``` -### 3. Доступ к UI +### 2. Доступ к UI Откройте в браузере: http://localhost:8088 @@ -63,13 +56,14 @@ make superset-dashboard - **🎯 Conversion to /confirmation** — доля просмотров `/confirmation` от просмотров `/home` KPI разложены в одну строку по 12-колоночной сетке Superset: четыре блока по 3 колонки. -`Unique Sessions` не вынесен отдельной KPI-плиткой, потому что в демо-данных -`user_domain_id` и `click_id` идут 1:1 и дают то же число, что `Unique Users`. +`Unique Sessions` не вынесен отдельной KPI-плиткой: в текущем дашборде важнее +развести события, пользователей и среднюю глубину визита. Генератор создаёт +повторные визиты, поэтому `user_domain_id` и `click_id` уже не идут 1:1. #### Динамика трафика - **📅 Events over Time** — линейный график событий с 5-минутными бакетами - (все события стенда укладываются в ~50 минут, поэтому часовая гранулярность - давала бы всего 2 точки и прямую линию) + (быстрый проверочный профиль покрывает 6 часов модельного времени, поэтому + 5-минутные бакеты дают видимую динамику без лишнего шума) - **📱 Traffic by Device** — pie chart распределения по устройствам #### География @@ -92,11 +86,9 @@ KPI разложены в одну строку по 12-колоночной с > **Почему именно одно зерно, а не сумма по слою.** Чарт берёт по одной > канонической таблице на слой (`browser_raw → browser_event → event → > v_events_enriched`). Если суммировать `total_rows` по всем таблицам слоя, -> в один столбец складываются таблицы разного зерна (события `1000` + визиты `99` -> + пустые error-таблицы) и получается **ложная «воронка потерь»**, которой нет. -> На одном зерне убывание становится настоящим: видимый шаг **1050 → 1000** — -> это дедупликация at-least-once потока по `event_id` в ODS -> (`ReplacingMergeTree`), а дальше число стабильно до витрины. +> в один столбец складываются таблицы разного зерна: события, визиты и +> error-таблицы. Получается **ложная «воронка потерь»**, которой нет. На одном +> зерне видно прохождение event-строк по слоям, а не сумму несравнимых таблиц. > > Настоящие сигналы качества (`rows_with_errors` в ODS, `orphan_events` в DDS) > на чистых демо-данных равны нулю и живут в `dm.dq_summary` отдельными @@ -110,7 +102,7 @@ KPI разложены в одну строку по 12-колоночной с | Фильтр | Поле | Тип | Применение | |--------|------|-----|------------| -| 📅 Date Range | `event_date` | Time Range | Charts с `event_date`; по умолчанию `No filter`, чтобы демо-данные 2022 года не скрывались | +| 📅 Date Range | `event_date` | Time Range | Charts с `event_date`; по умолчанию `No filter`, чтобы стартовая история не скрывалась фильтром даты | | 🌍 Country | `geo_country` | Multi-select | Charts на `dm.v_events_enriched` | | 📱 Device Type | `device_type` | Multi-select | Charts на `dm.v_events_enriched` | | 🌐 Browser | `browser_name` | Multi-select | Charts на `dm.v_events_enriched` | @@ -135,8 +127,10 @@ make clean # Остановка с удалением volumes make logs service=superset # Логи сервиса # ETL +make generated-history-analytics # Чистый прогон генерации до Superset +make generated-history-check # Проверка DM и Superset metadata make ddl # Применение DDL в ClickHouse -make data # Загрузка данных в Kafka +make data # Архивная загрузка data/*.jsonl в Kafka make transform # Запуск batch-процесса # Superset @@ -247,12 +241,7 @@ make superset-restart # Полная переинициализация docker compose down -v -docker compose up -d -make ddl -make data -make transform -make superset-init -make superset-dashboard +make generated-history-analytics ``` ### Нет данных в чартах diff --git a/generator/README.md b/generator/README.md index 0c862e5..f5a75b0 100644 --- a/generator/README.md +++ b/generator/README.md @@ -1,6 +1,7 @@ -# Генератор событий (MVP rev5) +# Генератор событий -Автономный генератор событий для Kafka с режимом `steady-stream`. +Автономный генератор событий для Kafka. В штатном стенде он создаёт стартовую +историю через `backfill`, а затем может продолжить поток в режиме `live`. Генератор строит поток по иерархии `пользователь → визит → событие`: один `click_id` живёт весь визит, события визита идут по страницам воронки с @@ -34,7 +35,7 @@ generator-service -> Kafka topics -> (потребители отдельно) | `src/clickstream_generator/service.py` | основной цикл сервиса | | `generator.py` | запуск сервиса и совместимый фасад | -## Режим работы: `steady-stream` +## Режим работы: `live` - Публикуем постепенно, **короткими тиками** (по умолчанию каждые 5 секунд) - На каждом тике отправляем небольшую порцию сообщений @@ -110,6 +111,25 @@ restart policy. Контейнерные значения `KAFKA_BOOTSTRAP_SERVERS` и `GEN_DATA_DIR` в compose оставлены безопасными внутренними значениями `kafka:29092` и `/data`. +Штатный чистый путь всего стенда запускается из корня репозитория: + +```bash +make generated-history-analytics +``` + +Эта команда очищает ClickHouse, Kafka-топики данных, state и manifest +генератора, создаёт стартовую историю, прогоняет STG -> ODS -> DDS -> DM и +проверяет Superset metadata. Файлы `data/*.jsonl` при этом не грузятся в Kafka: +они пока используются только как фактура для генератора. + +По умолчанию используется быстрый проверочный профиль на 6 часов модельного +времени (`GEN_MODEL_T_END=2026-01-01T06:00:00+00:00`). Суточный прогон доступен +явно: + +```bash +GEN_MODEL_T_END=2026-01-02T00:00:00+00:00 make generated-history-analytics +``` + ### Режим "раз в минуту" (для демо) Для контролируемых демо можно установить: @@ -141,6 +161,9 @@ make generator-restart # Запуск тестов make generator-test + +# Чистый аналитический прогон всего стенда +make generated-history-analytics ``` ## Метрики Prometheus diff --git a/scripts/check_generated_analytics.sh b/scripts/check_generated_analytics.sh new file mode 100644 index 0000000..83c87b0 --- /dev/null +++ b/scripts/check_generated_analytics.sh @@ -0,0 +1,421 @@ +#!/usr/bin/env bash +# +# Повторяемая проверка, что DM-витрины и Superset metadata работают на данных +# стартовой истории генератора, а не на архивном сиде 2022 года. + +set -euo pipefail + +COMPOSE_BIN="${COMPOSE_BIN:-docker compose}" +CLICKHOUSE_SERVICE="${CLICKHOUSE_SERVICE:-clickhouse}" +CLICKHOUSE_USER="${CLICKHOUSE_USER:-default}" +CLICKHOUSE_PASSWORD="${CLICKHOUSE_PASSWORD:-123456}" +GEN_MODEL_T0="${GEN_MODEL_T0:-2026-01-01T00:00:00+00:00}" +GEN_MODEL_T_END="${GEN_MODEL_T_END:-2026-01-01T06:00:00+00:00}" +REQUIRE_SUPERSET="${REQUIRE_SUPERSET:-1}" + +fail() { + echo "Ошибка: $*" >&2 + exit 1 +} + +clickhouse_datetime_literal() { + local value="$1" + value="${value/T/ }" + value="${value%Z}" + if [[ "${value}" =~ ^(.*)[+-][0-9]{2}:[0-9]{2}$ ]]; then + value="${BASH_REMATCH[1]}" + fi + echo "${value}" +} + +ch_query() { + ${COMPOSE_BIN} exec -T "${CLICKHOUSE_SERVICE}" clickhouse-client \ + --user="${CLICKHOUSE_USER}" \ + --password="${CLICKHOUSE_PASSWORD}" \ + --query "$1" +} + +pg_query() { + ${COMPOSE_BIN} exec -T postgres-metadata psql \ + -U airflow \ + -d superset \ + -At \ + -c "$1" | tr -d '\r' +} + +CH_MODEL_T0="$(clickhouse_datetime_literal "${GEN_MODEL_T0}")" +CH_MODEL_T_END="$(clickhouse_datetime_literal "${GEN_MODEL_T_END}")" + +echo "=== Проверка ClickHouse: данные генерации в DM ===" + +stats_query=" +WITH + toDateTime64('${CH_MODEL_T0}', 6) AS t0, + toDateTime64('${CH_MODEL_T_END}', 6) AS t_end +SELECT + count() AS events, + uniqExact(click_id) AS visits, + uniqExact(user_domain_id) AS users, + toString(min(event_ts)) AS min_event_ts, + toString(max(event_ts)) AS max_event_ts, + users < visits AND visits < events AS pyramid_ok, + min(event_ts) >= t0 AND max(event_ts) < t_end AS half_open_ok, + countIf(toYear(event_ts) = 2022) = 0 AS no_2022_rows, + hex(sipHash128(groupArray(tuple(event_id, click_id, user_domain_id, event_ts, page_url_path)))) AS digest +FROM +( + SELECT + event_id, + click_id, + user_domain_id, + event_ts, + page_url_path + FROM dm.v_events_enriched + WHERE event_ts >= t0 AND event_ts < t_end + ORDER BY event_id +) +FORMAT TabSeparated" + +stats="$(ch_query "${stats_query}")" +IFS=$'\t' read -r events visits users min_event_ts max_event_ts pyramid_ok half_open_ok no_2022_rows digest <<< "${stats}" + +[[ "${events}" =~ ^[0-9]+$ ]] || fail "не удалось прочитать число событий из dm.v_events_enriched" +(( events > 0 )) || fail "dm.v_events_enriched пустая в модельном диапазоне" +[[ "${pyramid_ok}" == "1" ]] || fail "нарушена пирамида users < visits < events" +[[ "${half_open_ok}" == "1" ]] || fail "данные вышли за диапазон [GEN_MODEL_T0, GEN_MODEL_T_END)" +[[ "${no_2022_rows}" == "1" ]] || fail "найдены строки 2022 года, похожие на архивный сид" + +echo "events=${events}" +echo "visits=${visits}" +echo "users=${users}" +echo "min_event_ts=${min_event_ts}" +echo "max_event_ts=${max_event_ts}" +echo "digest=${digest}" + +echo "" +echo "=== Проверка ClickHouse: возвраты пользователей ===" + +returns_query=" +WITH + toDateTime64('${CH_MODEL_T0}', 6) AS t0, + toDateTime64('${CH_MODEL_T_END}', 6) AS t_end, + users AS ( + SELECT user_domain_id, uniqExact(click_id) AS visits + FROM dm.v_events_enriched + WHERE event_ts >= t0 AND event_ts < t_end + AND user_domain_id IS NOT NULL + GROUP BY user_domain_id + ) +SELECT + count() AS users, + countIf(visits > 1) AS returning_users, + returning_users / users AS returning_share, + max(visits) AS max_visits_per_user +FROM users +FORMAT TabSeparated" + +returns="$(ch_query "${returns_query}")" +IFS=$'\t' read -r total_users returning_users returning_share max_visits_per_user <<< "${returns}" + +[[ "${total_users}" =~ ^[0-9]+$ ]] || fail "не удалось прочитать число пользователей" +(( total_users > 0 )) || fail "нет пользователей для проверки возвратов" +(( returning_users > 0 )) || fail "нет пользователей с повторными визитами" + +echo "users=${total_users}" +echo "returning_users=${returning_users}" +echo "returning_share=${returning_share}" +echo "max_visits_per_user=${max_visits_per_user}" + +echo "" +echo "=== Проверка ClickHouse: форма длины визита ===" + +visit_shape_query=" +WITH + toDateTime64('${CH_MODEL_T0}', 6) AS t0, + toDateTime64('${CH_MODEL_T_END}', 6) AS t_end, + 30 AS max_session_events, + sessions AS ( + SELECT + click_id, + count() AS events_count, + dateDiff('second', min(event_ts), max(event_ts)) AS duration_sec + FROM dm.v_events_enriched + WHERE event_ts >= t0 AND event_ts < t_end + GROUP BY click_id + ) +SELECT + count() AS visits, + countIf(events_count <= 2) / visits AS short_visit_share, + quantileExact(0.5)(events_count) AS median_events_per_visit, + avg(events_count) AS avg_events_per_visit, + countIf(events_count = max_session_events) / visits AS capped_visit_share, + quantileExact(0.5)(duration_sec) AS median_duration_sec, + quantileExact(0.95)(duration_sec) AS p95_duration_sec, + max(events_count) AS max_events_per_visit +FROM sessions +FORMAT TabSeparated" + +visit_shape="$(ch_query "${visit_shape_query}")" +IFS=$'\t' read -r shape_visits short_visit_share median_events_per_visit avg_events_per_visit capped_visit_share median_duration_sec p95_duration_sec max_events_per_visit <<< "${visit_shape}" + +[[ "${shape_visits}" =~ ^[0-9]+$ ]] || fail "не удалось прочитать число визитов" +(( shape_visits > 0 )) || fail "нет визитов для проверки формы" +(( max_events_per_visit <= 30 )) || fail "длина визита превысила GEN_MAX_SESSION_EVENTS" + +echo "visits=${shape_visits}" +echo "short_visit_share=${short_visit_share}" +echo "median_events_per_visit=${median_events_per_visit}" +echo "avg_events_per_visit=${avg_events_per_visit}" +echo "capped_visit_share=${capped_visit_share}" +echo "median_duration_sec=${median_duration_sec}" +echo "p95_duration_sec=${p95_duration_sec}" +echo "max_events_per_visit=${max_events_per_visit}" + +echo "" +echo "=== Проверка ClickHouse: ordered funnel ===" + +ordered_funnel_query=" +WITH + toDateTime64('${CH_MODEL_T0}', 6) AS t0, + toDateTime64('${CH_MODEL_T_END}', 6) AS t_end, + sessions AS ( + SELECT + click_id, + minIf(event_ts, page_url_path = '/home') AS home_ts, + minIf(event_ts, page_url_path IN ('/product_a', '/product_b')) AS product_ts, + minIf(event_ts, page_url_path = '/cart') AS cart_ts, + minIf(event_ts, page_url_path = '/payment') AS payment_ts, + minIf(event_ts, page_url_path = '/confirmation') AS confirmation_ts + FROM dm.v_events_enriched + WHERE event_ts >= t0 AND event_ts < t_end + GROUP BY click_id + ) +SELECT + countIf(home_ts IS NOT NULL) AS home, + countIf(home_ts IS NOT NULL AND product_ts > home_ts) AS products, + countIf(home_ts IS NOT NULL AND product_ts > home_ts AND cart_ts > product_ts) AS cart, + countIf(home_ts IS NOT NULL AND product_ts > home_ts AND cart_ts > product_ts AND payment_ts > cart_ts) AS payment, + countIf(home_ts IS NOT NULL AND product_ts > home_ts AND cart_ts > product_ts AND payment_ts > cart_ts AND confirmation_ts > payment_ts) AS confirmation, + products <= home AND cart <= products AND payment <= cart AND confirmation <= payment AS monotonic_ok, + confirmation / home AS confirmation_share +FROM sessions +FORMAT TabSeparated" + +ordered_funnel="$(ch_query "${ordered_funnel_query}")" +IFS=$'\t' read -r ordered_home ordered_products ordered_cart ordered_payment ordered_confirmation ordered_monotonic_ok ordered_confirmation_share <<< "${ordered_funnel}" + +[[ "${ordered_home}" =~ ^[0-9]+$ ]] || fail "не удалось прочитать ordered funnel" +(( ordered_home > 0 )) || fail "ordered funnel: нет /home" +[[ "${ordered_monotonic_ok}" == "1" ]] || fail "ordered funnel не монотонен" + +echo "home=${ordered_home}" +echo "products=${ordered_products}" +echo "cart=${ordered_cart}" +echo "payment=${ordered_payment}" +echo "confirmation=${ordered_confirmation}" +echo "monotonic_ok=${ordered_monotonic_ok}" +echo "confirmation_share=${ordered_confirmation_share}" + +echo "" +echo "=== Проверка ClickHouse: contains funnel ===" + +contains_funnel_query=" +WITH + toDateTime64('${CH_MODEL_T0}', 6) AS t0, + toDateTime64('${CH_MODEL_T_END}', 6) AS t_end, + sessions AS ( + SELECT + click_id, + countIf(page_url_path = '/home') > 0 AS has_home, + countIf(page_url_path IN ('/product_a', '/product_b')) > 0 AS has_product, + countIf(page_url_path = '/cart') > 0 AS has_cart, + countIf(page_url_path = '/payment') > 0 AS has_payment, + countIf(page_url_path = '/confirmation') > 0 AS has_confirmation + FROM dm.v_events_enriched + WHERE event_ts >= t0 AND event_ts < t_end + GROUP BY click_id + ) +SELECT + countIf(has_home) AS home, + countIf(has_home AND has_product) AS products, + countIf(has_home AND has_product AND has_cart) AS cart, + countIf(has_home AND has_product AND has_cart AND has_payment) AS payment, + countIf(has_home AND has_product AND has_cart AND has_payment AND has_confirmation) AS confirmation, + products <= home AND cart <= products AND payment <= cart AND confirmation <= payment AS monotonic_ok, + confirmation / home AS confirmation_share +FROM sessions +FORMAT TabSeparated" + +contains_funnel="$(ch_query "${contains_funnel_query}")" +IFS=$'\t' read -r contains_home contains_products contains_cart contains_payment contains_confirmation contains_monotonic_ok contains_confirmation_share <<< "${contains_funnel}" + +[[ "${contains_home}" =~ ^[0-9]+$ ]] || fail "не удалось прочитать contains funnel" +(( contains_home > 0 )) || fail "contains funnel: нет /home" +[[ "${contains_monotonic_ok}" == "1" ]] || fail "contains funnel не монотонен" + +echo "home=${contains_home}" +echo "products=${contains_products}" +echo "cart=${contains_cart}" +echo "payment=${contains_payment}" +echo "confirmation=${contains_confirmation}" +echo "monotonic_ok=${contains_monotonic_ok}" +echo "confirmation_share=${contains_confirmation_share}" + +echo "" +echo "=== Проверка ClickHouse: основные DM-витрины не пустые ===" + +views_query=" +SELECT source, rows +FROM +( + SELECT 'dm.v_events_enriched' AS source, count() AS rows FROM dm.v_events_enriched + UNION ALL SELECT 'dm.v_daily_traffic', count() FROM dm.v_daily_traffic + UNION ALL SELECT 'dm.v_top_pages_daily', count() FROM dm.v_top_pages_daily + UNION ALL SELECT 'dm.v_utm_effectiveness', count() FROM dm.v_utm_effectiveness + UNION ALL SELECT 'dm.v_session_overview', count() FROM dm.v_session_overview + UNION ALL SELECT 'dm.dq_summary', count() FROM dm.dq_summary +) +ORDER BY source +FORMAT TabSeparated" + +while IFS=$'\t' read -r source rows; do + [[ -n "${source}" ]] || continue + [[ "${rows}" =~ ^[0-9]+$ ]] || fail "не удалось прочитать число строк для ${source}" + (( rows > 0 )) || fail "${source} пустая" + echo "${source}=${rows}" +done < <(ch_query "${views_query}") + +if [[ "${REQUIRE_SUPERSET}" != "1" ]]; then + echo "" + echo "Проверка Superset пропущена: REQUIRE_SUPERSET=${REQUIRE_SUPERSET}" + exit 0 +fi + +echo "" +echo "=== Проверка Superset: datasets, charts и dashboard созданы ===" + +if ! ${COMPOSE_BIN} ps --services --filter "status=running" | grep -qx "superset"; then + fail "сервис superset не запущен" +fi + +${COMPOSE_BIN} exec -T superset curl -fsS http://localhost:8088/health >/dev/null \ + || fail "Superset health endpoint не отвечает" + +dataset_count="$(pg_query " +SELECT count(*) +FROM tables +WHERE schema = 'dm' + AND table_name IN ( + 'v_events_enriched', + 'v_daily_traffic', + 'v_utm_effectiveness', + 'v_top_pages_daily', + 'v_session_overview', + 'dq_summary' + );")" + +dashboard_count="$(pg_query " +SELECT count(*) +FROM dashboards +WHERE slug = 'ecommerce-analytics';")" + +chart_count="$(pg_query " +SELECT count(*) +FROM dashboard_slices ds +JOIN dashboards d ON d.id = ds.dashboard_id +WHERE d.slug = 'ecommerce-analytics';")" + +[[ "${dataset_count}" == "6" ]] || fail "ожидалось 6 Superset datasets, найдено ${dataset_count}" +[[ "${dashboard_count}" == "1" ]] || fail "dashboard ecommerce-analytics не найден" +(( chart_count > 0 )) || fail "dashboard ecommerce-analytics не связан с chart" + +echo "superset_datasets=${dataset_count}" +echo "superset_dashboards=${dashboard_count}" +echo "superset_dashboard_charts=${chart_count}" +echo "dashboard_url=http://localhost:8088/superset/dashboard/ecommerce-analytics/" + +echo "" +echo "=== Проверка Superset: dashboard открывается и читает ClickHouse ===" + +superset_probe="$( + ${COMPOSE_BIN} exec -T superset bash -s <<'PY' +set -euo pipefail +python - <<'PYTHON' +import json +import urllib.request + +from sqlalchemy import text + +from superset.app import create_app + +app = create_app() +with app.app_context(): + from superset.extensions import db + from superset.models.core import Database + from superset.models.dashboard import Dashboard + + dashboard = ( + db.session.query(Dashboard) + .filter_by(slug="ecommerce-analytics") + .one() + ) + + database = ( + db.session.query(Database) + .filter_by(database_name="clickhouse_dwh") + .one() + ) + with database.get_sqla_engine() as engine: + with engine.connect() as connection: + events = connection.execute( + text("SELECT count() FROM dm.v_events_enriched") + ).scalar() + +login_payload = json.dumps( + { + "username": "admin", + "password": "admin", + "provider": "db", + "refresh": True, + } +).encode("utf-8") + +login_request = urllib.request.Request( + "http://localhost:8088/api/v1/security/login", + data=login_payload, + headers={"Content-Type": "application/json"}, + method="POST", +) +with urllib.request.urlopen(login_request, timeout=10) as response: + login_status = response.status + token = json.loads(response.read().decode("utf-8"))["access_token"] + +dashboard_request = urllib.request.Request( + f"http://localhost:8088/api/v1/dashboard/{dashboard.id}", + headers={"Authorization": f"Bearer {token}"}, + method="GET", +) +with urllib.request.urlopen(dashboard_request, timeout=10) as response: + dashboard_status = response.status + dashboard_payload = json.loads(response.read().decode("utf-8")) + +dashboard_title = dashboard_payload["result"]["dashboard_title"] + +print(f"login_api_status={login_status}") +print(f"dashboard_api_status={dashboard_status}") +print(f"dashboard_title={dashboard_title}") +print(f"superset_clickhouse_events={events}") + +if login_status >= 400: + raise SystemExit("Superset login API failed") +if dashboard_status >= 400: + raise SystemExit("Superset dashboard API did not open") +if "E-commerce Analytics" not in dashboard_title: + raise SystemExit("Unexpected dashboard title") +if not isinstance(events, int) or events <= 0: + raise SystemExit("Superset ClickHouse query returned no data") +PYTHON +PY +)" + +echo "${superset_probe}" diff --git a/scripts/run_batch.sh b/scripts/run_batch.sh index 885e8da..8e44421 100755 --- a/scripts/run_batch.sh +++ b/scripts/run_batch.sh @@ -14,7 +14,7 @@ # # Требования: # - ClickHouse запущен (make up) -# - STG содержит данные (make data выполнен) +# - STG содержит данные генератора или ручной загрузки # # Стратегия: # Сейчас: полная перезагрузка (TRUNCATE + INSERT) — для демо @@ -81,7 +81,8 @@ STG_COUNT=$(${COMPOSE_BIN} exec -T "${CLICKHOUSE_SERVICE}" clickhouse-client \ if [[ "${STG_COUNT}" == "0" ]]; then echo "Предупреждение: Таблицы STG пусты." - echo "Сначала загрузите данные: make data" + echo "Сначала загрузите данные: make generated-history-analytics" + echo "Для ручной отладки можно выполнить backfill генератора или архивный make data." exit 1 fi diff --git a/scripts/run_generated_history_analytics.sh b/scripts/run_generated_history_analytics.sh new file mode 100644 index 0000000..26414a9 --- /dev/null +++ b/scripts/run_generated_history_analytics.sh @@ -0,0 +1,106 @@ +#!/usr/bin/env bash +# +# Чистый прогон аналитического контура от стартовой истории генератора до DM и Superset. +# +# Команда намеренно очищает volumes по умолчанию: так сбрасываются ClickHouse, +# Kafka-топики данных, state и manifest генератора. Это штатный повторяемый путь +# для проверки ADR-0006 на свежем стенде. + +set -euo pipefail + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +REPO_ROOT="$(cd "${SCRIPT_DIR}/.." && pwd)" + +COMPOSE_BIN="${COMPOSE_BIN:-docker compose}" +CLEAN_START="${CLEAN_START:-1}" +WAIT_CLICKHOUSE_SECONDS="${WAIT_CLICKHOUSE_SECONDS:-60}" +WAIT_STG_SECONDS="${WAIT_STG_SECONDS:-10}" + +GEN_SEED="${GEN_SEED:-4242}" +GEN_MODEL_T0="${GEN_MODEL_T0:-2026-01-01T00:00:00+00:00}" +GEN_MODEL_T_END="${GEN_MODEL_T_END:-2026-01-01T06:00:00+00:00}" +GEN_MODEL_TIMEZONE="${GEN_MODEL_TIMEZONE:-UTC}" +GEN_MODEL_TIME_SPEED="${GEN_MODEL_TIME_SPEED:-1}" +GEN_TICK_SECONDS="${GEN_TICK_SECONDS:-60}" +GEN_LAMBDA_BASE_PER_MIN="${GEN_LAMBDA_BASE_PER_MIN:-60}" +GEN_JITTER_PCT="${GEN_JITTER_PCT:-0}" +GEN_MIN_EVENTS_PER_TICK="${GEN_MIN_EVENTS_PER_TICK:-1}" +GEN_MAX_EVENTS_PER_TICK="${GEN_MAX_EVENTS_PER_TICK:-1000}" + +cd "${REPO_ROOT}" + +wait_for_clickhouse() { + local deadline + deadline=$((SECONDS + WAIT_CLICKHOUSE_SECONDS)) + + until ${COMPOSE_BIN} exec -T clickhouse clickhouse-client \ + --user=default \ + --password=123456 \ + --query "SELECT 1" >/dev/null 2>&1; do + if (( SECONDS >= deadline )); then + echo "Ошибка: ClickHouse не ответил за ${WAIT_CLICKHOUSE_SECONDS} сек." >&2 + exit 1 + fi + sleep 2 + done +} + +echo "=== Чистый прогон стартовой истории как источника аналитики ===" +echo "GEN_SEED=${GEN_SEED}" +echo "GEN_MODEL_T0=${GEN_MODEL_T0}" +echo "GEN_MODEL_T_END=${GEN_MODEL_T_END}" +echo "GEN_LAMBDA_BASE_PER_MIN=${GEN_LAMBDA_BASE_PER_MIN}" +echo "" + +if [[ "${CLEAN_START}" == "1" ]]; then + echo "Шаг 0: очистка volumes ClickHouse/Kafka/state" + ${COMPOSE_BIN} down -v --remove-orphans +else + echo "Шаг 0: CLEAN_START=0, очистка пропущена" +fi + +echo "Шаг 1: запуск ClickHouse и Kafka" +${COMPOSE_BIN} up -d clickhouse kafka +wait_for_clickhouse + +echo "Шаг 2: применение DDL" +bash "${SCRIPT_DIR}/apply_clickhouse_ddl.sh" + +echo "Шаг 3: сборка образа генератора" +${COMPOSE_BIN} build generator + +echo "Шаг 4: backfill стартовой истории в Kafka" +${COMPOSE_BIN} run --rm --no-deps \ + -e GEN_RUN_MODE=backfill \ + -e GEN_STATE_RESET=true \ + -e GEN_SEED="${GEN_SEED}" \ + -e GEN_MODEL_T0="${GEN_MODEL_T0}" \ + -e GEN_MODEL_T_END="${GEN_MODEL_T_END}" \ + -e GEN_MODEL_TIMEZONE="${GEN_MODEL_TIMEZONE}" \ + -e GEN_MODEL_TIME_SPEED="${GEN_MODEL_TIME_SPEED}" \ + -e GEN_TICK_SECONDS="${GEN_TICK_SECONDS}" \ + -e GEN_LAMBDA_BASE_PER_MIN="${GEN_LAMBDA_BASE_PER_MIN}" \ + -e GEN_JITTER_PCT="${GEN_JITTER_PCT}" \ + -e GEN_MIN_EVENTS_PER_TICK="${GEN_MIN_EVENTS_PER_TICK}" \ + -e GEN_MAX_EVENTS_PER_TICK="${GEN_MAX_EVENTS_PER_TICK}" \ + generator + +echo "Шаг 5: ожидание чтения Kafka Materialized View (${WAIT_STG_SECONDS} сек.)" +sleep "${WAIT_STG_SECONDS}" + +echo "Шаг 6: batch STG -> ODS -> DDS -> DM" +bash "${SCRIPT_DIR}/run_batch.sh" + +echo "Шаг 7: инициализация Superset metadata и dashboard" +${COMPOSE_BIN} up -d postgres-metadata clickhouse +${COMPOSE_BIN} up --abort-on-container-exit --exit-code-from superset-init superset-init +${COMPOSE_BIN} up -d --no-deps superset + +echo "Шаг 8: техническая проверка аналитического контура" +GEN_MODEL_T0="${GEN_MODEL_T0}" \ +GEN_MODEL_T_END="${GEN_MODEL_T_END}" \ +COMPOSE_BIN="${COMPOSE_BIN}" \ +bash "${SCRIPT_DIR}/check_generated_analytics.sh" + +echo "" +echo "Готово: стартовая история генератора доведена до DM и Superset."