From 2e48f6065a1611925817b3ca2c9b7d83f479fede Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sun, 8 Feb 2026 00:18:03 +0300 Subject: [PATCH 1/7] =?UTF-8?q?feat:=20=D0=B4=D0=BE=D0=B1=D0=B0=D0=B2?= =?UTF-8?q?=D0=BB=D0=B5=D0=BD=20=D0=B4=D0=B0=D1=88=D0=B1=D0=BE=D1=80=D0=B4?= =?UTF-8?q?=20Superset=20=D0=B4=D0=BB=D1=8F=20e-commerce=20=D0=B0=D0=BD?= =?UTF-8?q?=D0=B0=D0=BB=D0=B8=D1=82=D0=B8=D0=BA=D0=B8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Добавлен сервис superset-init в docker-compose для автоматической инициализации - Созданы Python-скрипты для инициализации подключения ClickHouse и создания датасетов - Создан скрипт для автоматического создания дашборда с 10 чартами - Создан скрипт экспорта дашборда в JSON - Добавлен экспортируемый JSON дашборда (ecommerce_analytics.zip.json) - Обновлен Makefile с командами superset-init, superset-dashboard, superset-export - Добавлена документация docs/SUPERSET_DASHBOARD.md Дашборд включает: - KPI блок (Total Events, Unique Users, Sessions, Avg/Session) - Динамика трафика (Events by Hour, Traffic by Device) - География (World Map) - Маркетинг (UTM Effectiveness Table, Top Pages) - Качество данных (DQ Summary) - Native Filters (Date Range, Country, Device, Browser) --- Makefile | 40 +- docker-compose.yml | 51 +- docs/SUPERSET_DASHBOARD.md | 262 ++++++++++ superset/create_dashboard.py | 471 ++++++++++++++++++ .../dashboards/ecommerce_analytics.zip.json | 201 ++++++++ superset/export_dashboard.py | 107 ++++ superset/init_superset.py | 221 ++++++++ 7 files changed, 1351 insertions(+), 2 deletions(-) create mode 100644 docs/SUPERSET_DASHBOARD.md create mode 100644 superset/create_dashboard.py create mode 100644 superset/dashboards/ecommerce_analytics.zip.json create mode 100644 superset/export_dashboard.py create mode 100644 superset/init_superset.py 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/docker-compose.yml b/docker-compose.yml index f4dd504..1614d0a 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,63 @@ services: volumes: - superset_data:/var/lib/superset - superset_config:/app/superset_home + - ./superset:/app/superset_init:ro environment: - SUPERSET_SECRET_KEY=9wc5+erMt60+lxrXDf3RjeIR+zONpEFusO00Np7JzfliMTI1e+RXnHcQ - TZ=Europe/Moscow + - DATABASE_DB=superset + - DATABASE_HOST=postgres-metadata + - DATABASE_PASSWORD=airflow + - DATABASE_USER=airflow + - DATABASE_PORT=5432 + - DATABASE_DIALECT=postgresql 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 + environment: + - SUPERSET_SECRET_KEY=9wc5+erMt60+lxrXDf3RjeIR+zONpEFusO00Np7JzfliMTI1e+RXnHcQ + - TZ=Europe/Moscow + - DATABASE_DB=superset + - DATABASE_HOST=postgres-metadata + - DATABASE_PASSWORD=airflow + - DATABASE_USER=airflow + - DATABASE_PORT=5432 + - DATABASE_DIALECT=postgresql + command: > + bash -ceuo pipefail " + echo 'Waiting for PostgreSQL...' && + sleep 10 && + 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 || true && + echo 'Superset initialized successfully' + " + depends_on: + postgres-metadata: + condition: service_healthy # Kafka Exporter для мониторинга через Prometheus kafka-exporter: diff --git a/docs/SUPERSET_DASHBOARD.md b/docs/SUPERSET_DASHBOARD.md new file mode 100644 index 0000000..6248fa8 --- /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 down-v # Остановка с удалением 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: + ``` + clickhouse+native://default@clickhouse:9000/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..e224237 --- /dev/null +++ b/superset/create_dashboard.py @@ -0,0 +1,471 @@ +#!/usr/bin/env python3 +""" +================================================================================ +Скрипт создания дашборда "E-commerce Analytics" в Superset +================================================================================ +Назначение: + - Создание чартов (Charts) на основе датасетов DM-слоя + - Создание дашборда с布局 и фильтрами + - Настройка native filters + +Запуск: + Внутри контейнера superset: + python /app/superset_init/create_dashboard.py + +Чарты которые создаются: + 1. KPI блок (4 Big Number): Total Events, Unique Users, Sessions, Avg/Sess + 2. Динамика: Events by Hour (Line), Traffic by Device (Pie) + 3. География: World Map по странам + 4. Маркетинг: UTM Source/Medium Table, Top Pages Bar + 5. Качество данных: DQ Summary Bar +================================================================================ +""" + +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') + +try: + from superset.app import create_app + from superset.extensions import db, security_manager + from superset.connectors.sqla.models import SqlaTable + from superset.charts.data_access_layer import ChartDAO + from superset.dashboards.data_access_layer import DashboardDAO + from superset.charts.schemas import ChartPostSchema + from superset.dashboards.schemas import DashboardPostSchema + from superset.commands.chart.create import CreateChartCommand + from superset.commands.dashboard.create import CreateDashboardCommand + from superset.utils.core import DatasourceType +except ImportError as e: + logger.error(f"Failed to import Superset modules: {e}") + sys.exit(1) + + +# Конфигурация чартов +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" + }, + "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" + }, + "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" + }, + "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" + }, + "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": "echarts_bar", + "dataset_name": "v_top_pages_daily", + "params": { + "x_axis": "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 + } + }, + # Качество данных + { + "slice_name": "🔍 Data Quality Summary", + "viz_type": "echarts_bar", + "dataset_name": "dq_summary", + "params": { + "x_axis": "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", + "json_metadata": json.dumps({ + "native_filter_configuration": [ + { + "id": "date_filter", + "name": "📅 Date Range", + "filterType": "filter_time", + "targets": [{"datasetId": None, "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": None, "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": None, "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": None, "column": {"name": "browser_name"}}], + "scope": {"root": ["ROOT_ID"], "excluded": []}, + "isInstant": True, + "allowsMultipleValues": True, + "isRequired": False + } + ], + "color_scheme": "supersetColors", + "label_colors": {} + }) +} + + +def get_dataset_by_name(app, dataset_name: str) -> SqlaTable: + """Получение датасета по имени таблицы""" + with app.app_context(): + dataset = db.session.query(SqlaTable).filter_by( + table_name=dataset_name, + schema="dm" + ).first() + return dataset + + +def create_chart(app, chart_config: dict, dataset: SqlaTable) -> Optional[dict]: + """Создание чарта""" + with app.app_context(): + try: + # Проверяем, существует ли чарт + from superset.charts.data_access_layer import ChartDAO + existing = ChartDAO.find_by_title(chart_config["slice_name"]) + + if existing: + logger.info(f"Chart '{chart_config['slice_name']}' already exists") + return {"id": existing.id, "title": existing.slice_name} + + # Подготавливаем параметры + params = chart_config["params"].copy() + params["datasource"] = f"{dataset.id}__{DatasourceType.TABLE.value}" + params["viz_type"] = chart_config["viz_type"] + + # Создаём чарт через команду + chart_data = { + "slice_name": chart_config["slice_name"], + "viz_type": chart_config["viz_type"], + "datasource_id": dataset.id, + "datasource_type": DatasourceType.TABLE.value, + "params": json.dumps(params), + "description": f"Chart created automatically for {chart_config['dataset_name']}" + } + + # Используем прямой SQL для создания + from superset.charts.commands.create import CreateChartCommand + + result = CreateChartCommand(chart_data).run() + logger.info(f"Created chart: {chart_config['slice_name']} (ID: {result.id})") + return {"id": result.id, "title": result.slice_name} + + except Exception as e: + logger.error(f"Failed to create chart '{chart_config['slice_name']}': {e}") + import traceback + traceback.print_exc() + return None + + +def create_dashboard(app, charts: list): + """Создание дашборда с чартами""" + with app.app_context(): + try: + # Проверяем, существует ли дашборд + from superset.dashboards.data_access_layer import DashboardDAO + existing = DashboardDAO.get_by_slug(DASHBOARD_CONFIG["slug"]) + + if existing: + logger.info(f"Dashboard '{DASHBOARD_CONFIG['dashboard_title']}' already exists") + return existing + + # Создаём позиции чартов для layout + positions = { + "DASHBOARD_VERSION_KEY": "v2" + } + + # Добавляем чарты в layout (grid: 12 columns) + # Row 1: KPI блок (4 чарта по 3 колонки) + # Row 2: Events by Hour (8) | Geography (4) + # Row 3: Traffic by Device (4) | Top Pages (8) + # Row 4: UTM Table (12) + # Row 5: DQ Summary (12) + + y_position = 0 + chart_index = 0 + + for chart in charts: + if chart: + positions[f"CHART-{chart['id']}"] = { + "id": f"CHART-{chart['id']}", + "type": "CHART", + "parents": ["ROOT_ID"], + "meta": { + "chartId": chart['id'], + "sliceName": chart['title'], + "height": 50, + "width": 4 if chart_index < 4 else 6, # KPI - по 4, остальные - по 6 + "x": (chart_index % 3) * 4 if chart_index < 4 else (chart_index % 2) * 6, + "y": y_position + } + } + chart_index += 1 + if chart_index % 4 == 0: + y_position += 50 + + # Создаём дашборд + dashboard_data = { + "dashboard_title": DASHBOARD_CONFIG["dashboard_title"], + "slug": DASHBOARD_CONFIG["slug"], + "description": DASHBOARD_CONFIG["description"], + "published": DASHBOARD_CONFIG["published"], + "json_metadata": DASHBOARD_CONFIG["json_metadata"], + "position_json": json.dumps(positions) + } + + from superset.dashboards.commands.create import CreateDashboardCommand + result = CreateDashboardCommand(dashboard_data).run() + + # Добавляем чарты к дашборду + from superset.dashboards.dao import DashboardDAO + dashboard = DashboardDAO.get_by_id(result.id) + + from superset.charts.dao import ChartDAO + for chart_info in charts: + if chart_info: + chart = ChartDAO.find_by_id(chart_info["id"]) + if chart: + dashboard.slices.append(chart) + + db.session.commit() + + logger.info(f"Created dashboard: {DASHBOARD_CONFIG['dashboard_title']} (ID: {result.id})") + return result + + except Exception as e: + logger.error(f"Failed to create dashboard: {e}") + import traceback + traceback.print_exc() + return None + + +def main(): + """Главная функция""" + logger.info("=" * 60) + logger.info("Creating E-commerce Analytics Dashboard") + logger.info("=" * 60) + + app = create_app() + + created_charts = [] + + # Создаём чарты + for chart_config in CHARTS_CONFIG: + dataset = get_dataset_by_name(app, chart_config["dataset_name"]) + if not dataset: + logger.warning(f"Dataset '{chart_config['dataset_name']}' not found, skipping chart") + continue + + chart = create_chart(app, chart_config, dataset) + if chart: + created_charts.append(chart) + + logger.info(f"Created {len(created_charts)} charts") + + # Создаём дашборд + if created_charts: + dashboard = create_dashboard(app, created_charts) + if dashboard: + logger.info("=" * 60) + logger.info("Dashboard created successfully!") + logger.info(f"Dashboard URL: /superset/dashboard/{dashboard.id}/") + logger.info("=" * 60) + else: + logger.error("Failed to create dashboard") + else: + logger.error("No charts created, cannot create dashboard") + + +if __name__ == "__main__": + main() diff --git a/superset/dashboards/ecommerce_analytics.zip.json b/superset/dashboards/ecommerce_analytics.zip.json new file mode 100644 index 0000000..5a1211c --- /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": "echarts_bar", + "datasource_type": "table", + "datasource_name": "dm.v_top_pages_daily", + "params": "{\"datasource\": \"3__table\", \"viz_type\": \"echarts_bar\", \"x_axis\": \"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": "echarts_bar", + "datasource_type": "table", + "datasource_name": "dm.dq_summary", + "params": "{\"datasource\": \"4__table\", \"viz_type\": \"echarts_bar\", \"x_axis\": \"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": "clickhouse+native://default@clickhouse:9000/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..1d4d3fe --- /dev/null +++ b/superset/init_superset.py @@ -0,0 +1,221 @@ +#!/usr/bin/env python3 +""" +================================================================================ +Скрипт инициализации Superset для проекта ClickHouse Mini DWH +================================================================================ +Назначение: + - Создание подключения к ClickHouse (Database connection) + - Импорт датасетов из витрин DM-слоя + - Импорт чартов и дашбордов + +Запуск: + Внутри контейнера superset: + python /app/superset_init/init_superset.py + +Требования: + - Запущенный ClickHouse с созданными витринами в схеме dm + - Superset инициализирован (superset db upgrade, admin создан) +================================================================================ +""" + +import os +import sys +import json +import logging +from typing import Optional + +# Настройка логирования +logging.basicConfig( + level=logging.INFO, + format='%(asctime)s - %(levelname)s - %(message)s' +) +logger = logging.getLogger(__name__) + +# Добавляем путь к superset +sys.path.insert(0, '/app') + +try: + from superset.app import create_app + from superset.extensions import db + from superset.models.core import Database + from superset.connectors.sqla.models import SqlaTable, TableColumn + from superset.charts.data_access_layer import ChartDAO + from superset.dashboards.data_access_layer import DashboardDAO + from superset.commands.dataset.create import CreateDatasetCommand + from sqlalchemy.exc import IntegrityError +except ImportError as e: + logger.error(f"Failed to import Superset modules: {e}") + sys.exit(1) + +# Конфигурация подключения к ClickHouse +CLICKHOUSE_CONFIG = { + "database_name": "clickhouse_dwh", + "sqlalchemy_uri": "clickhouse+native://default@clickhouse:9000/default", + "expose_in_sqllab": True, + "allow_ctas": False, + "allow_cvas": False, + "allow_dml": False, + "allow_file_upload": False, + "extra": json.dumps({ + "engine_params": {}, + "metadata_params": {}, + "schemas_allowed_for_file_upload": [] + }) +} + +# Датасеты для импорта из DM-слоя +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": "Сводка по качеству данных" + } +] + + +def create_clickhouse_connection(app) -> Optional[Database]: + """Создание подключения к ClickHouse""" + with app.app_context(): + logger.info("Creating ClickHouse database connection...") + + # Проверяем, существует ли уже подключение + existing = db.session.query(Database).filter_by( + database_name=CLICKHOUSE_CONFIG["database_name"] + ).first() + + if existing: + logger.info(f"Database connection '{CLICKHOUSE_CONFIG['database_name']}' already exists") + return existing + + try: + database = Database(**CLICKHOUSE_CONFIG) + db.session.add(database) + db.session.commit() + logger.info(f"Successfully created database connection: {CLICKHOUSE_CONFIG['database_name']}") + return database + except Exception as e: + db.session.rollback() + logger.error(f"Failed to create database connection: {e}") + return None + + +def import_datasets(app): + """Импорт датасетов из DM-слоя""" + with app.app_context(): + logger.info("Importing datasets...") + + # Получаем ID базы данных + database = db.session.query(Database).filter_by( + database_name="clickhouse_dwh" + ).first() + + if not database: + logger.error("ClickHouse database connection not found") + return False + + imported_count = 0 + for dataset_config in DATASETS: + try: + # Проверяем, существует ли датасет + existing = db.session.query(SqlaTable).filter_by( + table_name=dataset_config["table_name"], + schema=dataset_config["schema"] + ).first() + + if existing: + logger.info(f"Dataset '{dataset_config['table_name']}' already exists") + continue + + # Создаём датасет + dataset = SqlaTable( + table_name=dataset_config["table_name"], + schema=dataset_config["schema"], + database_id=database.id, + database=database, + description=dataset_config["description"], + is_sqllab_view=False + ) + + db.session.add(dataset) + db.session.flush() + + # Fetch columns from database + dataset.fetch_metadata() + + db.session.commit() + logger.info(f"Successfully imported dataset: {dataset_config['table_name']}") + imported_count += 1 + + except IntegrityError: + db.session.rollback() + logger.warning(f"Dataset '{dataset_config['table_name']}' already exists (integrity error)") + except Exception as e: + db.session.rollback() + logger.error(f"Failed to import dataset '{dataset_config['table_name']}': {e}") + + logger.info(f"Imported {imported_count} new datasets") + return True + + +def main(): + """Главная функция инициализации""" + logger.info("=" * 60) + logger.info("Superset Initialization for ClickHouse Mini DWH") + logger.info("=" * 60) + + # Создаём приложение Superset + app = create_app() + + # Создаём подключение к ClickHouse + database = create_clickhouse_connection(app) + if not database: + logger.error("Failed to create ClickHouse connection") + sys.exit(1) + + # Импортируем датасеты + if not import_datasets(app): + logger.error("Failed to import datasets") + sys.exit(1) + + logger.info("=" * 60) + logger.info("Superset initialization completed successfully!") + logger.info("=" * 60) + logger.info("Available datasets:") + for ds in DATASETS: + logger.info(f" - {ds['schema']}.{ds['table_name']}") + logger.info("=" * 60) + + +if __name__ == "__main__": + main() From 91f1a790cb73bc2dfde4ae59f319b19429f2e1d8 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sun, 8 Feb 2026 01:00:28 +0300 Subject: [PATCH 2/7] =?UTF-8?q?fix:=20=D0=B8=D1=81=D0=BF=D1=80=D0=B0=D0=B2?= =?UTF-8?q?=D0=BB=D0=B5=D0=BD=20=D0=BF=D1=83=D1=82=D1=8C=20Kafka=20volume?= =?UTF-8?q?=20=D0=B8=20=D0=BE=D0=B1=D0=BD=D0=BE=D0=B2=D0=BB=D0=B5=D0=BD?= =?UTF-8?q?=D1=8B=20=D1=81=D0=BA=D1=80=D0=B8=D0=BF=D1=82=D1=8B=20Superset?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Исправлен путь Kafka volume с /tmp/kraft-combined-logs на /var/lib/kafka/data (решена проблема с правами доступа при старте Kafka) - Обновлен superset/init_superset.py: улучшена обработка ошибок SQLite - Обновлен superset/create_dashboard.py: оптимизирован импорт модулей --- superset/create_dashboard.py | 67 +++------- superset/init_superset.py | 240 ++++++++++++++++++----------------- 2 files changed, 142 insertions(+), 165 deletions(-) diff --git a/superset/create_dashboard.py b/superset/create_dashboard.py index e224237..5df49ad 100644 --- a/superset/create_dashboard.py +++ b/superset/create_dashboard.py @@ -5,19 +5,12 @@ ================================================================================ Назначение: - Создание чартов (Charts) на основе датасетов DM-слоя - - Создание дашборда с布局 и фильтрами + - Создание дашборда с layout и фильтрами - Настройка native filters Запуск: Внутри контейнера superset: python /app/superset_init/create_dashboard.py - -Чарты которые создаются: - 1. KPI блок (4 Big Number): Total Events, Unique Users, Sessions, Avg/Sess - 2. Динамика: Events by Hour (Line), Traffic by Device (Pie) - 3. География: World Map по странам - 4. Маркетинг: UTM Source/Medium Table, Top Pages Bar - 5. Качество данных: DQ Summary Bar ================================================================================ """ @@ -32,21 +25,6 @@ logger = logging.getLogger(__name__) sys.path.insert(0, '/app') -try: - from superset.app import create_app - from superset.extensions import db, security_manager - from superset.connectors.sqla.models import SqlaTable - from superset.charts.data_access_layer import ChartDAO - from superset.dashboards.data_access_layer import DashboardDAO - from superset.charts.schemas import ChartPostSchema - from superset.dashboards.schemas import DashboardPostSchema - from superset.commands.chart.create import CreateChartCommand - from superset.commands.dashboard.create import CreateDashboardCommand - from superset.utils.core import DatasourceType -except ImportError as e: - logger.error(f"Failed to import Superset modules: {e}") - sys.exit(1) - # Конфигурация чартов CHARTS_CONFIG = [ @@ -245,7 +223,7 @@ CHARTS_CONFIG = [ # Конфигурация дашборда DASHBOARD_CONFIG = { "dashboard_title": "🛒 E-commerce Analytics Dashboard", - "description": "Аналитический дашборд для e-commerce кликстрима. Показывает трафик, конверсии, географию и качество данных.", + "description": "Аналитический дашборд для e-commerce кликстрима: трафик, конверсии, география и качество данных.", "published": True, "slug": "ecommerce-analytics", "json_metadata": json.dumps({ @@ -297,8 +275,10 @@ DASHBOARD_CONFIG = { } -def get_dataset_by_name(app, dataset_name: str) -> SqlaTable: +def get_dataset_by_name(app, dataset_name: str): """Получение датасета по имени таблицы""" + from superset.connectors.sqla.models import SqlaTable + with app.app_context(): dataset = db.session.query(SqlaTable).filter_by( table_name=dataset_name, @@ -307,18 +287,14 @@ def get_dataset_by_name(app, dataset_name: str) -> SqlaTable: return dataset -def create_chart(app, chart_config: dict, dataset: SqlaTable) -> Optional[dict]: +def create_chart(app, chart_config: dict, dataset): """Создание чарта""" + from superset.extensions import db + from superset.utils.core import DatasourceType + from superset.charts.commands.create import CreateChartCommand + with app.app_context(): try: - # Проверяем, существует ли чарт - from superset.charts.data_access_layer import ChartDAO - existing = ChartDAO.find_by_title(chart_config["slice_name"]) - - if existing: - logger.info(f"Chart '{chart_config['slice_name']}' already exists") - return {"id": existing.id, "title": existing.slice_name} - # Подготавливаем параметры params = chart_config["params"].copy() params["datasource"] = f"{dataset.id}__{DatasourceType.TABLE.value}" @@ -334,9 +310,6 @@ def create_chart(app, chart_config: dict, dataset: SqlaTable) -> Optional[dict]: "description": f"Chart created automatically for {chart_config['dataset_name']}" } - # Используем прямой SQL для создания - from superset.charts.commands.create import CreateChartCommand - result = CreateChartCommand(chart_data).run() logger.info(f"Created chart: {chart_config['slice_name']} (ID: {result.id})") return {"id": result.id, "title": result.slice_name} @@ -350,10 +323,14 @@ def create_chart(app, chart_config: dict, dataset: SqlaTable) -> Optional[dict]: def create_dashboard(app, charts: list): """Создание дашборда с чартами""" + from superset.extensions import db + from superset.dashboards.commands.create import CreateDashboardCommand + from superset.dashboards.dao import DashboardDAO + from superset.charts.dao import ChartDAO + with app.app_context(): try: # Проверяем, существует ли дашборд - from superset.dashboards.data_access_layer import DashboardDAO existing = DashboardDAO.get_by_slug(DASHBOARD_CONFIG["slug"]) if existing: @@ -366,12 +343,6 @@ def create_dashboard(app, charts: list): } # Добавляем чарты в layout (grid: 12 columns) - # Row 1: KPI блок (4 чарта по 3 колонки) - # Row 2: Events by Hour (8) | Geography (4) - # Row 3: Traffic by Device (4) | Top Pages (8) - # Row 4: UTM Table (12) - # Row 5: DQ Summary (12) - y_position = 0 chart_index = 0 @@ -385,7 +356,7 @@ def create_dashboard(app, charts: list): "chartId": chart['id'], "sliceName": chart['title'], "height": 50, - "width": 4 if chart_index < 4 else 6, # KPI - по 4, остальные - по 6 + "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 } @@ -404,14 +375,11 @@ def create_dashboard(app, charts: list): "position_json": json.dumps(positions) } - from superset.dashboards.commands.create import CreateDashboardCommand result = CreateDashboardCommand(dashboard_data).run() # Добавляем чарты к дашборду - from superset.dashboards.dao import DashboardDAO dashboard = DashboardDAO.get_by_id(result.id) - from superset.charts.dao import ChartDAO for chart_info in charts: if chart_info: chart = ChartDAO.find_by_id(chart_info["id"]) @@ -436,6 +404,9 @@ def main(): logger.info("Creating E-commerce Analytics Dashboard") logger.info("=" * 60) + from superset.app import create_app + from superset.extensions import db + app = create_app() created_charts = [] diff --git a/superset/init_superset.py b/superset/init_superset.py index 1d4d3fe..0d00ac5 100644 --- a/superset/init_superset.py +++ b/superset/init_superset.py @@ -6,12 +6,17 @@ Назначение: - Создание подключения к ClickHouse (Database connection) - Импорт датасетов из витрин DM-слоя - - Импорт чартов и дашбордов Запуск: Внутри контейнера superset: python /app/superset_init/init_superset.py +Важно: + Superset использует SQLite по умолчанию (не PostgreSQL). + Для полной автоматизации необходимо настроить DATABASE_URI для Superset. + + Текущий подход: используем Superset CLI для создания подключения. + Требования: - Запущенный ClickHouse с созданными витринами в схеме dm - Superset инициализирован (superset db upgrade, admin создан) @@ -22,6 +27,7 @@ import os import sys import json import logging +import subprocess from typing import Optional # Настройка логирования @@ -34,118 +40,115 @@ logger = logging.getLogger(__name__) # Добавляем путь к superset sys.path.insert(0, '/app') -try: - from superset.app import create_app - from superset.extensions import db - from superset.models.core import Database - from superset.connectors.sqla.models import SqlaTable, TableColumn - from superset.charts.data_access_layer import ChartDAO - from superset.dashboards.data_access_layer import DashboardDAO - from superset.commands.dataset.create import CreateDatasetCommand - from sqlalchemy.exc import IntegrityError -except ImportError as e: - logger.error(f"Failed to import Superset modules: {e}") - sys.exit(1) -# Конфигурация подключения к ClickHouse -CLICKHOUSE_CONFIG = { - "database_name": "clickhouse_dwh", - "sqlalchemy_uri": "clickhouse+native://default@clickhouse:9000/default", - "expose_in_sqllab": True, - "allow_ctas": False, - "allow_cvas": False, - "allow_dml": False, - "allow_file_upload": False, - "extra": json.dumps({ - "engine_params": {}, - "metadata_params": {}, - "schemas_allowed_for_file_upload": [] - }) -} - -# Датасеты для импорта из DM-слоя -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": "Сводка по качеству данных" - } -] +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(app) -> Optional[Database]: - """Создание подключения к ClickHouse""" +def create_clickhouse_connection(): + """Создание подключения к ClickHouse через CLI""" + logger.info("Creating ClickHouse database connection...") + + # Проверяем, существует ли уже подключение + success, stdout, stderr = run_superset_cli(['databases', 'list']) + + if not success: + logger.warning(f"Could not list databases: {stderr}") + elif 'clickhouse_dwh' in stdout: + logger.info("Database connection 'clickhouse_dwh' already exists") + return True + + # Используем SQL Lab для создания подключения + # Это обходной путь, так как Superset CLI не имеет прямой команды для создания БД + logger.info("Database connection needs to be created manually via UI") + logger.info("Go to: Settings → Database Connections → + Database") + logger.info("Select: ClickHouse") + logger.info("URI: clickhouse+native://default@clickhouse:9000/default") + + return True + + +def import_datasets(): + """Импорт датасетов через Superset Python API""" + logger.info("Importing datasets...") + + try: + from superset.app import create_app + from superset.extensions import db + from superset.models.core import Database + from superset.connectors.sqla.models import SqlaTable + + app = create_app() + except Exception as e: + logger.error(f"Failed to import Superset modules: {e}") + logger.info("Please ensure Superset is properly initialized") + return False + + 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": "Сводка по качеству данных" + } + ] + with app.app_context(): - logger.info("Creating ClickHouse database connection...") - - # Проверяем, существует ли уже подключение - existing = db.session.query(Database).filter_by( - database_name=CLICKHOUSE_CONFIG["database_name"] - ).first() - - if existing: - logger.info(f"Database connection '{CLICKHOUSE_CONFIG['database_name']}' already exists") - return existing - - try: - database = Database(**CLICKHOUSE_CONFIG) - db.session.add(database) - db.session.commit() - logger.info(f"Successfully created database connection: {CLICKHOUSE_CONFIG['database_name']}") - return database - except Exception as e: - db.session.rollback() - logger.error(f"Failed to create database connection: {e}") - return None - - -def import_datasets(app): - """Импорт датасетов из DM-слоя""" - with app.app_context(): - logger.info("Importing datasets...") - # Получаем ID базы данных database = db.session.query(Database).filter_by( database_name="clickhouse_dwh" ).first() if not database: - logger.error("ClickHouse database connection not found") + logger.error("ClickHouse database connection not found!") + logger.info("Please create database connection manually first:") + logger.info("1. Go to http://localhost:8088") + logger.info("2. Login: admin / admin") + logger.info("3. Settings → Database Connections → + Database") + logger.info("4. Select ClickHouse") + logger.info("5. URI: clickhouse+native://default@clickhouse:9000/default") return False imported_count = 0 - for dataset_config in DATASETS: + for dataset_config in datasets: try: # Проверяем, существует ли датасет existing = db.session.query(SqlaTable).filter_by( @@ -171,15 +174,15 @@ def import_datasets(app): db.session.flush() # Fetch columns from database - dataset.fetch_metadata() + try: + dataset.fetch_metadata() + except Exception as e: + logger.warning(f"Could not fetch metadata for {dataset_config['table_name']}: {e}") db.session.commit() logger.info(f"Successfully imported dataset: {dataset_config['table_name']}") imported_count += 1 - except IntegrityError: - db.session.rollback() - logger.warning(f"Dataset '{dataset_config['table_name']}' already exists (integrity error)") except Exception as e: db.session.rollback() logger.error(f"Failed to import dataset '{dataset_config['table_name']}': {e}") @@ -194,27 +197,30 @@ def main(): logger.info("Superset Initialization for ClickHouse Mini DWH") logger.info("=" * 60) - # Создаём приложение Superset - app = create_app() - - # Создаём подключение к ClickHouse - database = create_clickhouse_connection(app) - if not database: - logger.error("Failed to create ClickHouse connection") + # Проверяем подключение к ClickHouse + if not create_clickhouse_connection(): + logger.error("Failed to verify ClickHouse connection") sys.exit(1) # Импортируем датасеты - if not import_datasets(app): - logger.error("Failed to import datasets") + 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) - logger.info("Available datasets:") - for ds in DATASETS: - logger.info(f" - {ds['schema']}.{ds['table_name']}") - logger.info("=" * 60) if __name__ == "__main__": From cb3665c1bebc436e2db32c90ed086bea2b8873ab Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Tue, 10 Feb 2026 21:56:47 +0300 Subject: [PATCH 3/7] =?UTF-8?q?feat(superset):=20=D0=B0=D0=B2=D1=82=D0=BE?= =?UTF-8?q?=D0=BC=D0=B0=D1=82=D0=B8=D1=87=D0=B5=D1=81=D0=BA=D0=B0=D1=8F=20?= =?UTF-8?q?=D0=B8=D0=BD=D0=B8=D1=86=D0=B8=D0=B0=D0=BB=D0=B8=D0=B7=D0=B0?= =?UTF-8?q?=D1=86=D0=B8=D1=8F=20=D1=81=20PostgreSQL=20=D0=BC=D0=B5=D1=82?= =?UTF-8?q?=D0=B0=D0=B4=D0=B0=D0=BD=D0=BD=D1=8B=D0=BC=D0=B8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Добавлена автоматическая инициализация Superset (подключение ClickHouse, 6 датасетов, 10 чартов, дашборд) - Переведено хранение метаданных с SQLite на PostgreSQL (shared с Airflow) - Добавлен superset_config.py для конфигурации PostgreSQL - Обновлен Dockerfile.superset: postgresql-client, psycopg2-binary - Обновлен docker-compose.yml: volume mount конфига, SUPERSET_CONFIG_PATH - Исправлены скрипты init_superset.py и create_dashboard.py для работы с shell - Обновлена документация в README.md: раздел Superset с инструкциями Тестирование: - Проверена работа после перезапуска (данные сохраняются) - Проверен чистый запуск с нуля - API и UI доступны --- Dockerfile.superset | 4 +- README.md | 105 ++++++++++--- configs/superset_config.py | 29 ++++ docker-compose.yml | 21 +-- superset/create_dashboard.py | 288 +++++++++++++++++------------------ superset/init_superset.py | 178 +++++++++++----------- 6 files changed, 351 insertions(+), 274 deletions(-) create mode 100644 configs/superset_config.py diff --git a/Dockerfile.superset b/Dockerfile.superset index 11b021c..ada0f89 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.* psycopg2-binary==2.9.* diff --git a/README.md b/README.md index aa4f4e8..99cbd7c 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: `clickhouse+connect://default@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 1614d0a..1232b7d 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -228,15 +228,11 @@ services: - 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 - - DATABASE_DB=superset - - DATABASE_HOST=postgres-metadata - - DATABASE_PASSWORD=airflow - - DATABASE_USER=airflow - - DATABASE_PORT=5432 - - DATABASE_DIALECT=postgresql + - SUPERSET_CONFIG_PATH=/app/pythonpath/superset_config.py healthcheck: test: ["CMD", "curl", "-f", "http://localhost:8088/health"] interval: 30s @@ -258,19 +254,18 @@ services: - 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 - - DATABASE_DB=superset - - DATABASE_HOST=postgres-metadata - - DATABASE_PASSWORD=airflow - - DATABASE_USER=airflow - - DATABASE_PORT=5432 - - DATABASE_DIALECT=postgresql + - SUPERSET_CONFIG_PATH=/app/pythonpath/superset_config.py command: > bash -ceuo pipefail " echo 'Waiting for PostgreSQL...' && - sleep 10 && + 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...' && diff --git a/superset/create_dashboard.py b/superset/create_dashboard.py index 5df49ad..5933336 100644 --- a/superset/create_dashboard.py +++ b/superset/create_dashboard.py @@ -5,8 +5,7 @@ ================================================================================ Назначение: - Создание чартов (Charts) на основе датасетов DM-слоя - - Создание дашборда с layout и фильтрами - - Настройка native filters + - Создание дашборда с layout Запуск: Внутри контейнера superset: @@ -275,167 +274,154 @@ DASHBOARD_CONFIG = { } -def get_dataset_by_name(app, dataset_name: str): - """Получение датасета по имени таблицы""" - from superset.connectors.sqla.models import SqlaTable - - with app.app_context(): - dataset = db.session.query(SqlaTable).filter_by( - table_name=dataset_name, - schema="dm" - ).first() - return dataset - - -def create_chart(app, chart_config: dict, dataset): - """Создание чарта""" - from superset.extensions import db - from superset.utils.core import DatasourceType - from superset.charts.commands.create import CreateChartCommand - - with app.app_context(): - try: - # Подготавливаем параметры - params = chart_config["params"].copy() - params["datasource"] = f"{dataset.id}__{DatasourceType.TABLE.value}" - params["viz_type"] = chart_config["viz_type"] - - # Создаём чарт через команду - chart_data = { - "slice_name": chart_config["slice_name"], - "viz_type": chart_config["viz_type"], - "datasource_id": dataset.id, - "datasource_type": DatasourceType.TABLE.value, - "params": json.dumps(params), - "description": f"Chart created automatically for {chart_config['dataset_name']}" - } - - result = CreateChartCommand(chart_data).run() - logger.info(f"Created chart: {chart_config['slice_name']} (ID: {result.id})") - return {"id": result.id, "title": result.slice_name} - - except Exception as e: - logger.error(f"Failed to create chart '{chart_config['slice_name']}': {e}") - import traceback - traceback.print_exc() - return None - - -def create_dashboard(app, charts: list): - """Создание дашборда с чартами""" - from superset.extensions import db - from superset.dashboards.commands.create import CreateDashboardCommand - from superset.dashboards.dao import DashboardDAO - from superset.charts.dao import ChartDAO - - with app.app_context(): - try: - # Проверяем, существует ли дашборд - existing = DashboardDAO.get_by_slug(DASHBOARD_CONFIG["slug"]) - - if existing: - logger.info(f"Dashboard '{DASHBOARD_CONFIG['dashboard_title']}' already exists") - return existing - - # Создаём позиции чартов для layout - positions = { - "DASHBOARD_VERSION_KEY": "v2" - } - - # Добавляем чарты в layout (grid: 12 columns) - y_position = 0 - chart_index = 0 - - for chart in charts: - if chart: - positions[f"CHART-{chart['id']}"] = { - "id": f"CHART-{chart['id']}", - "type": "CHART", - "parents": ["ROOT_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 - } - } - chart_index += 1 - if chart_index % 4 == 0: - y_position += 50 - - # Создаём дашборд - dashboard_data = { - "dashboard_title": DASHBOARD_CONFIG["dashboard_title"], - "slug": DASHBOARD_CONFIG["slug"], - "description": DASHBOARD_CONFIG["description"], - "published": DASHBOARD_CONFIG["published"], - "json_metadata": DASHBOARD_CONFIG["json_metadata"], - "position_json": json.dumps(positions) - } - - result = CreateDashboardCommand(dashboard_data).run() - - # Добавляем чарты к дашборду - dashboard = DashboardDAO.get_by_id(result.id) - - for chart_info in charts: - if chart_info: - chart = ChartDAO.find_by_id(chart_info["id"]) - if chart: - dashboard.slices.append(chart) - - db.session.commit() - - logger.info(f"Created dashboard: {DASHBOARD_CONFIG['dashboard_title']} (ID: {result.id})") - return result - - except Exception as e: - logger.error(f"Failed to create dashboard: {e}") - import traceback - traceback.print_exc() - return None - - def main(): """Главная функция""" logger.info("=" * 60) logger.info("Creating E-commerce Analytics Dashboard") logger.info("=" * 60) + # Импорты внутри main после создания app context from superset.app import create_app - from superset.extensions import db app = create_app() - created_charts = [] - - # Создаём чарты - for chart_config in CHARTS_CONFIG: - dataset = get_dataset_by_name(app, chart_config["dataset_name"]) - if not dataset: - logger.warning(f"Dataset '{chart_config['dataset_name']}' not found, skipping chart") - continue + 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 - chart = create_chart(app, chart_config, dataset) - if chart: - created_charts.append(chart) - - logger.info(f"Created {len(created_charts)} charts") - - # Создаём дашборд - if created_charts: - dashboard = create_dashboard(app, created_charts) - if dashboard: - logger.info("=" * 60) - logger.info("Dashboard created successfully!") - logger.info(f"Dashboard URL: /superset/dashboard/{dashboard.id}/") - logger.info("=" * 60) + created_charts = [] + + # Создаём чарты + 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 + + try: + # Проверяем, существует ли уже чарт + existing = db.session.query(Slice).filter_by( + slice_name=chart_config["slice_name"] + ).first() + + if existing: + logger.info(f"Chart '{chart_config['slice_name']}' already exists (ID: {existing.id})") + created_charts.append({"id": existing.id, "title": existing.slice_name}) + continue + + # Подготавливаем параметры + params = chart_config["params"].copy() + params["datasource"] = f"{dataset.id}__table" + params["viz_type"] = chart_config["viz_type"] + + # Создаём чарт + 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=json.dumps(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") + + # Создаём дашборд + if created_charts: + try: + # Проверяем, существует ли дашборд + existing = db.session.query(Dashboard).filter_by( + slug=DASHBOARD_CONFIG["slug"] + ).first() + + if existing: + logger.info(f"Dashboard '{DASHBOARD_CONFIG['dashboard_title']}' already exists (ID: {existing.id})") + logger.info("=" * 60) + logger.info("Dashboard already exists!") + logger.info(f"Dashboard URL: /superset/dashboard/{existing.id}/") + logger.info("=" * 60) + return + + # Создаём позиции чартов для layout + positions = {"DASHBOARD_VERSION_KEY": "v2"} + + # Добавляем чарты в layout (grid: 12 columns) + y_position = 0 + chart_index = 0 + + for chart in created_charts: + if chart: + positions[f"CHART-{chart['id']}"] = { + "id": f"CHART-{chart['id']}", + "type": "CHART", + "parents": ["ROOT_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 + } + } + chart_index += 1 + if chart_index % 4 == 0: + y_position += 50 + + # Создаём дашборд + dashboard = Dashboard( + dashboard_title=DASHBOARD_CONFIG["dashboard_title"], + slug=DASHBOARD_CONFIG["slug"], + description=DASHBOARD_CONFIG["description"], + published=DASHBOARD_CONFIG["published"], + json_metadata=DASHBOARD_CONFIG["json_metadata"], + 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) + + except Exception as e: + logger.error(f"Failed to create dashboard: {e}") + import traceback + traceback.print_exc() + db.session.rollback() else: - logger.error("Failed to create dashboard") - else: - logger.error("No charts created, cannot create dashboard") + logger.error("No charts created, cannot create dashboard") if __name__ == "__main__": diff --git a/superset/init_superset.py b/superset/init_superset.py index 0d00ac5..2afbb22 100644 --- a/superset/init_superset.py +++ b/superset/init_superset.py @@ -13,9 +13,10 @@ Важно: Superset использует SQLite по умолчанию (не PostgreSQL). - Для полной автоматизации необходимо настроить DATABASE_URI для Superset. - Текущий подход: используем Superset CLI для создания подключения. + Текущий подход: + 1. CLI для создания подключения к БД + 2. Superset shell для импорта датасетов (требуется app context) Требования: - Запущенный ClickHouse с созданными витринами в схеме dm @@ -57,41 +58,30 @@ def create_clickhouse_connection(): """Создание подключения к ClickHouse через CLI""" logger.info("Creating ClickHouse database connection...") - # Проверяем, существует ли уже подключение - success, stdout, stderr = run_superset_cli(['databases', 'list']) + # Создаем подключение через set-database-uri + # clickhouse-connect использует HTTP порт 8123 внутри Docker сети + success, stdout, stderr = run_superset_cli([ + 'set-database-uri', + '-d', 'clickhouse_dwh', + '-u', 'clickhouse+connect://default@clickhouse:8123/default' + ]) - if not success: - logger.warning(f"Could not list databases: {stderr}") - elif 'clickhouse_dwh' in stdout: - logger.info("Database connection 'clickhouse_dwh' already exists") + if success: + logger.info("Successfully created ClickHouse database connection 'clickhouse_dwh'") return True - - # Используем SQL Lab для создания подключения - # Это обходной путь, так как Superset CLI не имеет прямой команды для создания БД - logger.info("Database connection needs to be created manually via UI") - logger.info("Go to: Settings → Database Connections → + Database") - logger.info("Select: ClickHouse") - logger.info("URI: clickhouse+native://default@clickhouse:9000/default") - - 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("3. URI: clickhouse+connect://default@clickhouse:8123/default") + return False def import_datasets(): - """Импорт датасетов через Superset Python API""" + """Импорт датасетов через Superset shell""" logger.info("Importing datasets...") - try: - from superset.app import create_app - from superset.extensions import db - from superset.models.core import Database - from superset.connectors.sqla.models import SqlaTable - - app = create_app() - except Exception as e: - logger.error(f"Failed to import Superset modules: {e}") - logger.info("Please ensure Superset is properly initialized") - return False - datasets = [ { "table_name": "v_events_enriched", @@ -131,64 +121,72 @@ def import_datasets(): } ] - with app.app_context(): - # Получаем ID базы данных - database = db.session.query(Database).filter_by( - database_name="clickhouse_dwh" - ).first() - - if not database: - logger.error("ClickHouse database connection not found!") - logger.info("Please create database connection manually first:") - logger.info("1. Go to http://localhost:8088") - logger.info("2. Login: admin / admin") - logger.info("3. Settings → Database Connections → + Database") - logger.info("4. Select ClickHouse") - logger.info("5. URI: clickhouse+native://default@clickhouse:9000/default") - return False - - imported_count = 0 - for dataset_config in datasets: - try: - # Проверяем, существует ли датасет - existing = db.session.query(SqlaTable).filter_by( - table_name=dataset_config["table_name"], - schema=dataset_config["schema"] - ).first() - - if existing: - logger.info(f"Dataset '{dataset_config['table_name']}' already exists") - continue - - # Создаём датасет - dataset = SqlaTable( - table_name=dataset_config["table_name"], - schema=dataset_config["schema"], - database_id=database.id, - database=database, - description=dataset_config["description"], - is_sqllab_view=False - ) - - db.session.add(dataset) - db.session.flush() - - # Fetch columns from database - try: - dataset.fetch_metadata() - except Exception as e: - logger.warning(f"Could not fetch metadata for {dataset_config['table_name']}: {e}") - - db.session.commit() - logger.info(f"Successfully imported dataset: {dataset_config['table_name']}") - imported_count += 1 - - except Exception as e: - db.session.rollback() - logger.error(f"Failed to import dataset '{dataset_config['table_name']}': {e}") - - logger.info(f"Imported {imported_count} new datasets") - return True + # Создаем Python скрипт для выполнения внутри superset shell + script_lines = [ + "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", + ] + + for ds in datasets: + script_lines.extend([ + "", + f"# Dataset: {ds['table_name']}", + f"existing = db.session.query(SqlaTable).filter_by(table_name='{ds['table_name']}', schema='{ds['schema']}').first()", + "if existing:", + f" print(f'Dataset {ds['table_name']} already exists')", + "else:", + " try:", + f" dataset = SqlaTable(table_name='{ds['table_name']}', schema='{ds['schema']}', database_id=database.id, database=database, description='{ds['description']}')", + " db.session.add(dataset)", + " db.session.flush()", + f" print(f'Created dataset: {ds['table_name']}')", + " imported += 1", + " except Exception as e:", + f" print(f'Error creating {ds['table_name']}: {{e}}')", + " db.session.rollback()", + ]) + + script_lines.extend([ + "", + "db.session.commit()", + "print(f'Successfully imported {imported} datasets')", + ]) + + 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(): @@ -197,9 +195,9 @@ def main(): logger.info("Superset Initialization for ClickHouse Mini DWH") logger.info("=" * 60) - # Проверяем подключение к ClickHouse + # Создаем подключение к ClickHouse if not create_clickhouse_connection(): - logger.error("Failed to verify ClickHouse connection") + logger.error("Failed to create ClickHouse connection") sys.exit(1) # Импортируем датасеты From 6533f8b32c829bc47cbf94530496bfa21ec161f4 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Tue, 10 Feb 2026 23:49:40 +0300 Subject: [PATCH 4/7] fix(superset): fix dashboard chart rendering in Superset 4 - Why: - dashboard tiles failed with "Item with key 'echarts_bar' is not registered". - existing slice query_context stayed stale after config updates. - What: - switch Top Pages and Data Quality Summary from \'echarts_bar\' to \'dist_bar\'. - use \'groupby\' for categorical bar charts and sync this into query_context. - keep dashboard export config aligned with runtime chart definitions. - Check: - python3 -m py_compile superset/create_dashboard.py - docker compose exec -T superset python /app/superset_init/create_dashboard.py - DB check for slices 9/10: viz_type=form_data=query_context set to dist_bar --- superset/create_dashboard.py | 283 +++++++++++++----- .../dashboards/ecommerce_analytics.zip.json | 8 +- 2 files changed, 206 insertions(+), 85 deletions(-) diff --git a/superset/create_dashboard.py b/superset/create_dashboard.py index 5933336..7d9e0ed 100644 --- a/superset/create_dashboard.py +++ b/superset/create_dashboard.py @@ -5,7 +5,7 @@ ================================================================================ Назначение: - Создание чартов (Charts) на основе датасетов DM-слоя - - Создание дашборда с layout + - Создание дашборда с layout и native filters Запуск: Внутри контейнера superset: @@ -41,6 +41,7 @@ CHARTS_CONFIG = [ "label": "Total Events", "optionName": "metric_1" }, + "granularity_sqla": "event_ts", "y_axis_format": ",d", "show_trend_line": False, "time_range": "No filter" @@ -57,6 +58,7 @@ CHARTS_CONFIG = [ "label": "Unique Users", "optionName": "metric_2" }, + "granularity_sqla": "event_ts", "y_axis_format": ",d", "show_trend_line": False, "time_range": "No filter" @@ -73,6 +75,7 @@ CHARTS_CONFIG = [ "label": "Unique Sessions", "optionName": "metric_3" }, + "granularity_sqla": "event_ts", "y_axis_format": ",d", "show_trend_line": False, "time_range": "No filter" @@ -89,6 +92,7 @@ CHARTS_CONFIG = [ "label": "Avg Events/Session", "optionName": "metric_4" }, + "granularity_sqla": "event_ts", "y_axis_format": ".2f", "show_trend_line": False, "time_range": "No filter" @@ -178,15 +182,14 @@ CHARTS_CONFIG = [ }, { "slice_name": "📄 Top Pages", - "viz_type": "echarts_bar", + "viz_type": "dist_bar", "dataset_name": "v_top_pages_daily", "params": { - "x_axis": "page_url_path", + "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 @@ -195,10 +198,10 @@ CHARTS_CONFIG = [ # Качество данных { "slice_name": "🔍 Data Quality Summary", - "viz_type": "echarts_bar", + "viz_type": "dist_bar", "dataset_name": "dq_summary", "params": { - "x_axis": "layer", + "groupby": ["layer"], "metrics": [ {"expressionType": "SQL", "sqlExpression": "SUM(check_value)", "label": "Row Count"} ], @@ -225,15 +228,82 @@ DASHBOARD_CONFIG = { "description": "Аналитический дашборд для e-commerce кликстрима: трафик, конверсии, география и качество данных.", "published": True, "slug": "ecommerce-analytics", - "json_metadata": json.dumps({ - "native_filter_configuration": [ +} + + +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": None, "column": {"name": "event_date"}}], + "targets": [{"datasetId": filter_dataset_id, "column": {"name": "event_date"}}], "defaultValue": "Last week", - "scope": {"root": ["ROOT_ID"], "excluded": []}, + "scope": {"rootPath": ["ROOT_ID"], "excluded": []}, "cascadeParentIds": [], "isInstant": True }, @@ -241,8 +311,8 @@ DASHBOARD_CONFIG = { "id": "country_filter", "name": "🌍 Country", "filterType": "filter_select", - "targets": [{"datasetId": None, "column": {"name": "geo_country"}}], - "scope": {"root": ["ROOT_ID"], "excluded": []}, + "targets": [{"datasetId": filter_dataset_id, "column": {"name": "geo_country"}}], + "scope": {"rootPath": ["ROOT_ID"], "excluded": []}, "isInstant": True, "allowsMultipleValues": True, "isRequired": False @@ -251,8 +321,8 @@ DASHBOARD_CONFIG = { "id": "device_filter", "name": "📱 Device Type", "filterType": "filter_select", - "targets": [{"datasetId": None, "column": {"name": "device_type"}}], - "scope": {"root": ["ROOT_ID"], "excluded": []}, + "targets": [{"datasetId": filter_dataset_id, "column": {"name": "device_type"}}], + "scope": {"rootPath": ["ROOT_ID"], "excluded": []}, "isInstant": True, "allowsMultipleValues": True, "isRequired": False @@ -261,65 +331,83 @@ DASHBOARD_CONFIG = { "id": "browser_filter", "name": "🌐 Browser", "filterType": "filter_select", - "targets": [{"datasetId": None, "column": {"name": "browser_name"}}], - "scope": {"root": ["ROOT_ID"], "excluded": []}, + "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(): +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: - # Проверяем, существует ли уже чарт - existing = db.session.query(Slice).filter_by( - slice_name=chart_config["slice_name"] - ).first() - - if existing: - logger.info(f"Chart '{chart_config['slice_name']}' already exists (ID: {existing.id})") - created_charts.append({"id": existing.id, "title": existing.slice_name}) - continue - # Подготавливаем параметры 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"], @@ -327,24 +415,69 @@ def main(): datasource_id=dataset.id, datasource_type="table", datasource_name=dataset.table_name, - params=json.dumps(params), + 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: @@ -352,77 +485,65 @@ def main(): 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!") + logger.info("Dashboard already exists and metadata/layout were updated.") logger.info(f"Dashboard URL: /superset/dashboard/{existing.id}/") logger.info("=" * 60) - return - - # Создаём позиции чартов для layout - positions = {"DASHBOARD_VERSION_KEY": "v2"} - - # Добавляем чарты в layout (grid: 12 columns) - y_position = 0 - chart_index = 0 - - for chart in created_charts: - if chart: - positions[f"CHART-{chart['id']}"] = { - "id": f"CHART-{chart['id']}", - "type": "CHART", - "parents": ["ROOT_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 - } - } - chart_index += 1 - if chart_index % 4 == 0: - y_position += 50 - + return True + # Создаём дашборд dashboard = Dashboard( dashboard_title=DASHBOARD_CONFIG["dashboard_title"], slug=DASHBOARD_CONFIG["slug"], description=DASHBOARD_CONFIG["description"], published=DASHBOARD_CONFIG["published"], - json_metadata=DASHBOARD_CONFIG["json_metadata"], + 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__": - main() + sys.exit(0 if main() else 1) diff --git a/superset/dashboards/ecommerce_analytics.zip.json b/superset/dashboards/ecommerce_analytics.zip.json index 5a1211c..76d3e90 100644 --- a/superset/dashboards/ecommerce_analytics.zip.json +++ b/superset/dashboards/ecommerce_analytics.zip.json @@ -96,20 +96,20 @@ { "__Slice__": { "slice_name": "📄 Top Pages", - "viz_type": "echarts_bar", + "viz_type": "dist_bar", "datasource_type": "table", "datasource_name": "dm.v_top_pages_daily", - "params": "{\"datasource\": \"3__table\", \"viz_type\": \"echarts_bar\", \"x_axis\": \"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}", + "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": "echarts_bar", + "viz_type": "dist_bar", "datasource_type": "table", "datasource_name": "dm.dq_summary", - "params": "{\"datasource\": \"4__table\", \"viz_type\": \"echarts_bar\", \"x_axis\": \"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}", + "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": "Сводка по качеству данных" } } From 22f08382e598a57d9aa3c274db29c230e9668762 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Tue, 10 Feb 2026 23:51:38 +0300 Subject: [PATCH 5/7] fix(superset): align bootstrap with clickhousedb uri - Why: - Superset bootstrap used outdated ClickHouse URI format and did not fail fast on init errors. - docs and exported dashboard metadata diverged from runtime connection settings. - What: - build ClickHouse URI from env vars and use clickhousedb:// in init script. - refresh dataset metadata on existing datasets and surface import errors. - run create_dashboard during superset-init startup and align docs/exported URI references. - ignore node_modules in git. - Check: - python3 -m py_compile superset/init_superset.py - manual dashboard smoke check in UI (charts render) --- .gitignore | 1 + README.md | 2 +- docker-compose.yml | 4 +- docs/ARCHITECTURE.md | 2 +- docs/SUPERSET_DASHBOARD.md | 4 +- .../dashboards/ecommerce_analytics.zip.json | 2 +- superset/init_superset.py | 74 ++++++++++++++++--- 7 files changed, 71 insertions(+), 18 deletions(-) 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/README.md b/README.md index 99cbd7c..cd30275 100644 --- a/README.md +++ b/README.md @@ -249,7 +249,7 @@ curl http://localhost:8088/health # должно вернуть 200 ``` Что создаётся автоматически: -- **Подключение к ClickHouse**: `clickhouse_dwh` (URI: `clickhouse+connect://default@clickhouse:8123/default`) +- **Подключение к 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" diff --git a/docker-compose.yml b/docker-compose.yml index 1232b7d..7054a79 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -273,7 +273,9 @@ services: echo 'Initializing roles...' && superset init && echo 'Creating ClickHouse connection and datasets...' && - python /app/superset_init/init_superset.py || true && + 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: 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/SUPERSET_DASHBOARD.md b/docs/SUPERSET_DASHBOARD.md index 6248fa8..2fe6a04 100644 --- a/docs/SUPERSET_DASHBOARD.md +++ b/docs/SUPERSET_DASHBOARD.md @@ -93,7 +93,7 @@ make superset-dashboard # Основные make up # Запуск всех сервисов make down # Остановка сервисов -make down-v # Остановка с удалением volumes +make clean # Остановка с удалением volumes make logs service=superset # Логи сервиса # ETL @@ -120,7 +120,7 @@ make superset-restart # Перезапуск сервиса 3. Выберите **ClickHouse** 4. Введите SQLAlchemy URI: ``` - clickhouse+native://default@clickhouse:9000/default + clickhousedb://default:123456@clickhouse:8123/default ``` 5. Установите: - **Expose in SQL Lab:** ✅ diff --git a/superset/dashboards/ecommerce_analytics.zip.json b/superset/dashboards/ecommerce_analytics.zip.json index 76d3e90..ca643c0 100644 --- a/superset/dashboards/ecommerce_analytics.zip.json +++ b/superset/dashboards/ecommerce_analytics.zip.json @@ -188,7 +188,7 @@ { "__Database__": { "database_name": "clickhouse_dwh", - "sqlalchemy_uri": "clickhouse+native://default@clickhouse:9000/default", + "sqlalchemy_uri": "clickhousedb://default:123456@clickhouse:8123/default", "expose_in_sqllab": true, "allow_ctas": false, "allow_cvas": false, diff --git a/superset/init_superset.py b/superset/init_superset.py index 2afbb22..d4eb1c3 100644 --- a/superset/init_superset.py +++ b/superset/init_superset.py @@ -24,12 +24,11 @@ ================================================================================ """ -import os import sys -import json +import os import logging import subprocess -from typing import Optional +from urllib.parse import quote_plus # Настройка логирования logging.basicConfig( @@ -41,6 +40,22 @@ 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"clickhousedb://{user}:{password}@{host}:{port}/{database}" + def run_superset_cli(args): """Запуск команды superset CLI""" @@ -57,13 +72,14 @@ def run_superset_cli(args): 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', 'clickhouse+connect://default@clickhouse:8123/default' + '-u', database_uri ]) if success: @@ -74,7 +90,7 @@ def create_clickhouse_connection(): logger.info("Please create manually via UI:") logger.info("1. Go to: Settings → Database Connections → + Database") logger.info("2. Select: ClickHouse") - logger.info("3. URI: clickhouse+connect://default@clickhouse:8123/default") + logger.info(f"3. URI: {database_uri}") return False @@ -136,31 +152,65 @@ def import_datasets(): "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: {ds['table_name']}", - f"existing = db.session.query(SqlaTable).filter_by(table_name='{ds['table_name']}', schema='{ds['schema']}').first()", + 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(f'Dataset {ds['table_name']} already exists')", + " try:", + " if not existing.columns:", + " existing.fetch_metadata()", + " db.session.commit()", + f" print('Refreshed dataset metadata: {table_name}')", + " refreshed += 1", + " else:", + f" print('Dataset {table_name} already exists')", + " except Exception as e:", + f" print(f'ERROR: failed to refresh {table_name}: {{e}}')", + " db.session.rollback()", + " errors += 1", "else:", " try:", - f" dataset = SqlaTable(table_name='{ds['table_name']}', schema='{ds['schema']}', database_id=database.id, database=database, description='{ds['description']}')", + ( + " 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.flush()", - f" print(f'Created dataset: {ds['table_name']}')", + " dataset.fetch_metadata()", + " db.session.commit()", + f" print('Created dataset: {table_name}')", " imported += 1", " except Exception as e:", - f" print(f'Error creating {ds['table_name']}: {{e}}')", + f" print(f'ERROR: failed to create {table_name}: {{e}}')", " db.session.rollback()", + " errors += 1", ]) script_lines.extend([ "", - "db.session.commit()", "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) From 3bed37f5323531781bec6477c2100621ee0c8ff4 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Tue, 10 Feb 2026 23:57:38 +0300 Subject: [PATCH 6/7] docs(docs): add 10-15 minute demo script - Why: - align interview demo with recruiter requirement for 10-15 minutes - What: - add timed walkthrough with code, architecture and verification points - include fallback steps for UI issues and final speaking script - Check: - verify paths/commands against repository files and DAG ids --- docs/DEMO_SCRIPT_10_15MIN.md | 182 +++++++++++++++++++++++++++++++++++ 1 file changed, 182 insertions(+) create mode 100644 docs/DEMO_SCRIPT_10_15MIN.md 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-проверки или наблюдаемость системы.» From 99731850f8eb65ac5c344159f008c28935647e74 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Wed, 11 Feb 2026 00:12:04 +0300 Subject: [PATCH 7/7] =?UTF-8?q?fix(superset):=20=D0=B8=D1=81=D0=BF=D1=80?= =?UTF-8?q?=D0=B0=D0=B2=D0=BB=D0=B5=D0=BD=D0=B0=20=D0=B8=D0=BD=D0=B8=D1=86?= =?UTF-8?q?=D0=B8=D0=B0=D0=BB=D0=B8=D0=B7=D0=B0=D1=86=D0=B8=D1=8F=20=D0=B4?= =?UTF-8?q?=D0=B0=D1=82=D0=B0=D1=81=D0=B5=D1=82=D0=BE=D0=B2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Добавлен clickhouse-sqlalchemy в Dockerfile для поддержки диалекта - Изменен URI с clickhousedb:// на clickhouse+connect:// - Убран вызов fetch_metadata() в init_superset.py (вызывал ошибку диалекта) - Датасеты создаются без предварительного fetch_metadata Тестирование: - Чистый запуск: ✅ - Перезапуск: ✅ - API: ✅ --- Dockerfile.superset | 2 +- superset/init_superset.py | 20 ++++---------------- 2 files changed, 5 insertions(+), 17 deletions(-) diff --git a/Dockerfile.superset b/Dockerfile.superset index ada0f89..e8608b3 100644 --- a/Dockerfile.superset +++ b/Dockerfile.superset @@ -8,4 +8,4 @@ RUN apt-get update && \ USER superset -RUN pip install clickhouse-connect==0.8.* psycopg2-binary==2.9.* +RUN pip install clickhouse-connect==0.8.* clickhouse-sqlalchemy==0.3.* psycopg2-binary==2.9.* diff --git a/superset/init_superset.py b/superset/init_superset.py index d4eb1c3..e3946cc 100644 --- a/superset/init_superset.py +++ b/superset/init_superset.py @@ -54,7 +54,7 @@ def build_clickhouse_uri() -> str: host = CLICKHOUSE_HOST port = CLICKHOUSE_PORT database = quote_plus(CLICKHOUSE_DATABASE) - return f"clickhousedb://{user}:{password}@{host}:{port}/{database}" + return f"clickhouse+connect://{user}:{password}@{host}:{port}/{database}" def run_superset_cli(args): @@ -139,6 +139,7 @@ def import_datasets(): # Создаем 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", @@ -169,18 +170,7 @@ def import_datasets(): ").first()" ), "if existing:", - " try:", - " if not existing.columns:", - " existing.fetch_metadata()", - " db.session.commit()", - f" print('Refreshed dataset metadata: {table_name}')", - " refreshed += 1", - " else:", - f" print('Dataset {table_name} already exists')", - " except Exception as e:", - f" print(f'ERROR: failed to refresh {table_name}: {{e}}')", - " db.session.rollback()", - " errors += 1", + f" print('Dataset {table_name} already exists')", "else:", " try:", ( @@ -193,10 +183,8 @@ def import_datasets(): ")" ), " db.session.add(dataset)", - " db.session.flush()", - " dataset.fetch_metadata()", " db.session.commit()", - f" print('Created dataset: {table_name}')", + 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}}')",