From d33cdb3fb0315f36da0a5e848a917c9ab1b027c4 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Fri, 6 Feb 2026 23:20:58 +0300 Subject: [PATCH 01/12] feat(infra): add airflow orchestration services Add Apache Airflow infrastructure with webserver, scheduler, and metadata database to enable DAG-based pipeline orchestration. Includes optimized requirements file and Docker configuration for Airflow 2.9.3. --- Dockerfile.airflow | 23 +++++++++ airflow/requirements.txt | 14 +++++ docker-compose.yml | 108 +++++++++++++++++++++++++++++++++++++++ 3 files changed, 145 insertions(+) create mode 100644 Dockerfile.airflow create mode 100644 airflow/requirements.txt diff --git a/Dockerfile.airflow b/Dockerfile.airflow new file mode 100644 index 0000000..0df3b0d --- /dev/null +++ b/Dockerfile.airflow @@ -0,0 +1,23 @@ +# Use the official Airflow image as base +FROM apache/airflow:2.9.3 + +# Set environment variables +ENV AIRFLOW_HOME=/opt/airflow + +# Copy the project requirements file +COPY --chown=airflow:0 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/airflow/requirements.txt b/airflow/requirements.txt new file mode 100644 index 0000000..ba5c858 --- /dev/null +++ b/airflow/requirements.txt @@ -0,0 +1,14 @@ +# Optimized requirements for Airflow Docker setup +# Only includes packages actually used in DAGs + +# Core database connector +psycopg2-binary==2.9.9 + +# Data processing (used in data_processing_dag.py and file_operations_dag.py) +pandas==2.1.4 + +# Data generation library (used in file_operations_dag.py) +mimesis==15.1.0 + +# Airflow PostgreSQL provider (used in sql_basic_dag.py and data_processing_dag.py) +apache-airflow-providers-postgres==5.11.1 diff --git a/docker-compose.yml b/docker-compose.yml index fc6b0a4..4cad047 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -1,3 +1,16 @@ +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: replace-me-with-random-string + POSTGRES_TRAINING_HOST: postgres-training + POSTGRES_TRAINING_PORT: "5432" + POSTGRES_TRAINING_DB: training + POSTGRES_TRAINING_USER: student + POSTGRES_TRAINING_PASSWORD: student + AIRFLOW_CONN_POSTGRES_TRAINING: postgresql://student:student@postgres-training:5432/training + + services: clickhouse: image: clickhouse/clickhouse-server:25.1 @@ -64,6 +77,100 @@ 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: . + image: airflow-optimized:2.9.2 + environment: + <<: *airflow-default-env + command: > + bash -c " + airflow webserver + " + ports: + - "8080:8080" + volumes: + - ./dags:/opt/airflow/dags + - ./data:/opt/airflow/data + networks: + - cs_dwh + depends_on: + airflow-init: + condition: service_completed_successfully + postgres-metadata: + condition: service_healthy + postgres-training: + condition: service_healthy + + # Airflow scheduler + airflow-scheduler: + build: . + # image: airflow-optimized:2.9.2 + environment: + <<: *airflow-default-env + command: > + bash -c " + airflow scheduler + " + volumes: + - ./dags:/opt/airflow/dags + - ./data:/opt/airflow/data + networks: + - cs_dwh + depends_on: + airflow-init: + condition: service_completed_successfully + postgres-metadata: + condition: service_healthy + postgres-training: + condition: service_healthy + + # Airflow init to create admin user + airflow-init: + build: . + user: "0:0" + # image: airflow-optimized:2.9.3 + environment: + <<: *airflow-default-env + volumes: + - ./dags:/opt/airflow/dags + - ./data:/opt/airflow/data + networks: + - cs_dwh + command: > + bash -ceuo pipefail " + mkdir -p /opt/airflow/data && + chmod -R 777 /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 postgres_training || true' && + su -s /bin/bash airflow -c \"airflow connections add 'postgres_training' --conn-uri \\\"$${AIRFLOW_CONN_POSTGRES_TRAINING}\\\" --conn-description 'Training exercises database'\" + " + depends_on: + postgres-metadata: + condition: service_healthy + prometheus: image: prom/prometheus:v2.53.4 volumes: @@ -127,6 +234,7 @@ networks: volumes: grafana_lib: kafka-data: + pgmeta: superset_data: superset_config: From fe9c15c0fecf374a788427c8985e21b8d08e0c34 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Fri, 6 Feb 2026 23:34:45 +0300 Subject: [PATCH 02/12] feat(airflow): configure ClickHouse connection and update infrastructure Update Airflow configuration to integrate with ClickHouse DWH instead of PostgreSQL training database. Changes include: - Switch Airflow dependencies from PostgreSQL to ClickHouse connector - Update docker-compose to use ClickHouse connection and correct Dockerfile - Refactor airflow/requirements.txt to include only essential packages - Add DAGs directory for ETL pipeline orchestration - Update documentation to reflect Airflow integration and access credentials - Adjust service dependencies to wait for ClickHouse startup --- AGENTS.md | 6 ++++++ Dockerfile.airflow | 2 +- README.md | 16 ++++++++++++-- airflow/requirements.txt | 17 ++++++--------- dags/.gitkeep | 0 dags/__init__.py | 0 dags/etl_pipeline_dag.py | 46 ++++++++++++++++++++++++++++++++++++++++ docker-compose.yml | 38 ++++++++++++++++----------------- docs/ARCHITECTURE.md | 11 ++++++++-- 9 files changed, 102 insertions(+), 34 deletions(-) create mode 100644 dags/.gitkeep create mode 100644 dags/__init__.py create mode 100644 dags/etl_pipeline_dag.py diff --git a/AGENTS.md b/AGENTS.md index 5a92a4b..d03d831 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,6 +20,7 @@ ## Ключевые артефакты ### Исполняемые файлы (текущая структура) +- `dags/` — Airflow DAGs для оркестрации ETL - `ddl/` — SQL для создания объектов БД: - `00_databases.sql` — создание БД stg/ods/dds/dm - `10_stg.sql` — STG слой (Kafka Engine + MV) @@ -32,6 +34,8 @@ - `apply_clickhouse_ddl.sh` — применение DDL - `load_kafka_data.sh` — загрузка в Kafka - `run_batch.sh` — запуск batch-процесса +- `airflow/` — конфигурация Airflow: + - `requirements.txt` — зависимости (clickhouse-connect и др.) ### Планы и документация (legacy) - `plans/clickhouse_ddl.md` — исходный план (inline DDL, legacy) @@ -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` @@ -89,6 +94,7 @@ - Kafka ingest: наличие данных в `stg.*` и типизированных строк в `ods.*`. - Мониторинг: доступность `/metrics` у ClickHouse и скрейп в Prometheus. +- **Airflow: `http://localhost:8080` должен показывать UI и DAG `etl_pipeline`.** - BI: витрина `dm.v_events_enriched` должна отвечать за разумное время при фильтре по дате. ## Связанная документация diff --git a/Dockerfile.airflow b/Dockerfile.airflow index 0df3b0d..0643813 100644 --- a/Dockerfile.airflow +++ b/Dockerfile.airflow @@ -5,7 +5,7 @@ FROM apache/airflow:2.9.3 ENV AIRFLOW_HOME=/opt/airflow # Copy the project requirements file -COPY --chown=airflow:0 requirements.txt ${AIRFLOW_HOME}/requirements.txt +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 diff --git a/README.md b/README.md index ac8eff7..6d73d44 100644 --- a/README.md +++ b/README.md @@ -1,6 +1,6 @@ # 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)]() @@ -57,6 +57,7 @@ docker compose exec clickhouse clickhouse-client \ |--------|-----|------------| | 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 | Визуализация метрик | @@ -85,7 +86,12 @@ flowchart TB DM["DM — витрины"] end + subgraph Airflow["⚙️ Airflow"] + DAG[ETL DAGs] + end + Sources -->|make data| Kafka -->|MV| STG -->|MV| ODS -->|Batch SQL| DDS -->|VIEW| DM + DAG -.->|оркестрация| ODS & DDS & DM ``` **Поток данных:** @@ -102,9 +108,12 @@ flowchart TB ``` . +├── dags/ # Airflow DAGs для оркестрации ├── ddl/ # SQL для создания объектов (00_databases → 40_dm) ├── jobs/ # Batch-трансформации (ODS→DDS, DDS→DM) ├── scripts/ # Автоматизация (apply ddl, load data, run batch) +├── airflow/ # Конфигурация Airflow +│ └── requirements.txt ├── docs/ # Документация │ └── ARCHITECTURE.md # Подробное описание слоёв ├── data/ # Исходные JSONL файлы @@ -187,7 +196,10 @@ flowchart LR ## 🔮 Развитие проекта -- [ ] **Airflow** — оркестрация batch-процесса +### ✅ Реализовано +- [x] **Airflow** — оркестрация batch-процесса (инфраструктура готова, DAGs в разработке) + +### 📋 В планах - [ ] **Инкрементальный batch** — watermark-based загрузка - [ ] **Материализация витрин** — для тяжёлых агрегаций - [ ] **DQ мониторинг** — алерты на ошибки парсинга diff --git a/airflow/requirements.txt b/airflow/requirements.txt index ba5c858..c6b6189 100644 --- a/airflow/requirements.txt +++ b/airflow/requirements.txt @@ -1,14 +1,11 @@ -# Optimized requirements for Airflow Docker setup -# Only includes packages actually used in DAGs +# Airflow requirements для ClickHouse DWH проекта -# Core database connector +# Core database connector для metadata psycopg2-binary==2.9.9 -# Data processing (used in data_processing_dag.py and file_operations_dag.py) +# ClickHouse provider для ETL +# Примечание: официальный провайдер deprecated, используем clickhouse-connect +clickhouse-connect==0.8.0 + +# Для работы с данными pandas==2.1.4 - -# Data generation library (used in file_operations_dag.py) -mimesis==15.1.0 - -# Airflow PostgreSQL provider (used in sql_basic_dag.py and data_processing_dag.py) -apache-airflow-providers-postgres==5.11.1 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/etl_pipeline_dag.py b/dags/etl_pipeline_dag.py new file mode 100644 index 0000000..feb0b1c --- /dev/null +++ b/dags/etl_pipeline_dag.py @@ -0,0 +1,46 @@ +""" +ETL Pipeline DAG для ClickHouse DWH + +Шаблон DAG для оркестрации пайплайна данных. +Полная реализация будет добавлена позже. + +Пайплайн: + 1. DDL - создание структуры БД + 2. Load - загрузка данных в Kafka + 3. Transform - batch трансформация ODS → DDS → DM +""" + +from datetime import datetime, timedelta +from airflow import DAG +from airflow.operators.bash import BashOperator +from airflow.operators.empty import EmptyOperator + +# Базовые настройки DAG +default_args = { + "owner": "airflow", + "depends_on_past": False, + "email_on_failure": False, + "email_on_retry": False, + "retries": 1, + "retry_delay": timedelta(minutes=5), +} + +with DAG( + dag_id="etl_pipeline", + default_args=default_args, + description="ETL pipeline для ClickHouse DWH", + schedule=None, # Запуск только вручную (пока) + start_date=datetime(2024, 1, 1), + catchup=False, + tags=["etl", "clickhouse", "dwh"], +) as dag: + + # TODO: добавить задачи пайплайна + # - ddl: создание структуры БД + # - load: загрузка данных в Kafka + # - transform: batch трансформация + + start = EmptyOperator(task_id="start") + end = EmptyOperator(task_id="end") + + start >> end diff --git a/docker-compose.yml b/docker-compose.yml index 4cad047..5da65f9 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -2,13 +2,9 @@ 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: replace-me-with-random-string - POSTGRES_TRAINING_HOST: postgres-training - POSTGRES_TRAINING_PORT: "5432" - POSTGRES_TRAINING_DB: training - POSTGRES_TRAINING_USER: student - POSTGRES_TRAINING_PASSWORD: student - AIRFLOW_CONN_POSTGRES_TRAINING: postgresql://student:student@postgres-training:5432/training + AIRFLOW__WEBSERVER__SECRET_KEY: ${AIRFLOW_SECRET_KEY:-replace-me-with-random-string} + # ClickHouse connection для ETL + AIRFLOW_CONN_CLICKHOUSE_DEFAULT: clickhouse://default:123456@clickhouse:8123/default services: @@ -98,8 +94,10 @@ services: # Airflow webserver airflow-webserver: - build: . - image: airflow-optimized:2.9.2 + build: + context: . + dockerfile: Dockerfile.airflow + image: airflow-optimized:2.9.3 environment: <<: *airflow-default-env command: > @@ -118,13 +116,14 @@ services: condition: service_completed_successfully postgres-metadata: condition: service_healthy - postgres-training: - condition: service_healthy + clickhouse: + condition: service_started # Airflow scheduler airflow-scheduler: - build: . - # image: airflow-optimized:2.9.2 + build: + context: . + dockerfile: Dockerfile.airflow environment: <<: *airflow-default-env command: > @@ -141,14 +140,15 @@ services: condition: service_completed_successfully postgres-metadata: condition: service_healthy - postgres-training: - condition: service_healthy + clickhouse: + condition: service_started # Airflow init to create admin user airflow-init: - build: . + build: + context: . + dockerfile: Dockerfile.airflow user: "0:0" - # image: airflow-optimized:2.9.3 environment: <<: *airflow-default-env volumes: @@ -164,8 +164,8 @@ services: 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 postgres_training || true' && - su -s /bin/bash airflow -c \"airflow connections add 'postgres_training' --conn-uri \\\"$${AIRFLOW_CONN_POSTGRES_TRAINING}\\\" --conn-description 'Training exercises database'\" + 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: diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 5400f5c..6369e23 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -586,9 +586,11 @@ INSERT INTO dm.daily_traffic SELECT * FROM dm.v_daily_traffic; ### Airflow-оркестрация +Инфраструктура Airflow развёрнута и готова к использованию: + ```python -# dag.py -with DAG('clickhouse_etl'): +# dags/etl_pipeline_dag.py +with DAG('etl_pipeline'): 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') @@ -596,6 +598,11 @@ with DAG('clickhouse_etl'): ddl >> load >> transform ``` +**Подключение к ClickHouse:** +- Connection: `clickhouse_default` +- URL: `clickhouse://default:123456@clickhouse:8123/default` +- Provider: `clickhouse-connect` (в `airflow/requirements.txt`) + --- ## Полезные запросы From cca315e4f6d8c56ab7b5b5de7560dbc4e33cb33e Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Fri, 6 Feb 2026 23:42:50 +0300 Subject: [PATCH 03/12] docs(airflow): add migration plan for ETL orchestration Detail the architecture for migrating ETL orchestration from make to Airflow. Define DAG structures for database initialization, Kafka data ingestion, batch transformation, and quality monitoring. Include technical specifications, operator details, and implementation phases. --- plans/airflow_dags_plan.md | 271 +++++++++++++++++++++++++++++++++++++ 1 file changed, 271 insertions(+) create mode 100644 plans/airflow_dags_plan.md diff --git a/plans/airflow_dags_plan.md b/plans/airflow_dags_plan.md new file mode 100644 index 0000000..f43d50a --- /dev/null +++ b/plans/airflow_dags_plan.md @@ -0,0 +1,271 @@ +# План миграции ETL-оркестрации на Airflow + +## Цель +Заменить `make`-команды на полноценную оркестрацию через Airflow DAGs с сохранением логики пайплайна STG → ODS → DDS → DM. + +--- + +## Архитектура потока данных (напоминание) + +``` +┌─────────────────────────────────────────────────────────────────────────────┐ +│ Источник: data/*.jsonl → Kafka → STG (Kafka Engine) → ODS (MV) │ +│ │ +│ STG → ODS: real-time через Materialized Views (не требует оркестрации) │ +│ ODS → DDS: batch через SQL (argMax + JOIN) — требует оркестрации │ +│ DDS → DM: batch через SQL (TRUNCATE + INSERT) — требует оркестрации │ +└─────────────────────────────────────────────────────────────────────────────┘ +``` + +--- + +## Структура DAG'ов + +### 1. `ddl_init` — Инициализация схемы БД + +**Назначение:** Создание всех объектов ClickHouse (базы, таблицы, MV, витрины). + +**Параметры:** +- `schedule`: `None` (только ручной запуск) +- `tags`: `["ddl", "init", "clickhouse"]` + +**Задачи (Tasks):** + +| Task ID | Описание | Тип оператора | +|---------|----------|---------------| +| `check_clickhouse` | Проверка доступности ClickHouse | `BashOperator` — `clickhouse-client --query="SELECT 1"` | +| `create_databases` | Создание БД: stg, ods, dds, dm | `SQLExecuteQueryOperator` — `ddl/00_databases.sql` | +| `create_stg` | Таблицы STG + Kafka Engine + MV | `SQLExecuteQueryOperator` — `ddl/10_stg.sql` | +| `create_ods` | Таблицы ODS + MV STG→ODS | `SQLExecuteQueryOperator` — `ddl/20_ods.sql` | +| `create_dds` | Таблицы DDS (batch-загрузка) | `SQLExecuteQueryOperator` — `ddl/30_dds.sql` | +| `create_dm` | Витрины DM (VIEW) | `SQLExecuteQueryOperator` — `ddl/40_dm.sql` | +| `verify_schema` | Проверка: список созданных таблиц | `BashOperator` — запрос `SHOW TABLES FROM each DB` | + +**Зависимости:** +``` +check_clickhouse >> create_databases >> [create_stg, create_ods, create_dds, create_dm] >> verify_schema +``` + +**Замечание:** Порядок важен — сначала `stg`+`ods` (MV работают сразу), потом `dds`+`dm`. + +--- + +### 2. `kafka_load` — Загрузка данных в Kafka + +**Назначение:** Загрузка JSON-данных из `data/*.jsonl` в Kafka-топики. + +**Параметры:** +- `schedule`: `None` (только ручной запуск) +- `tags`: `["kafka", "ingest", "demo"]` +- `params`: + - `limit`: int — количество строк для загрузки (default: 50, 0 = все) + - `full_load`: bool — загрузить полные файлы + - `reset_topics`: bool — пересоздать топики (default: true) + +**Задачи (Tasks):** + +| Task ID | Описание | Тип оператора | +|---------|----------|---------------| +| `check_kafka` | Проверка доступности Kafka | `BashOperator` — `kafka-topics.sh --list` | +| `reset_topics` | Удаление/создание топиков (conditional) | `BashOperator` — `kafka-topics.sh --delete/--create` | +| `load_browser` | Загрузка browser_events.jsonl | `BashOperator` — `kafka-console-producer.sh` | +| `load_location` | Загрузка location_events.jsonl | `BashOperator` — `kafka-console-producer.sh` | +| `load_device` | Загрузка device_events.jsonl | `BashOperator` — `kafka-console-producer.sh` | +| `load_geo` | Загрузка geo_events.jsonl | `BashOperator` — `kafka-console-producer.sh` | +| `verify_load` | Проверка: количество сообщений в топиках | `BashOperator` — `kafka-console-consumer.sh --from-beginning` или проверка через CH | + +**Зависимости:** +``` +check_kafka >> reset_topics >> [load_browser, load_location, load_device, load_geo] >> verify_load +``` + +**Особенности:** +- Загрузка файлов может идти параллельно (независимые топики). +- `head -n {{ params.limit }}` для среза данных. + +--- + +### 3. `etl_batch_transform` — Batch-трансформация ODS → DDS → DM + +**Назначение:** Основной ETL-пайплайн — сборка сущностей и обновление витрин. + +**Параметры:** +- `schedule`: `"@once"` для демо или `"*/15 * * * *"` (каждые 15 мин) +- `tags`: `["etl", "batch", "dds", "dm"]` +- `params`: + - `full_refresh`: bool — полная перезагрузка или инкремент (default: true для демо) + +**Задачи (Tasks):** + +| Task ID | Описание | Тип оператора | +|---------|----------|---------------| +| `wait_for_ods` | Ожидание появления данных в ODS | `BashOperator` — `SELECT count() FROM ods.browser_event` | +| `check_ods_quality` | Проверка качества ODS: ошибки парсинга | `SQLExecuteQueryOperator` — `SELECT layer, count() FROM ods.browser_event WHERE ...` | +| `truncate_dds` | Очистка DDS таблиц (conditional) | `SQLExecuteQueryOperator` — `TRUNCATE TABLE dds.click, dds.event` | +| `refresh_dds_click` | Загрузка dds.click (device + geo) | `SQLExecuteQueryOperator` — `jobs/30_dds_refresh.sql` (часть для click) | +| `refresh_dds_event` | Загрузка dds.event (browser + location) | `SQLExecuteQueryOperator` — `jobs/30_dds_refresh.sql` (часть для event) | +| `check_dds_integrity` | Проверка: orphan events (есть event, нет click) | `SQLExecuteQueryOperator` — `SELECT count() FROM dds.event WHERE click_id NOT IN (...)` | +| `refresh_dm_summary` | Обновление dm.dq_summary | `SQLExecuteQueryOperator` — `jobs/40_dm_refresh.sql` | +| `validate_dm` | Проверка: dq_summary не пустая | `BashOperator` — `SELECT * FROM dm.dq_summary` | + +**Зависимости:** +``` +wait_for_ods >> check_ods_quality >> truncate_dds >> [refresh_dds_click, refresh_dds_event] >> check_dds_integrity >> refresh_dm_summary >> validate_dm +``` + +**Особенности:** +- `refresh_dds_click` и `refresh_dds_event` независимы — можно параллельно. +- Для инкрементальной загрузки (в будущем) понадобится watermark (src_ingest_ts). + +--- + +### 4. `data_quality_monitor` — Мониторинг качества данных (опциональный) + +**Назначение:** Регулярная проверка DQ-метрик и алерты. + +**Параметры:** +- `schedule`: `"0 */1 * * *"` (каждый час) +- `tags`: `["dq", "monitoring", "alerts"]` + +**Задачи (Tasks):** + +| Task ID | Описание | Тип оператора | +|---------|----------|---------------| +| `check_stg_volume` | Проверка объёма STG | `SQLExecuteQueryOperator` | +| `check_ods_errors` | Проверка ошибок парсинга ODS | `SQLExecuteQueryOperator` — `SELECT count() FROM ods.*_errors` | +| `check_dds_orphans` | Проверка сиротских записей | `SQLExecuteQueryOperator` | +| `send_alert` | Отправка алерта (если проблемы) | `EmptyOperator` или callback | + +**Зависимости:** Линейная цепочка с условными переходами. + +--- + +## Технические детали реализации + +### Подключение к ClickHouse + +```python +# Connection в Airflow UI (Admin → Connections) +conn_id = "clickhouse_default" +conn_type = "generic" +host = "clickhouse" +port = 8123 # HTTP interface +login = "default" +password = "123456" +``` + +### SQL-операторы + +Для выполнения SQL использовать `SQLExecuteQueryOperator` с `clickhouse-connect`: + +```python +from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator + +refresh_dds = SQLExecuteQueryOperator( + task_id="refresh_dds_click", + conn_id="clickhouse_default", + sql=""" + INSERT INTO dds.click + SELECT ... -- SQL из jobs/30_dds_refresh.sql + """ +) +``` + +### Bash-операторы для Kafka + +```python +from airflow.operators.bash import BashOperator + +load_browser = BashOperator( + task_id="load_browser", + bash_command=""" + head -n {{ params.limit }} /opt/airflow/data/browser_events.jsonl | \ + kafka-console-producer.sh --bootstrap-server kafka:29092 --topic browser_events + """ +) +``` + +### Сенсоры (Sensors) + +Для ожидания данных использовать `SqlSensor`: + +```python +from airflow.providers.common.sql.sensors.sql import SqlSensor + +wait_for_ods = SqlSensor( + task_id="wait_for_ods", + conn_id="clickhouse_default", + sql="SELECT count() > 0 FROM ods.browser_event", + mode="poke", + poke_interval=30, + timeout=600 +) +``` + +--- + +## Последовательность внедрения + +1. **Этап 1: DDL и Batch** + - Создать `ddl_init` DAG + - Создать `etl_batch_transform` DAG + - Проверить полный цикл: DDL → load (ручной) → transform + +2. **Этап 2: Kafka Load** + - Создать `kafka_load` DAG с параметрами + - Интегрировать с `etl_batch_transform` через TriggerDagRunOperator + +3. **Этап 3: Мониторинг** + - Добавить `data_quality_monitor` DAG + - Настроить алерты (email/Slack) + +--- + +## Файловая структура + +``` +dags/ +├── __init__.py +├── ddl_init.py # DAG #1: Инициализация схемы +├── kafka_load.py # DAG #2: Загрузка в Kafka +├── etl_batch_transform.py # DAG #3: Batch ETL +├── data_quality_monitor.py # DAG #4: DQ мониторинг (опционально) +├── utils/ +│ ├── __init__.py +│ ├── clickhouse_helpers.py # Общие функции для CH +│ └── kafka_helpers.py # Общие функции для Kafka +└── sql/ # SQL-шаблоны (опционально) + ├── dds_click_insert.sql + ├── dds_event_insert.sql + └── dm_summary_insert.sql +``` + +--- + +## Особенности и ограничения + +1. **STG → ODS:** Работает через MV автоматически, не требует DAG. +2. **Очистка Kafka:** MV в ClickHouse запоминают offset'ы — для чистого старта нужно пересоздать MV. +3. **Полная перезагрузка:** Для демо используем `TRUNCATE + INSERT`. В продакшене — инкремент. +4. **Зависимости сервисов:** DAG'и должны проверять доступность ClickHouse/Kafka перед работой. +5. **Идемпотентность:** Batch-задачи должны быть идемпотентны (TRUNCATE перед INSERT). + +--- + +## Проверка после реализации + +```bash +# 1. Запуск Airflow +docker compose up -d airflow + +# 2. В UI должны появиться DAG'и: ddl_init, kafka_load, etl_batch_transform + +# 3. Тестовый прогон: +# - Trigger ddl_init → проверить таблицы в CH +# - Trigger kafka_load (limit=50) → проверить топики +# - Дождаться появления данных в ODS (автоматически через MV) +# - Trigger etl_batch_transform → проверить DDS и DM + +# 4. Проверка результатов: +docker compose exec clickhouse clickhouse-client -q "SELECT * FROM dm.dq_summary" +``` From 4da6a34e4c456a06e8b5893f070a9dadbd0396e1 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sat, 7 Feb 2026 00:03:32 +0300 Subject: [PATCH 04/12] docs(airflow): update dag implementation plan for mvp Consolidate the orchestration strategy by merging `kafka_load` and `etl_batch_transform` into a unified `etl_pipeline`. Replace BashOperator dependencies on Kafka CLI with PythonOperators utilizing `kafka-python`. Add detailed technical specifications for helper functions, MVP stages, and validation checks to align with current infrastructure constraints. --- plans/airflow_dags_plan.md | 391 +++++++++++++++---------------------- 1 file changed, 153 insertions(+), 238 deletions(-) diff --git a/plans/airflow_dags_plan.md b/plans/airflow_dags_plan.md index f43d50a..8bf242a 100644 --- a/plans/airflow_dags_plan.md +++ b/plans/airflow_dags_plan.md @@ -1,271 +1,186 @@ -# План миграции ETL-оркестрации на Airflow +# План развития Airflow DAG'ов ## Цель -Заменить `make`-команды на полноценную оркестрацию через Airflow DAGs с сохранением логики пайплайна STG → ODS → DDS → DM. +Перевести оркестрацию 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`. +- Файл `jobs/30_dds_refresh.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"]`. -``` -┌─────────────────────────────────────────────────────────────────────────────┐ -│ Источник: data/*.jsonl → Kafka → STG (Kafka Engine) → ODS (MV) │ -│ │ -│ STG → ODS: real-time через Materialized Views (не требует оркестрации) │ -│ ODS → DDS: batch через SQL (argMax + JOIN) — требует оркестрации │ -│ DDS → DM: batch через SQL (TRUNCATE + INSERT) — требует оркестрации │ -└─────────────────────────────────────────────────────────────────────────────┘ +### DAG 2 (обязательный): `etl_pipeline` +- `schedule`: `None` (ручной запуск для демо). +- `catchup`: `False`. +- `max_active_runs`: `1`. +- `tags`: `["etl", "clickhouse", "kafka", "demo"]`. + +### DAG 3 (опциональный): `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` | Создание БД | `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` | +| `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 DAG запускается вручную: при первом bootstrap, при изменении схемы, после `docker compose down -v`. -## Структура DAG'ов +## Дизайн DAG `etl_pipeline` +### Params (через Trigger DAG with config) +- `run_ingest`: bool, default `true`. +- `limit`: int, default `50`. +- `full_load`: bool, default `false`. +- `reset_topics`: bool, default `true`. +- `full_refresh`: bool, default `true`. -### 1. `ddl_init` — Инициализация схемы БД +### 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 | +| `check_input_files` | Проверка наличия `data/*_events.jsonl` | `PythonOperator` | -**Назначение:** Создание всех объектов ClickHouse (базы, таблицы, MV, витрины). +### TaskGroup `ingest` (выполняется только при `run_ingest=true`) +| Task ID | Что делает | Реализация | +|---------|------------|------------| +| `kafka_prepare_topics` | reset/create топиков по параметру `reset_topics` | `PythonOperator` + `kafka-python` AdminClient | +| `load_browser` | Публикация строк из `browser_events.jsonl` | `PythonOperator` + `KafkaProducer` | +| `load_location` | Публикация строк из `location_events.jsonl` | `PythonOperator` + `KafkaProducer` | +| `load_device` | Публикация строк из `device_events.jsonl` | `PythonOperator` + `KafkaProducer` | +| `load_geo` | Публикация строк из `geo_events.jsonl` | `PythonOperator` + `KafkaProducer` | +| `wait_for_stg_data` | Ожидание появления данных в `stg.*_raw` | `PythonSensor`/poll SQL | -**Параметры:** -- `schedule`: `None` (только ручной запуск) -- `tags`: `["ddl", "init", "clickhouse"]` - -**Задачи (Tasks):** - -| Task ID | Описание | Тип оператора | -|---------|----------|---------------| -| `check_clickhouse` | Проверка доступности ClickHouse | `BashOperator` — `clickhouse-client --query="SELECT 1"` | -| `create_databases` | Создание БД: stg, ods, dds, dm | `SQLExecuteQueryOperator` — `ddl/00_databases.sql` | -| `create_stg` | Таблицы STG + Kafka Engine + MV | `SQLExecuteQueryOperator` — `ddl/10_stg.sql` | -| `create_ods` | Таблицы ODS + MV STG→ODS | `SQLExecuteQueryOperator` — `ddl/20_ods.sql` | -| `create_dds` | Таблицы DDS (batch-загрузка) | `SQLExecuteQueryOperator` — `ddl/30_dds.sql` | -| `create_dm` | Витрины DM (VIEW) | `SQLExecuteQueryOperator` — `ddl/40_dm.sql` | -| `verify_schema` | Проверка: список созданных таблиц | `BashOperator` — запрос `SHOW TABLES FROM each DB` | - -**Зависимости:** -``` -check_clickhouse >> create_databases >> [create_stg, create_ods, create_dds, create_dm] >> verify_schema +Зависимости: +```text +kafka_prepare_topics >> [load_browser, load_location, load_device, load_geo] >> wait_for_stg_data ``` -**Замечание:** Порядок важен — сначала `stg`+`ods` (MV работают сразу), потом `dds`+`dm`. +Примечание: +- Для демо соблюдать ограничение на малый срез данных: `limit=20..50` по умолчанию. ---- +### 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 | +| `refresh_dds` | ODS → DDS | `jobs/30_dds_refresh.sql` | +| `check_dds_integrity` | Проверка orphan событий | inline SQL | +| `refresh_dm_summary` | DDS → DM DQ summary | `jobs/40_dm_refresh.sql` | +| `validate_dm_summary` | Проверка, что `dm.dq_summary` не пуста | SQL-check | -### 2. `kafka_load` — Загрузка данных в Kafka - -**Назначение:** Загрузка JSON-данных из `data/*.jsonl` в Kafka-топики. - -**Параметры:** -- `schedule`: `None` (только ручной запуск) -- `tags`: `["kafka", "ingest", "demo"]` -- `params`: - - `limit`: int — количество строк для загрузки (default: 50, 0 = все) - - `full_load`: bool — загрузить полные файлы - - `reset_topics`: bool — пересоздать топики (default: true) - -**Задачи (Tasks):** - -| Task ID | Описание | Тип оператора | -|---------|----------|---------------| -| `check_kafka` | Проверка доступности Kafka | `BashOperator` — `kafka-topics.sh --list` | -| `reset_topics` | Удаление/создание топиков (conditional) | `BashOperator` — `kafka-topics.sh --delete/--create` | -| `load_browser` | Загрузка browser_events.jsonl | `BashOperator` — `kafka-console-producer.sh` | -| `load_location` | Загрузка location_events.jsonl | `BashOperator` — `kafka-console-producer.sh` | -| `load_device` | Загрузка device_events.jsonl | `BashOperator` — `kafka-console-producer.sh` | -| `load_geo` | Загрузка geo_events.jsonl | `BashOperator` — `kafka-console-producer.sh` | -| `verify_load` | Проверка: количество сообщений в топиках | `BashOperator` — `kafka-console-consumer.sh --from-beginning` или проверка через CH | - -**Зависимости:** -``` -check_kafka >> reset_topics >> [load_browser, load_location, load_device, load_geo] >> verify_load +Зависимости: +```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 ``` -**Особенности:** -- Загрузка файлов может идти параллельно (независимые топики). -- `head -n {{ params.limit }}` для среза данных. - ---- - -### 3. `etl_batch_transform` — Batch-трансформация ODS → DDS → DM - -**Назначение:** Основной ETL-пайплайн — сборка сущностей и обновление витрин. - -**Параметры:** -- `schedule`: `"@once"` для демо или `"*/15 * * * *"` (каждые 15 мин) -- `tags`: `["etl", "batch", "dds", "dm"]` -- `params`: - - `full_refresh`: bool — полная перезагрузка или инкремент (default: true для демо) - -**Задачи (Tasks):** - -| Task ID | Описание | Тип оператора | -|---------|----------|---------------| -| `wait_for_ods` | Ожидание появления данных в ODS | `BashOperator` — `SELECT count() FROM ods.browser_event` | -| `check_ods_quality` | Проверка качества ODS: ошибки парсинга | `SQLExecuteQueryOperator` — `SELECT layer, count() FROM ods.browser_event WHERE ...` | -| `truncate_dds` | Очистка DDS таблиц (conditional) | `SQLExecuteQueryOperator` — `TRUNCATE TABLE dds.click, dds.event` | -| `refresh_dds_click` | Загрузка dds.click (device + geo) | `SQLExecuteQueryOperator` — `jobs/30_dds_refresh.sql` (часть для click) | -| `refresh_dds_event` | Загрузка dds.event (browser + location) | `SQLExecuteQueryOperator` — `jobs/30_dds_refresh.sql` (часть для event) | -| `check_dds_integrity` | Проверка: orphan events (есть event, нет click) | `SQLExecuteQueryOperator` — `SELECT count() FROM dds.event WHERE click_id NOT IN (...)` | -| `refresh_dm_summary` | Обновление dm.dq_summary | `SQLExecuteQueryOperator` — `jobs/40_dm_refresh.sql` | -| `validate_dm` | Проверка: dq_summary не пустая | `BashOperator` — `SELECT * FROM dm.dq_summary` | - -**Зависимости:** -``` -wait_for_ods >> check_ods_quality >> truncate_dds >> [refresh_dds_click, refresh_dds_event] >> check_dds_integrity >> refresh_dm_summary >> validate_dm +### Итоговая цепочка `etl_pipeline` +```text +precheck >> ingest(optional) >> transform ``` -**Особенности:** -- `refresh_dds_click` и `refresh_dds_event` независимы — можно параллельно. -- Для инкрементальной загрузки (в будущем) понадобится watermark (src_ingest_ts). +## Техническая реализация (приземленно) +### ClickHouse в Airflow +- Использовать `clickhouse-connect` напрямую в Python helper, а не `SQLExecuteQueryOperator`. +- Брать параметры подключения из `conn_id = clickhouse_default` через `BaseHook.get_connection`. ---- +### Kafka в Airflow +- Добавить зависимость `kafka-python` в `airflow/requirements.txt`. +- Использовать Python-код для: + - reset/create топиков; + - публикации строк из `.jsonl` (1 строка = 1 message value). -### 4. `data_quality_monitor` — Мониторинг качества данных (опциональный) +### Общие 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` -**Назначение:** Регулярная проверка DQ-метрик и алерты. - -**Параметры:** -- `schedule`: `"0 */1 * * *"` (каждый час) -- `tags`: `["dq", "monitoring", "alerts"]` - -**Задачи (Tasks):** - -| Task ID | Описание | Тип оператора | -|---------|----------|---------------| -| `check_stg_volume` | Проверка объёма STG | `SQLExecuteQueryOperator` | -| `check_ods_errors` | Проверка ошибок парсинга ODS | `SQLExecuteQueryOperator` — `SELECT count() FROM ods.*_errors` | -| `check_dds_orphans` | Проверка сиротских записей | `SQLExecuteQueryOperator` | -| `send_alert` | Отправка алерта (если проблемы) | `EmptyOperator` или callback | - -**Зависимости:** Линейная цепочка с условными переходами. - ---- - -## Технические детали реализации - -### Подключение к ClickHouse - -```python -# Connection в Airflow UI (Admin → Connections) -conn_id = "clickhouse_default" -conn_type = "generic" -host = "clickhouse" -port = 8123 # HTTP interface -login = "default" -password = "123456" -``` - -### SQL-операторы - -Для выполнения SQL использовать `SQLExecuteQueryOperator` с `clickhouse-connect`: - -```python -from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator - -refresh_dds = SQLExecuteQueryOperator( - task_id="refresh_dds_click", - conn_id="clickhouse_default", - sql=""" - INSERT INTO dds.click - SELECT ... -- SQL из jobs/30_dds_refresh.sql - """ -) -``` - -### Bash-операторы для Kafka - -```python -from airflow.operators.bash import BashOperator - -load_browser = BashOperator( - task_id="load_browser", - bash_command=""" - head -n {{ params.limit }} /opt/airflow/data/browser_events.jsonl | \ - kafka-console-producer.sh --bootstrap-server kafka:29092 --topic browser_events - """ -) -``` - -### Сенсоры (Sensors) - -Для ожидания данных использовать `SqlSensor`: - -```python -from airflow.providers.common.sql.sensors.sql import SqlSensor - -wait_for_ods = SqlSensor( - task_id="wait_for_ods", - conn_id="clickhouse_default", - sql="SELECT count() > 0 FROM ods.browser_event", - mode="poke", - poke_interval=30, - timeout=600 -) -``` - ---- - -## Последовательность внедрения - -1. **Этап 1: DDL и Batch** - - Создать `ddl_init` DAG - - Создать `etl_batch_transform` DAG - - Проверить полный цикл: DDL → load (ручной) → transform - -2. **Этап 2: Kafka Load** - - Создать `kafka_load` DAG с параметрами - - Интегрировать с `etl_batch_transform` через TriggerDagRunOperator - -3. **Этап 3: Мониторинг** - - Добавить `data_quality_monitor` DAG - - Настроить алерты (email/Slack) - ---- - -## Файловая структура - -``` +## Структура файлов +```text dags/ ├── __init__.py -├── ddl_init.py # DAG #1: Инициализация схемы -├── kafka_load.py # DAG #2: Загрузка в Kafka -├── etl_batch_transform.py # DAG #3: Batch ETL -├── data_quality_monitor.py # DAG #4: DQ мониторинг (опционально) -├── utils/ -│ ├── __init__.py -│ ├── clickhouse_helpers.py # Общие функции для CH -│ └── kafka_helpers.py # Общие функции для Kafka -└── sql/ # SQL-шаблоны (опционально) - ├── dds_click_insert.sql - ├── dds_event_insert.sql - └── dm_summary_insert.sql +├── ddl_init_dag.py # отдельный DAG для DDL (обязателен) +├── etl_pipeline_dag.py # основной ETL DAG (обязателен) +├── dq_monitor_dag.py # опциональный DAG мониторинга +└── utils/ + ├── __init__.py + ├── clickhouse_helpers.py + └── kafka_helpers.py ``` ---- +## Этапы внедрения +1. Этап 1 (MVP, обязательно): + - Реализовать `ddl_init` и `etl_pipeline`. + - В `etl_pipeline` оставить `precheck + transform`, ingest пока выполнять внешней командой `make data`. + - Проверить путь: `ddl_init` -> ODS -> DDS -> DM после ручной загрузки в Kafka. +2. Этап 2: + - Добавить в `etl_pipeline` ingest внутри DAG через `kafka-python`. + - Добавить ветвление `run_ingest=false` для сценария "только transform". +3. Этап 3: + - Добавить `dq_monitor` и alert callback (email/Slack/webhook). -## Особенности и ограничения - -1. **STG → ODS:** Работает через MV автоматически, не требует DAG. -2. **Очистка Kafka:** MV в ClickHouse запоминают offset'ы — для чистого старта нужно пересоздать MV. -3. **Полная перезагрузка:** Для демо используем `TRUNCATE + INSERT`. В продакшене — инкремент. -4. **Зависимости сервисов:** DAG'и должны проверять доступность ClickHouse/Kafka перед работой. -5. **Идемпотентность:** Batch-задачи должны быть идемпотентны (TRUNCATE перед INSERT). - ---- - -## Проверка после реализации +## Критерии готовности +- В Airflow UI видны DAG `ddl_init` и `etl_pipeline`. +- `etl_pipeline` падает с понятной ошибкой, если схема не применена. +- Ручной запуск DAG с `limit=50` завершает pipeline без падений. +- После прогона: + - в `ods.browser_event` есть строки; + - в `dds.click` и `dds.event` есть строки; + - `dm.dq_summary` заполнена. +- При повторном запуске с `full_refresh=true` нет неконтролируемых дублей в DDS. +## Минимальные smoke-checks ```bash -# 1. Запуск Airflow -docker compose up -d airflow +# 1) Запуск инфраструктуры +make up -# 2. В UI должны появиться DAG'и: ddl_init, kafka_load, etl_batch_transform +# 2) Открыть Airflow UI +# http://localhost:8080 (admin/admin) -# 3. Тестовый прогон: -# - Trigger ddl_init → проверить таблицы в CH -# - Trigger kafka_load (limit=50) → проверить топики -# - Дождаться появления данных в ODS (автоматически через MV) -# - Trigger etl_batch_transform → проверить DDS и DM +# 3) Один раз запустить ddl_init +# Trigger DAG ddl_init (без config или {"verify_only": false}) -# 4. Проверка результатов: -docker compose exec clickhouse clickhouse-client -q "SELECT * FROM dm.dq_summary" +# 4) Trigger DAG etl_pipeline с config: +# {"run_ingest": true, "limit": 50, "full_load": false, "reset_topics": true, "full_refresh": true} + +# 5) Проверка результатов +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" ``` + +## Следующий шаг после MVP +- Для ускорения можно разделить `jobs/30_dds_refresh.sql` на два файла и распараллелить `refresh_dds_click` и `refresh_dds_event` в DAG. +- Для продакшн-режима перейти с `full_refresh` на watermark-инкремент. From 6466921bda446b4ac4bc1ecc80e2574795058c58 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sat, 7 Feb 2026 19:44:37 +0300 Subject: [PATCH 05/12] docs(airflow): refactor dag implementation plan to separate kafka ingestion Split Kafka ingestion into a dedicated `kafka_load` DAG to enable independent experimentation with data loading without triggering the full ETL pipeline. Restructure implementation phases: Stage 1 uses `make data` for MVP, Stage 2 adds the standalone Kafka DAG, Stage 3 adds monitoring. Update DAG numbering, parameters, task groups, and acceptance criteria to reflect the new architecture. --- plans/airflow_dags_plan.md | 148 ++++++++++++++++++++++++------------- 1 file changed, 98 insertions(+), 50 deletions(-) diff --git a/plans/airflow_dags_plan.md b/plans/airflow_dags_plan.md index 8bf242a..4bc7716 100644 --- a/plans/airflow_dags_plan.md +++ b/plans/airflow_dags_plan.md @@ -1,11 +1,11 @@ # План развития Airflow DAG'ов ## Цель -Перевести оркестрацию ETL на Airflow так, чтобы пайплайн оставался устойчивым к "грязным" данным и соответствовал текущей цели проекта: Kafka → ClickHouse (STG → ODS → DDS → DM) → витрины для BI. +Перевести оркестрацию 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` использовать нельзя без донастройки образа. +- В 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. @@ -17,13 +17,21 @@ - `is_paused_upon_creation`: `True`. - `tags`: `["ddl", "bootstrap", "clickhouse"]`. -### DAG 2 (обязательный): `etl_pipeline` +### 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", "kafka", "demo"]`. +- `tags`: `["etl", "clickhouse", "demo"]`. -### DAG 3 (опциональный): `dq_monitor` +### DAG 4 (опциональный): `dq_monitor` - `schedule`: `0 * * * *`. - `catchup`: `False`. - `tags`: `["dq", "monitoring"]`. @@ -49,40 +57,54 @@ check_clickhouse >> ddl_00_databases >> ddl_10_stg >> ddl_20_ods >> ddl_30_dds > ``` Примечание: -- DDL DAG запускается вручную: при первом bootstrap, при изменении схемы, после `docker compose down -v`. +- `ddl_init` запускается вручную: при первом bootstrap, после `docker compose down -v`, после изменений схемы. -## Дизайн DAG `etl_pipeline` +## Дизайн DAG `kafka_load` (отдельный независимый контур) ### Params (через Trigger DAG with config) -- `run_ingest`: bool, default `true`. - `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 | -| `check_input_files` | Проверка наличия `data/*_events.jsonl` | `PythonOperator` | - -### TaskGroup `ingest` (выполняется только при `run_ingest=true`) -| Task ID | Что делает | Реализация | -|---------|------------|------------| -| `kafka_prepare_topics` | reset/create топиков по параметру `reset_topics` | `PythonOperator` + `kafka-python` AdminClient | -| `load_browser` | Публикация строк из `browser_events.jsonl` | `PythonOperator` + `KafkaProducer` | -| `load_location` | Публикация строк из `location_events.jsonl` | `PythonOperator` + `KafkaProducer` | -| `load_device` | Публикация строк из `device_events.jsonl` | `PythonOperator` + `KafkaProducer` | -| `load_geo` | Публикация строк из `geo_events.jsonl` | `PythonOperator` + `KafkaProducer` | -| `wait_for_stg_data` | Ожидание появления данных в `stg.*_raw` | `PythonSensor`/poll SQL | - -Зависимости: -```text -kafka_prepare_topics >> [load_browser, load_location, load_device, load_geo] >> wait_for_stg_data -``` - -Примечание: -- Для демо соблюдать ограничение на малый срез данных: `limit=20..50` по умолчанию. +| `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 | @@ -103,19 +125,29 @@ wait_for_ods_data >> check_ods_quality >> [truncate_dds_click, truncate_dds_even ### Итоговая цепочка `etl_pipeline` ```text -precheck >> ingest(optional) >> transform +precheck >> transform ``` +## Взаимодействие DAG'ов +Базовый сценарий: +```text +ddl_init -> kafka_load -> etl_pipeline +``` + +Экспериментальные сценарии: +- `kafka_load` отдельно: проверить разные наборы/параметры загрузки без запуска transform. +- `etl_pipeline` отдельно: повторно пересчитать DDS/DM по уже загруженным данным. + ## Техническая реализация (приземленно) ### ClickHouse в Airflow -- Использовать `clickhouse-connect` напрямую в Python helper, а не `SQLExecuteQueryOperator`. +- Использовать `clickhouse-connect` напрямую в Python helper. - Брать параметры подключения из `conn_id = clickhouse_default` через `BaseHook.get_connection`. ### Kafka в Airflow -- Добавить зависимость `kafka-python` в `airflow/requirements.txt`. +- Добавить `kafka-python` в `airflow/requirements.txt`. - Использовать Python-код для: - reset/create топиков; - - публикации строк из `.jsonl` (1 строка = 1 message value). + - публикации строк из `.jsonl` (`1 строка = 1 message value`). ### Общие helper-функции - `dags/utils/clickhouse_helpers.py`: @@ -125,13 +157,15 @@ precheck >> ingest(optional) >> transform - `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 (обязателен) -├── etl_pipeline_dag.py # основной ETL DAG (обязателен) +├── kafka_load_dag.py # отдельный DAG для ingest в Kafka (обязателен) +├── etl_pipeline_dag.py # основной DAG ODS -> DDS -> DM (обязателен) ├── dq_monitor_dag.py # опциональный DAG мониторинга └── utils/ ├── __init__.py @@ -142,23 +176,30 @@ dags/ ## Этапы внедрения 1. Этап 1 (MVP, обязательно): - Реализовать `ddl_init` и `etl_pipeline`. - - В `etl_pipeline` оставить `precheck + transform`, ingest пока выполнять внешней командой `make data`. - - Проверить путь: `ddl_init` -> ODS -> DDS -> DM после ручной загрузки в Kafka. + - Для загрузки данных использовать существующий сценарий `make data`. + - Проверить путь `ddl_init -> make data -> etl_pipeline`. 2. Этап 2: - - Добавить в `etl_pipeline` ingest внутри DAG через `kafka-python`. - - Добавить ветвление `run_ingest=false` для сценария "только transform". + - Реализовать отдельный 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). ## Критерии готовности -- В Airflow UI видны DAG `ddl_init` и `etl_pipeline`. -- `etl_pipeline` падает с понятной ошибкой, если схема не применена. -- Ручной запуск DAG с `limit=50` завершает pipeline без падений. -- После прогона: - - в `ods.browser_event` есть строки; - - в `dds.click` и `dds.event` есть строки; - - `dm.dq_summary` заполнена. -- При повторном запуске с `full_refresh=true` нет неконтролируемых дублей в DDS. +- Этап 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 @@ -171,16 +212,23 @@ make up # 3) Один раз запустить ddl_init # Trigger DAG ddl_init (без config или {"verify_only": false}) -# 4) Trigger DAG etl_pipeline с config: -# {"run_ingest": true, "limit": 50, "full_load": false, "reset_topics": true, "full_refresh": true} +# 4) Этап 1: загрузить данные текущим способом +make data -# 5) Проверка результатов +# 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 -- Для ускорения можно разделить `jobs/30_dds_refresh.sql` на два файла и распараллелить `refresh_dds_click` и `refresh_dds_event` в DAG. -- Для продакшн-режима перейти с `full_refresh` на watermark-инкремент. +- Разделить `jobs/30_dds_refresh.sql` на два файла и распараллелить `refresh_dds_click` и `refresh_dds_event`. +- Перейти с `full_refresh` на watermark-инкремент. From 9e340bb729bc5044def73a49766051e0e75a53ce Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sat, 7 Feb 2026 20:51:51 +0300 Subject: [PATCH 06/12] refactor(sql): reorganize sql files into structured directory hierarchy Move DDL files from flat ddl/ directory to sql/ddl/ with layer-based subdirectories (stg, ods, dds, dm). Move batch transformation SQL from jobs/ to sql/ layer directories. Update scripts and documentation to reflect new paths for improved organization and Airflow integration. --- AGENTS.md | 21 +++++----- README.md | 11 ++++- docs/ARCHITECTURE.md | 14 +++---- plans/airflow_dags_plan.md | 20 +++++----- plans/clickhouse_ddl.md | 26 ++++++------ scripts/apply_clickhouse_ddl.sh | 40 ++++++++++++------- scripts/run_batch.sh | 23 +++++++++-- {ddl => sql/ddl}/00_databases.sql | 0 {ddl => sql/ddl/dds}/30_dds.sql | 2 +- {ddl => sql/ddl/dm}/40_dm.sql | 2 +- {ddl => sql/ddl/ods}/20_ods.sql | 0 {ddl => sql/ddl/stg}/10_stg.sql | 0 .../dds/30_ods_to_dds.sql | 0 .../dm/40_dds_to_dm.sql | 0 14 files changed, 95 insertions(+), 64 deletions(-) rename {ddl => sql/ddl}/00_databases.sql (100%) rename {ddl => sql/ddl/dds}/30_dds.sql (99%) rename {ddl => sql/ddl/dm}/40_dm.sql (99%) rename {ddl => sql/ddl/ods}/20_ods.sql (100%) rename {ddl => sql/ddl/stg}/10_stg.sql (100%) rename jobs/30_dds_refresh.sql => sql/dds/30_ods_to_dds.sql (100%) rename jobs/40_dm_refresh.sql => sql/dm/40_dds_to_dm.sql (100%) 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 From 226807ecae9fa2aa9b798b23212220fc9a3ea3e0 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sat, 7 Feb 2026 20:55:14 +0300 Subject: [PATCH 07/12] chore(airflow): replace clickhouse-connect with airflow-clickhouse-plugin --- airflow/requirements.txt | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/airflow/requirements.txt b/airflow/requirements.txt index c6b6189..cdeea91 100644 --- a/airflow/requirements.txt +++ b/airflow/requirements.txt @@ -4,8 +4,7 @@ psycopg2-binary==2.9.9 # ClickHouse provider для ETL -# Примечание: официальный провайдер deprecated, используем clickhouse-connect -clickhouse-connect==0.8.0 +airflow-clickhouse-plugin==1.6.0 # Для работы с данными pandas==2.1.4 From de9d0429a243939400ef6ffaa79e16c0bbea20a1 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sat, 7 Feb 2026 21:36:55 +0300 Subject: [PATCH 08/12] fix(data): update dataset permissions for execution Adjusting file modes on jsonl and markdown files to allow execution within the ETL pipeline. --- data/2026-02-04_13-30-10.png | Bin data/DE-task.md | 0 data/browser_events.jsonl | 0 data/device_events.jsonl | 0 data/geo_events.jsonl | 0 data/location_events.jsonl | 0 6 files changed, 0 insertions(+), 0 deletions(-) mode change 100644 => 100755 data/2026-02-04_13-30-10.png mode change 100644 => 100755 data/DE-task.md mode change 100644 => 100755 data/browser_events.jsonl mode change 100644 => 100755 data/device_events.jsonl mode change 100644 => 100755 data/geo_events.jsonl mode change 100644 => 100755 data/location_events.jsonl 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 From 14d16c5caeaf6958b7c8b3f72ace9a3bb520791f Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sat, 7 Feb 2026 21:52:31 +0300 Subject: [PATCH 09/12] feat(airflow): implement dag orchestration for ddl and etl Add comprehensive DAG implementation for ClickHouse schema initialization and ETL pipeline orchestration. The ddl_init_dag manages database schema creation across stg/ods/dds/dm layers with verification capabilities. The etl_pipeline_dag implements full ODS to DDS to DM transformation flow with data quality checks, branching logic for full/incremental loads, and timeout handling for data availability. Additional changes: - Upgrade Airflow from 2.9.3 to 2.10.5 - Fix ClickHouse connection to use native protocol port 9000 - Mount SQL directory in docker-compose for DAG execution - Update project requirements and documentation comments - Remove unused pandas dependency --- Dockerfile.airflow | 2 +- airflow/requirements.txt | 9 +- dags/__pycache__/ddl_init_dag.cpython-312.pyc | Bin 0 -> 6608 bytes .../etl_pipeline_dag.cpython-312.pyc | Bin 0 -> 10146 bytes dags/ddl_init_dag.py | 188 ++++++++++++ dags/etl_pipeline_dag.py | 286 ++++++++++++++++-- docker-compose.yml | 9 +- 7 files changed, 459 insertions(+), 35 deletions(-) create mode 100644 dags/__pycache__/ddl_init_dag.cpython-312.pyc create mode 100644 dags/__pycache__/etl_pipeline_dag.cpython-312.pyc create mode 100644 dags/ddl_init_dag.py diff --git a/Dockerfile.airflow b/Dockerfile.airflow index 0643813..b857115 100644 --- a/Dockerfile.airflow +++ b/Dockerfile.airflow @@ -1,5 +1,5 @@ # Use the official Airflow image as base -FROM apache/airflow:2.9.3 +FROM apache/airflow:2.10.5 # Set environment variables ENV AIRFLOW_HOME=/opt/airflow diff --git a/airflow/requirements.txt b/airflow/requirements.txt index cdeea91..feae4fb 100644 --- a/airflow/requirements.txt +++ b/airflow/requirements.txt @@ -1,10 +1,7 @@ -# Airflow requirements для ClickHouse DWH проекта +# Airflow requirements для учебного ETL-проекта -# Core database connector для metadata +# Metadata DB для Airflow psycopg2-binary==2.9.9 -# ClickHouse provider для ETL +# ClickHouse operator/hook для DAG'ов airflow-clickhouse-plugin==1.6.0 - -# Для работы с данными -pandas==2.1.4 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 0000000000000000000000000000000000000000..9f0a80167f174f31c285def229b95b3bef23c7b4 GIT binary patch literal 6608 zcmb_gYj70TmA*aG^B!p)8oiJNjU*%uLL)F3Y$4;wdRQ_bY=rG%g15JsZb>uhdAPd= zJ!;q*V{eQX62L!#@~SACty)rOs}v_B@{4#~Tl=W|*d9byLzk&I6>nAiPee{-!>>K3 zXL?2g+dtOZQ+@m1d(J)Q+{gLuIrHm^3MYZ*$m_31pKc)J_t-EVZ=pi1?=%x~mGC4? zc#1cPRNNFcQCK&N=CB#Qv`EM4FdesqEqcF2w8m{=n_jny_P8VL(CaqQ8Fz(UdfhI% z;}ziwz3vbzLzSS4Z!@{q)SJ1i`W}oq*W>JB)OSG2FVvqy`x2wxRDY%30J1I+Itl_|;IDx6E{M1U zUyE5b{fS-vZg+DvAJU zZ{bh|#iM1*>N~ktU~X9u2aG^Nqu&F1?o}9bNqrAR1NjD+ybOdhf)pK}Ws?bURu=+m z@fKQIuoK6jsGq?Lrh|lw3>b~{i!g80TU3Q}en(wqU}k>{!r#kX2suB+LldMl7nevR z6;6pJ6Y^-#q}h2+5tL|L&>UFuf~au7MnB6blkhde9>K@k7nR1vCmg5sf15xqV|NB~Y-Jeo)=f*iDJmgICokTef3q=W=7!{V_iPF6Ip5a*&IHpk(c zmISS`)FTOsG^^S06BT65rMFp_IybAC<(I`s5h-!Whz(vqDr(&i%C8bd@5CZQ5dF?l z87fvpvJwd@271bsipA*Xf%kjl2PVU_&*4~<|G>kW7R?6FS%tXE7I+h}sucGk)}Vzp&Hqn5E5{9mT& ztUY6wsf%x9%rU&A9(Xb(ofNtFJ(^_9n|8?g;0`Hp=NFz#X-~#r&J?_7tTw|kgirDg{*d;Hk36Phh5v*S@wknAu>b4UWmvx+dI zXcjpoMitEjPG*+Vr!{I!bLw4UGy(o)k-^VW(4}yaB2Pz^NsW#rqY2H*rO;_Li;_+O z0ywk^34BnAo)N%@fy;r~3ROF}WkDH=QG~dVP-M*}3lmt=e2FA`T1-Z!*hn&7pwXyY zvlL9QMMQy12)vA5#2C;001r#i&GRA~O+*zITA|deRwX83W~i{NN;&`ocfl+F35qQF z#9KH2M9$lh9m%@_tFFeJtMRsL>qnJa)$M&B1V3m{w;f%n998Y3AA9OnJH@5tpe2E^AdNBZg=nU z$>qc9w&9gZ&^hvv%X@jtimOqj8~^f05AjV4^3lb}F)9pdj1RU&!&Rzd^T}fX4d|kpBI8_C0z*Ewu zcYg(a{Db-;7~2%|0FBj~|DWHd=ObJSIppu}Xh+O4jW1R5Zl`CTU&-rro;%mI~e> z2hEy!LQph1l1z+i^yy?$4APnnS_s3^<6yG%Bo;kT$e0+R7fSn~qIn{d$)qfxU^brQ z1?eg19YelnpvaPsE2}R}y)iW(T&Zl$K7FU5`TT)Py?@tx&6FMZ$X|0=xjJ)YW})KR z^wLbO`)eOG<+_ip)IPi7e=d9YPJQe79+mD@clJWL4${b`y~glI-|L%D!M)}w+yOms z`rMt`82n^RQcVW-1xH9er4(oxea8OrmMK7JX=AWu&6qX^TS_U^;%%=HL&88QR5N5r zuJUOrhbZT$Sx3;BJ_&}XKP(5@#n>D=Za_xxDhA^i2f$UYX8?FXZ%!$`EK1i*N>uX| zPlXH)73c6C!W$hk5jM)p61u>rxTfh5H01}N7lMXdR=40C&|S=o9yaJpH1{x@VqRQk z_(yn zuK6A3AHCf7&7PdI>o58VKj-mL=@f9NKPn_U7+7H<1PPVym*|_&rH4y~zuv?M2R7|90Dl)EW$semeX91++!0~l>Wt|=dAe9yflyF7h)duFkr+wh-K-m^ z`O9@k6Ls}~;)Mz#aM3v>M5#^80L=-U0cP2Hgrm8NIE)w!J|o}=PfE{}SvdS6hDFj7 z&^s4^73?^&58T7XE%hV6LWs#zy&X~%k)Qfv#^un@6`V}YkaYZyH? zI>-+7jT{*qV2_TBve+DZzQ2EPY)ncZ8=eLUflb1QLq{7uG`Vl@g%yYMa!7g&2pAH` z18{AunJJ38XR=ZBpZtVs`8{#`mUP@Db^lFj?~?Ajr0p)*^II}G|ic)!(k&yJj}It)Ep_SWRnWjG_*#Z2=*7 zYJ+)ibKcG58@ltpEuT3obyk%Q+^Zy2?RjrgzOMZUIL_9rres*}vjG z@X!E%Yu?+allQH-A1#sn*6cGN?yFFJEqQN)F8Pra_ui6ZpEWxIlD$m};oF|w+2P++ zF?sLSd^NLL-J7fK1;Yc@>=7`$zBBKIWoasntcN*P!_aCFye(hX39QxD?9nwRu~g5y zZqY4y+H?7_9NjRVxkYzEL#=Ab(Jc#eIeO1hI7dJEK~Ii8sy+wRu=-cATXZZ>d)2^> p9KGXO(=B@cpC?SX!hhD+`JP_hm(jOA+}mF<wt$WF%@KRFC{RTD zmPm2b5pa;cHR6oA0xr_GMch$Oz!NPAl#q9O#2YOQl#+f?q%2w~Q6_eZ8c4L~jYDh)QiO}8br>FqK4RB@88v?7j@<2UT5oq8l18cadz*??4(8$#Unz&VgX6|vW7I1vL`KuxOLn^ms;{OT7RUwY5j~mwGhm1D%B8c4${~JG9x$sBa}Wy?oO$co%{9k-0k`Fn1fl*~K?6o7)ZW60VkLSq>qC3AB34Sb!y?5{vdUfz7O88GJI06!hL&^FMvYYdV6*;eZzZN<)5c#3@%viwu|Bl)WQGx;i`&IHi==@at2 z-)58F1X5SzkAT<}7<`F=*UR#}`~keaE6)?~l}=kLGxX#hMt(c}lKgWZ^#M#c4ct!v zf8cUjej5NU$sfvB7%3=DFu-q~*%1kcCi>$^k>8i#g`gA{wyJbm)6WBk_dvpzRFdz) z+_!<*%Q=F`g8`%~Am(xT3N9AMfhysGt9TAjPs<-7#V01Ev3QIDyyL2H=~ME%>E{aK zoj{2J_&#ocwlQQU=7B8i6EHprtA3Axh2cH{Usy@=VIdd`jkO&OhNZTkDDr~j^V`0{ z%kQHVdoUJ@OF=0dkBNgmonqmF5-)|LykfIn-Yk@(TR zV`T;Ke>Hghxkt zVSh5h!>glk4E^!=1VB6mNS>xSsLHkTD>VFHp{)c=b2J=FO1$VZDaQEG7%wPJj!*D0 zPGsXTHW3sh#l=U1;RwEtz%wcEiX#sZcuAO4%=ifNqGBiSEbL5hQqhawjf8Se26vP) z!yg!M`mP6qe<*eM&bDO1}KZ88ob3>e<&o7J~#5slV|etyg$*M$VKnNI7>xQR9xnnAW>w z)*Q&RF=ZUj1rORNKDCwgQm0KJkS?Tqb`9`0rA*@`i}BRHhCN<-{}^Y{jo;NUvfgWbm-FQegtA?Ofuhjs`o_Usj~-grrl9T^MtuWm#R9|Z z{(bv~6?0f*N5TRxQ`<%{T^S%iw3uikr^*w6p%gF4e_2!e=FZo4p6{CX%s>6{_Fwh>YwwMk0ohWSt>197zB^stovDB9x2}=7HJ=#% z-S%&`Ul-l%XVd+x{2lH${roFl*)=k2geL1LJLfy=o9oMX>Ti0Q)1KyxXWgt>v3s&b zWjBjfrHfY0jsB*naiJL4-f>YyWkl|-D+jOamh1X64&XWPsoiy^Dr2va4K;uK+DW-b zX%Vf;FRW{N4W?iC4bVv{e29$so>$D$TQoRCQ%w~VtNL9nOHDN0EdP6QSfuD5IoxC zLr#FfE&{gv|9QnZab;fo5x7|@H^9krBj7)SAGe&}B0ez2{@L`SOP6y*r3i#$LO3+W zbLx_~D0|*lr0TiU^B`pst(@4U9BL_}Hj4)|?0mBG2-I5En!r)6KM# zGeUG0)jc<&PK9V7#x;jWnlW)mUFMs63J zmS4LYWH)V76~p@v4DRR|?qh~~w(o(<1}Gs616Ld&lcv6R-(J9i$s{(+F#|));J#sI z@W7ru6XSfQa-xj%+FIAh2DG&q0tiTi=w2whdx11T79-;F9;&P37+m9 z)xWlY7ln`-94aIh5d0HQh?&O~aGl@@(`V(6i0j030Gup%v-GL-%YH>Cg{SHg0uNdc zn(&jB72y#Lvgphv(cZrW-V;*=I5JVdaPuq7r&Sr=0P5hN(=#L&!Au3c_Iqk-f{{H0 z^B4qxpI@1HC$l89#I4{`P&uakzgdrizq3>2#>(yp0#U~kMYM6LibdonjS1!rN0XtWiTRk`}24S zm?nCm`91YtGt_lk>)hJeO=o(3*qOGq{!zpYIN@~D$w}zQKTf{@+W;|2ZsQ@GA<+Wl zL7fa73xp()+hL*x7?7!voCpI7jE^oem@0Bygbl~S5*#F^asZY;R1fTJB!wqQa2Gpf z|DeyHSdW2wXA{XtMA3ukQ3DZuFgEGaE9Oubl89RjUWjarQnJEs z2m^pI27jUm%?$OK)pf>ua@EZCY_aq7p_7N?hThpjnc}{gouAv?b7kl0zpMPm(1pf% z_glVyh+UztH*O)0lb1&#G9Q)eVi`wVw#2_Q+hFg`KYU^JkGI^gwcaqd0@{L(D(=_a z)$5$L1xjbNE$AtW|1PklmS`{P4AI{AL09ogjIX4#w5B_FS%f0UB=8O_`^D%mHLWY? z7UU5t(b9p1(4)R{I&8JVNuGnAGeB-p&}GOYAWz8gn>I@M++;kTTNDJ*GJt85m4b|- zxGMjkh@J!|NF;LkfzT)I{7_-6kW8}Y={{_rb z&VZ4-PKpyBpuK|I4z&ujQdm!TA4(6*eJ%`h2q+SuJ|Ph(IFNz>R3k9ycn|dvOAlB& zA^d)_7|KAQA<_EPI@nb20#V7|$1o5QAihKw%7h+@K@h$@% zJ{T2M;)*#pSMjbiPcf;;U*TO)Jj5}MVYLM^?2#ZSCCpblV8Rv;gCsF9yQE@4LcyuE z%RJ!9?p?Kr$#!5TFI)H3{I;vRjP)qkZrtm$u+z=NuA|ig!|wAlWcp8W)lgrWU+EOy zhPh|a|D42!mMhKsA$ie{@{(e}?Ns#Pm?YpuQVeS4B^rmCOgt72DQ5l{AA;J2YP*bL zgy$1Jo3IOzg+16Dz-jhH=b*R~@dz4Bm;}Zik47gj(Kr=%v*eGY&>4}pQbH0pY&CIS_ z)%CNRPj~&G>w<2k_ft>l8SdQJ*)h3p)BMpZ$I_eo%FeHzwPo$4bMz0ovqhdSJ*BU5=VE7LnX<+U<>|8ajHhGPoV9z-OkB6G{?t>S z^=!Yk^;+u!ML$9BrQzY;OW)P&OB@T7zS!}Njj|S>ww<)eC2MY2)-F_2cBk+-taD|n z`wEQwWmw&mUVu1kh~kVxl+Tnzk>yvg!hqF5)pE#}JGJrxIMVbfCLicx%~jK2QkKo% zR&cKAQ+{&Wl%f%ErV#y%fWI6Q27$v=jReI#5pD-0XkZ{hkZrS=n5gk`d0|wX2q#z# z#itjBVYYgG4Tt0Lh|jR>Rtp>;kZ1xgW6mXq;3k75lE!y3-x|GZjOblHr*>w_3NzcAso`Mt<}eXzxFC127=2_8N5l zG6H2SNLB;zyo=JOzzI)O!45ADmGcEL2$d~9Yw{Ub2C1B5T^db3DTcQ&=@;SBTxGb$ z8AC2Q13!XcGWsD1*r4DgK0NG9^xe+KWTSDFfFBYBDZVEzS=qU7}P$&v9yuwq^(=I3m zj6y#&;F#3>YS|^N&c%%dk}Q~WmdC2N5Qd4b!`}nr$DJ5I?)nq)YHIX6^rdd{|F9Cb6hvS!Db)hALIZPp#gIG&O%PhtG%C6`($S{y4Uu}B5q zV;=bVh2x{~+{YU9d~i4B!@q~{ZW!WdJx~Hyf@v%9+KK_O6*HsT0;Yok zX4}Hk*r53rFxgPdSa6Pn4+|y);a3j<(;dZ=`~KwTwGU&W8o7D30X&?D=qE9+fQgV& zt_{U=pBJH`Lret;h%4?NRj-Huj))Slz^deNl2RCGliNtHu7KH-Ql*i|4KF%C0YB*| z>v9OH&sKx)L?k&Hj`=Zg%+)Fajw#PZ@tgdP2TJ9e`>_6+RU-M{a^P@hs(fDXan-X3;;Ur+BrrL+KtfSrBv zK+m3m;e#r=7T>TC->9@!fcHJsH#n@67sw#;z~FG-uKh@e3_Jxyd!Zu8_U>he4(#3A zv;QEe1rV(iFa#l{UD1Qt5>oj3*VxQpa~zxJpi#>13#C~+E#Y~@cmbR5WAh?5v(PBj zgZl>i*quECd-{6Wy#s?Rz78GOv7>KjNH~dTFJW^En;&BHGB&5Nc?=s20@X5)r3xx`E9E0Hr05W>i8Ws zaGUbqrqCqtyq;UtIj%Iw@PZxKY631=}xiPrJE_bQ$-cI&RAaR zzG+{bwy)0E8*bVkO4}d0(37_N=S$P}hi5EzO&)sHLOWH_EK}YEgQ?0OQ)LS-Qzbl0 z7hF`;MtILwZa}EPPz})Xg(9kYfQHW9B0FunTT?|h-t9Ef2kE;m8~rp5Tus)jy)x_S z$+}l%i`}=J5bwn+XJFp9PD^NLk-+=a>;6&r3bYtM(~t=HDyV0PXy z>WfU@lsZhh1*j6xRSPvh?v~Xv_xKHK(*or+4$v1KxmDeo^)_Z*YqE74vZb}(SdBF{ z*-~-GLDg@}dRwxt#{57%4AjCvZA;eMoULj5#;C6aP}Kq?7?pC(`m|-e{BZwuOaHe% zqEcf66%m1 zOoLTz1_1~q+`KxuenZ-_LEilMb<5*f3w$~2NLxDOO}o;TUGnZj*DZ$tI5(ZPw8-r} z*DXC+OXb`%Y0FyK*Oj((U2~=_`{beVv}JsOqT=*+-PaUl-LCugX&u@2g@BGS7N4&ZnR*Wue*JSFuW-21jg#Z^!0-nYBy0^&Up|Ie~t*4g^2^}qBz)UP*w>~%x` o 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 index feb0b1c..586f0dc 100644 --- a/dags/etl_pipeline_dag.py +++ b/dags/etl_pipeline_dag.py @@ -1,46 +1,282 @@ """ -ETL Pipeline DAG для ClickHouse DWH +DAG ETL-процесса ODS -> DDS -> DM для учебного проекта. -Шаблон DAG для оркестрации пайплайна данных. -Полная реализация будет добавлена позже. - -Пайплайн: - 1. DDL - создание структуры БД - 2. Load - загрузка данных в Kafka - 3. Transform - batch трансформация 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 airflow import DAG -from airflow.operators.bash import BashOperator -from airflow.operators.empty import EmptyOperator +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=5), + "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, - description="ETL pipeline для ClickHouse DWH", - schedule=None, # Запуск только вручную (пока) + schedule=None, start_date=datetime(2024, 1, 1), catchup=False, - tags=["etl", "clickhouse", "dwh"], + 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: - - # TODO: добавить задачи пайплайна - # - ddl: создание структуры БД - # - load: загрузка данных в Kafka - # - transform: batch трансформация - - start = EmptyOperator(task_id="start") - end = EmptyOperator(task_id="end") - - start >> end + 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/docker-compose.yml b/docker-compose.yml index 5da65f9..976d03f 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -4,7 +4,7 @@ x-airflow-env: &airflow-default-env 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:8123/default + AIRFLOW_CONN_CLICKHOUSE_DEFAULT: clickhouse://default:123456@clickhouse:9000/default services: @@ -97,7 +97,7 @@ services: build: context: . dockerfile: Dockerfile.airflow - image: airflow-optimized:2.9.3 + image: airflow-optimized:2.10.5 environment: <<: *airflow-default-env command: > @@ -108,6 +108,7 @@ services: - "8080:8080" volumes: - ./dags:/opt/airflow/dags + - ./sql:/opt/airflow/sql:ro - ./data:/opt/airflow/data networks: - cs_dwh @@ -132,6 +133,7 @@ services: " volumes: - ./dags:/opt/airflow/dags + - ./sql:/opt/airflow/sql:ro - ./data:/opt/airflow/data networks: - cs_dwh @@ -153,13 +155,14 @@ services: <<: *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 777 /opt/airflow/data || true && + 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' && From 4d8f9d42f43940ee181293e4df38bdb5e84c20c8 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sat, 7 Feb 2026 21:53:52 +0300 Subject: [PATCH 10/12] feat(infra): add clickhouse data persistence volume Add persistent volume for ClickHouse to preserve data across container restarts. The volume `clickhouse-data` is mounted to `/var/lib/clickhouse`, ensuring data remains when containers are recreated. --- README.md | 3 +++ docker-compose.yml | 3 ++- 2 files changed, 5 insertions(+), 1 deletion(-) diff --git a/README.md b/README.md index f57eb32..a9bf2fc 100644 --- a/README.md +++ b/README.md @@ -140,6 +140,9 @@ flowchart TB | `FULL=1 make data` | Загрузить полный датасет | | `make transform` | Запустить batch-процесс | +> Примечание: данные ClickHouse теперь сохраняются в Docker volume `clickhouse-data`. +> `docker compose down` сохраняет данные, `docker compose down -v` удаляет все volume (включая ClickHouse). + --- ## 🔗 Ключи данных diff --git a/docker-compose.yml b/docker-compose.yml index 976d03f..8a39681 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -20,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 @@ -235,9 +236,9 @@ networks: driver: bridge volumes: + clickhouse-data: grafana_lib: kafka-data: pgmeta: superset_data: superset_config: - From d9b8a3f909699c3ee7811f182ac8d32a7b3e78c5 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sat, 7 Feb 2026 21:57:15 +0300 Subject: [PATCH 11/12] docs(airflow): update plugin and DAG docs --- AGENTS.md | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index 0b50e4c..20becea 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -34,7 +34,7 @@ - `load_kafka_data.sh` — загрузка в Kafka - `run_batch.sh` — запуск batch-процесса - `airflow/` — конфигурация Airflow: - - `requirements.txt` — зависимости (clickhouse-connect и др.) + - `requirements.txt` — зависимости Airflow/ClickHouse plugin ### Планы и документация (legacy) - `plans/clickhouse_ddl.md` — исходный план (inline DDL, legacy) @@ -62,6 +62,7 @@ - `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`): @@ -84,6 +85,7 @@ - Держать изменения минимальными и по теме задания (инфра, схема, ingest, витрины). - Не коммитить секреты. Если требуется пароль/ключи — использовать `.env` и примеры `.env.example`. - README/планы обновлять вместе с изменениями инфраструктуры/DDL. +- Для спорных или меняющихся API (особенно Airflow/operators/providers) проверять актуальную документацию через `context7` и фиксировать решение в коде/документации. - **Комментарии в коде — на русском языке**: - SQL: заголовочный блок с описанием файла, комментарии к каждому логическому блоку - Bash: шапка с назначением/запуском/требованиями, секции разделены `# -----` @@ -93,7 +95,7 @@ - Kafka ingest: наличие данных в `stg.*` и типизированных строк в `ods.*`. - Мониторинг: доступность `/metrics` у ClickHouse и скрейп в Prometheus. -- **Airflow: `http://localhost:8080` должен показывать UI и DAG `etl_pipeline`.** +- **Airflow: `http://localhost:8080` должен показывать UI и DAG `ddl_init` и `etl_pipeline`.** - BI: витрина `dm.v_events_enriched` должна отвечать за разумное время при фильтре по дате. ## Связанная документация From 225ae8bedba6c1d5329e4efd974c2d1d28cd1c28 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sat, 7 Feb 2026 22:03:13 +0300 Subject: [PATCH 12/12] docs: update README and architecture docs for Airflow orchestration workflow --- README.md | 189 +++++++++++++++++++++++-------------------- docker-compose.yml | 2 +- docs/ARCHITECTURE.md | 21 ++--- 3 files changed, 113 insertions(+), 99 deletions(-) diff --git a/README.md b/README.md index a9bf2fc..02867ff 100644 --- a/README.md +++ b/README.md @@ -4,54 +4,62 @@ [![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 | Назначение | |--------|-----|------------| @@ -64,47 +72,43 @@ docker compose exec clickhouse clickhouse-client \ --- -## 🏗️ Архитектура +## Архитектура (в двух словах) ```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[ETL DAGs] + 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) --- -## 📁 Структура проекта +## Структура проекта ``` . @@ -130,22 +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`. -> `docker compose down` сохраняет данные, `docker compose down -v` удаляет все volume (включая ClickHouse). +Примечания про сохранность данных: +- Данные ClickHouse сохраняются в Docker volume `clickhouse-data`. +- Данные Kafka сохраняются в Docker volume `kafka-data`. +- `docker compose down` сохраняет named volumes, `docker compose down -v` удаляет их (и данные пропадут). --- -## 🔗 Ключи данных +## Ключи данных (как джойним) ```mermaid flowchart LR @@ -181,41 +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` — воронка страниц - ---- - -## 🔮 Развитие проекта - -### ✅ Реализовано -- [x] **Airflow** — оркестрация batch-процесса (инфраструктура готова, DAGs в разработке) - -### 📋 В планах -- [ ] **Инкрементальный batch** — watermark-based загрузка -- [ ] **Материализация витрин** — для тяжёлых агрегаций -- [ ] **DQ мониторинг** — алерты на ошибки парсинга - ---- - -## 📝 Лицензия - -Проект создан для образовательных целей в рамках DE-тестового задания. diff --git a/docker-compose.yml b/docker-compose.yml index 8a39681..c3bad80 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -37,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: diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 381bd0b..6e5c37d 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -589,19 +589,22 @@ INSERT INTO dm.daily_traffic SELECT * FROM dm.v_daily_traffic; Инфраструктура Airflow развёрнута и готова к использованию: ```python -# dags/etl_pipeline_dag.py -with DAG('etl_pipeline'): - 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:8123/default` -- Provider: `clickhouse-connect` (в `airflow/requirements.txt`) +- 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/...`). ---