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) # Импортируем датасеты