diff --git a/.scratch/docs-accuracy-audit/findings.md b/.scratch/docs-accuracy-audit/findings.md index 3f0c893..b1782fa 100644 --- a/.scratch/docs-accuracy-audit/findings.md +++ b/.scratch/docs-accuracy-audit/findings.md @@ -1,6 +1,15 @@ # Аудит точности технической документации против кода -Дата: 2026-06-06 · Статус: открыто (починка не начата) +Дата: 2026-06-06 · Статус: **закрыто** (все 12 находок починены 2026-06-06) + +> Резолюция. 🔴 #1–#3 правлены в `ARCHITECTURE.md` (схема `*_errors`, DQ-split как +> пересекающийся, `dds.event` как browser-driven LEFT JOIN + маркер `location_not_found`) — +> формулировки подтянуты к урокам 2–3. #4: владелец решил **убрать** `make superset-export` +> (раздел доки + target в Makefile + битый `superset/export_dashboard.py` удалены). +> 🟡 #6–#10 и 🟢 #11–#12 закрыты в `ARCHITECTURE.md`/`OPERATIONS.md`/`REPO_MAP.md` +> (добавлен перечень DQ-маркеров, параметр `wait_stg_timeout_sec`, порты Superset/креды CH, +> пропущенные артефакты). Заодно зафиксирована рамка: основной путь — Airflow, +> `scripts/`/`make` — запасной. Контекст: при переработке корневого `README.md` (mentee-first) встал вопрос, можно ли смело отправлять читателя в профильные доки — не устарели ли они сами. Прогнали сверку diff --git a/Makefile b/Makefile index 2470dec..3c7b388 100644 --- a/Makefile +++ b/Makefile @@ -1,6 +1,6 @@ .PHONY: up down clean ddl data transform logs \ reload-monitoring recover-monitoring \ - superset-init superset-dashboard superset-export superset-ui superset-restart + superset-init superset-dashboard superset-ui superset-restart COMPOSE ?= docker compose @@ -77,10 +77,6 @@ superset-init: 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" diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 4630e71..c61195c 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -146,7 +146,7 @@ CREATE TABLE stg.browser_raw ( | `location_event` | event_id | ReplacingMergeTree(src_ingest_ts) | Данные страниц | | `device_by_click` | click_id | ReplacingMergeTree(src_ingest_ts) | Устройства | | `geo_by_click` | click_id | ReplacingMergeTree(src_ingest_ts) | Гео-данные | -| `*_errors` | — | MergeTree | Строки с критичными ошибками парсинга | +| `*_errors` | — | MergeTree | Копия строк с любой ошибкой разбора: Kafka-метаданные + `raw` + `error_reason` (не типизированная копия события) | **Batch шаг наполнения ODS:** @@ -155,11 +155,12 @@ CREATE TABLE stg.browser_raw ( | `load_ods` (`sql/ods/20_stg_to_ods.sql`) | Полная пересборка ODS из STG в рамках DAG `etl_pipeline` | | `TRUNCATE ods.*` | Очистка перед пересборкой для детерминированного результата | | `INSERT ... SELECT` | Типизация валидных строк в основные ODS таблицы | -| `INSERT ... SELECT` в `*_errors` | Сохранение строк с критичными ошибками парсинга | +| `INSERT ... SELECT` в `*_errors` | Сохранение копии строк с любой ошибкой разбора (Kafka-метаданные + `raw` + `error_reason`) | -**Логика разделения:** -- **Основная таблица**: строки с валидным business key (`WHERE key IS NOT NULL`). -- **Таблица ошибок**: строки с критичными ошибками (невалидный key, timestamp/координаты и т.п.) через отдельный `INSERT ... SELECT`. +**Логика разделения (DQ-split):** +- **Основная таблица** `ods.*`: строки с валидным бизнес-ключом (`WHERE key IS NOT NULL`). Ошибка по *неключевому* полю не выкидывает строку — она остаётся, но помечается в массиве `parse_errors`. +- **Таблица ошибок** `ods.*_errors`: **копия** строк, где при разборе случилась *любая* ошибка (отдельный `INSERT ... SELECT`). У неё своя схема — не типизированное событие, а сырьё для разбора. +- Условия пересекаются нарочно: строка с валидным ключом, но битым неключевым полем попадает **и в основную таблицу, и в `*_errors`**. Подробный разбор — в разделе [«Обработка ошибок в ODS»](#обработка-ошибок-в-ods-dq-split) и в уроке 2 курса. **Пример структуры:** ```sql @@ -202,7 +203,7 @@ GROUP BY error; | Таблица | PK | Источники | JOIN-ключ | |---------|-----|-----------|-----------| -| `event` | event_id | browser_event + location_event | click_id → click | +| `event` | event_id | browser_event (ведущая) + location_event (LEFT) | click_id → click | | `click` | click_id | device_by_click + geo_by_click | — | **Структура:** @@ -301,7 +302,7 @@ LEFT JOIN ( | `v_daily_traffic` | Агрегация трафика | День × страна × устройство × браузер × UTM | | `v_top_pages_daily` | Популярность страниц | День × URL path | | `v_utm_effectiveness` | Маркетинговая аналитика | День × UTM source/medium/campaign | -| `v_session_overview` | Сессионная аналитика | День × пользователь × сессия | +| `v_session_overview` | Сессионная аналитика (только идентифицированные пользователи, `user_domain_id IS NOT NULL`) | День × пользователь × сессия | | `v_dq_errors_daily` | Мониторинг качества | День × тип ошибки | **Пример:** @@ -328,7 +329,7 @@ FROM ... ``` - `TRUNCATE` предотвращает накопление дубликатов при повторных запусках -- Хранит статистику по всем слоям (stg/ods/dds) для быстрой проверки +- Хранит статистику по всем слоям (`stg`/`ods`/`dds`/`dm`): `total_rows` по таблицам, `rows_with_errors` в ODS, `orphan_events` в DDS (события без своего клика) и `total_rows` финальной витрины `dm.v_events_enriched` — чтобы lineage замыкался на одном (event) зерне вплоть до DM **Почему VIEW:** - Для демо: достаточно производительности @@ -395,7 +396,7 @@ sequenceDiagram ```mermaid erDiagram - BROWSER_EVENT ||--|| LOCATION_EVENT : "event_id" + BROWSER_EVENT ||--o| LOCATION_EVENT : "event_id (LEFT)" BROWSER_EVENT ||--o| DEVICE_BY_CLICK : "click_id" BROWSER_EVENT ||--o| GEO_BY_CLICK : "click_id" @@ -424,6 +425,7 @@ erDiagram UUID click_id PK String os String os_name + String os_timezone String device_type UInt8 device_is_mobile String user_custom_id @@ -443,11 +445,11 @@ erDiagram ### Сборка DDS-сущностей -**event** (browser + location): +**event** (browser + location), `browser_event` — ведущая таблица: ```mermaid flowchart LR subgraph ODS["ODS"] - B["browser_event"] + B["browser_event
(ведущая)"] L["location_event"] end @@ -455,10 +457,12 @@ flowchart LR EV["event"] end - B -->|JOIN event_id| EV - L -->|JOIN event_id| EV + B -->|все event_id| EV + L -.->|LEFT JOIN event_id| EV ``` +**Важно:** сборка идёт **от browser** через `LEFT JOIN location`. `event_id`, которые есть только в `location_event` (без browser), в `dds.event` **не попадают**. Если у события нет своей location-строки — событие остаётся, поля страницы/UTM пустые (`NULL`), и в `ods_parse_errors` ставится маркер `location_not_found`. Тот же принцип, что и в `dds.click`: не теряем, а оставляем видимый след. Подробнее — урок 3 курса (`docs/course/lessons/03_ods_to_dds.md`). + **click** (device + geo) с поддержкой partial data: ```mermaid flowchart LR @@ -512,24 +516,37 @@ flowchart LR | **MV + JOIN** | Реалтайм | Eventual consistency, дубли при late arrival | | **Batch (выбрано)** | Согласованность, контроль | Задержка до следующего запуска | -### Обработка ошибок в ODS +### Обработка ошибок в ODS (DQ-split) **Проблема:** Грязные данные могут содержать не только невалидные ключи, но и невалидные timestamp/координаты/ID. -**Решение:** -1. **Основная таблица**: строки с валидным business key (`WHERE key IS NOT NULL`) -2. **Таблица ошибок**: строки с критичными ошибками парсинга через отдельные `INSERT ... SELECT` -3. **DQ-метрики**: массив `parse_errors` для аудита +**Решение — DQ-split:** +1. **Основная таблица** `ods.*`: строки с валидным бизнес-ключом (`WHERE key IS NOT NULL`). Ошибки по *неключевым* полям не выкидывают строку — она остаётся, но помечается массивом `parse_errors`. +2. **Таблица ошибок** `ods.*_errors`: **копия** строк, где при разборе случилась *любая* ошибка. Это не типизированная копия события, а сырьё для разбора — метаданные доставки из Kafka, исходный JSON и причина ошибки: + - `ingest_ts`, `kafka_topic`, `kafka_partition`, `kafka_offset`, `kafka_ts` — координаты сообщения в Kafka; + - `raw` — исходный JSON «как пришёл»; + - `error_reason` — список несработавших полей одной строкой (`arrayStringConcat(parse_errors, ',')`). +3. **DQ-метрики**: массив `parse_errors` в основной таблице для аудита. + +**Одна строка может попасть в оба места — это не баг, а замысел.** Строка с валидным ключом, но битым неключевым полем (например, валидный `event_id`, но `event_ts IS NULL`) и **остаётся** в основной таблице (с меткой в `parse_errors`), и **копируется** в `*_errors`. Основная таблица отвечает на вопрос «что есть для работы», таблица ошибок — «что пришло битым и требует разбора». Подробный разбор — в уроке 2 курса (`docs/course/lessons/02_stg_to_ods.md`). ```sql --- Основная таблица -INSERT INTO ods.browser_event -SELECT ... FROM stg.browser_raw WHERE event_id IS NOT NULL; +-- event_id, event_ts, click_id, parse_errors — это не колонки stg.browser_raw, +-- а алиасы из блока WITH, где raw (сырой JSON) разбирается через JSONExtract*/*OrNull. +-- Здесь WITH опущен для краткости; полная версия — в sql/ods/20_stg_to_ods.sql. --- Таблица ошибок +-- Основная таблица: валидный ключ; parse_errors помечает битые неключевые поля +INSERT INTO ods.browser_event +SELECT ..., parse_errors FROM stg.browser_raw WHERE event_id IS NOT NULL; + +-- Таблица ошибок: другая схема (Kafka-метаданные + raw + error_reason); +-- сюда едет копия любой строки с хотя бы одной ошибкой разбора INSERT INTO ods.browser_event_errors -SELECT ... FROM stg.browser_raw -WHERE event_id IS NULL OR event_ts IS NULL OR click_id IS NULL; +SELECT ingest_ts, kafka_topic, kafka_partition, kafka_offset, kafka_ts, + raw, arrayStringConcat(parse_errors, ',') AS error_reason +FROM stg.browser_raw +WHERE length(parse_errors) > 0 + AND (event_id IS NULL OR event_ts IS NULL OR click_id IS NULL); ``` ### Partial data в DDS @@ -539,7 +556,22 @@ WHERE event_id IS NULL OR event_ts IS NULL OR click_id IS NULL; **Решение:** 1. **UNION DISTINCT** всех click_id из обоих источников 2. **LEFT JOIN** для получения данных (обрабатываем device-only и geo-only) -3. **DQ-маркеры**: `device_not_found`, `geo_not_found` в `parse_errors` +3. **DQ-маркеры** в `ods_parse_errors`: `device_not_found`, `geo_not_found` (нет соответствующего источника по `click_id`), `geo_country_missing` (гео есть, но страна не определена) + +### Перечень DQ-маркеров + +Маркеры качества копятся в массивах `parse_errors` (ODS) и `ods_parse_errors` (DDS). Полный список: + +| Слой / таблица | Маркеры | Когда ставится | +|----------------|---------|----------------| +| ODS `browser_event` | `bad_event_id`, `bad_event_timestamp`, `bad_click_id` | поле не разобралось в нужный тип | +| ODS `location_event` | `bad_event_id` | не разобрался `event_id` | +| ODS `device_by_click` | `bad_click_id`, `bad_user_domain_id` | не разобрались `click_id` / `user_domain_id` | +| ODS `geo_by_click` | `bad_click_id`, `bad_geo_latitude`, `bad_geo_longitude` | не разобрались ключ или координаты | +| DDS `click` | `device_not_found`, `geo_not_found`, `geo_country_missing` | нет источника по `click_id` либо страна не определена | +| DDS `event` | `location_not_found` | у события нет своей location-строки | + +В DDS наследуются `parse_errors` только **ведущего** источника (device → `dds.click`, browser → `dds.event`); к ним добавляются маркеры стыковки (`*_not_found`, `*_missing`). `parse_errors` из geo/location в DDS не переносятся. --- diff --git a/docs/OPERATIONS.md b/docs/OPERATIONS.md index 9470e80..375df1c 100644 --- a/docs/OPERATIONS.md +++ b/docs/OPERATIONS.md @@ -22,13 +22,14 @@ Порты задаются в `docker-compose.yml`: -- ClickHouse native: `localhost:8002` -- ClickHouse HTTP: `localhost:9123` +- ClickHouse native: `localhost:8002` (пользователь `default`, пароль `123456`) +- ClickHouse HTTP / play-консоль: `http://localhost:9123/play` (`default` / `123456`) - Kafka: `localhost:9092` - Kafka UI: `http://localhost:8082` - Airflow: `http://localhost:8080` (`admin/admin`) +- Superset: `http://localhost:8088` (`admin/admin`) - Prometheus: `http://localhost:9090` -- Grafana: `http://localhost:3000` +- Grafana: `http://localhost:3000` (`admin/admin`) ## Airflow DAGs @@ -57,7 +58,9 @@ ### `etl_pipeline` - Запуск: ручной (`Trigger DAG with config`) -- Параметр: `full_refresh` (`bool`, default `true`) — очистить DDS перед загрузкой +- Параметры: + - `full_refresh` (`bool`, default `true`) — очистить DDS перед загрузкой + - `wait_stg_timeout_sec` (`int`, default `600`, minimum `30`) — сколько секунд задача `wait_for_stg_data` ждёт появления данных в STG, прежде чем упасть по таймауту - Зависимость: требует наличия данных в STG (от `kafka_load` или `make data`) - Гейт целостности DDS: `check_dds_integrity` считает события без клика, а `assert_dds_integrity` роняет DAG при `orphan_events > 0`. Проверка идёт после diff --git a/docs/REPO_MAP.md b/docs/REPO_MAP.md index d313b2b..05bebb1 100644 --- a/docs/REPO_MAP.md +++ b/docs/REPO_MAP.md @@ -4,25 +4,40 @@ ## Исполняемые файлы -### Airflow +### Airflow (основной путь запуска) - `airflow/dags/ddl_init_dag.py` — инициализация схемы ClickHouse - `airflow/dags/kafka_load_dag.py` — загрузка в Kafka из JSONL - `airflow/dags/etl_pipeline_dag.py` — ETL процесс STG -> ODS -> DDS -> DM - `airflow/dags/utils/kafka_helpers.py` — helper-функции для Kafka +- `airflow/dags/utils/sql_helpers.py` — чтение и подготовка SQL-файлов для DAG +- `airflow/dags/utils/airflow_params.py` — разбор и валидация параметров DAG - `airflow/requirements.txt` — зависимости Airflow/ClickHouse plugin ### SQL +DDL (форма таблиц): + - `sql/ddl/00_databases.sql` — создание БД `stg`/`ods`/`dds`/`dm` - `sql/ddl/stg/10_stg.sql` — STG (Kafka Engine + MV) -- `sql/ddl/ods/20_ods.sql` — ODS (типизация + MV для ошибок) +- `sql/ddl/ods/20_ods.sql` — ODS: типизированные таблицы и `*_errors` (наполняются batch, не MV) - `sql/ddl/dds/30_dds.sql` — DDS (таблицы для batch-загрузки) - `sql/ddl/dm/40_dm.sql` — DM (витрины VIEW) -- `sql/dds/30_ods_to_dds.sql` — ODS -> DDS (argMax + JOIN) -- `sql/dm/40_dds_to_dm.sql` — обновление `dq_summary` -### Скрипты +Трансформации (наполнение, шаги `etl_pipeline`): + +- `sql/ods/20_stg_to_ods.sql` — STG -> ODS: типизация + DQ-split (валидный ключ → `ods.*`, любая ошибка → `ods.*_errors`) +- `sql/dds/30_ods_to_dds.sql` — ODS -> DDS (argMax + LEFT JOIN) +- `sql/dm/40_dds_to_dm.sql` — DDS -> DM: пересборка `dm.dq_summary` (TRUNCATE+INSERT) по всем слоям; сами витрины `dm.v_*` — это VIEW из DDL + +### Superset + +- `superset/init_superset.py` — подключение к ClickHouse + создание датасетов +- `superset/create_dashboard.py` — сборка дашборда с чартами + +### Скрипты (запасной путь, не основной) + +Shell-скрипты `scripts/*` (и обёртки `make ddl`/`make data`/`make transform`) — это локальный fallback в обход Airflow. Основной путь запуска — DAG-и Airflow (см. выше). - `scripts/apply_clickhouse_ddl.sh` — применение DDL - `scripts/load_kafka_data.sh` — загрузка в Kafka @@ -48,8 +63,12 @@ - `README.md` — быстрый старт и обзор проекта - `docs/ARCHITECTURE.md` — техническая архитектура - `docs/OPERATIONS.md` — запуск, проверки, troubleshooting +- `docs/SUPERSET_DASHBOARD.md` — настройка и использование дашборда Superset - `docs/DE-task.md` — исходное задание - `docs/COMMIT_RULES.md` — правила коммитов +- `docs/course/` — продвинутый учебный курс на базе стенда (PRD, план, уроки) +- `docs/adr/` — архитектурные решения (ADR) +- `docs/agents/` — контракты для агентских скиллов (issue-tracker, triage, domain) ## Legacy-планы diff --git a/docs/SUPERSET_DASHBOARD.md b/docs/SUPERSET_DASHBOARD.md index d778227..78b4541 100644 --- a/docs/SUPERSET_DASHBOARD.md +++ b/docs/SUPERSET_DASHBOARD.md @@ -142,7 +142,6 @@ make transform # Запуск batch-процесса # Superset make superset-init # Инициализация (подключение + датасеты) make superset-dashboard # Создание дашборда -make superset-export # Экспорт дашборда в JSON make superset-ui # Показать URL и логин make superset-restart # Перезапуск сервиса ``` @@ -195,22 +194,16 @@ python /app/superset_init/init_superset.py --- -## Экспорт и импорт дашборда +## Импорт дашборда -### Экспорт - -```bash -# Автоматический экспорт в JSON -make superset-export - -# Результат: superset/dashboards/ecommerce_analytics.json -``` - -### Импорт +Основной способ собрать дашборд — `make superset-dashboard` (скрипт `create_dashboard.py`). +Готовый экспорт дашборда лежит в репозитории на случай ручного импорта: +`superset/dashboards/ecommerce_analytics.zip.json` (внутри контейнера — +`/app/superset_init/dashboards/ecommerce_analytics.zip.json`). ```bash # Импорт через CLI -docker compose exec superset superset import-dashboards -p /app/superset_init/dashboards/ecommerce_analytics.json +docker compose exec superset superset import-dashboards -p /app/superset_init/dashboards/ecommerce_analytics.zip.json # Или через UI: Settings → Import Dashboards ``` @@ -289,7 +282,7 @@ docker compose exec superset bash -c "curl clickhouse:8123" | Сервис | URL | Логин/Пароль | |--------|-----|--------------| | Superset | http://localhost:8088 | admin / admin | -| ClickHouse HTTP | http://localhost:9123 | default / (пустой) | +| ClickHouse HTTP | http://localhost:9123 | default / 123456 | | Airflow | http://localhost:8080 | admin / admin | | Grafana | http://localhost:3000 | admin / admin | | Prometheus | http://localhost:9090 | - | diff --git a/superset/export_dashboard.py b/superset/export_dashboard.py deleted file mode 100644 index 92395cc..0000000 --- a/superset/export_dashboard.py +++ /dev/null @@ -1,107 +0,0 @@ -#!/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)