diff --git a/AGENTS.md b/AGENTS.md index 5a92a4b..20becea 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -12,6 +12,7 @@ - Kafka (источник событий, 1 JSON message = 1 event) - ClickHouse (STG → ODS → DDS → DM) +- **Airflow** (оркестрация ETL-пайплайна) - Superset (BI поверх витрин) - Prometheus + Grafana (мониторинг) - простой инструмент/скрипт, который читает `.jsonl` и пишет события в Kafka @@ -19,19 +20,21 @@ ## Ключевые артефакты ### Исполняемые файлы (текущая структура) -- `ddl/` — SQL для создания объектов БД: - - `00_databases.sql` — создание БД stg/ods/dds/dm - - `10_stg.sql` — STG слой (Kafka Engine + MV) - - `20_ods.sql` — ODS слой (типизация + MV для ошибок) - - `30_dds.sql` — DDS слой (таблицы для batch-загрузки) - - `40_dm.sql` — DM слой (витрины VIEW) -- `jobs/` — batch-трансформации: - - `30_dds_refresh.sql` — ODS → DDS (argMax + JOIN) - - `40_dm_refresh.sql` — обновление DQ_summary +- `dags/` — Airflow DAGs для оркестрации ETL +- `sql/` — SQL по слоям: + - `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/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 - `scripts/` — скрипты автоматизации: - `apply_clickhouse_ddl.sh` — применение DDL - `load_kafka_data.sh` — загрузка в Kafka - `run_batch.sh` — запуск batch-процесса +- `airflow/` — конфигурация Airflow: + - `requirements.txt` — зависимости Airflow/ClickHouse plugin ### Планы и документация (legacy) - `plans/clickhouse_ddl.md` — исходный план (inline DDL, legacy) @@ -53,12 +56,13 @@ Базовые команды: - `make up` (или `docker compose up -d`) -- `make ddl` (применяет SQL из `ddl/*.sql` в ClickHouse) +- `make ddl` (применяет SQL из `sql/ddl/00_databases.sql` и `sql/ddl/*/*.sql` в ClickHouse) - `make data` (пересоздаёт топики и заливает небольшой срез данных в Kafka; полный режим — `FULL=1 make data`) - `make transform` (запускает batch-процесс ODS → DDS → DM) - `docker compose up -d` - `docker compose ps` - `docker compose logs -f --tail=200 ` +- `docker compose down` (сохраняет named volumes, включая `clickhouse-data`) - `docker compose down -v` (удалит volumes; используйте осознанно) Порты (см. `docker-compose.yml`): @@ -67,6 +71,7 @@ - ClickHouse HTTP: `localhost:9123` - Kafka: `localhost:9092` - Kafka UI: `http://localhost:8082` +- **Airflow: `http://localhost:8080` (admin/admin)** - Prometheus: `http://localhost:9090` - Grafana: `http://localhost:3000` @@ -80,15 +85,17 @@ - Держать изменения минимальными и по теме задания (инфра, схема, ingest, витрины). - Не коммитить секреты. Если требуется пароль/ключи — использовать `.env` и примеры `.env.example`. - README/планы обновлять вместе с изменениями инфраструктуры/DDL. +- Для спорных или меняющихся API (особенно Airflow/operators/providers) проверять актуальную документацию через `context7` и фиксировать решение в коде/документации. - **Комментарии в коде — на русском языке**: - SQL: заголовочный блок с описанием файла, комментарии к каждому логическому блоку - Bash: шапка с назначением/запуском/требованиями, секции разделены `# -----` - - См. существующие файлы как пример (`ddl/20_ods.sql`, `jobs/30_dds_refresh.sql`, `scripts/run_batch.sh`) + - См. существующие файлы как пример (`sql/ddl/ods/20_ods.sql`, `sql/dds/30_ods_to_dds.sql`, `scripts/run_batch.sh`) ## Быстрые проверки - Kafka ingest: наличие данных в `stg.*` и типизированных строк в `ods.*`. - Мониторинг: доступность `/metrics` у ClickHouse и скрейп в Prometheus. +- **Airflow: `http://localhost:8080` должен показывать UI и DAG `ddl_init` и `etl_pipeline`.** - BI: витрина `dm.v_events_enriched` должна отвечать за разумное время при фильтре по дате. ## Связанная документация diff --git a/Dockerfile.airflow b/Dockerfile.airflow new file mode 100644 index 0000000..b857115 --- /dev/null +++ b/Dockerfile.airflow @@ -0,0 +1,23 @@ +# Use the official Airflow image as base +FROM apache/airflow:2.10.5 + +# Set environment variables +ENV AIRFLOW_HOME=/opt/airflow + +# Copy the project requirements file +COPY --chown=airflow:0 airflow/requirements.txt ${AIRFLOW_HOME}/requirements.txt + +# Install the required Python packages +RUN pip install --no-cache-dir -r ${AIRFLOW_HOME}/requirements.txt + +# Create necessary directories +RUN mkdir -p /opt/airflow/dags /opt/airflow/logs /opt/airflow/config /opt/airflow/data + +# Set working directory +WORKDIR ${AIRFLOW_HOME} + +# Expose the webserver port +EXPOSE 8080 + +# The image will use Airflow's default entrypoint +# Commands will be passed at runtime diff --git a/README.md b/README.md index ac8eff7..02867ff 100644 --- a/README.md +++ b/README.md @@ -1,110 +1,130 @@ # ClickHouse Mini DWH для кликстрима -[![Stack](https://img.shields.io/badge/stack-Kafka%20%7C%20ClickHouse%20%7C%20Superset-blue)](./docker-compose.yml) +[![Stack](https://img.shields.io/badge/stack-Kafka%20%7C%20ClickHouse%20%7C%20Airflow%20%7C%20Superset-blue)](./docker-compose.yml) [![Layers](https://img.shields.io/badge/layers-STG%20→%20ODS%20→%20DDS%20→%20DM-green)](./docs/ARCHITECTURE.md) [![License](https://img.shields.io/badge/license-Educational-orange)]() -Многослойное хранилище данных (STG → ODS → DDS → DM) для анализа кликстрима e-commerce. +Мини-демо для решения задания [DE-task.md](./data/DE-task.md): развернуть инфраструктуру на своей машине, прогнать кликстрим через Kafka в ClickHouse, сделать регулярный расчёт в Airflow и подготовить витрины под дашборд. -Данные поступают из Kafka, проходят типизацию и обогащение, формируя витрины для BI-аналитики. +Фокус проекта: быстро показать работающий end-to-end сценарий и понятным языком объяснить, как устроены слои и почему пайплайн не падает на "грязных" данных. -> **Соответствие заданию:** Реализован полный цикл Data Engineering: ingestion → хранилище со слоями → регулярный процесс трансформации → витрины для дашборда. +Коротко про поток: +`data/*.jsonl` -> Kafka (1 строка = 1 сообщение) -> ClickHouse `stg` (сырые JSON) -> `ods` (типизация + DQ) -> Airflow batch -> `dds` (сущности) -> `dm` (витрины VIEW) -> Superset. --- -## 🚀 Быстрый старт +## Быстрый старт (демо-сценарий) ```bash -# 1. Поднять инфраструктуру (Kafka + ClickHouse + Superset) +# 1) Поднять инфраструктуру make up -# 2. Создать структуру БД -make ddl - -# 3. Загрузить данные (автоматически потекут STG → ODS) -make data # первые 50 строк -# или: FULL=1 make data # полный датасет (1000 строк) - -# 4. Подождать 5-10 сек (данные проходят через Kafka) -sleep 10 - -# 5. Запустить batch-трансформацию (ODS → DDS → DM) -make transform +# Проверить статусы контейнеров +docker compose ps ``` -**Проверка:** -```bash -# Статистика по слоям -docker compose exec clickhouse clickhouse-client \ - --user=default --password=123456 --query=" - SELECT database, countDistinct(table) AS tables, sum(rows) AS rows - FROM system.parts WHERE database IN ('stg','ods','dds','dm') - GROUP BY database ORDER BY database -" +Дальше основной путь идёт через Airflow (как в задании). -# Пример запроса к витрине -docker compose exec clickhouse clickhouse-client \ - --user=default --password=123456 --query=" - SELECT * FROM dm.v_utm_effectiveness ORDER BY clicks DESC LIMIT 5 -" +1. Открыть Airflow UI: `http://localhost:8080` (admin/admin) +2. Включить (unpause) и запустить `ddl_init` (создаёт базы/таблицы/VIEW в ClickHouse) + +Опционально можно триггернуть DAG из CLI (удобно для CI/скрипта): +```bash +docker compose exec -T airflow-webserver airflow dags trigger ddl_init +``` + +Загрузка небольшого среза данных в Kafka: +```bash +make data # по умолчанию первые 50 строк +# или: FULL=1 make data # полный датасет (1000 строк) +``` + +Запуск batch-трансформации (ODS -> DDS -> DM) в Airflow (если DAG выключен, сначала unpause): +```bash +docker compose exec -T airflow-webserver airflow dags trigger etl_pipeline \ + --conf '{"full_refresh": true}' +``` + +Smoke-check результата в ClickHouse: +```bash +docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query \ + "SELECT 'ods.browser_event' AS t, count() AS rows FROM ods.browser_event" +docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query \ + "SELECT 'dds.click' AS t, count() AS rows FROM dds.click" +docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query \ + "SELECT 'dds.event' AS t, count() AS rows FROM dds.event" +docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query \ + "SELECT 'dm.dq_summary' AS t, count() AS rows FROM dm.dq_summary" ``` --- -## 📊 Доступные сервисы +## Доступные сервисы | Сервис | 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 | Визуализация метрик | --- -## 🏗️ Архитектура +## Архитектура (в двух словах) ```mermaid flowchart TB - subgraph Sources["📁 JSON файлы"] + subgraph Sources["JSONL файлы"] BE[browser_events.jsonl] LE[location_events.jsonl] DE[device_events.jsonl] GE[geo_events.jsonl] end - subgraph Kafka["🚀 Kafka"] + subgraph Kafka["Kafka"] KT[Топики] end - subgraph CH["🗄️ ClickHouse"] - STG["STG — сырые JSON"] - ODS["ODS — типизированные"] - DDS["DDS — сущности"] - DM["DM — витрины"] + subgraph CH["ClickHouse"] + STG["stg: сырьё + Kafka MV"] + ODS["ods: типизация + DQ"] + DDS["dds: сущности"] + DM["dm: витрины (VIEW)"] + end + + subgraph Airflow["Airflow"] + DAG[DAG: ddl_init / etl_pipeline] end Sources -->|make data| Kafka -->|MV| STG -->|MV| ODS -->|Batch SQL| DDS -->|VIEW| DM + DAG -.->|оркестрация| ODS & DDS & DM ``` -**Поток данных:** -1. **STG** — сырые JSON из Kafka (MergeTree) -2. **ODS** — типизированные данные + DQ (ReplacingMergeTree) -3. **DDS** — собранные сущности event + click (Batch SQL) -4. **DM** — витрины для BI (VIEW) +Особенность задания про "грязные данные": парсинг не валит pipeline, ошибки фиксируются в `ods.*_errors` и в поле `parse_errors`. [Подробное описание архитектуры →](./docs/ARCHITECTURE.md) --- -## 📁 Структура проекта +## Структура проекта ``` . -├── ddl/ # SQL для создания объектов (00_databases → 40_dm) -├── jobs/ # Batch-трансформации (ODS→DDS, DDS→DM) +├── dags/ # Airflow DAGs для оркестрации +├── sql/ +│ ├── ddl/ # DDL по слоям +│ │ ├── 00_databases.sql +│ │ ├── stg/10_stg.sql +│ │ ├── ods/20_ods.sql +│ │ ├── dds/30_dds.sql +│ │ └── dm/40_dm.sql +│ ├── dds/ # Batch SQL: ODS -> DDS +│ └── dm/ # Batch SQL: DDS -> DM ├── scripts/ # Автоматизация (apply ddl, load data, run batch) +├── airflow/ # Конфигурация Airflow +│ └── requirements.txt ├── docs/ # Документация │ └── ARCHITECTURE.md # Подробное описание слоёв ├── data/ # Исходные JSONL файлы @@ -114,19 +134,24 @@ flowchart TB --- -## 🛠️ Команды Makefile +## Команды Makefile | Команда | Описание | |---------|----------| | `make up` | Поднять инфраструктуру | -| `make ddl` | Создать структуру БД | +| `make ddl` | Применить DDL в ClickHouse (вне Airflow) | | `make data` | Загрузить данные в Kafka (50 строк) | | `FULL=1 make data` | Загрузить полный датасет | -| `make transform` | Запустить batch-процесс | +| `make transform` | Запустить batch-процесс (вне Airflow) | + +Примечания про сохранность данных: +- Данные ClickHouse сохраняются в Docker volume `clickhouse-data`. +- Данные Kafka сохраняются в Docker volume `kafka-data`. +- `docker compose down` сохраняет named volumes, `docker compose down -v` удаляет их (и данные пропадут). --- -## 🔗 Ключи данных +## Ключи данных (как джойним) ```mermaid flowchart LR @@ -162,38 +187,46 @@ flowchart LR --- -## 📚 Документация +## Дашборд в Superset (опционально, но полезно) + +1. Открыть `http://localhost:8088` +2. Database -> Add: + - URI: `clickhouse+connect://default:123456@clickhouse:8123/default` +3. Создать datasets из `dm.v_*` (VIEW) и собрать несколько графиков + +Идеи графиков под задание: +- Трафик по дням: `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) + +--- + +## Частые проблемы + +- `etl_pipeline` падает с сообщением про схему: сначала запустите `ddl_init`. +- После `docker compose down -v` схема и данные исчезнут: нужно заново `ddl_init` и `make data`. +- Подключения используют разные протоколы: + - Airflow (ClickHouseOperator) ходит в ClickHouse по native TCP (порт `9000` внутри сети Docker). + - Superset (clickhouse-connect) ходит по HTTP (порт `8123` внутри сети Docker). + +--- + +## Статус проекта + +Реализовано (Этап 1): +- DAG `ddl_init`: последовательное применение DDL + проверка схемы. +- DAG `etl_pipeline`: precheck, ожидание данных в ODS, пересчёт DDS/DM, базовые проверки. +- Устойчивость к "грязным" данным: ошибки парсинга сохраняются в ODS, а не валят ingest. + +В планах (не требуется для MVP задания): +- DAG `kafka_load` (чистый ingest из `.jsonl` в Kafka средствами Airflow). +- Инкрементальный batch (watermark вместо `full_refresh`). +- DQ мониторинг по расписанию. + +--- + +## Документация - [Архитектура и слои](./docs/ARCHITECTURE.md) — подробное описание STG/ODS/DDS/DM, ER-диаграммы, обоснование решений - [DE-task.md](./data/DE-task.md) — исходное задание - ---- - -## 🎯 Дашборд в Superset - -1. Открыть http://localhost:8088 -2. Database → Add: - - **URI:** `clickhouse+connect://default:123456@clickhouse:8123/default` -3. Datasets → Add from `dm.v_*` -4. Charts & Dashboard - -Основные витрины: -- `v_events_enriched` — полное обогащение -- `v_daily_traffic` — агрегация по дням -- `v_utm_effectiveness` — эффективность кампаний -- `v_top_pages_daily` — воронка страниц - ---- - -## 🔮 Развитие проекта - -- [ ] **Airflow** — оркестрация batch-процесса -- [ ] **Инкрементальный batch** — watermark-based загрузка -- [ ] **Материализация витрин** — для тяжёлых агрегаций -- [ ] **DQ мониторинг** — алерты на ошибки парсинга - ---- - -## 📝 Лицензия - -Проект создан для образовательных целей в рамках DE-тестового задания. diff --git a/airflow/requirements.txt b/airflow/requirements.txt new file mode 100644 index 0000000..feae4fb --- /dev/null +++ b/airflow/requirements.txt @@ -0,0 +1,7 @@ +# Airflow requirements для учебного ETL-проекта + +# Metadata DB для Airflow +psycopg2-binary==2.9.9 + +# ClickHouse operator/hook для DAG'ов +airflow-clickhouse-plugin==1.6.0 diff --git a/dags/.gitkeep b/dags/.gitkeep new file mode 100644 index 0000000..e69de29 diff --git a/dags/__init__.py b/dags/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/dags/__pycache__/ddl_init_dag.cpython-312.pyc b/dags/__pycache__/ddl_init_dag.cpython-312.pyc new file mode 100644 index 0000000..9f0a801 Binary files /dev/null and b/dags/__pycache__/ddl_init_dag.cpython-312.pyc differ diff --git a/dags/__pycache__/etl_pipeline_dag.cpython-312.pyc b/dags/__pycache__/etl_pipeline_dag.cpython-312.pyc new file mode 100644 index 0000000..aa2e353 Binary files /dev/null and b/dags/__pycache__/etl_pipeline_dag.cpython-312.pyc differ diff --git a/dags/ddl_init_dag.py b/dags/ddl_init_dag.py new file mode 100644 index 0000000..9529177 --- /dev/null +++ b/dags/ddl_init_dag.py @@ -0,0 +1,188 @@ +""" +DAG инициализации DDL в ClickHouse. + +Учебный формат: +- каждая операция DDL выполняется отдельной SQL-task; +- SQL-файлы вызываются явно по фиксированным путям; +- режим verify_only позволяет прогонять только проверки схемы. +""" + +from __future__ import annotations + +import re +from datetime import datetime, timedelta +from pathlib import Path + +from airflow import DAG +from airflow.exceptions import AirflowException +from airflow.models.param import Param +from airflow.operators.empty import EmptyOperator +from airflow.operators.python import BranchPythonOperator, PythonOperator +from airflow.utils.trigger_rule import TriggerRule +from airflow_clickhouse_plugin.operators.clickhouse import ClickHouseOperator + + +# ----------------------------------------------------------------------------- +# Базовые настройки DAG +# ----------------------------------------------------------------------------- +default_args = { + "owner": "airflow", + "depends_on_past": False, + "email_on_failure": False, + "email_on_retry": False, + "retries": 1, + "retry_delay": timedelta(minutes=2), +} + + +# ----------------------------------------------------------------------------- +# SQL-файлы проекта +# ----------------------------------------------------------------------------- +SQL_ROOT = Path(__file__).resolve().parents[1] / "sql" + + +def load_sql_statements(relative_path: str) -> tuple[str, ...]: + """Читает SQL-файл и делит его на отдельные команды по ';'.""" + file_path = SQL_ROOT / relative_path + if not file_path.is_file(): + raise AirflowException(f"SQL-файл не найден: {file_path}") + + sql_text = file_path.read_text(encoding="utf-8") + statements: list[str] = [] + for segment in sql_text.split(";"): + # Убираем блочные и строковые комментарии, чтобы не отправлять "пустые" запросы. + no_block_comments = re.sub(r"/\*.*?\*/", "", segment, flags=re.S) + lines = [line for line in no_block_comments.splitlines() if not line.strip().startswith("--")] + cleaned = "\n".join(lines).strip() + if cleaned: + statements.append(cleaned) + + if not statements: + raise AirflowException(f"SQL-файл пустой: {file_path}") + return tuple(statements) + + +# ----------------------------------------------------------------------------- +# SQL-проверки +# ----------------------------------------------------------------------------- +SQL_CHECK_CLICKHOUSE = "SELECT 1 AS ok" + +SQL_VERIFY_SCHEMA = """ +SELECT + (SELECT count() FROM system.tables WHERE database = 'stg' AND name = 'browser_raw') AS stg_browser_raw, + (SELECT count() FROM system.tables WHERE database = 'ods' AND name = 'browser_event') AS ods_browser_event, + (SELECT count() FROM system.tables WHERE database = 'dds' AND name = 'click') AS dds_click, + (SELECT count() FROM system.tables WHERE database = 'dds' AND name = 'event') AS dds_event, + (SELECT count() FROM system.tables WHERE database = 'dm' AND name = 'v_events_enriched') AS dm_v_events_enriched +""" + + +# ----------------------------------------------------------------------------- +# Управляющие функции +# ----------------------------------------------------------------------------- +def choose_ddl_mode(**context) -> str: + """Выбирает ветку выполнения: full DDL или только verify.""" + dag_run = context.get("dag_run") + conf = dag_run.conf if dag_run else {} + verify_only = bool(conf.get("verify_only", context["params"]["verify_only"])) + return "skip_ddl" if verify_only else "ddl_00_databases" + + +def assert_schema_ready(**context) -> None: + """Проверяет результат финальной SQL-проверки схемы.""" + ti = context["ti"] + result = ti.xcom_pull(task_ids="verify_schema_sql") + + if not result or not result[0] or len(result[0]) != 5: + raise AirflowException(f"Некорректный результат проверки схемы: {result}") + + if any(value == 0 for value in result[0]): + raise AirflowException( + "Схема применена не полностью. Проверьте таблицы/VIEW stg, ods, dds, dm." + ) + + +with DAG( + dag_id="ddl_init", + description="Инициализация схемы ClickHouse (stg/ods/dds/dm)", + default_args=default_args, + schedule=None, + start_date=datetime(2024, 1, 1), + catchup=False, + max_active_runs=1, + is_paused_upon_creation=True, + tags=["ddl", "bootstrap", "clickhouse"], + params={ + "verify_only": Param(False, type="boolean"), + }, +) as dag: + check_clickhouse = ClickHouseOperator( + task_id="check_clickhouse", + sql=SQL_CHECK_CLICKHOUSE, + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + choose_mode = BranchPythonOperator( + task_id="choose_mode", + python_callable=choose_ddl_mode, + ) + + ddl_00_databases = ClickHouseOperator( + task_id="ddl_00_databases", + sql=load_sql_statements("ddl/00_databases.sql"), + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + ddl_10_stg = ClickHouseOperator( + task_id="ddl_10_stg", + sql=load_sql_statements("ddl/stg/10_stg.sql"), + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + ddl_20_ods = ClickHouseOperator( + task_id="ddl_20_ods", + sql=load_sql_statements("ddl/ods/20_ods.sql"), + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + ddl_30_dds = ClickHouseOperator( + task_id="ddl_30_dds", + sql=load_sql_statements("ddl/dds/30_dds.sql"), + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + ddl_40_dm = ClickHouseOperator( + task_id="ddl_40_dm", + sql=load_sql_statements("ddl/dm/40_dm.sql"), + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + skip_ddl = EmptyOperator(task_id="skip_ddl") + + ddl_complete = EmptyOperator( + task_id="ddl_complete", + trigger_rule=TriggerRule.NONE_FAILED_MIN_ONE_SUCCESS, + ) + + verify_schema_sql = ClickHouseOperator( + task_id="verify_schema_sql", + sql=SQL_VERIFY_SCHEMA, + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + verify_schema = PythonOperator( + task_id="verify_schema", + python_callable=assert_schema_ready, + ) + + check_clickhouse >> choose_mode + choose_mode >> skip_ddl >> ddl_complete + choose_mode >> ddl_00_databases >> ddl_10_stg >> ddl_20_ods >> ddl_30_dds >> ddl_40_dm >> ddl_complete + ddl_complete >> verify_schema_sql >> verify_schema diff --git a/dags/etl_pipeline_dag.py b/dags/etl_pipeline_dag.py new file mode 100644 index 0000000..586f0dc --- /dev/null +++ b/dags/etl_pipeline_dag.py @@ -0,0 +1,282 @@ +""" +DAG ETL-процесса ODS -> DDS -> DM для учебного проекта. + +Принципы реализации: +- SQL выполняется явными task на ClickHouseOperator; +- SQL-файлы вызываются по фиксированным путям; +- Python используется только для управляющей логики (branch/wait/assert). +""" + +from __future__ import annotations + +import re +import time +from datetime import datetime, timedelta +from pathlib import Path + +from airflow import DAG +from airflow.exceptions import AirflowException +from airflow.models.param import Param +from airflow.operators.empty import EmptyOperator +from airflow.operators.python import BranchPythonOperator, PythonOperator +from airflow.utils.task_group import TaskGroup +from airflow.utils.trigger_rule import TriggerRule +from airflow_clickhouse_plugin.hooks.clickhouse import ClickHouseHook +from airflow_clickhouse_plugin.operators.clickhouse import ClickHouseOperator + + +# ----------------------------------------------------------------------------- +# Базовые настройки DAG +# ----------------------------------------------------------------------------- +default_args = { + "owner": "airflow", + "depends_on_past": False, + "email_on_failure": False, + "email_on_retry": False, + "retries": 1, + "retry_delay": timedelta(minutes=2), +} + + +# ----------------------------------------------------------------------------- +# SQL-файлы проекта +# ----------------------------------------------------------------------------- +SQL_ROOT = Path(__file__).resolve().parents[1] / "sql" + + +def load_sql_statements(relative_path: str) -> tuple[str, ...]: + """Читает SQL-файл и делит его на отдельные команды по ';'.""" + file_path = SQL_ROOT / relative_path + if not file_path.is_file(): + raise AirflowException(f"SQL-файл не найден: {file_path}") + + sql_text = file_path.read_text(encoding="utf-8") + statements: list[str] = [] + for segment in sql_text.split(";"): + # Убираем блочные и строковые комментарии, чтобы не отправлять "пустые" запросы. + no_block_comments = re.sub(r"/\*.*?\*/", "", segment, flags=re.S) + lines = [line for line in no_block_comments.splitlines() if not line.strip().startswith("--")] + cleaned = "\n".join(lines).strip() + if cleaned: + statements.append(cleaned) + + if not statements: + raise AirflowException(f"SQL-файл пустой: {file_path}") + return tuple(statements) + + +# ----------------------------------------------------------------------------- +# SQL для проверок и технических шагов +# ----------------------------------------------------------------------------- +SQL_CHECK_CLICKHOUSE = "SELECT 1 AS ok" + +SQL_CHECK_SCHEMA_READY = """ +SELECT + (SELECT count() FROM system.tables WHERE database = 'stg' AND name = 'browser_raw') AS stg_browser_raw, + (SELECT count() FROM system.tables WHERE database = 'ods' AND name = 'browser_event') AS ods_browser_event, + (SELECT count() FROM system.tables WHERE database = 'dds' AND name = 'event') AS dds_event, + (SELECT count() FROM system.tables WHERE database = 'dm' AND name = 'v_events_enriched') AS dm_v_events_enriched +""" + +SQL_CHECK_ODS_QUALITY = """ +SELECT + count() AS total_rows, + countIf(length(parse_errors) > 0) AS rows_with_errors, + round(if(count() = 0, 0, countIf(length(parse_errors) > 0) / count() * 100), 2) AS error_pct +FROM ods.browser_event +""" + +SQL_TRUNCATE_DDS_CLICK = "TRUNCATE TABLE dds.click" +SQL_TRUNCATE_DDS_EVENT = "TRUNCATE TABLE dds.event" + +SQL_CHECK_DDS_INTEGRITY = """ +SELECT + countIf(click_id IS NOT NULL AND click_id NOT IN (SELECT click_id FROM dds.click)) AS orphan_events +FROM dds.event +""" + +SQL_VALIDATE_DM_SUMMARY = "SELECT count() AS dq_rows FROM dm.dq_summary" + + +# ----------------------------------------------------------------------------- +# Управляющие функции +# ----------------------------------------------------------------------------- +def assert_schema_ready(**context) -> None: + """Падает, если DDL не применён полностью.""" + ti = context["ti"] + result = ti.xcom_pull(task_ids="precheck.check_schema_ready_sql") + + if not result or not result[0] or len(result[0]) != 4: + raise AirflowException(f"Некорректный результат check_schema_ready_sql: {result}") + + if any(value == 0 for value in result[0]): + raise AirflowException( + "Схема не готова: сначала запустите DAG ddl_init, затем повторите etl_pipeline." + ) + + +def wait_for_ods_data(**context) -> None: + """ + Ожидает появления строк в ods.browser_event до заданного таймаута. + Таймаут берётся из dag_run.conf.wait_ods_timeout_sec или из params. + """ + dag_run = context.get("dag_run") + conf = dag_run.conf if dag_run else {} + timeout_sec = int(conf.get("wait_ods_timeout_sec", context["params"]["wait_ods_timeout_sec"])) + poll_interval_sec = 10 + + hook = ClickHouseHook(clickhouse_conn_id="clickhouse_default", database="default") + started = time.monotonic() + + while True: + rows = hook.execute("SELECT count() FROM ods.browser_event") + count_rows = int(rows[0][0]) if rows else 0 + if count_rows > 0: + return + + elapsed = int(time.monotonic() - started) + if elapsed >= timeout_sec: + raise AirflowException( + f"Таймаут ожидания ODS истёк ({timeout_sec} сек). " + "Таблица ods.browser_event всё ещё пуста." + ) + + time.sleep(poll_interval_sec) + + +def choose_full_refresh(**context) -> str: + """Ветвление: делать TRUNCATE DDS или пропустить.""" + dag_run = context.get("dag_run") + conf = dag_run.conf if dag_run else {} + full_refresh = bool(conf.get("full_refresh", context["params"]["full_refresh"])) + return "transform.truncate_dds_click" if full_refresh else "transform.skip_truncate" + + +def assert_dm_summary_not_empty(**context) -> None: + """Проверяет, что dm.dq_summary заполнена после загрузки.""" + ti = context["ti"] + result = ti.xcom_pull(task_ids="transform.validate_dm_summary_sql") + + if not result or not result[0] or len(result[0]) != 1: + raise AirflowException(f"Некорректный результат validate_dm_summary_sql: {result}") + + dq_rows = int(result[0][0]) + if dq_rows <= 0: + raise AirflowException("dm.dq_summary пуста после load_dm_summary.") + + +with DAG( + dag_id="etl_pipeline", + description="ETL ODS -> DDS -> DM для demo-проекта", + default_args=default_args, + schedule=None, + start_date=datetime(2024, 1, 1), + catchup=False, + max_active_runs=1, + is_paused_upon_creation=True, + tags=["etl", "clickhouse", "demo"], + params={ + "full_refresh": Param(True, type="boolean"), + "wait_ods_timeout_sec": Param(600, type="integer", minimum=30), + }, +) as dag: + with TaskGroup(group_id="precheck") as precheck: + check_clickhouse = ClickHouseOperator( + task_id="check_clickhouse", + sql=SQL_CHECK_CLICKHOUSE, + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + check_schema_ready_sql = ClickHouseOperator( + task_id="check_schema_ready_sql", + sql=SQL_CHECK_SCHEMA_READY, + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + check_schema_ready = PythonOperator( + task_id="check_schema_ready", + python_callable=assert_schema_ready, + ) + + check_clickhouse >> check_schema_ready_sql >> check_schema_ready + + with TaskGroup(group_id="transform") as transform: + wait_for_ods_data_task = PythonOperator( + task_id="wait_for_ods_data", + python_callable=wait_for_ods_data, + ) + + check_ods_quality = ClickHouseOperator( + task_id="check_ods_quality", + sql=SQL_CHECK_ODS_QUALITY, + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + choose_refresh_mode = BranchPythonOperator( + task_id="choose_refresh_mode", + python_callable=choose_full_refresh, + ) + + truncate_dds_click = ClickHouseOperator( + task_id="truncate_dds_click", + sql=SQL_TRUNCATE_DDS_CLICK, + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + truncate_dds_event = ClickHouseOperator( + task_id="truncate_dds_event", + sql=SQL_TRUNCATE_DDS_EVENT, + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + skip_truncate = EmptyOperator(task_id="skip_truncate") + + truncate_complete = EmptyOperator( + task_id="truncate_complete", + trigger_rule=TriggerRule.NONE_FAILED_MIN_ONE_SUCCESS, + ) + + load_dds = ClickHouseOperator( + task_id="load_dds", + sql=load_sql_statements("dds/30_ods_to_dds.sql"), + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + check_dds_integrity = ClickHouseOperator( + task_id="check_dds_integrity", + sql=SQL_CHECK_DDS_INTEGRITY, + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + load_dm_summary = ClickHouseOperator( + task_id="load_dm_summary", + sql=load_sql_statements("dm/40_dds_to_dm.sql"), + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + validate_dm_summary_sql = ClickHouseOperator( + task_id="validate_dm_summary_sql", + sql=SQL_VALIDATE_DM_SUMMARY, + clickhouse_conn_id="clickhouse_default", + database="default", + ) + + validate_dm_summary = PythonOperator( + task_id="validate_dm_summary", + python_callable=assert_dm_summary_not_empty, + ) + + wait_for_ods_data_task >> check_ods_quality >> choose_refresh_mode + choose_refresh_mode >> truncate_dds_click >> truncate_dds_event >> truncate_complete + choose_refresh_mode >> skip_truncate >> truncate_complete + truncate_complete >> load_dds >> check_dds_integrity >> load_dm_summary >> validate_dm_summary_sql >> validate_dm_summary + + precheck >> transform diff --git a/data/2026-02-04_13-30-10.png b/data/2026-02-04_13-30-10.png old mode 100644 new mode 100755 diff --git a/data/DE-task.md b/data/DE-task.md old mode 100644 new mode 100755 diff --git a/data/browser_events.jsonl b/data/browser_events.jsonl old mode 100644 new mode 100755 diff --git a/data/device_events.jsonl b/data/device_events.jsonl old mode 100644 new mode 100755 diff --git a/data/geo_events.jsonl b/data/geo_events.jsonl old mode 100644 new mode 100755 diff --git a/data/location_events.jsonl b/data/location_events.jsonl old mode 100644 new mode 100755 diff --git a/docker-compose.yml b/docker-compose.yml index fc6b0a4..c3bad80 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -1,3 +1,12 @@ +x-airflow-env: &airflow-default-env + AIRFLOW__CORE__LOAD_EXAMPLES: "False" + AIRFLOW__DATABASE__SQL_ALCHEMY_CONN: postgresql+psycopg2://airflow:airflow@postgres-metadata:5432/airflow + AIRFLOW__CORE__EXECUTOR: LocalExecutor + AIRFLOW__WEBSERVER__SECRET_KEY: ${AIRFLOW_SECRET_KEY:-replace-me-with-random-string} + # ClickHouse connection для ETL + AIRFLOW_CONN_CLICKHOUSE_DEFAULT: clickhouse://default:123456@clickhouse:9000/default + + services: clickhouse: image: clickhouse/clickhouse-server:25.1 @@ -11,6 +20,7 @@ services: networks: - cs_dwh volumes: + - clickhouse-data:/var/lib/clickhouse - ./configs/default_user.xml:/etc/clickhouse-server/users.d/default_user.xml # - ./configs/z_config.xml:/etc/clickhouse-server/config.d/z_config.xml # - ./configs/macros_ch1.xml:/etc/clickhouse-server/config.d/macros.xml @@ -27,7 +37,7 @@ services: # container_name: kafka ports: - 9092:9092 - # Если хочешь сохранять топики/сообщения между `docker compose down/up` — раскомментируй: + # Данные топиков/сообщений сохраняются в named volume `kafka-data` (между `docker compose down/up`). volumes: - kafka-data:/tmp/kraft-combined-logs environment: @@ -64,6 +74,107 @@ services: networks: - cs_dwh + # PostgreSQL for Airflow metadata + postgres-metadata: + image: postgres:16 + environment: + POSTGRES_USER: airflow + POSTGRES_PASSWORD: airflow + POSTGRES_DB: airflow + ports: + - "5434:5432" + volumes: + - pgmeta:/var/lib/postgresql/data + networks: + - cs_dwh + healthcheck: + test: ["CMD-SHELL", "pg_isready -U airflow -d airflow"] + interval: 5s + timeout: 5s + retries: 20 + + # Airflow webserver + airflow-webserver: + build: + context: . + dockerfile: Dockerfile.airflow + image: airflow-optimized:2.10.5 + environment: + <<: *airflow-default-env + command: > + bash -c " + airflow webserver + " + ports: + - "8080:8080" + volumes: + - ./dags:/opt/airflow/dags + - ./sql:/opt/airflow/sql:ro + - ./data:/opt/airflow/data + networks: + - cs_dwh + depends_on: + airflow-init: + condition: service_completed_successfully + postgres-metadata: + condition: service_healthy + clickhouse: + condition: service_started + + # Airflow scheduler + airflow-scheduler: + build: + context: . + dockerfile: Dockerfile.airflow + environment: + <<: *airflow-default-env + command: > + bash -c " + airflow scheduler + " + volumes: + - ./dags:/opt/airflow/dags + - ./sql:/opt/airflow/sql:ro + - ./data:/opt/airflow/data + networks: + - cs_dwh + depends_on: + airflow-init: + condition: service_completed_successfully + postgres-metadata: + condition: service_healthy + clickhouse: + condition: service_started + + # Airflow init to create admin user + airflow-init: + build: + context: . + dockerfile: Dockerfile.airflow + user: "0:0" + environment: + <<: *airflow-default-env + volumes: + - ./dags:/opt/airflow/dags + - ./sql:/opt/airflow/sql:ro + - ./data:/opt/airflow/data + networks: + - cs_dwh + command: > + bash -ceuo pipefail " + mkdir -p /opt/airflow/data && + chmod -R a+rX /opt/airflow/data || true && + chown -R airflow:0 /opt/airflow/data || true && + umask 000 && + su -s /bin/bash airflow -c 'airflow db migrate' && + su -s /bin/bash airflow -c 'airflow users create --username admin --password admin --firstname Admin --lastname User --role Admin --email admin@example.org' && + su -s /bin/bash airflow -c 'airflow connections delete clickhouse_default || true' && + su -s /bin/bash airflow -c \"airflow connections add 'clickhouse_default' --conn-uri \\\"$${AIRFLOW_CONN_CLICKHOUSE_DEFAULT}\\\" --conn-description 'ClickHouse DWH'\" + " + depends_on: + postgres-metadata: + condition: service_healthy + prometheus: image: prom/prometheus:v2.53.4 volumes: @@ -125,8 +236,9 @@ networks: driver: bridge volumes: + clickhouse-data: grafana_lib: kafka-data: + pgmeta: superset_data: superset_config: - diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 5400f5c..6e5c37d 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -369,11 +369,11 @@ sequenceDiagram K-->>User: ✅ Инфраструктура готова User->>Make: make ddl - Make->>CH: ddl/00_databases.sql - Make->>CH: ddl/10_stg.sql (Kafka Engine) - Make->>CH: ddl/20_ods.sql (MV) - Make->>CH: ddl/30_dds.sql - Make->>CH: ddl/40_dm.sql + Make->>CH: sql/ddl/00_databases.sql + Make->>CH: sql/ddl/stg/10_stg.sql (Kafka Engine) + Make->>CH: sql/ddl/ods/20_ods.sql (MV) + Make->>CH: sql/ddl/dds/30_dds.sql + Make->>CH: sql/ddl/dm/40_dm.sql CH-->>User: ✅ Структура БД создана User->>Make: make data @@ -389,10 +389,10 @@ sequenceDiagram CH-->>User: ✅ Данные в STG/ODS User->>Make: make transform - Make->>CH: jobs/30_dds_refresh.sql + Make->>CH: sql/dds/30_ods_to_dds.sql CH->>ODS: argMax() — снапшот CH->>DDS: JOIN + INSERT - Make->>CH: jobs/40_dm_refresh.sql + Make->>CH: sql/dm/40_dds_to_dm.sql CH->>DM: DQ summary CH-->>User: ✅ DDS/DM обновлены ``` @@ -586,16 +586,26 @@ INSERT INTO dm.daily_traffic SELECT * FROM dm.v_daily_traffic; ### Airflow-оркестрация +Инфраструктура Airflow развёрнута и готова к использованию: + ```python -# dag.py -with DAG('clickhouse_etl'): - ddl = BashOperator(task_id='ddl', bash_command='make ddl') - load = BashOperator(task_id='load', bash_command='make data') - transform = BashOperator(task_id='transform', bash_command='make transform') - - ddl >> load >> transform +# dags/ddl_init_dag.py и dags/etl_pipeline_dag.py +# +# Учебный формат: +# - DDL и трансформации выполняются явными SQL-task через ClickHouseOperator; +# - SQL-файлы вызываются по фиксированным путям; +# - загрузка данных в Kafka (Этап 1) выполняется через `make data`. +# +# Основной demo-сценарий: +# ddl_init -> make data -> etl_pipeline ``` +**Подключение к ClickHouse:** +- Connection: `clickhouse_default` +- URL: `clickhouse://default:123456@clickhouse:9000/default` (native TCP для Airflow plugin) +- Provider/интеграция: `airflow-clickhouse-plugin` (в `airflow/requirements.txt`), задачи выполняются через `ClickHouseOperator`. +- Примечание: Superset подключается к ClickHouse по HTTP (обычно `clickhouse+connect://...:8123/...`). + --- ## Полезные запросы diff --git a/plans/airflow_dags_plan.md b/plans/airflow_dags_plan.md new file mode 100644 index 0000000..c12a8eb --- /dev/null +++ b/plans/airflow_dags_plan.md @@ -0,0 +1,234 @@ +# План развития Airflow DAG'ов + +## Цель +Перевести оркестрацию ETL на Airflow так, чтобы пайплайн оставался устойчивым к "грязным" данным и соответствовал целям проекта: Kafka → ClickHouse (STG → ODS → DDS → DM) → витрины для BI. + +## Что важно учесть в текущем репозитории +- В `AGENTS.md` как quick check ожидается DAG `etl_pipeline`. +- В Airflow-контейнере сейчас нет Kafka CLI, поэтому `kafka-topics.sh` и `kafka-console-producer.sh` из `BashOperator` не используем. +- DDL должен выполняться строго последовательно: `00 -> 10 -> 20 -> 30 -> 40`. +- Файл `sql/dds/30_ods_to_dds.sql` уже включает обе загрузки (`dds.click` и `dds.event`), поэтому в MVP это одна task. + +## Архитектура оркестрации +### DAG 1 (обязательный): `ddl_init` +- `schedule`: `None` (только ручной запуск). +- `catchup`: `False`. +- `max_active_runs`: `1`. +- `is_paused_upon_creation`: `True`. +- `tags`: `["ddl", "bootstrap", "clickhouse"]`. + +### DAG 2 (обязательный): `kafka_load` +- `schedule`: `None` (ручной/экспериментальный запуск). +- `catchup`: `False`. +- `max_active_runs`: `1`. +- `is_paused_upon_creation`: `True`. +- `tags`: `["kafka", "ingest", "experiments"]`. +- Реализация: этап 2 (после MVP). + +### DAG 3 (обязательный): `etl_pipeline` +- `schedule`: `None` (ручной запуск для демо). +- `catchup`: `False`. +- `max_active_runs`: `1`. +- `tags`: `["etl", "clickhouse", "demo"]`. + +### DAG 4 (опциональный): `dq_monitor` +- `schedule`: `0 * * * *`. +- `catchup`: `False`. +- `tags`: `["dq", "monitoring"]`. + +## Дизайн DAG `ddl_init` +### Params +- `verify_only`: bool, default `false` (прогон только проверок без применения DDL). + +### Tasks +| Task ID | Что делает | Источник SQL/реализация | +|---------|------------|--------------------------| +| `check_clickhouse` | Проверка доступности CH (`SELECT 1`) | `PythonOperator` + `clickhouse-connect` | +| `ddl_00_databases` | Создание БД | `sql/ddl/00_databases.sql` | +| `ddl_10_stg` | STG + Kafka Engine + MV | `sql/ddl/stg/10_stg.sql` | +| `ddl_20_ods` | ODS + MV STG→ODS + *_errors | `sql/ddl/ods/20_ods.sql` | +| `ddl_30_dds` | Таблицы DDS | `sql/ddl/dds/30_dds.sql` | +| `ddl_40_dm` | VIEW витрины DM | `sql/ddl/dm/40_dm.sql` | +| `verify_schema` | Проверка ключевых таблиц/VIEW | SQL-check | + +Зависимости: +```text +check_clickhouse >> ddl_00_databases >> ddl_10_stg >> ddl_20_ods >> ddl_30_dds >> ddl_40_dm >> verify_schema +``` + +Примечание: +- `ddl_init` запускается вручную: при первом bootstrap, после `docker compose down -v`, после изменений схемы. + +## Дизайн DAG `kafka_load` (отдельный независимый контур) +### Params (через Trigger DAG with config) +- `limit`: int, default `50`. +- `full_load`: bool, default `false`. +- `reset_topics`: bool, default `true`. +- `load_browser`: bool, default `true`. +- `load_location`: bool, default `true`. +- `load_device`: bool, default `true`. +- `load_geo`: bool, default `true`. + +### TaskGroup `precheck` +| Task ID | Что делает | Реализация | +|---------|------------|------------| +| `check_kafka` | Проверка доступности Kafka broker | `PythonOperator` + `kafka-python` | +| `check_input_files` | Проверка наличия `data/*_events.jsonl` | `PythonOperator` | +| `validate_load_params` | Валидация параметров загрузки (`limit`, `full_load`, флаги потоков) | `PythonOperator` | + +### TaskGroup `ingest` +| Task ID | Что делает | Реализация | +|---------|------------|------------| +| `prepare_topics` | reset/create топиков по `reset_topics` | `PythonOperator` + AdminClient | +| `load_browser_events` | Публикация `browser_events.jsonl` | `PythonOperator` + KafkaProducer | +| `load_location_events` | Публикация `location_events.jsonl` | `PythonOperator` + KafkaProducer | +| `load_device_events` | Публикация `device_events.jsonl` | `PythonOperator` + KafkaProducer | +| `load_geo_events` | Публикация `geo_events.jsonl` | `PythonOperator` + KafkaProducer | +| `verify_publish_counts` | Проверка, что отправлено > 0 сообщений в выбранные потоки | `PythonOperator` (по XCom) | + +Зависимости: +```text +precheck >> prepare_topics >> [load_browser_events, load_location_events, load_device_events, load_geo_events] >> verify_publish_counts +``` + +Примечания: +- DAG намеренно независим от `etl_pipeline`: можно запускать ingest отдельно для экспериментов. +- Авто-триггер `etl_pipeline` не включаем по умолчанию; при необходимости добавляется отдельным параметром позже. + +## Дизайн DAG `etl_pipeline` +### Params (через Trigger DAG with config) +- `full_refresh`: bool, default `true`. +- `wait_ods_timeout_sec`: int, default `600`. + +### TaskGroup `precheck` +| Task ID | Что делает | Реализация | +|---------|------------|------------| +| `check_clickhouse` | Проверка доступности CH (`SELECT 1`) | `PythonOperator` + `clickhouse-connect` | +| `check_schema_ready` | Проверка, что DDL уже применён (`stg.browser_raw`, `ods.browser_event`, `dds.event`, `dm.v_events_enriched`) | SQL-check, fail fast | + +### TaskGroup `transform` +| Task ID | Что делает | Источник SQL | +|---------|------------|--------------| +| `wait_for_ods_data` | Ожидание строк в `ods.browser_event` | SQL-check | +| `check_ods_quality` | Базовые DQ-метрики ODS (ошибки/total) | SQL-check | +| `truncate_dds_click` | Очистка `dds.click` при `full_refresh=true` | inline SQL | +| `truncate_dds_event` | Очистка `dds.event` при `full_refresh=true` | inline SQL | +| `load_dds` | ODS → DDS | `sql/dds/30_ods_to_dds.sql` | +| `check_dds_integrity` | Проверка orphan событий | inline SQL | +| `load_dm_summary` | DDS → DM DQ summary | `sql/dm/40_dds_to_dm.sql` | +| `validate_dm_summary` | Проверка, что `dm.dq_summary` не пуста | SQL-check | + +Зависимости: +```text +wait_for_ods_data >> check_ods_quality >> [truncate_dds_click, truncate_dds_event] >> load_dds >> check_dds_integrity >> load_dm_summary >> validate_dm_summary +``` + +### Итоговая цепочка `etl_pipeline` +```text +precheck >> transform +``` + +## Взаимодействие DAG'ов +Базовый сценарий: +```text +ddl_init -> kafka_load -> etl_pipeline +``` + +Экспериментальные сценарии: +- `kafka_load` отдельно: проверить разные наборы/параметры загрузки без запуска transform. +- `etl_pipeline` отдельно: повторно пересчитать DDS/DM по уже загруженным данным. + +## Техническая реализация (приземленно) +### ClickHouse в Airflow +- Использовать `clickhouse-connect` напрямую в Python helper. +- Брать параметры подключения из `conn_id = clickhouse_default` через `BaseHook.get_connection`. + +### Kafka в Airflow +- Добавить `kafka-python` в `airflow/requirements.txt`. +- Использовать Python-код для: + - reset/create топиков; + - публикации строк из `.jsonl` (`1 строка = 1 message value`). + +### Общие helper-функции +- `dags/utils/clickhouse_helpers.py`: + - `execute_sql(sql: str) -> None` + - `execute_sql_file(path: str) -> None` + - `fetch_one(sql: str) -> tuple` +- `dags/utils/kafka_helpers.py`: + - `prepare_topics(reset: bool) -> None` + - `load_jsonl(file_path: str, topic: str, limit: int, full_load: bool) -> int` + - `check_kafka_ready() -> None` + +## Структура файлов +```text +dags/ +├── __init__.py +├── ddl_init_dag.py # отдельный DAG для DDL (обязателен) +├── kafka_load_dag.py # отдельный DAG для ingest в Kafka (обязателен) +├── etl_pipeline_dag.py # основной DAG ODS -> DDS -> DM (обязателен) +├── dq_monitor_dag.py # опциональный DAG мониторинга +└── utils/ + ├── __init__.py + ├── clickhouse_helpers.py + └── kafka_helpers.py +``` + +## Этапы внедрения +1. Этап 1 (MVP, обязательно): + - Реализовать `ddl_init` и `etl_pipeline`. + - Для загрузки данных использовать существующий сценарий `make data`. + - Проверить путь `ddl_init -> make data -> etl_pipeline`. +2. Этап 2: + - Реализовать отдельный DAG `kafka_load` на `kafka-python`. + - Перенести загрузку из `make data` в `kafka_load` (функциональный паритет). + - Добавить в `kafka_load` расширенные параметры экспериментов (выбор потоков, сценарии reset/no-reset). + - Добавить опциональный параметр автотриггера `etl_pipeline` (по умолчанию `false`). +3. Этап 3: + - Добавить `dq_monitor` и alert callback (email/Slack/webhook). + +## Критерии готовности +- Этап 1: + - В Airflow UI видны DAG `ddl_init` и `etl_pipeline`. + - `etl_pipeline` падает с понятной ошибкой, если схема не применена или ODS пуста. + - После прогона `make data -> etl_pipeline`: + - в `ods.browser_event` есть строки; + - в `dds.click` и `dds.event` есть строки; + - `dm.dq_summary` заполнена. +- Этап 2: + - В Airflow UI дополнительно виден DAG `kafka_load`. + - `kafka_load` с `limit=50` завершает отправку сообщений без падений. + - После прогона `kafka_load -> etl_pipeline` результаты совпадают с `make data -> etl_pipeline`. +- Для всех этапов: + - При повторном запуске `etl_pipeline` с `full_refresh=true` нет неконтролируемых дублей в DDS. + +## Минимальные smoke-checks +```bash +# 1) Запуск инфраструктуры +make up + +# 2) Открыть Airflow UI +# http://localhost:8080 (admin/admin) + +# 3) Один раз запустить ddl_init +# Trigger DAG ddl_init (без config или {"verify_only": false}) + +# 4) Этап 1: загрузить данные текущим способом +make data + +# 5) Запустить etl_pipeline +# {"full_refresh": true} + +# 6) Проверка результатов +docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query "SELECT count() FROM ods.browser_event" +docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query "SELECT count() FROM dds.click" +docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query "SELECT count() FROM dds.event" +docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query "SELECT count() FROM dm.dq_summary" + +# 7) Этап 2: после реализации kafka_load +# Trigger DAG kafka_load {"limit": 50, "full_load": false, "reset_topics": true} +# Trigger DAG etl_pipeline {"full_refresh": true} +``` + +## Следующий шаг после MVP +- Разделить `sql/dds/30_ods_to_dds.sql` на два файла и распараллелить `load_dds_click` и `load_dds_event`. +- Перейти с `full_refresh` на watermark-инкремент. diff --git a/plans/clickhouse_ddl.md b/plans/clickhouse_ddl.md index 685682c..fdf80ff 100644 --- a/plans/clickhouse_ddl.md +++ b/plans/clickhouse_ddl.md @@ -114,32 +114,32 @@ flowchart LR ## План актуализации DDL (target state репозитория) -Цель: перестать исполнять DDL из markdown и хранить **исполняемые** DDL в отдельных `ddl/*.sql` (по слоям), чтобы: +Цель: перестать исполнять DDL из markdown и хранить **исполняемые** DDL в отдельных `sql/*/*.sql` (по слоям), чтобы: - применять их “тонким раннером” через `clickhouse-client` (через `make ddl`); - в будущем легко перенести выполнение в Airflow (1 файл = 1 task, линейные зависимости). -Важно: Kafka-объекты STG включаем **по умолчанию** (как часть `ddl/10_stg.sql`). +Важно: Kafka-объекты STG включаем **по умолчанию** (как часть `sql/ddl/stg/10_stg.sql`). ### Артефакты DDL (планируемые файлы) -- `ddl/00_databases.sql` — базы `stg/ods/dds/dm`. -- `ddl/10_stg.sql` — STG raw (`stg.*_raw`) + Kafka source tables (`ENGINE = Kafka`) + MV `Kafka → STG`. -- `ddl/20_ods.sql` — ODS таблицы типизации + DQ (`parse_errors`) + MV `STG → ODS` + таблицы `ods_*_errors` для строк с битыми ключами. -- `ddl/30_dds.sql` — DDS таблицы (`dds.event`, `dds.click`) **без MV** (только `CREATE TABLE`). -- `ddl/40_dm.sql` — витрины `VIEW` для Superset (`dm.v_*`). +- `sql/ddl/00_databases.sql` — базы `stg/ods/dds/dm`. +- `sql/ddl/stg/10_stg.sql` — STG raw (`stg.*_raw`) + Kafka source tables (`ENGINE = Kafka`) + MV `Kafka → STG`. +- `sql/ddl/ods/20_ods.sql` — ODS таблицы типизации + DQ (`parse_errors`) + MV `STG → ODS` + таблицы `ods_*_errors` для строк с битыми ключами. +- `sql/ddl/dds/30_dds.sql` — DDS таблицы (`dds.event`, `dds.click`) **без MV** (только `CREATE TABLE`). +- `sql/ddl/dm/40_dm.sql` — витрины `VIEW` для Superset (`dm.v_*`). -BI-ограничения (ресурсы/пользователь) **не выносим в `ddl/*.sql`**: оставляем это только как текст/пример в этом плане, чтобы не смешивать инфраструктуру доступа с DDL витрин. +BI-ограничения (ресурсы/пользователь) **не выносим в `sql/*/*.sql`**: оставляем это только как текст/пример в этом плане, чтобы не смешивать инфраструктуру доступа с DDL витрин. ### Артефакты batch-трансформаций (планируемые файлы) -- `jobs/30_dds_refresh.sql` — регулярная батч‑сборка DDS из ODS: +- `sql/dds/30_ods_to_dds.sql` — регулярная батч‑сборка DDS из ODS: - получить “последнюю версию” строк по ключам (`event_id`/`click_id`) через `argMax(..., src_ingest_ts)` (или эквивалент); - выполнить join snapshot’ов и загрузить в `dds.event`/`dds.click` (для демо возможно “full rebuild”; позже — инкрементально). ### Исполнение DDL (make сейчас / Airflow потом) -Требования к файлам `ddl/*.sql`: +Требования к файлам `sql/*/*.sql`: - идемпотентность (`IF NOT EXISTS`), чтобы повторные прогоны были безопасны; - строгий порядок исполнения: `00 → 10 → 20 → 30 → 40` (из‑за зависимостей MV); @@ -163,11 +163,11 @@ BI-ограничения (ресурсы/пользователь) **не вы Текущее “как запускаем” (целевое, для реализации следующим шагом): - `make ddl` вызывает `scripts/apply_clickhouse_ddl.sh`; -- скрипт прогоняет `ddl/*.sql` по порядку через `clickhouse-client --multiquery` внутри контейнера ClickHouse. +- скрипт прогоняет `sql/*/*.sql` по порядку через `clickhouse-client --multiquery` внутри контейнера ClickHouse. Batch‑трансформации (целевое, для реализации следующим шагом): -- `make transform` (или аналогичная команда) запускает `jobs/30_dds_refresh.sql` через `clickhouse-client`; +- `make transform` (или аналогичная команда) запускает `sql/dds/30_ods_to_dds.sql` через `clickhouse-client`; - в будущем Airflow будет делать то же самое по расписанию (один job‑SQL = один task). ### Параметры окружения (docker compose) @@ -178,7 +178,7 @@ Batch‑трансформации (целевое, для реализации --- -## Приложение A: текущий inline DDL (legacy; будет вынесен в `ddl/*.sql`) +## Приложение A: текущий inline DDL (legacy; будет вынесен в `sql/*/*.sql`) ### 0) Базы данных diff --git a/scripts/apply_clickhouse_ddl.sh b/scripts/apply_clickhouse_ddl.sh index 8188f90..1f21349 100755 --- a/scripts/apply_clickhouse_ddl.sh +++ b/scripts/apply_clickhouse_ddl.sh @@ -3,8 +3,8 @@ # Скрипт применения DDL в ClickHouse # # Назначение: -# Последовательно применяет SQL-файлы из ddl/*.sql в базу ClickHouse. -# Файлы применяются в алфавитном порядке (00 → 10 → 20 → 30 → 40). +# Последовательно применяет SQL-файлы из sql/ddl/* в базу ClickHouse. +# Порядок фиксированный: 00 → 10 → 20 → 30 → 40. # # Как запускать: # make ddl @@ -22,7 +22,7 @@ set -euo pipefail # Директория со скриптом SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" -DDL_DIR="${SCRIPT_DIR}/../ddl" +SQL_ROOT_DIR="${SCRIPT_DIR}/../sql" # Параметры подключения (можно переопределить через переменные окружения) COMPOSE_BIN="${COMPOSE_BIN:-docker compose}" @@ -31,7 +31,7 @@ CLICKHOUSE_DB="${CLICKHOUSE_DB:-default}" CLICKHOUSE_USER="${CLICKHOUSE_USER:-default}" CLICKHOUSE_PASSWORD="${CLICKHOUSE_PASSWORD:-123456}" -echo "Применение DDL из ${DDL_DIR}..." +echo "Применение DDL из ${SQL_ROOT_DIR}..." # ----------------------------------------------------------------------------- # Проверка: ClickHouse запущен? @@ -45,18 +45,28 @@ fi # ----------------------------------------------------------------------------- # Применение SQL-файлов по порядку # ----------------------------------------------------------------------------- -# shellcheck disable=SC2044 -for sql_file in "${DDL_DIR}"/*.sql; do - if [[ -f "$sql_file" ]]; then - echo "Применение: $(basename "$sql_file")" - ${COMPOSE_BIN} exec -T "${CLICKHOUSE_SERVICE}" clickhouse-client \ - --user="${CLICKHOUSE_USER}" \ - --password="${CLICKHOUSE_PASSWORD}" \ - --database="${CLICKHOUSE_DB}" \ - --multiquery \ - < "$sql_file" - echo " ✓ OK" +DDL_FILES=( + "${SQL_ROOT_DIR}/ddl/00_databases.sql" + "${SQL_ROOT_DIR}/ddl/stg/10_stg.sql" + "${SQL_ROOT_DIR}/ddl/ods/20_ods.sql" + "${SQL_ROOT_DIR}/ddl/dds/30_dds.sql" + "${SQL_ROOT_DIR}/ddl/dm/40_dm.sql" +) + +for sql_file in "${DDL_FILES[@]}"; do + if [[ ! -f "$sql_file" ]]; then + echo "Ошибка: не найден SQL-файл: $sql_file" >&2 + exit 1 fi + + echo "Применение: $(basename "$sql_file")" + ${COMPOSE_BIN} exec -T "${CLICKHOUSE_SERVICE}" clickhouse-client \ + --user="${CLICKHOUSE_USER}" \ + --password="${CLICKHOUSE_PASSWORD}" \ + --database="${CLICKHOUSE_DB}" \ + --multiquery \ + < "$sql_file" + echo " ✓ OK" done echo "" diff --git a/scripts/run_batch.sh b/scripts/run_batch.sh index 6b84730..9c36b60 100755 --- a/scripts/run_batch.sh +++ b/scripts/run_batch.sh @@ -3,7 +3,7 @@ # Скрипт batch-трансформации данных: ODS → DDS → DM # # Назначение: -# Запускает SQL-скрипты из jobs/ для преобразования данных между слоями: +# Запускает SQL-скрипты из sql/dds и sql/dm для преобразования данных между слоями: # 1. ODS → DDS : Сборка сущностей из типизированных данных # 2. DDS → DM : Обновление сводки по качеству данных (dq_summary) # @@ -24,7 +24,9 @@ set -euo pipefail # Директория со скриптом SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" -JOBS_DIR="${SCRIPT_DIR}/../jobs" +SQL_ROOT_DIR="${SCRIPT_DIR}/../sql" +DDS_TRANSFORM_SQL="${SQL_ROOT_DIR}/dds/30_ods_to_dds.sql" +DM_TRANSFORM_SQL="${SQL_ROOT_DIR}/dm/40_dds_to_dm.sql" # Параметры подключения COMPOSE_BIN="${COMPOSE_BIN:-docker compose}" @@ -42,6 +44,19 @@ if ! ${COMPOSE_BIN} ps | grep -q "${CLICKHOUSE_SERVICE}"; then exit 1 fi +# ----------------------------------------------------------------------------- +# Проверка: SQL-файлы batch существуют? +# ----------------------------------------------------------------------------- +if [[ ! -f "${DDS_TRANSFORM_SQL}" ]]; then + echo "Ошибка: Не найден SQL-файл: ${DDS_TRANSFORM_SQL}" + exit 1 +fi + +if [[ ! -f "${DM_TRANSFORM_SQL}" ]]; then + echo "Ошибка: Не найден SQL-файл: ${DM_TRANSFORM_SQL}" + exit 1 +fi + # ----------------------------------------------------------------------------- # Проверка: в ODS есть данные? # ----------------------------------------------------------------------------- @@ -83,7 +98,7 @@ ${COMPOSE_BIN} exec -T "${CLICKHOUSE_SERVICE}" clickhouse-client \ --user="${CLICKHOUSE_USER}" \ --password="${CLICKHOUSE_PASSWORD}" \ --database="${CLICKHOUSE_DB}" \ - --multiquery < "${JOBS_DIR}/30_dds_refresh.sql" + --multiquery < "${DDS_TRANSFORM_SQL}" echo " ✓ DDS обновлён" @@ -105,7 +120,7 @@ ${COMPOSE_BIN} exec -T "${CLICKHOUSE_SERVICE}" clickhouse-client \ --user="${CLICKHOUSE_USER}" \ --password="${CLICKHOUSE_PASSWORD}" \ --database="${CLICKHOUSE_DB}" \ - --multiquery < "${JOBS_DIR}/40_dm_refresh.sql" + --multiquery < "${DM_TRANSFORM_SQL}" echo " ✓ DM обновлён" diff --git a/ddl/00_databases.sql b/sql/ddl/00_databases.sql similarity index 100% rename from ddl/00_databases.sql rename to sql/ddl/00_databases.sql diff --git a/ddl/30_dds.sql b/sql/ddl/dds/30_dds.sql similarity index 99% rename from ddl/30_dds.sql rename to sql/ddl/dds/30_dds.sql index fc60469..a9e732f 100644 --- a/ddl/30_dds.sql +++ b/sql/ddl/dds/30_dds.sql @@ -8,7 +8,7 @@ -- -- Загрузка: -- Batch SQL (не MV!) — для согласованности при late arrivals --- См. jobs/30_dds_refresh.sql +-- См. sql/dds/30_ods_to_dds.sql -- -- Почему не MV: -- - MV с JOIN даёт eventual consistency (данные приходят в разное время) diff --git a/ddl/40_dm.sql b/sql/ddl/dm/40_dm.sql similarity index 99% rename from ddl/40_dm.sql rename to sql/ddl/dm/40_dm.sql index ad985d3..0bd89a5 100644 --- a/ddl/40_dm.sql +++ b/sql/ddl/dm/40_dm.sql @@ -13,7 +13,7 @@ -- -- Для продакшена: -- - Если тяжёлые агрегации тормозят — материализовать в таблицы --- - См. пример закомментированный в jobs/40_dm_refresh.sql +-- - См. пример закомментированный в sql/dm/40_dds_to_dm.sql -- ============================================================================ -- ---------------------------------------------------------------------------- diff --git a/ddl/20_ods.sql b/sql/ddl/ods/20_ods.sql similarity index 100% rename from ddl/20_ods.sql rename to sql/ddl/ods/20_ods.sql diff --git a/ddl/10_stg.sql b/sql/ddl/stg/10_stg.sql similarity index 100% rename from ddl/10_stg.sql rename to sql/ddl/stg/10_stg.sql diff --git a/jobs/30_dds_refresh.sql b/sql/dds/30_ods_to_dds.sql similarity index 100% rename from jobs/30_dds_refresh.sql rename to sql/dds/30_ods_to_dds.sql diff --git a/jobs/40_dm_refresh.sql b/sql/dm/40_dds_to_dm.sql similarity index 100% rename from jobs/40_dm_refresh.sql rename to sql/dm/40_dds_to_dm.sql