feat(stand): переведён аналитический путь на стартовую историю

- Зачем:
  - чистый стенд должен строить аналитику из генерации, а не из архивного JSONL-сида.
- Что:
  - добавлены команды generated-history-analytics и generated-history-check.
  - обновлены README, operations и Superset-документы под путь generator backfill -> DM -> Superset.
  - создан follow-up на миграцию учебных материалов с архивного сида.
- Проверка:
  - make generated-history-analytics.
  - make generated-history-check.
  - reviewer gate issue 06 пройден без блокирующих находок.
This commit is contained in:
2026-06-14 20:22:17 +03:00
parent d5408f28e9
commit d35254ead6
11 changed files with 909 additions and 118 deletions
@@ -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.
@@ -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`.
+9
View File
@@ -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 "=== Перезагрузка сервисов мониторинга ==="
+26 -31
View File
@@ -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` и в
+62 -36
View File
@@ -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` использовать только для ручных экспериментов.
+9 -6
View File
@@ -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
+19 -30
View File
@@ -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
```
### Нет данных в чартах
+26 -3
View File
@@ -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
+421
View File
@@ -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}"
+3 -2
View File
@@ -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
+106
View File
@@ -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."