diff --git a/AGENTS.md b/AGENTS.md index d03d831..0b50e4c 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -21,15 +21,14 @@ ### Исполняемые файлы (текущая структура) - `dags/` — Airflow DAGs для оркестрации ETL -- `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 +- `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 @@ -57,7 +56,7 @@ Базовые команды: - `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` @@ -88,7 +87,7 @@ - **Комментарии в коде — на русском языке**: - 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`) ## Быстрые проверки diff --git a/README.md b/README.md index 6d73d44..f57eb32 100644 --- a/README.md +++ b/README.md @@ -109,8 +109,15 @@ flowchart TB ``` . ├── dags/ # Airflow DAGs для оркестрации -├── ddl/ # SQL для создания объектов (00_databases → 40_dm) -├── jobs/ # Batch-трансформации (ODS→DDS, DDS→DM) +├── 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 diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 6369e23..381bd0b 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 обновлены ``` diff --git a/plans/airflow_dags_plan.md b/plans/airflow_dags_plan.md index 4bc7716..c12a8eb 100644 --- a/plans/airflow_dags_plan.md +++ b/plans/airflow_dags_plan.md @@ -7,7 +7,7 @@ - В `AGENTS.md` как quick check ожидается DAG `etl_pipeline`. - В Airflow-контейнере сейчас нет Kafka CLI, поэтому `kafka-topics.sh` и `kafka-console-producer.sh` из `BashOperator` не используем. - DDL должен выполняться строго последовательно: `00 -> 10 -> 20 -> 30 -> 40`. -- Файл `jobs/30_dds_refresh.sql` уже включает обе загрузки (`dds.click` и `dds.event`), поэтому в MVP это одна task. +- Файл `sql/dds/30_ods_to_dds.sql` уже включает обе загрузки (`dds.click` и `dds.event`), поэтому в MVP это одна task. ## Архитектура оркестрации ### DAG 1 (обязательный): `ddl_init` @@ -44,11 +44,11 @@ | Task ID | Что делает | Источник SQL/реализация | |---------|------------|--------------------------| | `check_clickhouse` | Проверка доступности CH (`SELECT 1`) | `PythonOperator` + `clickhouse-connect` | -| `ddl_00_databases` | Создание БД | `ddl/00_databases.sql` | -| `ddl_10_stg` | STG + Kafka Engine + MV | `ddl/10_stg.sql` | -| `ddl_20_ods` | ODS + MV STG→ODS + *_errors | `ddl/20_ods.sql` | -| `ddl_30_dds` | Таблицы DDS | `ddl/30_dds.sql` | -| `ddl_40_dm` | VIEW витрины DM | `ddl/40_dm.sql` | +| `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 | Зависимости: @@ -113,14 +113,14 @@ precheck >> prepare_topics >> [load_browser_events, load_location_events, load_d | `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 | -| `refresh_dds` | ODS → DDS | `jobs/30_dds_refresh.sql` | +| `load_dds` | ODS → DDS | `sql/dds/30_ods_to_dds.sql` | | `check_dds_integrity` | Проверка orphan событий | inline SQL | -| `refresh_dm_summary` | DDS → DM DQ summary | `jobs/40_dm_refresh.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] >> refresh_dds >> check_dds_integrity >> refresh_dm_summary >> validate_dm_summary +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` @@ -230,5 +230,5 @@ docker compose exec -T clickhouse clickhouse-client --user=default --password=12 ``` ## Следующий шаг после MVP -- Разделить `jobs/30_dds_refresh.sql` на два файла и распараллелить `refresh_dds_click` и `refresh_dds_event`. +- Разделить `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