diff --git a/.gitignore b/.gitignore index d4c6094..7775eb3 100644 --- a/.gitignore +++ b/.gitignore @@ -171,3 +171,4 @@ $RECYCLE.BIN/ *.msm *.msp *.lnk +node_modules diff --git a/Dockerfile.superset b/Dockerfile.superset index 11b021c..e8608b3 100644 --- a/Dockerfile.superset +++ b/Dockerfile.superset @@ -3,9 +3,9 @@ FROM apache/superset:4.1.2-dev USER root RUN apt-get update && \ - apt-get install -y bind9-host bat less iputils-ping curl && \ + apt-get install -y bind9-host bat less iputils-ping curl postgresql-client && \ apt-get clean USER superset -RUN pip install clickhouse-connect==0.8.* +RUN pip install clickhouse-connect==0.8.* clickhouse-sqlalchemy==0.3.* psycopg2-binary==2.9.* diff --git a/Makefile b/Makefile index 8b066ef..db375a0 100644 --- a/Makefile +++ b/Makefile @@ -1,7 +1,13 @@ -.PHONY: up down clean ddl data transform reload-monitoring recover-monitoring +.PHONY: up down clean ddl data transform logs \ + reload-monitoring recover-monitoring \ + superset-init superset-dashboard superset-export superset-ui superset-restart COMPOSE ?= docker compose +# ============================================================================ +# Основные команды +# ============================================================================ + up: $(COMPOSE) up -d @@ -13,6 +19,13 @@ down: clean: $(COMPOSE) down -v --remove-orphans +logs: + $(COMPOSE) logs -f --tail=200 $(service) + +# ============================================================================ +# ETL Pipeline +# ============================================================================ + ddl: bash ./scripts/apply_clickhouse_ddl.sh @@ -49,3 +62,28 @@ recover-monitoring: @echo "=== Проверка targets ===" @curl -s http://localhost:9090/api/v1/targets | grep -o '"job":"[^"]*"' | sort | uniq @curl -s http://localhost:9090/api/v1/targets | grep -o '"health":"[^"]*"' | sort | uniq -c + +# ============================================================================ +# Superset команды +# ============================================================================ + +# Инициализация Superset (подключение к ClickHouse + датасеты) +superset-init: + $(COMPOSE) exec -T superset bash -c "python /app/superset_init/init_superset.py" + +# Создание дашборда с чартами +superset-dashboard: + $(COMPOSE) exec -T superset bash -c "python /app/superset_init/create_dashboard.py" + +# Экспорт дашборда в JSON +superset-export: + $(COMPOSE) exec -T superset bash -c "python /app/superset_init/export_dashboard.py" + +# Открыть Superset UI +superset-ui: + @echo "Superset доступен по адресу: http://localhost:8088" + @echo "Логин: admin / Пароль: admin" + +# Перезапуск Superset +superset-restart: + $(COMPOSE) restart superset diff --git a/README.md b/README.md index aa4f4e8..cd30275 100644 --- a/README.md +++ b/README.md @@ -66,14 +66,16 @@ docker compose exec -T clickhouse clickhouse-client --user=default --password=12 ## Доступные сервисы -| Сервис | URL | Назначение | -|--------|-----|------------| -| ClickHouse HTTP | http://localhost:9123/play | SQL-запросы | -| Kafka UI | http://localhost:8082 | Просмотр топиков | -| Airflow | http://localhost:8080 | Оркестрация ETL (admin/admin) | -| Superset | http://localhost:8088 | BI-дашборды | -| Prometheus | http://localhost:9090 | Метрики | -| Grafana | http://localhost:3000 | Визуализация метрик | +| Сервис | URL | Назначение | Логин/Пароль | +|--------|-----|------------|--------------| +| ClickHouse HTTP | http://localhost:9123/play | SQL-запросы | default/123456 | +| Kafka UI | http://localhost:8082 | Просмотр топиков | — | +| Airflow | http://localhost:8080 | Оркестрация ETL | admin/admin | +| Superset | http://localhost:8088 | BI-дашборды | admin/admin | +| Prometheus | http://localhost:9090 | Метрики | — | +| Grafana | http://localhost:3000 | Визуализация метрик | admin/admin | + +Superset: после `make up` готовый дашборд доступен по адресу http://localhost:8088/superset/dashboard/1/ --- @@ -149,15 +151,22 @@ flowchart LR │ ├── ods/ # Batch SQL: STG -> ODS │ ├── dds/ # Batch SQL: ODS -> DDS │ └── dm/ # Batch SQL: DDS -> DM +├── configs/ # Конфигурации сервисов +│ ├── superset_config.py # Конфиг Superset (PostgreSQL metadata) +│ ├── prometheus/ # Prometheus конфигурация +│ └── grafana/ # Grafana dashboards & datasources ├── scripts/ # Служебные shell-скрипты (legacy fallback, не основной путь) ├── airflow/ # Конфигурация Airflow │ ├── dags/ # Airflow DAGs для оркестрации │ └── requirements.txt +├── superset/ # Скрипты инициализации Superset +│ ├── init_superset.py # Подключение к ClickHouse + датасеты +│ └── create_dashboard.py # Создание дашборда с чартами ├── docs/ # Документация │ └── ARCHITECTURE.md # Подробное описание слоёв ├── data/ # Исходные JSONL файлы ├── docker-compose.yml -└── Makefile # Команды: up, ddl, transform +└── Makefile # Команды: up, ddl, transform, superset-* ``` --- @@ -169,6 +178,10 @@ flowchart LR | `make up` | Поднять инфраструктуру | | `make ddl` | Применить DDL в ClickHouse (вне Airflow) | | `make transform` | Запустить batch-процесс `STG -> ODS -> DDS -> DM` (вне Airflow) | +| `make superset-init` | Подключение к ClickHouse + импорт датасетов | +| `make superset-dashboard` | Создание дашборда с чартами | +| `make superset-ui` | Показать URL Superset | +| `make superset-restart` | Перезапуск Superset | Примечания про сохранность данных: - Данные ClickHouse сохраняются в Docker volume `clickhouse-data`. @@ -215,16 +228,69 @@ flowchart LR ## Дашборд в Superset (опционально, но полезно) -1. Открыть `http://localhost:8088` -2. Database -> Add: - - URI: `clickhouse+connect://default:123456@clickhouse:8123/default` -3. Создать datasets из `dm.v_*` (VIEW) и собрать несколько графиков +Superset развёрнут с автоматической инициализацией: подключение к ClickHouse, датасеты и дашборд создаются автоматически при первом запуске. -Идеи графиков под задание: -- Трафик по дням: `dm.v_daily_traffic` (events, uniq_users) -- Эффективность UTM: `dm.v_utm_effectiveness` (clicks, purchases) -- Популярные страницы: `dm.v_top_pages_daily` (pageviews) -- Качество данных: `dm.v_dq_errors_daily` (rows_cnt по error_code) +### Быстрый доступ + +| URL | Назначение | Логин/Пароль | +|-----|------------|--------------| +| http://localhost:8088 | Superset UI | admin/admin | +| http://localhost:8088/superset/dashboard/1/ | Готовый дашборд | — | + +### Автоматическая инициализация (рекомендуется) + +```bash +# При первом запуске инфраструктуры +make up + +# Дашборд создаётся автоматически через 30-60 секунд +# Проверить готовность: +curl http://localhost:8088/health # должно вернуть 200 +``` + +Что создаётся автоматически: +- **Подключение к ClickHouse**: `clickhouse_dwh` (URI: `clickhousedb://default:123456@clickhouse:8123/default`) +- **Датасеты** (6 шт.): `v_events_enriched`, `v_daily_traffic`, `v_utm_effectiveness`, `v_top_pages_daily`, `v_session_overview`, `dq_summary` +- **Чарты** (10 шт.): KPI метрики, графики трафика, география, UTM-эффективность, качество данных +- **Дашборд**: "🛒 E-commerce Analytics Dashboard" + +### Ручная инициализация (если автоматика не сработала) + +```bash +# Подключение к ClickHouse + датасеты +make superset-init + +# Создание дашборда с чартами +make superset-dashboard +``` + +### Структура дашборда + +Дашборд "E-commerce Analytics Dashboard" включает: + +| Блок | Чарты | Датасет | +|------|-------|---------| +| **KPI** | Total Events, Unique Users, Unique Sessions, Avg Events/Session | `v_events_enriched` | +| **Динамика** | Events by Hour (timeline), Traffic by Device (pie) | `v_events_enriched` | +| **География** | World Map по странам | `v_events_enriched` | +| **Маркетинг** | UTM Effectiveness Table, Top Pages | `v_utm_effectiveness`, `v_top_pages_daily` | +| **Quality** | Data Quality Summary | `dq_summary` | + +### Архитектура Superset + +``` +┌─────────────────┐ ┌──────────────────┐ ┌─────────────────┐ +│ Superset UI │────▶│ PostgreSQL │────▶│ ClickHouse │ +│ (localhost) │ │ (metadata) │ │ (данные) │ +│ :8088 │ │ dashboards, │ │ dm.v_* VIEW │ +└─────────────────┘ │ datasets, charts│ └─────────────────┘ + └──────────────────┘ +``` + +Особенности конфигурации: +- **Metadata**: PostgreSQL (shared с Airflow) — данные сохраняются при перезапуске +- **Data**: ClickHouse через `clickhouse-connect` (HTTP порт 8123) +- **Config**: `configs/superset_config.py` (PostgreSQL URI, секретный ключ) --- @@ -257,6 +323,8 @@ flowchart LR - Подключения используют разные протоколы: - Airflow (ClickHouseOperator) ходит в ClickHouse по native TCP (порт `9000` внутри сети Docker). - Superset (clickhouse-connect) ходит по HTTP (порт `8123` внутри сети Docker). +- **Superset**: дашборд не появился сразу — подождите 30-60 секунд после `make up`, затем проверьте `curl http://localhost:8088/health`. +- **Superset**: при полном сбросе (`docker compose down -v`) метаданные Superset пропадут т.к. используется общая PostgreSQL. Для чистого перезапуска Superset удалите только БД `superset` в PostgreSQL и перезапустите контейнеры. --- @@ -267,6 +335,7 @@ flowchart LR - **Ingest**: DAG `ddl_init` (DDL + проверка схемы), DAG `kafka_load` (параметры `limit`, `reset_topics`) - **Трансформации**: DAG `etl_pipeline` (pre-check, batch STG→ODS→DDS→DM, валидация) - **Витрины**: VIEW в DM для бизнес-дашбордов (трафик, UTM, качество данных) +- **Superset**: Автоматическая инициализация (подключение ClickHouse, 6 датасетов, 10 чартов, дашборд), метаданные в PostgreSQL - **Надёжность**: ошибки парсинга сохраняются в ODS, пайплайн не падает на "грязных" данных - **Мониторинг**: Prometheus скрейпит ClickHouse метрики, Grafana дашборд и alert rules diff --git a/configs/superset_config.py b/configs/superset_config.py new file mode 100644 index 0000000..761e52d --- /dev/null +++ b/configs/superset_config.py @@ -0,0 +1,29 @@ +# Superset configuration for PostgreSQL metadata store +import os + +# Database URI for PostgreSQL +SQLALCHEMY_DATABASE_URI = 'postgresql://airflow:airflow@postgres-metadata:5432/superset' + +# Secret key (should match docker-compose) +SECRET_KEY = os.getenv('SUPERSET_SECRET_KEY', '9wc5+erMt60+lxrXDf3RjeIR+zONpEFusO00Np7JzfliMTI1e+RXnHcQ') + +# Disable debug mode +DEBUG = False + +# Enable CSRF protection +WTF_CSRF_ENABLED = True + +# Session configuration +SESSION_TYPE = 'filesystem' +SESSION_COOKIE_SECURE = False +SESSION_COOKIE_HTTPONLY = True +SESSION_COOKIE_SAMESITE = 'Lax' + +# Cache configuration (optional, using simple cache) +CACHE_CONFIG = { + 'CACHE_TYPE': 'SimpleCache', + 'CACHE_DEFAULT_TIMEOUT': 300 +} + +# Timezone +DEFAULT_TIMEZONE = 'Europe/Moscow' diff --git a/docker-compose.yml b/docker-compose.yml index f4dd504..7054a79 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -216,7 +216,7 @@ services: superset: # image: apache/superset - build: + build: context: . dockerfile: Dockerfile.superset restart: unless-stopped @@ -227,14 +227,60 @@ services: volumes: - superset_data:/var/lib/superset - superset_config:/app/superset_home + - ./superset:/app/superset_init:ro + - ./configs/superset_config.py:/app/pythonpath/superset_config.py:ro environment: - SUPERSET_SECRET_KEY=9wc5+erMt60+lxrXDf3RjeIR+zONpEFusO00Np7JzfliMTI1e+RXnHcQ - TZ=Europe/Moscow + - SUPERSET_CONFIG_PATH=/app/pythonpath/superset_config.py healthcheck: test: ["CMD", "curl", "-f", "http://localhost:8088/health"] interval: 30s timeout: 10s retries: 5 + depends_on: + superset-init: + condition: service_completed_successfully + + # Superset initialization service + superset-init: + build: + context: . + dockerfile: Dockerfile.superset + user: "0:0" + networks: + - cs_dwh + volumes: + - superset_data:/var/lib/superset + - superset_config:/app/superset_home + - ./superset:/app/superset_init:ro + - ./configs/superset_config.py:/app/pythonpath/superset_config.py:ro + environment: + - SUPERSET_SECRET_KEY=9wc5+erMt60+lxrXDf3RjeIR+zONpEFusO00Np7JzfliMTI1e+RXnHcQ + - TZ=Europe/Moscow + - SUPERSET_CONFIG_PATH=/app/pythonpath/superset_config.py + command: > + bash -ceuo pipefail " + echo 'Waiting for PostgreSQL...' && + sleep 5 && + echo 'Creating superset database if not exists...' && + PGPASSWORD=airflow psql -h postgres-metadata -U airflow -d airflow -tc \"SELECT 1 FROM pg_database WHERE datname='superset'\" | grep -q 1 || \\ + PGPASSWORD=airflow psql -h postgres-metadata -U airflow -d airflow -c \"CREATE DATABASE superset;\" && + echo 'Initializing Superset DB...' && + superset db upgrade && + echo 'Creating admin user...' && + superset fab create-admin --username admin --password admin --firstname Superset --lastname Admin --email admin@example.org || true && + echo 'Initializing roles...' && + superset init && + echo 'Creating ClickHouse connection and datasets...' && + python /app/superset_init/init_superset.py && + echo 'Creating dashboard and charts...' && + python /app/superset_init/create_dashboard.py && + echo 'Superset initialized successfully' + " + depends_on: + postgres-metadata: + condition: service_healthy # Kafka Exporter для мониторинга через Prometheus kafka-exporter: diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 56b7362..4630e71 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -602,7 +602,7 @@ INSERT INTO dm.daily_traffic SELECT * FROM dm.v_daily_traffic; - URL: `clickhouse://default:123456@clickhouse:9000/default` (native TCP для Airflow plugin) - Provider/интеграция: `airflow-clickhouse-plugin` (в `airflow/requirements.txt`), задачи выполняются через `ClickHouseOperator`. - Дополнительно: `kafka-python==2.0.6` для работы с Kafka из DAG. -- Примечание: Superset подключается к ClickHouse по HTTP (обычно `clickhouse+connect://...:8123/...`). +- Примечание: Superset подключается к ClickHouse по HTTP (обычно `clickhousedb://...:8123/...`). --- diff --git a/docs/DEMO_SCRIPT_10_15MIN.md b/docs/DEMO_SCRIPT_10_15MIN.md new file mode 100644 index 0000000..16a82c3 --- /dev/null +++ b/docs/DEMO_SCRIPT_10_15MIN.md @@ -0,0 +1,182 @@ +# Сценарий демо на 10-15 минут (код + архитектура) + +Цель: дать студенту готовый сценарий, который показывает не только запуск стенда, но и инженерные решения в коде. + +Формат: запись экрана + голос. +Ориентир по времени: 12 минут (допуск 10-15). + +--- + +## 1) Подготовка перед записью (за 10-20 минут) + +```bash +make up +docker compose ps +docker compose exec -T airflow-webserver airflow dags trigger ddl_init +docker compose exec -T airflow-webserver airflow dags trigger kafka_load --conf '{"limit": 50, "reset_topics": true}' +docker compose exec -T airflow-webserver airflow dags trigger etl_pipeline --conf '{"full_refresh": true}' +``` + +Проверить, что открываются: +- Kafka UI: `http://localhost:8082` +- Grafana: `http://localhost:3000/d/clickhouse-overview/clickhouse-overview` +- Airflow: `http://localhost:8080/dags/ddl_init/grid?tab=details` +- Superset (если используете): `http://localhost:8088/login/?next=/` +- ClickHouse Play: `http://localhost:9123/play` + +Подготовить вкладки заранее: +- `Makefile` +- `docker-compose.yml` +- `docs/ARCHITECTURE.md` +- `airflow/dags/ddl_init_dag.py` +- `airflow/dags/kafka_load_dag.py` +- `airflow/dags/etl_pipeline_dag.py` +- `sql/ddl/stg/10_stg.sql` +- `sql/ods/20_stg_to_ods.sql` +- `sql/dds/30_ods_to_dds.sql` +- `sql/dm/40_dds_to_dm.sql` + +Опционально подготовить DBeaver (если хотите показывать не через Play): +- Host: `localhost` +- Port: `9123` (HTTP) или `8002` (native) +- User: `default` +- Password: `123456` + +--- + +## 2) Поминутный план выступления + +### 0:00-1:30 Инфраструктура и цель проекта + +Что показывать: +- Терминал с `docker compose ps` +- `Makefile` +- `docker-compose.yml` + +Что говорить: +- «Это учебный mini DWH для кликстрима: Kafka, ClickHouse, Airflow, Superset, Prometheus, Grafana.» +- «Инфраструктура поднимается одной командой `make up`; внутри это `docker compose up -d`.» +- «В `Makefile` также есть команды для остановки, очистки, перезагрузки мониторинга и recovery.» +- «Сервисная цель проекта: быстро и повторяемо показать end-to-end поток данных до витрин.» + +Что подчеркнуть в коде: +- В `Makefile` показать цели `up/down/clean/reload-monitoring/recover-monitoring`. +- В `docker-compose.yml` бегло показать ключевые сервисы и порты. + +### 1:30-3:30 Архитектура и логика выбора + +Что показывать: +- `docs/ARCHITECTURE.md` (диаграммы потока, слои STG/ODS/DDS/DM). + +Что говорить: +- «Управление сделано через 3 DAG: `ddl_init`, `kafka_load`, `etl_pipeline`.» +- «STG нужен для сырых событий как есть, чтобы сохранять воспроизводимость.» +- «ODS типизирует и валидирует данные, включая фиксацию ошибок парсинга.» +- «DDS собирает бизнес-сущности `event` и `click` для аналитики.» +- «DM отдает витрины и агрегаты для BI и интервью-демо.» + +Объяснение решений: +- «Разделение на слои уменьшает связность и ускоряет диагностику проблем.» +- «Грязные данные не останавливают пайплайн: ошибки уходят в `ods.*_errors` и DQ-слой.» + +### 3:30-6:30 Показ кода DAG-ов + +Что показывать: +- `airflow/dags/ddl_init_dag.py` +- `airflow/dags/kafka_load_dag.py` +- `airflow/dags/etl_pipeline_dag.py` + +Что говорить: +- «В `ddl_init` код разворачивает DDL в ClickHouse и подготавливает структуру слоев.» +- «В `kafka_load` есть управляемые параметры `limit` и `reset_topics` для быстрого smoke-прогона.» +- «В `etl_pipeline` выполняются шаги STG->ODS->DDS->DM с конфигурацией `full_refresh`; итоговый DM-блок здесь — загрузка `dm.dq_summary`.» +- «Логика запуска ручная: это удобно для демонстрации на собеседовании и для отладки.» + +Что обязательно назвать: +- «Почему `limit=50` в демо: скорость и повторяемость важнее полноты.» +- «Почему DAG-и разделены: проще локализовать сбой и перезапустить только нужный этап.» + +### 6:30-8:30 Показ SQL и модели данных + +Что показывать: +- `sql/ddl/stg/10_stg.sql` +- `sql/ods/20_stg_to_ods.sql` +- `sql/dds/30_ods_to_dds.sql` +- `sql/dm/40_dds_to_dm.sql` + +Что говорить: +- «В STG используется связка Kafka Engine + Materialized View + MergeTree таблицы.» +- «ODS делает типизацию, нормализацию и отправку проблемных строк в таблицы ошибок.» +- «DDS собирает сущности по ключам (`event_id`, `click_id`), чтобы упростить аналитику.» +- «Витрины DM строятся поверх DDS и готовы для BI.» + +Короткий акцент на DQ: +- «Вместо падения на невалидном JSON сохраняем ошибку и продолжаем обработку потока.» + +### 8:30-10:30 Прогон в Airflow + проверка результата + +Что показывать: +- Airflow UI: последний `Success` у `ddl_init`, `kafka_load`, `etl_pipeline` +- ClickHouse Play или DBeaver + +Что выполнять: + +```sql +SELECT count() AS rows FROM stg.browser_raw; +SELECT count() AS rows FROM ods.browser_event; +SELECT count() AS rows FROM dds.event; +SELECT * FROM dm.v_daily_traffic ORDER BY event_date DESC LIMIT 10; +SELECT * FROM dm.dq_summary ORDER BY layer, table_name, check_name LIMIT 20; +``` + +Что говорить: +- «На экране видно прохождение данных по слоям и непустые витрины.» +- «DQ summary подтверждает контроль качества и обработку проблемных записей.» + +### 10:30-12:00 Мониторинг и финал + +Что показывать: +- Grafana: ClickHouse/Kafka/Airflow dashboards +- Prometheus targets (опционально) +- Superset dashboard (если подготовлен) + +Что говорить: +- «Мониторинг показывает здоровье стенда и ключевые технические метрики.» +- «На BI-слое уже можно отвечать на базовые бизнес-вопросы по трафику и UTM.» +- «Итог: решение покрывает инфраструктуру, ingestion, трансформации, DQ, витрины и observability.» + +--- + +## 3) План Б, если что-то сломалось на записи + +Если не открывается UI: + +```bash +docker compose ps +docker compose logs -f --tail=100 airflow-webserver +docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query "SELECT count() FROM dds.event" +curl -s http://localhost:9090/api/v1/targets | grep -o '"health":"[^"]*"' +``` + +Если Grafana пустая: + +```bash +make reload-monitoring +``` + +Если мониторинг завис: + +```bash +make recover-monitoring +``` + +Короткая реплика: +- «Даже при проблемах UI я показываю проверку через CLI и SQL, чтобы подтвердить работоспособность пайплайна.» + +--- + +## 4) Готовый текст финала (20-30 секунд) + +«Я реализовал end-to-end mini DWH для кликстрима: от инфраструктуры и ingestion до витрин и мониторинга. +Архитектура послойная STG-ODS-DDS-DM, orchestration через Airflow DAG-и, а невалидные данные фиксируются без падения пайплайна. +Если нужно, могу детально разобрать любой уровень: DAG-код, SQL-трансформации, DQ-проверки или наблюдаемость системы.» diff --git a/docs/SUPERSET_DASHBOARD.md b/docs/SUPERSET_DASHBOARD.md new file mode 100644 index 0000000..2fe6a04 --- /dev/null +++ b/docs/SUPERSET_DASHBOARD.md @@ -0,0 +1,262 @@ +# Дашборд Superset для E-commerce Analytics + +Документация по настройке и использованию Superset дашборда для анализа кликстрима. + +--- + +## Быстрый старт + +### 1. Запуск инфраструктуры + +```bash +# Запуск всех сервисов +make up + +# Применение DDL в ClickHouse +make ddl + +# Загрузка данных в Kafka +make data + +# Запуск ETL-пайплайна (ODS → DDS → DM) +make transform +``` + +### 2. Инициализация Superset + +```bash +# Автоматическая инициализация (создание подключения и датасетов) +make superset-init + +# Создание дашборда с чартами +make superset-dashboard +``` + +### 3. Доступ к UI + +Откройте в браузере: http://localhost:8088 + +**Логин:** `admin` +**Пароль:** `admin` + +--- + +## Структура дашборда + +### Витрины данных (Datasets) + +| Витрина | Таблица ClickHouse | Описание | +|---------|-------------------|------------| +| **Events Enriched** | `dm.v_events_enriched` | Полная обогащённая витрина событий | +| **Daily Traffic** | `dm.v_daily_traffic` | Агрегаты по дням | +| **UTM Effectiveness** | `dm.v_utm_effectiveness` | Эффективность маркетинговых каналов | +| **Top Pages** | `dm.v_top_pages_daily` | Популярность страниц | +| **Session Overview** | `dm.v_session_overview` | Анализ сессий | +| **DQ Summary** | `dm.dq_summary` | Качество данных | + +### Чарты (Charts) + +#### KPI-блок (верх дашборда) +- **📊 Total Events** — общее количество событий +- **👤 Unique Users** — уникальные пользователи +- **🎯 Unique Sessions** — уникальные сессии (click_id) +- **📈 Avg Events/Session** — среднее количество событий на сессию + +#### Динамика трафика +- **📅 Events by Hour** — линейный график событий по часам +- **📱 Traffic by Device** — pie chart распределения по устройствам + +#### География +- **🌍 Geography Map** — world map с распределением по странам + +#### Маркетинг +- **🔗 UTM Effectiveness Table** — таблица эффективности UTM-меток +- **📄 Top Pages** — bar chart топ-20 страниц + +#### Качество данных +- **🔍 Data Quality Summary** — статистика по слоям STG/ODS/DDS + +### Фильтры (Native Filters) + +| Фильтр | Поле | Тип | Применение | +|--------|------|-----|------------| +| 📅 Date Range | `event_date` | Time Range | Все чарты | +| 🌍 Country | `geo_country` | Multi-select | Все чарты | +| 📱 Device Type | `device_type` | Multi-select | Все чарты | +| 🌐 Browser | `browser_name` | Multi-select | Все чарты | + +--- + +## Команды Makefile + +```bash +# Основные +make up # Запуск всех сервисов +make down # Остановка сервисов +make clean # Остановка с удалением volumes +make logs service=superset # Логи сервиса + +# ETL +make ddl # Применение DDL в ClickHouse +make data # Загрузка данных в Kafka +make transform # Запуск batch-процесса + +# Superset +make superset-init # Инициализация (подключение + датасеты) +make superset-dashboard # Создание дашборда +make superset-export # Экспорт дашборда в JSON +make superset-ui # Показать URL и логин +make superset-restart # Перезапуск сервиса +``` + +--- + +## Ручная настройка (если автоматика не сработала) + +### Создание подключения к ClickHouse + +1. Откройте **Settings → Database Connections** +2. Нажмите **+ Database** +3. Выберите **ClickHouse** +4. Введите SQLAlchemy URI: + ``` + clickhousedb://default:123456@clickhouse:8123/default + ``` +5. Установите: + - **Expose in SQL Lab:** ✅ + - **Allow DDL:** ❌ +6. Нажмите **Connect** + +### Импорт датасетов + +```bash +# Внутри контейнера +docker compose exec superset bash +python /app/superset_init/init_superset.py +``` + +### Создание чартов вручную + +1. Перейдите в **Charts → + Chart** +2. Выберите датасет (например, `dm.v_events_enriched`) +3. Настройте визуализацию: + - **Viz Type:** Big Number / Line Chart / Pie Chart / World Map / Table + - **Metrics:** COUNT(*), COUNT(DISTINCT ...) + - **Dimensions:** группировки + - **Filters:** фильтры +4. Нажмите **Create Chart** + +### Создание дашборда + +1. **Dashboards → + Dashboard** +2. Назовите: "E-commerce Analytics Dashboard" +3. Добавьте чарты из списка +4. Настройте layout (drag-and-drop) +5. Добавьте Native Filters (фильтры вверху) +6. Сохраните + +--- + +## Экспорт и импорт дашборда + +### Экспорт + +```bash +# Автоматический экспорт в JSON +make superset-export + +# Результат: superset/dashboards/ecommerce_analytics.json +``` + +### Импорт + +```bash +# Импорт через CLI +docker compose exec superset superset import-dashboards -p /app/superset_init/dashboards/ecommerce_analytics.json + +# Или через UI: Settings → Import Dashboards +``` + +--- + +## Расширение дашборда + +### Добавление нового чарта + +1. Отредактируйте `superset/create_dashboard.py` +2. Добавьте конфигурацию в `CHARTS_CONFIG` +3. Запустите: `make superset-dashboard` + +Пример нового чарта: +```python +{ + "slice_name": "📊 My New Chart", + "viz_type": "echarts_bar", + "dataset_name": "v_events_enriched", + "params": { + "x_axis": "event_type", + "metrics": [{"sqlExpression": "COUNT(*)", "label": "Count"}], + "time_range": "No filter" + } +} +``` + +--- + +## Troubleshooting + +### Superset не стартует + +```bash +# Проверить логи +make logs service=superset + +# Перезапуск +make superset-restart + +# Полная переинициализация +docker compose down -v +docker compose up -d +make superset-init +``` + +### Нет данных в чартах + +```bash +# Проверить данные в ClickHouse +docker compose exec clickhouse clickhouse-client -q "SELECT count() FROM dm.v_events_enriched" + +# Перезапустить ETL +make transform +``` + +### Ошибка подключения к ClickHouse + +```bash +# Проверить доступность ClickHouse +docker compose exec superset bash -c "ping clickhouse" + +# Проверить порт +docker compose exec superset bash -c "curl clickhouse:8123" +``` + +--- + +## Порты сервисов + +| Сервис | URL | Логин/Пароль | +|--------|-----|--------------| +| Superset | http://localhost:8088 | admin / admin | +| ClickHouse HTTP | http://localhost:9123 | default / (пустой) | +| Airflow | http://localhost:8080 | admin / admin | +| Grafana | http://localhost:3000 | admin / admin | +| Prometheus | http://localhost:9090 | - | +| Kafka UI | http://localhost:8082 | - | + +--- + +## Дополнительные ресурсы + +- [Superset Documentation](https://superset.apache.org/docs/intro) +- [ClickHouse SQL Reference](https://clickhouse.com/docs/en/sql-reference) +- [ARCHITECTURE.md](./ARCHITECTURE.md) — архитектура хранилища diff --git a/superset/create_dashboard.py b/superset/create_dashboard.py new file mode 100644 index 0000000..7d9e0ed --- /dev/null +++ b/superset/create_dashboard.py @@ -0,0 +1,549 @@ +#!/usr/bin/env python3 +""" +================================================================================ +Скрипт создания дашборда "E-commerce Analytics" в Superset +================================================================================ +Назначение: + - Создание чартов (Charts) на основе датасетов DM-слоя + - Создание дашборда с layout и native filters + +Запуск: + Внутри контейнера superset: + python /app/superset_init/create_dashboard.py +================================================================================ +""" + +import os +import sys +import json +import logging +from datetime import datetime + +logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s') +logger = logging.getLogger(__name__) + +sys.path.insert(0, '/app') + + +# Конфигурация чартов +CHARTS_CONFIG = [ + # KPI блок + { + "slice_name": "📊 Total Events", + "viz_type": "big_number", + "dataset_name": "v_events_enriched", + "params": { + "metric": { + "expressionType": "SQL", + "sqlExpression": "COUNT(*)", + "column": None, + "aggregate": None, + "label": "Total Events", + "optionName": "metric_1" + }, + "granularity_sqla": "event_ts", + "y_axis_format": ",d", + "show_trend_line": False, + "time_range": "No filter" + } + }, + { + "slice_name": "👤 Unique Users", + "viz_type": "big_number", + "dataset_name": "v_events_enriched", + "params": { + "metric": { + "expressionType": "SQL", + "sqlExpression": "COUNT(DISTINCT user_domain_id)", + "label": "Unique Users", + "optionName": "metric_2" + }, + "granularity_sqla": "event_ts", + "y_axis_format": ",d", + "show_trend_line": False, + "time_range": "No filter" + } + }, + { + "slice_name": "🎯 Unique Sessions", + "viz_type": "big_number", + "dataset_name": "v_events_enriched", + "params": { + "metric": { + "expressionType": "SQL", + "sqlExpression": "COUNT(DISTINCT click_id)", + "label": "Unique Sessions", + "optionName": "metric_3" + }, + "granularity_sqla": "event_ts", + "y_axis_format": ",d", + "show_trend_line": False, + "time_range": "No filter" + } + }, + { + "slice_name": "📈 Avg Events/Session", + "viz_type": "big_number", + "dataset_name": "v_events_enriched", + "params": { + "metric": { + "expressionType": "SQL", + "sqlExpression": "COUNT(*) / COUNT(DISTINCT click_id)", + "label": "Avg Events/Session", + "optionName": "metric_4" + }, + "granularity_sqla": "event_ts", + "y_axis_format": ".2f", + "show_trend_line": False, + "time_range": "No filter" + } + }, + # Динамика + { + "slice_name": "📅 Events by Hour", + "viz_type": "echarts_timeseries_line", + "dataset_name": "v_events_enriched", + "params": { + "granularity_sqla": "event_ts", + "time_grain_sqla": "PT1H", + "metrics": [ + { + "expressionType": "SQL", + "sqlExpression": "COUNT(*)", + "label": "Events" + } + ], + "groupby": [], + "time_range": "Last week", + "adhoc_filters": [], + "row_limit": 10000 + } + }, + { + "slice_name": "📱 Traffic by Device", + "viz_type": "pie", + "dataset_name": "v_events_enriched", + "params": { + "groupby": ["device_type"], + "metric": { + "expressionType": "SQL", + "sqlExpression": "COUNT(*)", + "label": "Count" + }, + "row_limit": 100, + "donut": True, + "show_legend": True, + "labels_outside": True, + "time_range": "No filter" + } + }, + # География + { + "slice_name": "🌍 Geography Map", + "viz_type": "world_map", + "dataset_name": "v_events_enriched", + "params": { + "entity": "geo_country", + "metric": { + "expressionType": "SQL", + "sqlExpression": "COUNT(*)", + "label": "Events" + }, + "row_limit": 500, + "linear_color_scheme": "blue_white_yellow", + "time_range": "No filter" + } + }, + # Маркетинг + { + "slice_name": "🔗 UTM Effectiveness Table", + "viz_type": "table", + "dataset_name": "v_utm_effectiveness", + "params": { + "groupby": ["utm_source", "utm_medium", "utm_campaign"], + "metrics": [ + {"expressionType": "SQL", "sqlExpression": "SUM(clicks)", "label": "Clicks"}, + {"expressionType": "SQL", "sqlExpression": "SUM(uniq_users)", "label": "Users"}, + {"expressionType": "SQL", "sqlExpression": "SUM(uniq_sessions)", "label": "Sessions"} + ], + "row_limit": 100, + "time_range": "No filter", + "adhoc_filters": [ + { + "clause": "WHERE", + "expressionType": "SQL", + "sqlExpression": "utm_source IS NOT NULL", + "subject": None, + "operator": None, + "comparator": None + } + ] + } + }, + { + "slice_name": "📄 Top Pages", + "viz_type": "dist_bar", + "dataset_name": "v_top_pages_daily", + "params": { + "groupby": ["page_url_path"], + "metrics": [ + {"expressionType": "SQL", "sqlExpression": "SUM(pageviews)", "label": "Pageviews"} + ], + "row_limit": 20, + "time_range": "No filter", + "orientation": "vertical", + "show_legend": False + } + }, + # Качество данных + { + "slice_name": "🔍 Data Quality Summary", + "viz_type": "dist_bar", + "dataset_name": "dq_summary", + "params": { + "groupby": ["layer"], + "metrics": [ + {"expressionType": "SQL", "sqlExpression": "SUM(check_value)", "label": "Row Count"} + ], + "adhoc_filters": [ + { + "clause": "WHERE", + "expressionType": "SQL", + "sqlExpression": "check_name = 'total_rows'", + "subject": None, + "operator": None, + "comparator": None + } + ], + "row_limit": 100, + "time_range": "No filter", + "show_legend": False + } + } +] + +# Конфигурация дашборда +DASHBOARD_CONFIG = { + "dashboard_title": "🛒 E-commerce Analytics Dashboard", + "description": "Аналитический дашборд для e-commerce кликстрима: трафик, конверсии, география и качество данных.", + "published": True, + "slug": "ecommerce-analytics", +} + + +def sync_query_context(chart, params: dict, dataset_id: int) -> None: + """ + Синхронизирует сохраненный query_context с обновленными params чарта. + + Для Superset 4.x у `big_number` запрос валидируется как time-series и + ожидает `granularity` в query_context (проверено по актуальной документации). + """ + if not chart.query_context: + return + + try: + query_context = json.loads(chart.query_context) + except (TypeError, json.JSONDecodeError): + logger.warning("Chart ID %s has invalid query_context, skip sync", chart.id) + return + + query_context["datasource"] = {"id": dataset_id, "type": "table"} + query_context["form_data"] = { + **params, + "datasource": f"{dataset_id}__table", + "viz_type": chart.viz_type, + "slice_id": chart.id, + } + + queries = query_context.get("queries") + if not isinstance(queries, list) or not queries: + chart.query_context = json.dumps(query_context) + return + + if chart.viz_type == "big_number": + query = queries[0] + granularity = params.get("granularity_sqla") + if granularity: + query["granularity"] = granularity + query["is_timeseries"] = True + query["time_range"] = params.get("time_range", query.get("time_range")) + if "metric" in params: + query["metrics"] = [params["metric"]] + extras = query.get("extras") if isinstance(query.get("extras"), dict) else {} + if "time_grain_sqla" in params: + extras["time_grain_sqla"] = params.get("time_grain_sqla") + query["extras"] = extras + elif params.get("x_axis") or params.get("groupby"): + # Для категориальных графиков синхронизируем колонки измерений. + dimensions = params.get("groupby") + if not dimensions and params.get("x_axis"): + dimensions = [params["x_axis"]] + query = queries[0] + query["columns"] = dimensions + if "metrics" in params: + query["metrics"] = params["metrics"] + elif "metric" in params: + query["metrics"] = [params["metric"]] + query["row_limit"] = params.get("row_limit", query.get("row_limit")) + query["time_range"] = params.get("time_range", query.get("time_range")) + query["is_timeseries"] = False + + chart.query_context = json.dumps(query_context) + + +def build_dashboard_metadata(filter_dataset_id: int | None) -> str: + """Формирует json_metadata с валидными datasetId для native filters.""" + native_filters = [] + + if filter_dataset_id is not None: + native_filters = [ + { + "id": "date_filter", + "name": "📅 Date Range", + "filterType": "filter_time", + "targets": [{"datasetId": filter_dataset_id, "column": {"name": "event_date"}}], + "defaultValue": "Last week", + "scope": {"rootPath": ["ROOT_ID"], "excluded": []}, + "cascadeParentIds": [], + "isInstant": True + }, + { + "id": "country_filter", + "name": "🌍 Country", + "filterType": "filter_select", + "targets": [{"datasetId": filter_dataset_id, "column": {"name": "geo_country"}}], + "scope": {"rootPath": ["ROOT_ID"], "excluded": []}, + "isInstant": True, + "allowsMultipleValues": True, + "isRequired": False + }, + { + "id": "device_filter", + "name": "📱 Device Type", + "filterType": "filter_select", + "targets": [{"datasetId": filter_dataset_id, "column": {"name": "device_type"}}], + "scope": {"rootPath": ["ROOT_ID"], "excluded": []}, + "isInstant": True, + "allowsMultipleValues": True, + "isRequired": False + }, + { + "id": "browser_filter", + "name": "🌐 Browser", + "filterType": "filter_select", + "targets": [{"datasetId": filter_dataset_id, "column": {"name": "browser_name"}}], + "scope": {"rootPath": ["ROOT_ID"], "excluded": []}, + "isInstant": True, + "allowsMultipleValues": True, + "isRequired": False + } + ] + + metadata = { + "native_filter_configuration": native_filters, + "color_scheme": "supersetColors", + "label_colors": {} + } + return json.dumps(metadata) + + +def main() -> bool: + """Главная функция""" + logger.info("=" * 60) + logger.info("Creating E-commerce Analytics Dashboard") + logger.info("=" * 60) + + # Импорты внутри main после создания app context + from superset.app import create_app + + app = create_app() + + with app.app_context(): + from superset.extensions import db + from superset.models.slice import Slice + from superset.models.dashboard import Dashboard + from superset.connectors.sqla.models import SqlaTable + + created_charts = [] + datasets_by_name = {} + + # Создаём чарты + for chart_config in CHARTS_CONFIG: + dataset = db.session.query(SqlaTable).filter_by( + table_name=chart_config["dataset_name"], + schema="dm" + ).first() + + if not dataset: + logger.warning(f"Dataset '{chart_config['dataset_name']}' not found, skipping chart") + continue + datasets_by_name[chart_config["dataset_name"]] = dataset.id + + try: + # Подготавливаем параметры + params = chart_config["params"].copy() + params["datasource"] = f"{dataset.id}__table" + params["viz_type"] = chart_config["viz_type"] + serialized_params = json.dumps(params) + + # Проверяем, существует ли уже чарт + existing = db.session.query(Slice).filter_by( + slice_name=chart_config["slice_name"] + ).first() + + if existing: + # Синхронизируем параметры существующего чарта с конфигом. + existing.viz_type = chart_config["viz_type"] + existing.datasource_id = dataset.id + existing.datasource_type = "table" + existing.datasource_name = dataset.table_name + existing.params = serialized_params + sync_query_context(existing, params, dataset.id) + existing.description = f"Chart created automatically for {chart_config['dataset_name']}" + db.session.flush() + logger.info( + f"Chart '{chart_config['slice_name']}' already exists (ID: {existing.id}), " + "params synced" + ) + created_charts.append({"id": existing.id, "title": existing.slice_name}) + continue + + # Создаём чарт + chart = Slice( + slice_name=chart_config["slice_name"], + viz_type=chart_config["viz_type"], + datasource_id=dataset.id, + datasource_type="table", + datasource_name=dataset.table_name, + params=serialized_params, + description=f"Chart created automatically for {chart_config['dataset_name']}" + ) + + db.session.add(chart) + db.session.flush() + + logger.info(f"Created chart: {chart_config['slice_name']} (ID: {chart.id})") + created_charts.append({"id": chart.id, "title": chart.slice_name}) + + except Exception as e: + logger.error(f"Failed to create chart '{chart_config['slice_name']}': {e}") + import traceback + traceback.print_exc() + db.session.rollback() + + logger.info(f"Created/Found {len(created_charts)} charts") + metadata_json = build_dashboard_metadata(datasets_by_name.get("v_events_enriched")) + + # Создаём позиции чартов для layout. + # Обязательные блоки ROOT_ID/GRID_ID нужны для корректной работы /tabs. + positions = { + "DASHBOARD_VERSION_KEY": "v2", + "ROOT_ID": { + "id": "ROOT_ID", + "type": "ROOT", + "children": ["GRID_ID"], + }, + "GRID_ID": { + "id": "GRID_ID", + "type": "GRID", + "children": [], + "parents": ["ROOT_ID"], + "meta": {"background": "BACKGROUND_TRANSPARENT"}, + }, + } + + # Добавляем чарты в layout (grid: 12 columns) + y_position = 0 + chart_index = 0 + + for chart in created_charts: + if chart: + chart_component_id = f"CHART-{chart['id']}" + positions[chart_component_id] = { + "id": chart_component_id, + "type": "CHART", + "children": [], + "parents": ["ROOT_ID", "GRID_ID"], + "meta": { + "chartId": chart['id'], + "sliceName": chart['title'], + "height": 50, + "width": 4 if chart_index < 4 else 6, + "x": (chart_index % 3) * 4 if chart_index < 4 else (chart_index % 2) * 6, + "y": y_position, + }, + } + positions["GRID_ID"]["children"].append(chart_component_id) + chart_index += 1 + if chart_index % 4 == 0: + y_position += 50 + + # Создаём дашборд + if created_charts: + try: + # Проверяем, существует ли дашборд + existing = db.session.query(Dashboard).filter_by( + slug=DASHBOARD_CONFIG["slug"] + ).first() + + if existing: + existing.description = DASHBOARD_CONFIG["description"] + existing.published = DASHBOARD_CONFIG["published"] + existing.json_metadata = metadata_json + existing.position_json = json.dumps(positions) + existing.slices = [] + for chart_info in created_charts: + chart = db.session.query(Slice).filter_by(id=chart_info["id"]).first() + if chart: + existing.slices.append(chart) + db.session.commit() + logger.info(f"Dashboard '{DASHBOARD_CONFIG['dashboard_title']}' already exists (ID: {existing.id})") + logger.info("=" * 60) + logger.info("Dashboard already exists and metadata/layout were updated.") + logger.info(f"Dashboard URL: /superset/dashboard/{existing.id}/") + logger.info("=" * 60) + return True + + # Создаём дашборд + dashboard = Dashboard( + dashboard_title=DASHBOARD_CONFIG["dashboard_title"], + slug=DASHBOARD_CONFIG["slug"], + description=DASHBOARD_CONFIG["description"], + published=DASHBOARD_CONFIG["published"], + json_metadata=metadata_json, + position_json=json.dumps(positions) + ) + + db.session.add(dashboard) + db.session.flush() + + # Добавляем чарты к дашборду + for chart_info in created_charts: + if chart_info: + chart = db.session.query(Slice).filter_by(id=chart_info["id"]).first() + if chart: + dashboard.slices.append(chart) + + db.session.commit() + + logger.info(f"Created dashboard: {DASHBOARD_CONFIG['dashboard_title']} (ID: {dashboard.id})") + logger.info("=" * 60) + logger.info("Dashboard created successfully!") + logger.info(f"Dashboard URL: /superset/dashboard/{dashboard.id}/") + logger.info("=" * 60) + return True + + except Exception as e: + logger.error(f"Failed to create dashboard: {e}") + import traceback + traceback.print_exc() + db.session.rollback() + return False + else: + logger.error("No charts created, cannot create dashboard") + return False + + return False + +if __name__ == "__main__": + sys.exit(0 if main() else 1) diff --git a/superset/dashboards/ecommerce_analytics.zip.json b/superset/dashboards/ecommerce_analytics.zip.json new file mode 100644 index 0000000..ca643c0 --- /dev/null +++ b/superset/dashboards/ecommerce_analytics.zip.json @@ -0,0 +1,201 @@ +{ + "dashboards": [ + { + "__Dashboard__": { + "dashboard_title": "🛒 E-commerce Analytics Dashboard", + "description": "Аналитический дашборд для e-commerce кликстрима: трафик, конверсии, география и качество данных.", + "slug": "ecommerce-analytics", + "published": true, + "json_metadata": "{\"native_filter_configuration\": [{\"id\": \"date_filter\", \"name\": \"📅 Date Range\", \"filterType\": \"filter_time\", \"targets\": [{\"datasetId\": null, \"column\": {\"name\": \"event_date\"}}], \"defaultValue\": \"Last week\", \"scope\": {\"root\": [\"ROOT_ID\"], \"excluded\": []}, \"cascadeParentIds\": [], \"isInstant\": true}, {\"id\": \"country_filter\", \"name\": \"🌍 Country\", \"filterType\": \"filter_select\", \"targets\": [{\"datasetId\": null, \"column\": {\"name\": \"geo_country\"}}], \"scope\": {\"root\": [\"ROOT_ID\"], \"excluded\": []}, \"isInstant\": true, \"allowsMultipleValues\": true, \"isRequired\": false}, {\"id\": \"device_filter\", \"name\": \"📱 Device Type\", \"filterType\": \"filter_select\", \"targets\": [{\"datasetId\": null, \"column\": {\"name\": \"device_type\"}}], \"scope\": {\"root\": [\"ROOT_ID\"], \"excluded\": []}, \"isInstant\": true, \"allowsMultipleValues\": true, \"isRequired\": false}, {\"id\": \"browser_filter\", \"name\": \"🌐 Browser\", \"filterType\": \"filter_select\", \"targets\": [{\"datasetId\": null, \"column\": {\"name\": \"browser_name\"}}], \"scope\": {\"root\": [\"ROOT_ID\"], \"excluded\": []}, \"isInstant\": true, \"allowsMultipleValues\": true, \"isRequired\": false}], \"color_scheme\": \"supersetColors\", \"label_colors\": {}}", + "position_json": "{\"DASHBOARD_VERSION_KEY\": \"v2\", \"CHART-1\": {\"id\": \"CHART-1\", \"type\": \"CHART\", \"parents\": [\"ROOT_ID\"], \"meta\": {\"chartId\": 1, \"sliceName\": \"📊 Total Events\", \"height\": 50, \"width\": 4, \"x\": 0, \"y\": 0}}, \"CHART-2\": {\"id\": \"CHART-2\", \"type\": \"CHART\", \"parents\": [\"ROOT_ID\"], \"meta\": {\"chartId\": 2, \"sliceName\": \"👤 Unique Users\", \"height\": 50, \"width\": 4, \"x\": 4, \"y\": 0}}, \"CHART-3\": {\"id\": \"CHART-3\", \"type\": \"CHART\", \"parents\": [\"ROOT_ID\"], \"meta\": {\"chartId\": 3, \"sliceName\": \"🎯 Unique Sessions\", \"height\": 50, \"width\": 4, \"x\": 8, \"y\": 0}}, \"CHART-4\": {\"id\": \"CHART-4\", \"type\": \"CHART\", \"parents\": [\"ROOT_ID\"], \"meta\": {\"chartId\": 4, \"sliceName\": \"📈 Avg Events/Session\", \"height\": 50, \"width\": 4, \"x\": 0, \"y\": 50}}, \"CHART-5\": {\"id\": \"CHART-5\", \"type\": \"CHART\", \"parents\": [\"ROOT_ID\"], \"meta\": {\"chartId\": 5, \"sliceName\": \"📅 Events by Hour\", \"height\": 50, \"width\": 8, \"x\": 0, \"y\": 100}}, \"CHART-6\": {\"id\": \"CHART-6\", \"type\": \"CHART\", \"parents\": [\"ROOT_ID\"], \"meta\": {\"chartId\": 6, \"sliceName\": \"📱 Traffic by Device\", \"height\": 50, \"width\": 4, \"x\": 8, \"y\": 100}}, \"CHART-7\": {\"id\": \"CHART-7\", \"type\": \"CHART\", \"parents\": [\"ROOT_ID\"], \"meta\": {\"chartId\": 7, \"sliceName\": \"🌍 Geography Map\", \"height\": 50, \"width\": 6, \"x\": 0, \"y\": 150}}, \"CHART-8\": {\"id\": \"CHART-8\", \"type\": \"CHART\", \"parents\": [\"ROOT_ID\"], \"meta\": {\"chartId\": 8, \"sliceName\": \"🔗 UTM Effectiveness Table\", \"height\": 50, \"width\": 6, \"x\": 6, \"y\": 150}}, \"CHART-9\": {\"id\": \"CHART-9\", \"type\": \"CHART\", \"parents\": [\"ROOT_ID\"], \"meta\": {\"chartId\": 9, \"sliceName\": \"📄 Top Pages\", \"height\": 50, \"width\": 6, \"x\": 0, \"y\": 200}}, \"CHART-10\": {\"id\": \"CHART-10\", \"type\": \"CHART\", \"parents\": [\"ROOT_ID\"], \"meta\": {\"chartId\": 10, \"sliceName\": \"🔍 Data Quality Summary\", \"height\": 50, \"width\": 6, \"x\": 6, \"y\": 200}}}", + "slices": [1, 2, 3, 4, 5, 6, 7, 8, 9, 10] + } + } + ], + "charts": [ + { + "__Slice__": { + "slice_name": "📊 Total Events", + "viz_type": "big_number", + "datasource_type": "table", + "datasource_name": "dm.v_events_enriched", + "params": "{\"datasource\": \"1__table\", \"viz_type\": \"big_number\", \"metric\": {\"expressionType\": \"SQL\", \"sqlExpression\": \"COUNT(*)\", \"column\": null, \"aggregate\": null, \"label\": \"Total Events\", \"optionName\": \"metric_1\"}, \"y_axis_format\": \",d\", \"show_trend_line\": false, \"time_range\": \"No filter\"}", + "description": "Общее количество событий" + } + }, + { + "__Slice__": { + "slice_name": "👤 Unique Users", + "viz_type": "big_number", + "datasource_type": "table", + "datasource_name": "dm.v_events_enriched", + "params": "{\"datasource\": \"1__table\", \"viz_type\": \"big_number\", \"metric\": {\"expressionType\": \"SQL\", \"sqlExpression\": \"COUNT(DISTINCT user_domain_id)\", \"label\": \"Unique Users\", \"optionName\": \"metric_2\"}, \"y_axis_format\": \",d\", \"show_trend_line\": false, \"time_range\": \"No filter\"}", + "description": "Уникальные пользователи" + } + }, + { + "__Slice__": { + "slice_name": "🎯 Unique Sessions", + "viz_type": "big_number", + "datasource_type": "table", + "datasource_name": "dm.v_events_enriched", + "params": "{\"datasource\": \"1__table\", \"viz_type\": \"big_number\", \"metric\": {\"expressionType\": \"SQL\", \"sqlExpression\": \"COUNT(DISTINCT click_id)\", \"label\": \"Unique Sessions\", \"optionName\": \"metric_3\"}, \"y_axis_format\": \",d\", \"show_trend_line\": false, \"time_range\": \"No filter\"}", + "description": "Уникальные сессии" + } + }, + { + "__Slice__": { + "slice_name": "📈 Avg Events/Session", + "viz_type": "big_number", + "datasource_type": "table", + "datasource_name": "dm.v_events_enriched", + "params": "{\"datasource\": \"1__table\", \"viz_type\": \"big_number\", \"metric\": {\"expressionType\": \"SQL\", \"sqlExpression\": \"COUNT(*) / COUNT(DISTINCT click_id)\", \"label\": \"Avg Events/Session\", \"optionName\": \"metric_4\"}, \"y_axis_format\": \".2f\", \"show_trend_line\": false, \"time_range\": \"No filter\"}", + "description": "Среднее количество событий на сессию" + } + }, + { + "__Slice__": { + "slice_name": "📅 Events by Hour", + "viz_type": "echarts_timeseries_line", + "datasource_type": "table", + "datasource_name": "dm.v_events_enriched", + "params": "{\"datasource\": \"1__table\", \"viz_type\": \"echarts_timeseries_line\", \"granularity_sqla\": \"event_ts\", \"time_grain_sqla\": \"PT1H\", \"metrics\": [{\"expressionType\": \"SQL\", \"sqlExpression\": \"COUNT(*)\", \"label\": \"Events\"}], \"groupby\": [], \"time_range\": \"Last week\", \"adhoc_filters\": [], \"row_limit\": 10000}", + "description": "События по часам" + } + }, + { + "__Slice__": { + "slice_name": "📱 Traffic by Device", + "viz_type": "pie", + "datasource_type": "table", + "datasource_name": "dm.v_events_enriched", + "params": "{\"datasource\": \"1__table\", \"viz_type\": \"pie\", \"groupby\": [\"device_type\"], \"metric\": {\"expressionType\": \"SQL\", \"sqlExpression\": \"COUNT(*)\", \"label\": \"Count\"}, \"row_limit\": 100, \"donut\": true, \"show_legend\": true, \"labels_outside\": true, \"time_range\": \"No filter\"}", + "description": "Распределение трафика по устройствам" + } + }, + { + "__Slice__": { + "slice_name": "🌍 Geography Map", + "viz_type": "world_map", + "datasource_type": "table", + "datasource_name": "dm.v_events_enriched", + "params": "{\"datasource\": \"1__table\", \"viz_type\": \"world_map\", \"entity\": \"geo_country\", \"metric\": {\"expressionType\": \"SQL\", \"sqlExpression\": \"COUNT(*)\", \"label\": \"Events\"}, \"row_limit\": 500, \"linear_color_scheme\": \"blue_white_yellow\", \"time_range\": \"No filter\"}", + "description": "География посетителей" + } + }, + { + "__Slice__": { + "slice_name": "🔗 UTM Effectiveness Table", + "viz_type": "table", + "datasource_type": "table", + "datasource_name": "dm.v_utm_effectiveness", + "params": "{\"datasource\": \"2__table\", \"viz_type\": \"table\", \"groupby\": [\"utm_source\", \"utm_medium\", \"utm_campaign\"], \"metrics\": [{\"expressionType\": \"SQL\", \"sqlExpression\": \"SUM(clicks)\", \"label\": \"Clicks\"}, {\"expressionType\": \"SQL\", \"sqlExpression\": \"SUM(uniq_users)\", \"label\": \"Users\"}, {\"expressionType\": \"SQL\", \"sqlExpression\": \"SUM(uniq_sessions)\", \"label\": \"Sessions\"}], \"row_limit\": 100, \"time_range\": \"No filter\", \"adhoc_filters\": [{\"clause\": \"WHERE\", \"expressionType\": \"SQL\", \"sqlExpression\": \"utm_source IS NOT NULL\", \"subject\": null, \"operator\": null, \"comparator\": null}]}\n", + "description": "Эффективность UTM-кампаний" + } + }, + { + "__Slice__": { + "slice_name": "📄 Top Pages", + "viz_type": "dist_bar", + "datasource_type": "table", + "datasource_name": "dm.v_top_pages_daily", + "params": "{\"datasource\": \"3__table\", \"viz_type\": \"dist_bar\", \"groupby\": [\"page_url_path\"], \"metrics\": [{\"expressionType\": \"SQL\", \"sqlExpression\": \"SUM(pageviews)\", \"label\": \"Pageviews\"}], \"row_limit\": 20, \"order_by_cols\": [[\"SUM(pageviews)\", false]], \"time_range\": \"No filter\", \"orientation\": \"vertical\", \"show_legend\": false}", + "description": "Топ страниц по просмотрам" + } + }, + { + "__Slice__": { + "slice_name": "🔍 Data Quality Summary", + "viz_type": "dist_bar", + "datasource_type": "table", + "datasource_name": "dm.dq_summary", + "params": "{\"datasource\": \"4__table\", \"viz_type\": \"dist_bar\", \"groupby\": [\"layer\"], \"metrics\": [{\"expressionType\": \"SQL\", \"sqlExpression\": \"SUM(check_value)\", \"label\": \"Row Count\"}], \"adhoc_filters\": [{\"clause\": \"WHERE\", \"expressionType\": \"SQL\", \"sqlExpression\": \"check_name = 'total_rows'\", \"subject\": null, \"operator\": null, \"comparator\": null}], \"row_limit\": 100, \"time_range\": \"No filter\", \"show_legend\": false}", + "description": "Сводка по качеству данных" + } + } + ], + "datasets": [ + { + "__SqlaTable__": { + "table_name": "v_events_enriched", + "schema": "dm", + "database": "clickhouse_dwh", + "description": "Полная обогащённая витрина событий (event + click)", + "columns": [ + {"column_name": "event_id", "type": "UUID", "description": "UUID события"}, + {"column_name": "event_ts", "type": "DateTime64(6)", "description": "Время события"}, + {"column_name": "event_date", "type": "Date", "description": "Дата события"}, + {"column_name": "event_type", "type": "String", "description": "Тип события"}, + {"column_name": "click_id", "type": "UUID", "description": "ID сессии"}, + {"column_name": "user_domain_id", "type": "UUID", "description": "ID пользователя"}, + {"column_name": "device_type", "type": "String", "description": "Тип устройства"}, + {"column_name": "geo_country", "type": "String", "description": "Страна"}, + {"column_name": "browser_name", "type": "String", "description": "Браузер"}, + {"column_name": "utm_source", "type": "String", "description": "UTM Source"}, + {"column_name": "utm_medium", "type": "String", "description": "UTM Medium"}, + {"column_name": "page_url_path", "type": "String", "description": "Путь URL"} + ] + } + }, + { + "__SqlaTable__": { + "table_name": "v_utm_effectiveness", + "schema": "dm", + "database": "clickhouse_dwh", + "description": "Эффективность UTM-кампаний", + "columns": [ + {"column_name": "event_date", "type": "Date", "description": "Дата"}, + {"column_name": "utm_source", "type": "String", "description": "UTM Source"}, + {"column_name": "utm_medium", "type": "String", "description": "UTM Medium"}, + {"column_name": "utm_campaign", "type": "String", "description": "UTM Campaign"}, + {"column_name": "clicks", "type": "UInt64", "description": "Клики"}, + {"column_name": "uniq_users", "type": "UInt64", "description": "Уникальные пользователи"}, + {"column_name": "uniq_sessions", "type": "UInt64", "description": "Уникальные сессии"} + ] + } + }, + { + "__SqlaTable__": { + "table_name": "v_top_pages_daily", + "schema": "dm", + "database": "clickhouse_dwh", + "description": "Популярность страниц по дням", + "columns": [ + {"column_name": "event_date", "type": "Date", "description": "Дата"}, + {"column_name": "page_url_path", "type": "String", "description": "Путь URL"}, + {"column_name": "pageviews", "type": "UInt64", "description": "Просмотры"}, + {"column_name": "uniq_clicks", "type": "UInt64", "description": "Уникальные клики"} + ] + } + }, + { + "__SqlaTable__": { + "table_name": "dq_summary", + "schema": "dm", + "database": "clickhouse_dwh", + "description": "Сводка по качеству данных", + "columns": [ + {"column_name": "check_date", "type": "Date", "description": "Дата проверки"}, + {"column_name": "layer", "type": "String", "description": "Слой (stg/ods/dds)"}, + {"column_name": "table_name", "type": "String", "description": "Имя таблицы"}, + {"column_name": "check_name", "type": "String", "description": "Тип проверки"}, + {"column_name": "check_value", "type": "UInt64", "description": "Значение"} + ] + } + } + ], + "databases": [ + { + "__Database__": { + "database_name": "clickhouse_dwh", + "sqlalchemy_uri": "clickhousedb://default:123456@clickhouse:8123/default", + "expose_in_sqllab": true, + "allow_ctas": false, + "allow_cvas": false, + "allow_dml": false, + "allow_file_upload": false, + "extra": "{}" + } + } + ] +} diff --git a/superset/export_dashboard.py b/superset/export_dashboard.py new file mode 100644 index 0000000..92395cc --- /dev/null +++ b/superset/export_dashboard.py @@ -0,0 +1,107 @@ +#!/usr/bin/env python3 +""" +================================================================================ +Экспорт дашборда Superset в JSON-формат +================================================================================ +Назначение: + - Экспорт созданного дашборда в JSON для версионирования + - Формат совместимый с superset import-dashboards + +Запуск: + docker compose exec superset python /app/superset_init/export_dashboard.py +================================================================================ +""" + +import json +import sys +import os +sys.path.insert(0, '/app') + +try: + from superset.app import create_app + from superset.dashboards.data_access_layer import DashboardDAO + from superset.charts.data_access_layer import ChartDAO +except ImportError as e: + print(f"Error importing: {e}") + sys.exit(1) + + +def export_dashboard(slug: str, output_path: str): + """Экспорт дашборда в JSON""" + app = create_app() + + with app.app_context(): + dashboard = DashboardDAO.get_by_slug(slug) + + if not dashboard: + print(f"Dashboard with slug '{slug}' not found") + return False + + # Собираем данные дашборда + dashboard_data = { + "dashboards": [ + { + "__Dashboard__": { + "dashboard_title": dashboard.dashboard_title, + "description": dashboard.description, + "slug": dashboard.slug, + "json_metadata": dashboard.json_metadata, + "position_json": dashboard.position_json, + "published": dashboard.published, + "slices": [] + } + } + ], + "charts": [], + "datasets": [] + } + + # Добавляем чарты + for slice_obj in dashboard.slices: + chart_data = { + "__Slice__": { + "slice_name": slice_obj.slice_name, + "viz_type": slice_obj.viz_type, + "params": slice_obj.params, + "description": slice_obj.description, + "datasource_type": slice_obj.datasource_type, + "datasource_name": slice_obj.datasource.name if slice_obj.datasource else None + } + } + dashboard_data["dashboards"][0]["__Dashboard__"]["slices"].append(slice_obj.id) + dashboard_data["charts"].append(chart_data) + + # Добавляем датасет + if slice_obj.datasource: + ds = slice_obj.datasource + dataset_data = { + "__SqlaTable__": { + "table_name": ds.table_name, + "schema": ds.schema, + "database": ds.database.database_name if ds.database else None, + "description": ds.description, + "columns": [ + { + "column_name": col.column_name, + "type": col.type, + "description": col.description + } + for col in ds.columns + ] + } + } + # Добавляем уникальные датасеты + if dataset_data not in dashboard_data["datasets"]: + dashboard_data["datasets"].append(dataset_data) + + # Сохраняем в файл + with open(output_path, 'w', encoding='utf-8') as f: + json.dump(dashboard_data, f, indent=2, ensure_ascii=False) + + print(f"Dashboard exported to: {output_path}") + return True + + +if __name__ == "__main__": + output_file = "/app/superset_init/dashboards/ecommerce_analytics.json" + export_dashboard("ecommerce-analytics", output_file) diff --git a/superset/init_superset.py b/superset/init_superset.py new file mode 100644 index 0000000..e3946cc --- /dev/null +++ b/superset/init_superset.py @@ -0,0 +1,263 @@ +#!/usr/bin/env python3 +""" +================================================================================ +Скрипт инициализации Superset для проекта ClickHouse Mini DWH +================================================================================ +Назначение: + - Создание подключения к ClickHouse (Database connection) + - Импорт датасетов из витрин DM-слоя + +Запуск: + Внутри контейнера superset: + python /app/superset_init/init_superset.py + +Важно: + Superset использует SQLite по умолчанию (не PostgreSQL). + + Текущий подход: + 1. CLI для создания подключения к БД + 2. Superset shell для импорта датасетов (требуется app context) + +Требования: + - Запущенный ClickHouse с созданными витринами в схеме dm + - Superset инициализирован (superset db upgrade, admin создан) +================================================================================ +""" + +import sys +import os +import logging +import subprocess +from urllib.parse import quote_plus + +# Настройка логирования +logging.basicConfig( + level=logging.INFO, + format='%(asctime)s - %(levelname)s - %(message)s' +) +logger = logging.getLogger(__name__) + +# Добавляем путь к superset +sys.path.insert(0, '/app') + +CLICKHOUSE_USER = os.getenv('CLICKHOUSE_USER', 'default') +CLICKHOUSE_PASSWORD = os.getenv('CLICKHOUSE_PASSWORD', '123456') +CLICKHOUSE_HOST = os.getenv('CLICKHOUSE_HOST', 'clickhouse') +CLICKHOUSE_PORT = os.getenv('CLICKHOUSE_HTTP_PORT', '8123') +CLICKHOUSE_DATABASE = os.getenv('CLICKHOUSE_DATABASE', 'default') + + +def build_clickhouse_uri() -> str: + """Собирает URI подключения к ClickHouse для Superset.""" + user = quote_plus(CLICKHOUSE_USER) + password = quote_plus(CLICKHOUSE_PASSWORD) + host = CLICKHOUSE_HOST + port = CLICKHOUSE_PORT + database = quote_plus(CLICKHOUSE_DATABASE) + return f"clickhouse+connect://{user}:{password}@{host}:{port}/{database}" + + +def run_superset_cli(args): + """Запуск команды superset CLI""" + cmd = ['superset'] + args + logger.info(f"Running: {' '.join(cmd)}") + result = subprocess.run(cmd, capture_output=True, text=True) + if result.returncode != 0: + logger.error(f"Command failed: {result.stderr}") + else: + logger.info(f"Command output: {result.stdout}") + return result.returncode == 0, result.stdout, result.stderr + + +def create_clickhouse_connection(): + """Создание подключения к ClickHouse через CLI""" + logger.info("Creating ClickHouse database connection...") + database_uri = build_clickhouse_uri() + + # Создаем подключение через set-database-uri + # clickhouse-connect использует HTTP порт 8123 внутри Docker сети + success, stdout, stderr = run_superset_cli([ + 'set-database-uri', + '-d', 'clickhouse_dwh', + '-u', database_uri + ]) + + if success: + logger.info("Successfully created ClickHouse database connection 'clickhouse_dwh'") + return True + else: + logger.error(f"Failed to create database connection: {stderr}") + logger.info("Please create manually via UI:") + logger.info("1. Go to: Settings → Database Connections → + Database") + logger.info("2. Select: ClickHouse") + logger.info(f"3. URI: {database_uri}") + return False + + +def import_datasets(): + """Импорт датасетов через Superset shell""" + logger.info("Importing datasets...") + + datasets = [ + { + "table_name": "v_events_enriched", + "schema": "dm", + "database_name": "clickhouse_dwh", + "description": "Полная обогащённая витрина событий (event + click)" + }, + { + "table_name": "v_daily_traffic", + "schema": "dm", + "database_name": "clickhouse_dwh", + "description": "Агрегация трафика по дням и измерениям" + }, + { + "table_name": "v_utm_effectiveness", + "schema": "dm", + "database_name": "clickhouse_dwh", + "description": "Эффективность UTM-кампаний" + }, + { + "table_name": "v_top_pages_daily", + "schema": "dm", + "database_name": "clickhouse_dwh", + "description": "Популярность страниц по дням" + }, + { + "table_name": "v_session_overview", + "schema": "dm", + "database_name": "clickhouse_dwh", + "description": "Обзор сессий пользователей" + }, + { + "table_name": "dq_summary", + "schema": "dm", + "database_name": "clickhouse_dwh", + "description": "Сводка по качеству данных" + } + ] + + # Создаем Python скрипт для выполнения внутри superset shell + script_lines = [ + "import clickhouse_connect # Регистрирует диалект clickhouse+connect", + "from superset.extensions import db", + "from superset.models.core import Database", + "from superset.connectors.sqla.models import SqlaTable", + "", + "# Получаем базу данных", + "database = db.session.query(Database).filter_by(database_name='clickhouse_dwh').first()", + "if not database:", + " print('ERROR: Database clickhouse_dwh not found')", + " exit(1)", + "", + "print(f'Found database: {database.database_name} (id={database.id})')", + "", + "imported = 0", + "refreshed = 0", + "errors = 0", + ] + + for ds in datasets: + table_name = ds["table_name"] + schema_name = ds["schema"] + description = ds["description"] + script_lines.extend([ + "", + f"# Dataset: {table_name}", + ( + "existing = db.session.query(SqlaTable).filter_by(" + f"table_name={table_name!r}, schema={schema_name!r}" + ").first()" + ), + "if existing:", + f" print('Dataset {table_name} already exists')", + "else:", + " try:", + ( + " dataset = SqlaTable(" + f"table_name={table_name!r}, " + f"schema={schema_name!r}, " + "database_id=database.id, " + "database=database, " + f"description={description!r}" + ")" + ), + " db.session.add(dataset)", + " db.session.commit()", + f" print('Created dataset: {table_name} (metadata will be fetched on first use)')", + " imported += 1", + " except Exception as e:", + f" print(f'ERROR: failed to create {table_name}: {{e}}')", + " db.session.rollback()", + " errors += 1", + ]) + + script_lines.extend([ + "", + "print(f'Successfully imported {imported} datasets')", + "print(f'Refreshed metadata for {refreshed} datasets')", + "print(f'Errors: {errors}')", + "if errors > 0:", + " raise SystemExit(1)", + ]) + + script_content = '\n'.join(script_lines) + + # Запускаем через superset shell + cmd = ['superset', 'shell'] + logger.info("Running datasets import via superset shell...") + + result = subprocess.run( + cmd, + input=script_content, + capture_output=True, + text=True + ) + + if result.returncode != 0: + logger.error(f"Shell command failed: {result.stderr}") + return False + + logger.info(f"Shell output:\n{result.stdout}") + if "ERROR" in result.stdout: + logger.error("Failed to import some datasets") + return False + + logger.info("Datasets imported successfully") + return True + + +def main(): + """Главная функция инициализации""" + logger.info("=" * 60) + logger.info("Superset Initialization for ClickHouse Mini DWH") + logger.info("=" * 60) + + # Создаем подключение к ClickHouse + if not create_clickhouse_connection(): + logger.error("Failed to create ClickHouse connection") + sys.exit(1) + + # Импортируем датасеты + try: + if not import_datasets(): + logger.error("Failed to import datasets") + logger.info("\nTo create datasets manually:") + logger.info("1. Go to http://localhost:8088") + logger.info("2. Datasets → + Dataset") + logger.info("3. Select 'clickhouse_dwh' database") + logger.info("4. Select schema 'dm' and desired table") + sys.exit(1) + except Exception as e: + logger.error(f"Error importing datasets: {e}") + import traceback + traceback.print_exc() + sys.exit(1) + + logger.info("=" * 60) + logger.info("Superset initialization completed successfully!") + logger.info("=" * 60) + + +if __name__ == "__main__": + main()