Merge branch 'feature/airflow-orchestration'

This commit is contained in:
2026-02-07 22:03:41 +03:00
29 changed files with 1067 additions and 146 deletions
+18 -11
View File
@@ -12,6 +12,7 @@
- Kafka (источник событий, 1 JSON message = 1 event) - Kafka (источник событий, 1 JSON message = 1 event)
- ClickHouse (STG → ODS → DDS → DM) - ClickHouse (STG → ODS → DDS → DM)
- **Airflow** (оркестрация ETL-пайплайна)
- Superset (BI поверх витрин) - Superset (BI поверх витрин)
- Prometheus + Grafana (мониторинг) - Prometheus + Grafana (мониторинг)
- простой инструмент/скрипт, который читает `.jsonl` и пишет события в Kafka - простой инструмент/скрипт, который читает `.jsonl` и пишет события в Kafka
@@ -19,19 +20,21 @@
## Ключевые артефакты ## Ключевые артефакты
### Исполняемые файлы (текущая структура) ### Исполняемые файлы (текущая структура)
- `ddl/`SQL для создания объектов БД: - `dags/`Airflow DAGs для оркестрации ETL
- `00_databases.sql`создание БД stg/ods/dds/dm - `sql/`SQL по слоям:
- `10_stg.sql` — STG слой (Kafka Engine + MV) - `sql/ddl/00_databases.sql` — создание БД stg/ods/dds/dm
- `20_ods.sql`ODS слой (типизация + MV для ошибок) - `sql/ddl/stg/10_stg.sql`STG слой (Kafka Engine + MV)
- `30_dds.sql`DDS слой (таблицы для batch-загрузки) - `sql/ddl/ods/20_ods.sql`ODS слой (типизация + MV для ошибок)
- `40_dm.sql` — DM слой (витрины VIEW) - `sql/ddl/dds/30_dds.sql` — DDS слой (таблицы для batch-загрузки)
- `jobs/` — batch-трансформации: - `sql/ddl/dm/40_dm.sql` — DM слой (витрины VIEW)
- `30_dds_refresh.sql` — ODS → DDS (argMax + JOIN) - `sql/dds/30_ods_to_dds.sql` — ODS → DDS (argMax + JOIN)
- `40_dm_refresh.sql` — обновление DQ_summary - `sql/dm/40_dds_to_dm.sql` — обновление DQ_summary
- `scripts/` — скрипты автоматизации: - `scripts/` — скрипты автоматизации:
- `apply_clickhouse_ddl.sh` — применение DDL - `apply_clickhouse_ddl.sh` — применение DDL
- `load_kafka_data.sh` — загрузка в Kafka - `load_kafka_data.sh` — загрузка в Kafka
- `run_batch.sh` — запуск batch-процесса - `run_batch.sh` — запуск batch-процесса
- `airflow/` — конфигурация Airflow:
- `requirements.txt` — зависимости Airflow/ClickHouse plugin
### Планы и документация (legacy) ### Планы и документация (legacy)
- `plans/clickhouse_ddl.md` — исходный план (inline DDL, legacy) - `plans/clickhouse_ddl.md` — исходный план (inline DDL, legacy)
@@ -53,12 +56,13 @@
Базовые команды: Базовые команды:
- `make up` (или `docker compose up -d`) - `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 data` (пересоздаёт топики и заливает небольшой срез данных в Kafka; полный режим — `FULL=1 make data`)
- `make transform` (запускает batch-процесс ODS → DDS → DM) - `make transform` (запускает batch-процесс ODS → DDS → DM)
- `docker compose up -d` - `docker compose up -d`
- `docker compose ps` - `docker compose ps`
- `docker compose logs -f --tail=200 <service>` - `docker compose logs -f --tail=200 <service>`
- `docker compose down` (сохраняет named volumes, включая `clickhouse-data`)
- `docker compose down -v` (удалит volumes; используйте осознанно) - `docker compose down -v` (удалит volumes; используйте осознанно)
Порты (см. `docker-compose.yml`): Порты (см. `docker-compose.yml`):
@@ -67,6 +71,7 @@
- ClickHouse HTTP: `localhost:9123` - ClickHouse HTTP: `localhost:9123`
- Kafka: `localhost:9092` - Kafka: `localhost:9092`
- Kafka UI: `http://localhost:8082` - Kafka UI: `http://localhost:8082`
- **Airflow: `http://localhost:8080` (admin/admin)**
- Prometheus: `http://localhost:9090` - Prometheus: `http://localhost:9090`
- Grafana: `http://localhost:3000` - Grafana: `http://localhost:3000`
@@ -80,15 +85,17 @@
- Держать изменения минимальными и по теме задания (инфра, схема, ingest, витрины). - Держать изменения минимальными и по теме задания (инфра, схема, ingest, витрины).
- Не коммитить секреты. Если требуется пароль/ключи — использовать `.env` и примеры `.env.example`. - Не коммитить секреты. Если требуется пароль/ключи — использовать `.env` и примеры `.env.example`.
- README/планы обновлять вместе с изменениями инфраструктуры/DDL. - README/планы обновлять вместе с изменениями инфраструктуры/DDL.
- Для спорных или меняющихся API (особенно Airflow/operators/providers) проверять актуальную документацию через `context7` и фиксировать решение в коде/документации.
- **Комментарии в коде — на русском языке**: - **Комментарии в коде — на русском языке**:
- SQL: заголовочный блок с описанием файла, комментарии к каждому логическому блоку - SQL: заголовочный блок с описанием файла, комментарии к каждому логическому блоку
- Bash: шапка с назначением/запуском/требованиями, секции разделены `# -----` - Bash: шапка с назначением/запуском/требованиями, секции разделены `# -----`
- См. существующие файлы как пример (`ddl/20_ods.sql`, `jobs/30_dds_refresh.sql`, `scripts/run_batch.sh`) - См. существующие файлы как пример (`sql/ddl/ods/20_ods.sql`, `sql/dds/30_ods_to_dds.sql`, `scripts/run_batch.sh`)
## Быстрые проверки ## Быстрые проверки
- Kafka ingest: наличие данных в `stg.*` и типизированных строк в `ods.*`. - Kafka ingest: наличие данных в `stg.*` и типизированных строк в `ods.*`.
- Мониторинг: доступность `/metrics` у ClickHouse и скрейп в Prometheus. - Мониторинг: доступность `/metrics` у ClickHouse и скрейп в Prometheus.
- **Airflow: `http://localhost:8080` должен показывать UI и DAG `ddl_init` и `etl_pipeline`.**
- BI: витрина `dm.v_events_enriched` должна отвечать за разумное время при фильтре по дате. - BI: витрина `dm.v_events_enriched` должна отвечать за разумное время при фильтре по дате.
## Связанная документация ## Связанная документация
+23
View File
@@ -0,0 +1,23 @@
# Use the official Airflow image as base
FROM apache/airflow:2.10.5
# Set environment variables
ENV AIRFLOW_HOME=/opt/airflow
# Copy the project requirements file
COPY --chown=airflow:0 airflow/requirements.txt ${AIRFLOW_HOME}/requirements.txt
# Install the required Python packages
RUN pip install --no-cache-dir -r ${AIRFLOW_HOME}/requirements.txt
# Create necessary directories
RUN mkdir -p /opt/airflow/dags /opt/airflow/logs /opt/airflow/config /opt/airflow/data
# Set working directory
WORKDIR ${AIRFLOW_HOME}
# Expose the webserver port
EXPOSE 8080
# The image will use Airflow's default entrypoint
# Commands will be passed at runtime
+118 -85
View File
@@ -1,110 +1,130 @@
# ClickHouse Mini DWH для кликстрима # 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) [![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)]() [![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 ```bash
# 1. Поднять инфраструктуру (Kafka + ClickHouse + Superset) # 1) Поднять инфраструктуру
make up make up
# 2. Создать структуру БД # Проверить статусы контейнеров
make ddl docker compose ps
# 3. Загрузить данные (автоматически потекут STG → ODS)
make data # первые 50 строк
# или: FULL=1 make data # полный датасет (1000 строк)
# 4. Подождать 5-10 сек (данные проходят через Kafka)
sleep 10
# 5. Запустить batch-трансформацию (ODS → DDS → DM)
make transform
``` ```
**Проверка:** Дальше основной путь идёт через Airflow (как в задании).
```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
"
# Пример запроса к витрине 1. Открыть Airflow UI: `http://localhost:8080` (admin/admin)
docker compose exec clickhouse clickhouse-client \ 2. Включить (unpause) и запустить `ddl_init` (создаёт базы/таблицы/VIEW в ClickHouse)
--user=default --password=123456 --query="
SELECT * FROM dm.v_utm_effectiveness ORDER BY clicks DESC LIMIT 5 Опционально можно триггернуть 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 | Назначение | | Сервис | URL | Назначение |
|--------|-----|------------| |--------|-----|------------|
| ClickHouse HTTP | http://localhost:9123/play | SQL-запросы | | ClickHouse HTTP | http://localhost:9123/play | SQL-запросы |
| Kafka UI | http://localhost:8082 | Просмотр топиков | | Kafka UI | http://localhost:8082 | Просмотр топиков |
| Airflow | http://localhost:8080 | Оркестрация ETL (admin/admin) |
| Superset | http://localhost:8088 | BI-дашборды | | Superset | http://localhost:8088 | BI-дашборды |
| Prometheus | http://localhost:9090 | Метрики | | Prometheus | http://localhost:9090 | Метрики |
| Grafana | http://localhost:3000 | Визуализация метрик | | Grafana | http://localhost:3000 | Визуализация метрик |
--- ---
## 🏗️ Архитектура ## Архитектура (в двух словах)
```mermaid ```mermaid
flowchart TB flowchart TB
subgraph Sources["📁 JSON файлы"] subgraph Sources["JSONL файлы"]
BE[browser_events.jsonl] BE[browser_events.jsonl]
LE[location_events.jsonl] LE[location_events.jsonl]
DE[device_events.jsonl] DE[device_events.jsonl]
GE[geo_events.jsonl] GE[geo_events.jsonl]
end end
subgraph Kafka["🚀 Kafka"] subgraph Kafka["Kafka"]
KT[Топики] KT[Топики]
end end
subgraph CH["🗄️ ClickHouse"] subgraph CH["ClickHouse"]
STG["STG — сырые JSON"] STG["stg: сырьё + Kafka MV"]
ODS["ODS — типизированные"] ODS["ods: типизация + DQ"]
DDS["DDS — сущности"] DDS["dds: сущности"]
DM["DM — витрины"] DM["dm: витрины (VIEW)"]
end
subgraph Airflow["Airflow"]
DAG[DAG: ddl_init / etl_pipeline]
end end
Sources -->|make data| Kafka -->|MV| STG -->|MV| ODS -->|Batch SQL| DDS -->|VIEW| DM Sources -->|make data| Kafka -->|MV| STG -->|MV| ODS -->|Batch SQL| DDS -->|VIEW| DM
DAG -.->|оркестрация| ODS & DDS & DM
``` ```
**Поток данных:** Особенность задания про "грязные данные": парсинг не валит pipeline, ошибки фиксируются в `ods.*_errors` и в поле `parse_errors`.
1. **STG** — сырые JSON из Kafka (MergeTree)
2. **ODS** — типизированные данные + DQ (ReplacingMergeTree)
3. **DDS** — собранные сущности event + click (Batch SQL)
4. **DM** — витрины для BI (VIEW)
[Подробное описание архитектуры →](./docs/ARCHITECTURE.md) [Подробное описание архитектуры →](./docs/ARCHITECTURE.md)
--- ---
## 📁 Структура проекта ## Структура проекта
``` ```
. .
├── ddl/ # SQL для создания объектов (00_databases → 40_dm) ├── dags/ # Airflow DAGs для оркестрации
├── 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) ├── scripts/ # Автоматизация (apply ddl, load data, run batch)
├── airflow/ # Конфигурация Airflow
│ └── requirements.txt
├── docs/ # Документация ├── docs/ # Документация
│ └── ARCHITECTURE.md # Подробное описание слоёв │ └── ARCHITECTURE.md # Подробное описание слоёв
├── data/ # Исходные JSONL файлы ├── data/ # Исходные JSONL файлы
@@ -114,19 +134,24 @@ flowchart TB
--- ---
## 🛠️ Команды Makefile ## Команды Makefile
| Команда | Описание | | Команда | Описание |
|---------|----------| |---------|----------|
| `make up` | Поднять инфраструктуру | | `make up` | Поднять инфраструктуру |
| `make ddl` | Создать структуру БД | | `make ddl` | Применить DDL в ClickHouse (вне Airflow) |
| `make data` | Загрузить данные в Kafka (50 строк) | | `make data` | Загрузить данные в Kafka (50 строк) |
| `FULL=1 make data` | Загрузить полный датасет | | `FULL=1 make data` | Загрузить полный датасет |
| `make transform` | Запустить batch-процесс | | `make transform` | Запустить batch-процесс (вне Airflow) |
Примечания про сохранность данных:
- Данные ClickHouse сохраняются в Docker volume `clickhouse-data`.
- Данные Kafka сохраняются в Docker volume `kafka-data`.
- `docker compose down` сохраняет named volumes, `docker compose down -v` удаляет их (и данные пропадут).
--- ---
## 🔗 Ключи данных ## Ключи данных (как джойним)
```mermaid ```mermaid
flowchart LR flowchart LR
@@ -162,38 +187,46 @@ flowchart LR
--- ---
## 📚 Документация ## Дашборд в Superset (опционально, но полезно)
1. Открыть `http://localhost:8088`
2. Database -> Add:
- URI: `clickhouse+connect://default:123456@clickhouse:8123/default`
3. Создать datasets из `dm.v_*` (VIEW) и собрать несколько графиков
Идеи графиков под задание:
- Трафик по дням: `dm.v_daily_traffic` (events, uniq_users)
- Эффективность UTM: `dm.v_utm_effectiveness` (clicks, purchases)
- Популярные страницы: `dm.v_top_pages_daily` (pageviews)
- Качество данных: `dm.v_dq_errors_daily` (rows_cnt по error_code)
---
## Частые проблемы
- `etl_pipeline` падает с сообщением про схему: сначала запустите `ddl_init`.
- После `docker compose down -v` схема и данные исчезнут: нужно заново `ddl_init` и `make data`.
- Подключения используют разные протоколы:
- Airflow (ClickHouseOperator) ходит в ClickHouse по native TCP (порт `9000` внутри сети Docker).
- Superset (clickhouse-connect) ходит по HTTP (порт `8123` внутри сети Docker).
---
## Статус проекта
Реализовано (Этап 1):
- DAG `ddl_init`: последовательное применение DDL + проверка схемы.
- DAG `etl_pipeline`: precheck, ожидание данных в ODS, пересчёт DDS/DM, базовые проверки.
- Устойчивость к "грязным" данным: ошибки парсинга сохраняются в ODS, а не валят ingest.
В планах (не требуется для MVP задания):
- DAG `kafka_load` (чистый ingest из `.jsonl` в Kafka средствами Airflow).
- Инкрементальный batch (watermark вместо `full_refresh`).
- DQ мониторинг по расписанию.
---
## Документация
- [Архитектура и слои](./docs/ARCHITECTURE.md) — подробное описание STG/ODS/DDS/DM, ER-диаграммы, обоснование решений - [Архитектура и слои](./docs/ARCHITECTURE.md) — подробное описание STG/ODS/DDS/DM, ER-диаграммы, обоснование решений
- [DE-task.md](./data/DE-task.md) — исходное задание - [DE-task.md](./data/DE-task.md) — исходное задание
---
## 🎯 Дашборд в Superset
1. Открыть http://localhost:8088
2. Database → Add:
- **URI:** `clickhouse+connect://default:123456@clickhouse:8123/default`
3. Datasets → Add from `dm.v_*`
4. Charts & Dashboard
Основные витрины:
- `v_events_enriched` — полное обогащение
- `v_daily_traffic` — агрегация по дням
- `v_utm_effectiveness` — эффективность кампаний
- `v_top_pages_daily` — воронка страниц
---
## 🔮 Развитие проекта
- [ ] **Airflow** — оркестрация batch-процесса
- [ ] **Инкрементальный batch** — watermark-based загрузка
- [ ] **Материализация витрин** — для тяжёлых агрегаций
- [ ] **DQ мониторинг** — алерты на ошибки парсинга
---
## 📝 Лицензия
Проект создан для образовательных целей в рамках DE-тестового задания.
+7
View File
@@ -0,0 +1,7 @@
# Airflow requirements для учебного ETL-проекта
# Metadata DB для Airflow
psycopg2-binary==2.9.9
# ClickHouse operator/hook для DAG'ов
airflow-clickhouse-plugin==1.6.0
View File
View File
Binary file not shown.
Binary file not shown.
+188
View File
@@ -0,0 +1,188 @@
"""
DAG инициализации DDL в ClickHouse.
Учебный формат:
- каждая операция DDL выполняется отдельной SQL-task;
- SQL-файлы вызываются явно по фиксированным путям;
- режим verify_only позволяет прогонять только проверки схемы.
"""
from __future__ import annotations
import re
from datetime import datetime, timedelta
from pathlib import Path
from airflow import DAG
from airflow.exceptions import AirflowException
from airflow.models.param import Param
from airflow.operators.empty import EmptyOperator
from airflow.operators.python import BranchPythonOperator, PythonOperator
from airflow.utils.trigger_rule import TriggerRule
from airflow_clickhouse_plugin.operators.clickhouse import ClickHouseOperator
# -----------------------------------------------------------------------------
# Базовые настройки DAG
# -----------------------------------------------------------------------------
default_args = {
"owner": "airflow",
"depends_on_past": False,
"email_on_failure": False,
"email_on_retry": False,
"retries": 1,
"retry_delay": timedelta(minutes=2),
}
# -----------------------------------------------------------------------------
# SQL-файлы проекта
# -----------------------------------------------------------------------------
SQL_ROOT = Path(__file__).resolve().parents[1] / "sql"
def load_sql_statements(relative_path: str) -> tuple[str, ...]:
"""Читает SQL-файл и делит его на отдельные команды по ';'."""
file_path = SQL_ROOT / relative_path
if not file_path.is_file():
raise AirflowException(f"SQL-файл не найден: {file_path}")
sql_text = file_path.read_text(encoding="utf-8")
statements: list[str] = []
for segment in sql_text.split(";"):
# Убираем блочные и строковые комментарии, чтобы не отправлять "пустые" запросы.
no_block_comments = re.sub(r"/\*.*?\*/", "", segment, flags=re.S)
lines = [line for line in no_block_comments.splitlines() if not line.strip().startswith("--")]
cleaned = "\n".join(lines).strip()
if cleaned:
statements.append(cleaned)
if not statements:
raise AirflowException(f"SQL-файл пустой: {file_path}")
return tuple(statements)
# -----------------------------------------------------------------------------
# SQL-проверки
# -----------------------------------------------------------------------------
SQL_CHECK_CLICKHOUSE = "SELECT 1 AS ok"
SQL_VERIFY_SCHEMA = """
SELECT
(SELECT count() FROM system.tables WHERE database = 'stg' AND name = 'browser_raw') AS stg_browser_raw,
(SELECT count() FROM system.tables WHERE database = 'ods' AND name = 'browser_event') AS ods_browser_event,
(SELECT count() FROM system.tables WHERE database = 'dds' AND name = 'click') AS dds_click,
(SELECT count() FROM system.tables WHERE database = 'dds' AND name = 'event') AS dds_event,
(SELECT count() FROM system.tables WHERE database = 'dm' AND name = 'v_events_enriched') AS dm_v_events_enriched
"""
# -----------------------------------------------------------------------------
# Управляющие функции
# -----------------------------------------------------------------------------
def choose_ddl_mode(**context) -> str:
"""Выбирает ветку выполнения: full DDL или только verify."""
dag_run = context.get("dag_run")
conf = dag_run.conf if dag_run else {}
verify_only = bool(conf.get("verify_only", context["params"]["verify_only"]))
return "skip_ddl" if verify_only else "ddl_00_databases"
def assert_schema_ready(**context) -> None:
"""Проверяет результат финальной SQL-проверки схемы."""
ti = context["ti"]
result = ti.xcom_pull(task_ids="verify_schema_sql")
if not result or not result[0] or len(result[0]) != 5:
raise AirflowException(f"Некорректный результат проверки схемы: {result}")
if any(value == 0 for value in result[0]):
raise AirflowException(
"Схема применена не полностью. Проверьте таблицы/VIEW stg, ods, dds, dm."
)
with DAG(
dag_id="ddl_init",
description="Инициализация схемы ClickHouse (stg/ods/dds/dm)",
default_args=default_args,
schedule=None,
start_date=datetime(2024, 1, 1),
catchup=False,
max_active_runs=1,
is_paused_upon_creation=True,
tags=["ddl", "bootstrap", "clickhouse"],
params={
"verify_only": Param(False, type="boolean"),
},
) as dag:
check_clickhouse = ClickHouseOperator(
task_id="check_clickhouse",
sql=SQL_CHECK_CLICKHOUSE,
clickhouse_conn_id="clickhouse_default",
database="default",
)
choose_mode = BranchPythonOperator(
task_id="choose_mode",
python_callable=choose_ddl_mode,
)
ddl_00_databases = ClickHouseOperator(
task_id="ddl_00_databases",
sql=load_sql_statements("ddl/00_databases.sql"),
clickhouse_conn_id="clickhouse_default",
database="default",
)
ddl_10_stg = ClickHouseOperator(
task_id="ddl_10_stg",
sql=load_sql_statements("ddl/stg/10_stg.sql"),
clickhouse_conn_id="clickhouse_default",
database="default",
)
ddl_20_ods = ClickHouseOperator(
task_id="ddl_20_ods",
sql=load_sql_statements("ddl/ods/20_ods.sql"),
clickhouse_conn_id="clickhouse_default",
database="default",
)
ddl_30_dds = ClickHouseOperator(
task_id="ddl_30_dds",
sql=load_sql_statements("ddl/dds/30_dds.sql"),
clickhouse_conn_id="clickhouse_default",
database="default",
)
ddl_40_dm = ClickHouseOperator(
task_id="ddl_40_dm",
sql=load_sql_statements("ddl/dm/40_dm.sql"),
clickhouse_conn_id="clickhouse_default",
database="default",
)
skip_ddl = EmptyOperator(task_id="skip_ddl")
ddl_complete = EmptyOperator(
task_id="ddl_complete",
trigger_rule=TriggerRule.NONE_FAILED_MIN_ONE_SUCCESS,
)
verify_schema_sql = ClickHouseOperator(
task_id="verify_schema_sql",
sql=SQL_VERIFY_SCHEMA,
clickhouse_conn_id="clickhouse_default",
database="default",
)
verify_schema = PythonOperator(
task_id="verify_schema",
python_callable=assert_schema_ready,
)
check_clickhouse >> choose_mode
choose_mode >> skip_ddl >> ddl_complete
choose_mode >> ddl_00_databases >> ddl_10_stg >> ddl_20_ods >> ddl_30_dds >> ddl_40_dm >> ddl_complete
ddl_complete >> verify_schema_sql >> verify_schema
+282
View File
@@ -0,0 +1,282 @@
"""
DAG ETL-процесса ODS -> DDS -> DM для учебного проекта.
Принципы реализации:
- SQL выполняется явными task на ClickHouseOperator;
- SQL-файлы вызываются по фиксированным путям;
- Python используется только для управляющей логики (branch/wait/assert).
"""
from __future__ import annotations
import re
import time
from datetime import datetime, timedelta
from pathlib import Path
from airflow import DAG
from airflow.exceptions import AirflowException
from airflow.models.param import Param
from airflow.operators.empty import EmptyOperator
from airflow.operators.python import BranchPythonOperator, PythonOperator
from airflow.utils.task_group import TaskGroup
from airflow.utils.trigger_rule import TriggerRule
from airflow_clickhouse_plugin.hooks.clickhouse import ClickHouseHook
from airflow_clickhouse_plugin.operators.clickhouse import ClickHouseOperator
# -----------------------------------------------------------------------------
# Базовые настройки DAG
# -----------------------------------------------------------------------------
default_args = {
"owner": "airflow",
"depends_on_past": False,
"email_on_failure": False,
"email_on_retry": False,
"retries": 1,
"retry_delay": timedelta(minutes=2),
}
# -----------------------------------------------------------------------------
# SQL-файлы проекта
# -----------------------------------------------------------------------------
SQL_ROOT = Path(__file__).resolve().parents[1] / "sql"
def load_sql_statements(relative_path: str) -> tuple[str, ...]:
"""Читает SQL-файл и делит его на отдельные команды по ';'."""
file_path = SQL_ROOT / relative_path
if not file_path.is_file():
raise AirflowException(f"SQL-файл не найден: {file_path}")
sql_text = file_path.read_text(encoding="utf-8")
statements: list[str] = []
for segment in sql_text.split(";"):
# Убираем блочные и строковые комментарии, чтобы не отправлять "пустые" запросы.
no_block_comments = re.sub(r"/\*.*?\*/", "", segment, flags=re.S)
lines = [line for line in no_block_comments.splitlines() if not line.strip().startswith("--")]
cleaned = "\n".join(lines).strip()
if cleaned:
statements.append(cleaned)
if not statements:
raise AirflowException(f"SQL-файл пустой: {file_path}")
return tuple(statements)
# -----------------------------------------------------------------------------
# SQL для проверок и технических шагов
# -----------------------------------------------------------------------------
SQL_CHECK_CLICKHOUSE = "SELECT 1 AS ok"
SQL_CHECK_SCHEMA_READY = """
SELECT
(SELECT count() FROM system.tables WHERE database = 'stg' AND name = 'browser_raw') AS stg_browser_raw,
(SELECT count() FROM system.tables WHERE database = 'ods' AND name = 'browser_event') AS ods_browser_event,
(SELECT count() FROM system.tables WHERE database = 'dds' AND name = 'event') AS dds_event,
(SELECT count() FROM system.tables WHERE database = 'dm' AND name = 'v_events_enriched') AS dm_v_events_enriched
"""
SQL_CHECK_ODS_QUALITY = """
SELECT
count() AS total_rows,
countIf(length(parse_errors) > 0) AS rows_with_errors,
round(if(count() = 0, 0, countIf(length(parse_errors) > 0) / count() * 100), 2) AS error_pct
FROM ods.browser_event
"""
SQL_TRUNCATE_DDS_CLICK = "TRUNCATE TABLE dds.click"
SQL_TRUNCATE_DDS_EVENT = "TRUNCATE TABLE dds.event"
SQL_CHECK_DDS_INTEGRITY = """
SELECT
countIf(click_id IS NOT NULL AND click_id NOT IN (SELECT click_id FROM dds.click)) AS orphan_events
FROM dds.event
"""
SQL_VALIDATE_DM_SUMMARY = "SELECT count() AS dq_rows FROM dm.dq_summary"
# -----------------------------------------------------------------------------
# Управляющие функции
# -----------------------------------------------------------------------------
def assert_schema_ready(**context) -> None:
"""Падает, если DDL не применён полностью."""
ti = context["ti"]
result = ti.xcom_pull(task_ids="precheck.check_schema_ready_sql")
if not result or not result[0] or len(result[0]) != 4:
raise AirflowException(f"Некорректный результат check_schema_ready_sql: {result}")
if any(value == 0 for value in result[0]):
raise AirflowException(
"Схема не готова: сначала запустите DAG ddl_init, затем повторите etl_pipeline."
)
def wait_for_ods_data(**context) -> None:
"""
Ожидает появления строк в ods.browser_event до заданного таймаута.
Таймаут берётся из dag_run.conf.wait_ods_timeout_sec или из params.
"""
dag_run = context.get("dag_run")
conf = dag_run.conf if dag_run else {}
timeout_sec = int(conf.get("wait_ods_timeout_sec", context["params"]["wait_ods_timeout_sec"]))
poll_interval_sec = 10
hook = ClickHouseHook(clickhouse_conn_id="clickhouse_default", database="default")
started = time.monotonic()
while True:
rows = hook.execute("SELECT count() FROM ods.browser_event")
count_rows = int(rows[0][0]) if rows else 0
if count_rows > 0:
return
elapsed = int(time.monotonic() - started)
if elapsed >= timeout_sec:
raise AirflowException(
f"Таймаут ожидания ODS истёк ({timeout_sec} сек). "
"Таблица ods.browser_event всё ещё пуста."
)
time.sleep(poll_interval_sec)
def choose_full_refresh(**context) -> str:
"""Ветвление: делать TRUNCATE DDS или пропустить."""
dag_run = context.get("dag_run")
conf = dag_run.conf if dag_run else {}
full_refresh = bool(conf.get("full_refresh", context["params"]["full_refresh"]))
return "transform.truncate_dds_click" if full_refresh else "transform.skip_truncate"
def assert_dm_summary_not_empty(**context) -> None:
"""Проверяет, что dm.dq_summary заполнена после загрузки."""
ti = context["ti"]
result = ti.xcom_pull(task_ids="transform.validate_dm_summary_sql")
if not result or not result[0] or len(result[0]) != 1:
raise AirflowException(f"Некорректный результат validate_dm_summary_sql: {result}")
dq_rows = int(result[0][0])
if dq_rows <= 0:
raise AirflowException("dm.dq_summary пуста после load_dm_summary.")
with DAG(
dag_id="etl_pipeline",
description="ETL ODS -> DDS -> DM для demo-проекта",
default_args=default_args,
schedule=None,
start_date=datetime(2024, 1, 1),
catchup=False,
max_active_runs=1,
is_paused_upon_creation=True,
tags=["etl", "clickhouse", "demo"],
params={
"full_refresh": Param(True, type="boolean"),
"wait_ods_timeout_sec": Param(600, type="integer", minimum=30),
},
) as dag:
with TaskGroup(group_id="precheck") as precheck:
check_clickhouse = ClickHouseOperator(
task_id="check_clickhouse",
sql=SQL_CHECK_CLICKHOUSE,
clickhouse_conn_id="clickhouse_default",
database="default",
)
check_schema_ready_sql = ClickHouseOperator(
task_id="check_schema_ready_sql",
sql=SQL_CHECK_SCHEMA_READY,
clickhouse_conn_id="clickhouse_default",
database="default",
)
check_schema_ready = PythonOperator(
task_id="check_schema_ready",
python_callable=assert_schema_ready,
)
check_clickhouse >> check_schema_ready_sql >> check_schema_ready
with TaskGroup(group_id="transform") as transform:
wait_for_ods_data_task = PythonOperator(
task_id="wait_for_ods_data",
python_callable=wait_for_ods_data,
)
check_ods_quality = ClickHouseOperator(
task_id="check_ods_quality",
sql=SQL_CHECK_ODS_QUALITY,
clickhouse_conn_id="clickhouse_default",
database="default",
)
choose_refresh_mode = BranchPythonOperator(
task_id="choose_refresh_mode",
python_callable=choose_full_refresh,
)
truncate_dds_click = ClickHouseOperator(
task_id="truncate_dds_click",
sql=SQL_TRUNCATE_DDS_CLICK,
clickhouse_conn_id="clickhouse_default",
database="default",
)
truncate_dds_event = ClickHouseOperator(
task_id="truncate_dds_event",
sql=SQL_TRUNCATE_DDS_EVENT,
clickhouse_conn_id="clickhouse_default",
database="default",
)
skip_truncate = EmptyOperator(task_id="skip_truncate")
truncate_complete = EmptyOperator(
task_id="truncate_complete",
trigger_rule=TriggerRule.NONE_FAILED_MIN_ONE_SUCCESS,
)
load_dds = ClickHouseOperator(
task_id="load_dds",
sql=load_sql_statements("dds/30_ods_to_dds.sql"),
clickhouse_conn_id="clickhouse_default",
database="default",
)
check_dds_integrity = ClickHouseOperator(
task_id="check_dds_integrity",
sql=SQL_CHECK_DDS_INTEGRITY,
clickhouse_conn_id="clickhouse_default",
database="default",
)
load_dm_summary = ClickHouseOperator(
task_id="load_dm_summary",
sql=load_sql_statements("dm/40_dds_to_dm.sql"),
clickhouse_conn_id="clickhouse_default",
database="default",
)
validate_dm_summary_sql = ClickHouseOperator(
task_id="validate_dm_summary_sql",
sql=SQL_VALIDATE_DM_SUMMARY,
clickhouse_conn_id="clickhouse_default",
database="default",
)
validate_dm_summary = PythonOperator(
task_id="validate_dm_summary",
python_callable=assert_dm_summary_not_empty,
)
wait_for_ods_data_task >> check_ods_quality >> choose_refresh_mode
choose_refresh_mode >> truncate_dds_click >> truncate_dds_event >> truncate_complete
choose_refresh_mode >> skip_truncate >> truncate_complete
truncate_complete >> load_dds >> check_dds_integrity >> load_dm_summary >> validate_dm_summary_sql >> validate_dm_summary
precheck >> transform
Regular → Executable
View File

Before

Width:  |  Height:  |  Size: 526 KiB

After

Width:  |  Height:  |  Size: 526 KiB

Regular → Executable
View File
Regular → Executable
View File
Regular → Executable
View File
Regular → Executable
View File
Regular → Executable
View File
+114 -2
View File
@@ -1,3 +1,12 @@
x-airflow-env: &airflow-default-env
AIRFLOW__CORE__LOAD_EXAMPLES: "False"
AIRFLOW__DATABASE__SQL_ALCHEMY_CONN: postgresql+psycopg2://airflow:airflow@postgres-metadata:5432/airflow
AIRFLOW__CORE__EXECUTOR: LocalExecutor
AIRFLOW__WEBSERVER__SECRET_KEY: ${AIRFLOW_SECRET_KEY:-replace-me-with-random-string}
# ClickHouse connection для ETL
AIRFLOW_CONN_CLICKHOUSE_DEFAULT: clickhouse://default:123456@clickhouse:9000/default
services: services:
clickhouse: clickhouse:
image: clickhouse/clickhouse-server:25.1 image: clickhouse/clickhouse-server:25.1
@@ -11,6 +20,7 @@ services:
networks: networks:
- cs_dwh - cs_dwh
volumes: volumes:
- clickhouse-data:/var/lib/clickhouse
- ./configs/default_user.xml:/etc/clickhouse-server/users.d/default_user.xml - ./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/z_config.xml:/etc/clickhouse-server/config.d/z_config.xml
# - ./configs/macros_ch1.xml:/etc/clickhouse-server/config.d/macros.xml # - ./configs/macros_ch1.xml:/etc/clickhouse-server/config.d/macros.xml
@@ -27,7 +37,7 @@ services:
# container_name: kafka # container_name: kafka
ports: ports:
- 9092:9092 - 9092:9092
# Если хочешь сохранять топики/сообщения между `docker compose down/up` — раскомментируй: # Данные топиков/сообщений сохраняются в named volume `kafka-data` (между `docker compose down/up`).
volumes: volumes:
- kafka-data:/tmp/kraft-combined-logs - kafka-data:/tmp/kraft-combined-logs
environment: environment:
@@ -64,6 +74,107 @@ services:
networks: networks:
- cs_dwh - cs_dwh
# PostgreSQL for Airflow metadata
postgres-metadata:
image: postgres:16
environment:
POSTGRES_USER: airflow
POSTGRES_PASSWORD: airflow
POSTGRES_DB: airflow
ports:
- "5434:5432"
volumes:
- pgmeta:/var/lib/postgresql/data
networks:
- cs_dwh
healthcheck:
test: ["CMD-SHELL", "pg_isready -U airflow -d airflow"]
interval: 5s
timeout: 5s
retries: 20
# Airflow webserver
airflow-webserver:
build:
context: .
dockerfile: Dockerfile.airflow
image: airflow-optimized:2.10.5
environment:
<<: *airflow-default-env
command: >
bash -c "
airflow webserver
"
ports:
- "8080:8080"
volumes:
- ./dags:/opt/airflow/dags
- ./sql:/opt/airflow/sql:ro
- ./data:/opt/airflow/data
networks:
- cs_dwh
depends_on:
airflow-init:
condition: service_completed_successfully
postgres-metadata:
condition: service_healthy
clickhouse:
condition: service_started
# Airflow scheduler
airflow-scheduler:
build:
context: .
dockerfile: Dockerfile.airflow
environment:
<<: *airflow-default-env
command: >
bash -c "
airflow scheduler
"
volumes:
- ./dags:/opt/airflow/dags
- ./sql:/opt/airflow/sql:ro
- ./data:/opt/airflow/data
networks:
- cs_dwh
depends_on:
airflow-init:
condition: service_completed_successfully
postgres-metadata:
condition: service_healthy
clickhouse:
condition: service_started
# Airflow init to create admin user
airflow-init:
build:
context: .
dockerfile: Dockerfile.airflow
user: "0:0"
environment:
<<: *airflow-default-env
volumes:
- ./dags:/opt/airflow/dags
- ./sql:/opt/airflow/sql:ro
- ./data:/opt/airflow/data
networks:
- cs_dwh
command: >
bash -ceuo pipefail "
mkdir -p /opt/airflow/data &&
chmod -R a+rX /opt/airflow/data || true &&
chown -R airflow:0 /opt/airflow/data || true &&
umask 000 &&
su -s /bin/bash airflow -c 'airflow db migrate' &&
su -s /bin/bash airflow -c 'airflow users create --username admin --password admin --firstname Admin --lastname User --role Admin --email admin@example.org' &&
su -s /bin/bash airflow -c 'airflow connections delete clickhouse_default || true' &&
su -s /bin/bash airflow -c \"airflow connections add 'clickhouse_default' --conn-uri \\\"$${AIRFLOW_CONN_CLICKHOUSE_DEFAULT}\\\" --conn-description 'ClickHouse DWH'\"
"
depends_on:
postgres-metadata:
condition: service_healthy
prometheus: prometheus:
image: prom/prometheus:v2.53.4 image: prom/prometheus:v2.53.4
volumes: volumes:
@@ -125,8 +236,9 @@ networks:
driver: bridge driver: bridge
volumes: volumes:
clickhouse-data:
grafana_lib: grafana_lib:
kafka-data: kafka-data:
pgmeta:
superset_data: superset_data:
superset_config: superset_config:
+24 -14
View File
@@ -369,11 +369,11 @@ sequenceDiagram
K-->>User: ✅ Инфраструктура готова K-->>User: ✅ Инфраструктура готова
User->>Make: make ddl User->>Make: make ddl
Make->>CH: ddl/00_databases.sql Make->>CH: sql/ddl/00_databases.sql
Make->>CH: ddl/10_stg.sql (Kafka Engine) Make->>CH: sql/ddl/stg/10_stg.sql (Kafka Engine)
Make->>CH: ddl/20_ods.sql (MV) Make->>CH: sql/ddl/ods/20_ods.sql (MV)
Make->>CH: ddl/30_dds.sql Make->>CH: sql/ddl/dds/30_dds.sql
Make->>CH: ddl/40_dm.sql Make->>CH: sql/ddl/dm/40_dm.sql
CH-->>User: ✅ Структура БД создана CH-->>User: ✅ Структура БД создана
User->>Make: make data User->>Make: make data
@@ -389,10 +389,10 @@ sequenceDiagram
CH-->>User: ✅ Данные в STG/ODS CH-->>User: ✅ Данные в STG/ODS
User->>Make: make transform User->>Make: make transform
Make->>CH: jobs/30_dds_refresh.sql Make->>CH: sql/dds/30_ods_to_dds.sql
CH->>ODS: argMax() — снапшот CH->>ODS: argMax() — снапшот
CH->>DDS: JOIN + INSERT 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->>DM: DQ summary
CH-->>User: ✅ DDS/DM обновлены CH-->>User: ✅ DDS/DM обновлены
``` ```
@@ -586,16 +586,26 @@ INSERT INTO dm.daily_traffic SELECT * FROM dm.v_daily_traffic;
### Airflow-оркестрация ### Airflow-оркестрация
Инфраструктура Airflow развёрнута и готова к использованию:
```python ```python
# dag.py # dags/ddl_init_dag.py и dags/etl_pipeline_dag.py
with DAG('clickhouse_etl'): #
ddl = BashOperator(task_id='ddl', bash_command='make ddl') # Учебный формат:
load = BashOperator(task_id='load', bash_command='make data') # - DDL и трансформации выполняются явными SQL-task через ClickHouseOperator;
transform = BashOperator(task_id='transform', bash_command='make transform') # - SQL-файлы вызываются по фиксированным путям;
# - загрузка данных в Kafka (Этап 1) выполняется через `make data`.
ddl >> load >> transform #
# Основной demo-сценарий:
# ddl_init -> make data -> etl_pipeline
``` ```
**Подключение к ClickHouse:**
- Connection: `clickhouse_default`
- URL: `clickhouse://default:123456@clickhouse:9000/default` (native TCP для Airflow plugin)
- Provider/интеграция: `airflow-clickhouse-plugin``airflow/requirements.txt`), задачи выполняются через `ClickHouseOperator`.
- Примечание: Superset подключается к ClickHouse по HTTP (обычно `clickhouse+connect://...:8123/...`).
--- ---
## Полезные запросы ## Полезные запросы
+234
View File
@@ -0,0 +1,234 @@
# План развития Airflow DAG'ов
## Цель
Перевести оркестрацию ETL на Airflow так, чтобы пайплайн оставался устойчивым к "грязным" данным и соответствовал целям проекта: Kafka → ClickHouse (STG → ODS → DDS → DM) → витрины для BI.
## Что важно учесть в текущем репозитории
- В `AGENTS.md` как quick check ожидается DAG `etl_pipeline`.
- В Airflow-контейнере сейчас нет Kafka CLI, поэтому `kafka-topics.sh` и `kafka-console-producer.sh` из `BashOperator` не используем.
- DDL должен выполняться строго последовательно: `00 -> 10 -> 20 -> 30 -> 40`.
- Файл `sql/dds/30_ods_to_dds.sql` уже включает обе загрузки (`dds.click` и `dds.event`), поэтому в MVP это одна task.
## Архитектура оркестрации
### DAG 1 (обязательный): `ddl_init`
- `schedule`: `None` (только ручной запуск).
- `catchup`: `False`.
- `max_active_runs`: `1`.
- `is_paused_upon_creation`: `True`.
- `tags`: `["ddl", "bootstrap", "clickhouse"]`.
### DAG 2 (обязательный): `kafka_load`
- `schedule`: `None` (ручной/экспериментальный запуск).
- `catchup`: `False`.
- `max_active_runs`: `1`.
- `is_paused_upon_creation`: `True`.
- `tags`: `["kafka", "ingest", "experiments"]`.
- Реализация: этап 2 (после MVP).
### DAG 3 (обязательный): `etl_pipeline`
- `schedule`: `None` (ручной запуск для демо).
- `catchup`: `False`.
- `max_active_runs`: `1`.
- `tags`: `["etl", "clickhouse", "demo"]`.
### DAG 4 (опциональный): `dq_monitor`
- `schedule`: `0 * * * *`.
- `catchup`: `False`.
- `tags`: `["dq", "monitoring"]`.
## Дизайн DAG `ddl_init`
### Params
- `verify_only`: bool, default `false` (прогон только проверок без применения DDL).
### Tasks
| Task ID | Что делает | Источник SQL/реализация |
|---------|------------|--------------------------|
| `check_clickhouse` | Проверка доступности CH (`SELECT 1`) | `PythonOperator` + `clickhouse-connect` |
| `ddl_00_databases` | Создание БД | `sql/ddl/00_databases.sql` |
| `ddl_10_stg` | STG + Kafka Engine + MV | `sql/ddl/stg/10_stg.sql` |
| `ddl_20_ods` | ODS + MV STG→ODS + *_errors | `sql/ddl/ods/20_ods.sql` |
| `ddl_30_dds` | Таблицы DDS | `sql/ddl/dds/30_dds.sql` |
| `ddl_40_dm` | VIEW витрины DM | `sql/ddl/dm/40_dm.sql` |
| `verify_schema` | Проверка ключевых таблиц/VIEW | SQL-check |
Зависимости:
```text
check_clickhouse >> ddl_00_databases >> ddl_10_stg >> ddl_20_ods >> ddl_30_dds >> ddl_40_dm >> verify_schema
```
Примечание:
- `ddl_init` запускается вручную: при первом bootstrap, после `docker compose down -v`, после изменений схемы.
## Дизайн DAG `kafka_load` (отдельный независимый контур)
### Params (через Trigger DAG with config)
- `limit`: int, default `50`.
- `full_load`: bool, default `false`.
- `reset_topics`: bool, default `true`.
- `load_browser`: bool, default `true`.
- `load_location`: bool, default `true`.
- `load_device`: bool, default `true`.
- `load_geo`: bool, default `true`.
### TaskGroup `precheck`
| Task ID | Что делает | Реализация |
|---------|------------|------------|
| `check_kafka` | Проверка доступности Kafka broker | `PythonOperator` + `kafka-python` |
| `check_input_files` | Проверка наличия `data/*_events.jsonl` | `PythonOperator` |
| `validate_load_params` | Валидация параметров загрузки (`limit`, `full_load`, флаги потоков) | `PythonOperator` |
### TaskGroup `ingest`
| Task ID | Что делает | Реализация |
|---------|------------|------------|
| `prepare_topics` | reset/create топиков по `reset_topics` | `PythonOperator` + AdminClient |
| `load_browser_events` | Публикация `browser_events.jsonl` | `PythonOperator` + KafkaProducer |
| `load_location_events` | Публикация `location_events.jsonl` | `PythonOperator` + KafkaProducer |
| `load_device_events` | Публикация `device_events.jsonl` | `PythonOperator` + KafkaProducer |
| `load_geo_events` | Публикация `geo_events.jsonl` | `PythonOperator` + KafkaProducer |
| `verify_publish_counts` | Проверка, что отправлено > 0 сообщений в выбранные потоки | `PythonOperator` (по XCom) |
Зависимости:
```text
precheck >> prepare_topics >> [load_browser_events, load_location_events, load_device_events, load_geo_events] >> verify_publish_counts
```
Примечания:
- DAG намеренно независим от `etl_pipeline`: можно запускать ingest отдельно для экспериментов.
- Авто-триггер `etl_pipeline` не включаем по умолчанию; при необходимости добавляется отдельным параметром позже.
## Дизайн DAG `etl_pipeline`
### Params (через Trigger DAG with config)
- `full_refresh`: bool, default `true`.
- `wait_ods_timeout_sec`: int, default `600`.
### TaskGroup `precheck`
| Task ID | Что делает | Реализация |
|---------|------------|------------|
| `check_clickhouse` | Проверка доступности CH (`SELECT 1`) | `PythonOperator` + `clickhouse-connect` |
| `check_schema_ready` | Проверка, что DDL уже применён (`stg.browser_raw`, `ods.browser_event`, `dds.event`, `dm.v_events_enriched`) | SQL-check, fail fast |
### TaskGroup `transform`
| Task ID | Что делает | Источник SQL |
|---------|------------|--------------|
| `wait_for_ods_data` | Ожидание строк в `ods.browser_event` | SQL-check |
| `check_ods_quality` | Базовые DQ-метрики ODS (ошибки/total) | SQL-check |
| `truncate_dds_click` | Очистка `dds.click` при `full_refresh=true` | inline SQL |
| `truncate_dds_event` | Очистка `dds.event` при `full_refresh=true` | inline SQL |
| `load_dds` | ODS → DDS | `sql/dds/30_ods_to_dds.sql` |
| `check_dds_integrity` | Проверка orphan событий | inline SQL |
| `load_dm_summary` | DDS → DM DQ summary | `sql/dm/40_dds_to_dm.sql` |
| `validate_dm_summary` | Проверка, что `dm.dq_summary` не пуста | SQL-check |
Зависимости:
```text
wait_for_ods_data >> check_ods_quality >> [truncate_dds_click, truncate_dds_event] >> load_dds >> check_dds_integrity >> load_dm_summary >> validate_dm_summary
```
### Итоговая цепочка `etl_pipeline`
```text
precheck >> transform
```
## Взаимодействие DAG'ов
Базовый сценарий:
```text
ddl_init -> kafka_load -> etl_pipeline
```
Экспериментальные сценарии:
- `kafka_load` отдельно: проверить разные наборы/параметры загрузки без запуска transform.
- `etl_pipeline` отдельно: повторно пересчитать DDS/DM по уже загруженным данным.
## Техническая реализация (приземленно)
### ClickHouse в Airflow
- Использовать `clickhouse-connect` напрямую в Python helper.
- Брать параметры подключения из `conn_id = clickhouse_default` через `BaseHook.get_connection`.
### Kafka в Airflow
- Добавить `kafka-python` в `airflow/requirements.txt`.
- Использовать Python-код для:
- reset/create топиков;
- публикации строк из `.jsonl` (`1 строка = 1 message value`).
### Общие helper-функции
- `dags/utils/clickhouse_helpers.py`:
- `execute_sql(sql: str) -> None`
- `execute_sql_file(path: str) -> None`
- `fetch_one(sql: str) -> tuple`
- `dags/utils/kafka_helpers.py`:
- `prepare_topics(reset: bool) -> None`
- `load_jsonl(file_path: str, topic: str, limit: int, full_load: bool) -> int`
- `check_kafka_ready() -> None`
## Структура файлов
```text
dags/
├── __init__.py
├── ddl_init_dag.py # отдельный DAG для DDL (обязателен)
├── kafka_load_dag.py # отдельный DAG для ingest в Kafka (обязателен)
├── etl_pipeline_dag.py # основной DAG ODS -> DDS -> DM (обязателен)
├── dq_monitor_dag.py # опциональный DAG мониторинга
└── utils/
├── __init__.py
├── clickhouse_helpers.py
└── kafka_helpers.py
```
## Этапы внедрения
1. Этап 1 (MVP, обязательно):
- Реализовать `ddl_init` и `etl_pipeline`.
- Для загрузки данных использовать существующий сценарий `make data`.
- Проверить путь `ddl_init -> make data -> etl_pipeline`.
2. Этап 2:
- Реализовать отдельный DAG `kafka_load` на `kafka-python`.
- Перенести загрузку из `make data` в `kafka_load` (функциональный паритет).
- Добавить в `kafka_load` расширенные параметры экспериментов (выбор потоков, сценарии reset/no-reset).
- Добавить опциональный параметр автотриггера `etl_pipeline` (по умолчанию `false`).
3. Этап 3:
- Добавить `dq_monitor` и alert callback (email/Slack/webhook).
## Критерии готовности
- Этап 1:
- В Airflow UI видны DAG `ddl_init` и `etl_pipeline`.
- `etl_pipeline` падает с понятной ошибкой, если схема не применена или ODS пуста.
- После прогона `make data -> etl_pipeline`:
- в `ods.browser_event` есть строки;
- в `dds.click` и `dds.event` есть строки;
- `dm.dq_summary` заполнена.
- Этап 2:
- В Airflow UI дополнительно виден DAG `kafka_load`.
- `kafka_load` с `limit=50` завершает отправку сообщений без падений.
- После прогона `kafka_load -> etl_pipeline` результаты совпадают с `make data -> etl_pipeline`.
- Для всех этапов:
- При повторном запуске `etl_pipeline` с `full_refresh=true` нет неконтролируемых дублей в DDS.
## Минимальные smoke-checks
```bash
# 1) Запуск инфраструктуры
make up
# 2) Открыть Airflow UI
# http://localhost:8080 (admin/admin)
# 3) Один раз запустить ddl_init
# Trigger DAG ddl_init (без config или {"verify_only": false})
# 4) Этап 1: загрузить данные текущим способом
make data
# 5) Запустить etl_pipeline
# {"full_refresh": true}
# 6) Проверка результатов
docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query "SELECT count() FROM ods.browser_event"
docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query "SELECT count() FROM dds.click"
docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query "SELECT count() FROM dds.event"
docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query "SELECT count() FROM dm.dq_summary"
# 7) Этап 2: после реализации kafka_load
# Trigger DAG kafka_load {"limit": 50, "full_load": false, "reset_topics": true}
# Trigger DAG etl_pipeline {"full_refresh": true}
```
## Следующий шаг после MVP
- Разделить `sql/dds/30_ods_to_dds.sql` на два файла и распараллелить `load_dds_click` и `load_dds_event`.
- Перейти с `full_refresh` на watermark-инкремент.
+13 -13
View File
@@ -114,32 +114,32 @@ flowchart LR
## План актуализации DDL (target state репозитория) ## План актуализации DDL (target state репозитория)
Цель: перестать исполнять DDL из markdown и хранить **исполняемые** DDL в отдельных `ddl/*.sql` (по слоям), чтобы: Цель: перестать исполнять DDL из markdown и хранить **исполняемые** DDL в отдельных `sql/*/*.sql` (по слоям), чтобы:
- применять их “тонким раннером” через `clickhouse-client` (через `make ddl`); - применять их “тонким раннером” через `clickhouse-client` (через `make ddl`);
- в будущем легко перенести выполнение в Airflow (1 файл = 1 task, линейные зависимости). - в будущем легко перенести выполнение в Airflow (1 файл = 1 task, линейные зависимости).
Важно: Kafka-объекты STG включаем **по умолчанию** (как часть `ddl/10_stg.sql`). Важно: Kafka-объекты STG включаем **по умолчанию** (как часть `sql/ddl/stg/10_stg.sql`).
### Артефакты DDL (планируемые файлы) ### Артефакты DDL (планируемые файлы)
- `ddl/00_databases.sql` — базы `stg/ods/dds/dm`. - `sql/ddl/00_databases.sql` — базы `stg/ods/dds/dm`.
- `ddl/10_stg.sql` — STG raw (`stg.*_raw`) + Kafka source tables (`ENGINE = Kafka`) + MV `Kafka → STG`. - `sql/ddl/stg/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` для строк с битыми ключами. - `sql/ddl/ods/20_ods.sql` — ODS таблицы типизации + DQ (`parse_errors`) + MV `STG → ODS` + таблицы `ods_*_errors` для строк с битыми ключами.
- `ddl/30_dds.sql` — DDS таблицы (`dds.event`, `dds.click`) **без MV** (только `CREATE TABLE`). - `sql/ddl/dds/30_dds.sql` — DDS таблицы (`dds.event`, `dds.click`) **без MV** (только `CREATE TABLE`).
- `ddl/40_dm.sql` — витрины `VIEW` для Superset (`dm.v_*`). - `sql/ddl/dm/40_dm.sql` — витрины `VIEW` для Superset (`dm.v_*`).
BI-ограничения (ресурсы/пользователь) **не выносим в `ddl/*.sql`**: оставляем это только как текст/пример в этом плане, чтобы не смешивать инфраструктуру доступа с DDL витрин. BI-ограничения (ресурсы/пользователь) **не выносим в `sql/*/*.sql`**: оставляем это только как текст/пример в этом плане, чтобы не смешивать инфраструктуру доступа с DDL витрин.
### Артефакты batch-трансформаций (планируемые файлы) ### Артефакты 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)` (или эквивалент); - получить “последнюю версию” строк по ключам (`event_id`/`click_id`) через `argMax(..., src_ingest_ts)` (или эквивалент);
- выполнить join snapshot’ов и загрузить в `dds.event`/`dds.click` (для демо возможно “full rebuild”; позже — инкрементально). - выполнить join snapshot’ов и загрузить в `dds.event`/`dds.click` (для демо возможно “full rebuild”; позже — инкрементально).
### Исполнение DDL (make сейчас / Airflow потом) ### Исполнение DDL (make сейчас / Airflow потом)
Требования к файлам `ddl/*.sql`: Требования к файлам `sql/*/*.sql`:
- идемпотентность (`IF NOT EXISTS`), чтобы повторные прогоны были безопасны; - идемпотентность (`IF NOT EXISTS`), чтобы повторные прогоны были безопасны;
- строгий порядок исполнения: `00 → 10 → 20 → 30 → 40` (из‑за зависимостей MV); - строгий порядок исполнения: `00 → 10 → 20 → 30 → 40` (из‑за зависимостей MV);
@@ -163,11 +163,11 @@ BI-ограничения (ресурсы/пользователь) **не вы
Текущее “как запускаем” (целевое, для реализации следующим шагом): Текущее “как запускаем” (целевое, для реализации следующим шагом):
- `make ddl` вызывает `scripts/apply_clickhouse_ddl.sh`; - `make ddl` вызывает `scripts/apply_clickhouse_ddl.sh`;
- скрипт прогоняет `ddl/*.sql` по порядку через `clickhouse-client --multiquery` внутри контейнера ClickHouse. - скрипт прогоняет `sql/*/*.sql` по порядку через `clickhouse-client --multiquery` внутри контейнера ClickHouse.
Batch‑трансформации (целевое, для реализации следующим шагом): 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). - в будущем Airflow будет делать то же самое по расписанию (один job‑SQL = один task).
### Параметры окружения (docker compose) ### Параметры окружения (docker compose)
@@ -178,7 +178,7 @@ Batch‑трансформации (целевое, для реализации
--- ---
## Приложение A: текущий inline DDL (legacy; будет вынесен в `ddl/*.sql`) ## Приложение A: текущий inline DDL (legacy; будет вынесен в `sql/*/*.sql`)
### 0) Базы данных ### 0) Базы данных
+25 -15
View File
@@ -3,8 +3,8 @@
# Скрипт применения DDL в ClickHouse # Скрипт применения DDL в ClickHouse
# #
# Назначение: # Назначение:
# Последовательно применяет SQL-файлы из ddl/*.sql в базу ClickHouse. # Последовательно применяет SQL-файлы из sql/ddl/* в базу ClickHouse.
# Файлы применяются в алфавитном порядке (00 → 10 → 20 → 30 → 40). # Порядок фиксированный: 00 → 10 → 20 → 30 → 40.
# #
# Как запускать: # Как запускать:
# make ddl # make ddl
@@ -22,7 +22,7 @@ set -euo pipefail
# Директория со скриптом # Директория со скриптом
SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" 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}" COMPOSE_BIN="${COMPOSE_BIN:-docker compose}"
@@ -31,7 +31,7 @@ CLICKHOUSE_DB="${CLICKHOUSE_DB:-default}"
CLICKHOUSE_USER="${CLICKHOUSE_USER:-default}" CLICKHOUSE_USER="${CLICKHOUSE_USER:-default}"
CLICKHOUSE_PASSWORD="${CLICKHOUSE_PASSWORD:-123456}" CLICKHOUSE_PASSWORD="${CLICKHOUSE_PASSWORD:-123456}"
echo "Применение DDL из ${DDL_DIR}..." echo "Применение DDL из ${SQL_ROOT_DIR}..."
# ----------------------------------------------------------------------------- # -----------------------------------------------------------------------------
# Проверка: ClickHouse запущен? # Проверка: ClickHouse запущен?
@@ -45,18 +45,28 @@ fi
# ----------------------------------------------------------------------------- # -----------------------------------------------------------------------------
# Применение SQL-файлов по порядку # Применение SQL-файлов по порядку
# ----------------------------------------------------------------------------- # -----------------------------------------------------------------------------
# shellcheck disable=SC2044 DDL_FILES=(
for sql_file in "${DDL_DIR}"/*.sql; do "${SQL_ROOT_DIR}/ddl/00_databases.sql"
if [[ -f "$sql_file" ]]; then "${SQL_ROOT_DIR}/ddl/stg/10_stg.sql"
echo "Применение: $(basename "$sql_file")" "${SQL_ROOT_DIR}/ddl/ods/20_ods.sql"
${COMPOSE_BIN} exec -T "${CLICKHOUSE_SERVICE}" clickhouse-client \ "${SQL_ROOT_DIR}/ddl/dds/30_dds.sql"
--user="${CLICKHOUSE_USER}" \ "${SQL_ROOT_DIR}/ddl/dm/40_dm.sql"
--password="${CLICKHOUSE_PASSWORD}" \ )
--database="${CLICKHOUSE_DB}" \
--multiquery \ for sql_file in "${DDL_FILES[@]}"; do
< "$sql_file" if [[ ! -f "$sql_file" ]]; then
echo " ✓ OK" echo "Ошибка: не найден SQL-файл: $sql_file" >&2
exit 1
fi 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 done
echo "" echo ""
+19 -4
View File
@@ -3,7 +3,7 @@
# Скрипт batch-трансформации данных: ODS → DDS → DM # Скрипт batch-трансформации данных: ODS → DDS → DM
# #
# Назначение: # Назначение:
# Запускает SQL-скрипты из jobs/ для преобразования данных между слоями: # Запускает SQL-скрипты из sql/dds и sql/dm для преобразования данных между слоями:
# 1. ODS → DDS : Сборка сущностей из типизированных данных # 1. ODS → DDS : Сборка сущностей из типизированных данных
# 2. DDS → DM : Обновление сводки по качеству данных (dq_summary) # 2. DDS → DM : Обновление сводки по качеству данных (dq_summary)
# #
@@ -24,7 +24,9 @@ set -euo pipefail
# Директория со скриптом # Директория со скриптом
SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" 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}" COMPOSE_BIN="${COMPOSE_BIN:-docker compose}"
@@ -42,6 +44,19 @@ if ! ${COMPOSE_BIN} ps | grep -q "${CLICKHOUSE_SERVICE}"; then
exit 1 exit 1
fi 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 есть данные? # Проверка: в ODS есть данные?
# ----------------------------------------------------------------------------- # -----------------------------------------------------------------------------
@@ -83,7 +98,7 @@ ${COMPOSE_BIN} exec -T "${CLICKHOUSE_SERVICE}" clickhouse-client \
--user="${CLICKHOUSE_USER}" \ --user="${CLICKHOUSE_USER}" \
--password="${CLICKHOUSE_PASSWORD}" \ --password="${CLICKHOUSE_PASSWORD}" \
--database="${CLICKHOUSE_DB}" \ --database="${CLICKHOUSE_DB}" \
--multiquery < "${JOBS_DIR}/30_dds_refresh.sql" --multiquery < "${DDS_TRANSFORM_SQL}"
echo " ✓ DDS обновлён" echo " ✓ DDS обновлён"
@@ -105,7 +120,7 @@ ${COMPOSE_BIN} exec -T "${CLICKHOUSE_SERVICE}" clickhouse-client \
--user="${CLICKHOUSE_USER}" \ --user="${CLICKHOUSE_USER}" \
--password="${CLICKHOUSE_PASSWORD}" \ --password="${CLICKHOUSE_PASSWORD}" \
--database="${CLICKHOUSE_DB}" \ --database="${CLICKHOUSE_DB}" \
--multiquery < "${JOBS_DIR}/40_dm_refresh.sql" --multiquery < "${DM_TRANSFORM_SQL}"
echo " ✓ DM обновлён" echo " ✓ DM обновлён"
+1 -1
View File
@@ -8,7 +8,7 @@
-- --
-- Загрузка: -- Загрузка:
-- Batch SQL (не MV!) — для согласованности при late arrivals -- Batch SQL (не MV!) — для согласованности при late arrivals
-- См. jobs/30_dds_refresh.sql -- См. sql/dds/30_ods_to_dds.sql
-- --
-- Почему не MV: -- Почему не MV:
-- - MV с JOIN даёт eventual consistency (данные приходят в разное время) -- - MV с JOIN даёт eventual consistency (данные приходят в разное время)
+1 -1
View File
@@ -13,7 +13,7 @@
-- --
-- Для продакшена: -- Для продакшена:
-- - Если тяжёлые агрегации тормозят — материализовать в таблицы -- - Если тяжёлые агрегации тормозят — материализовать в таблицы
-- - См. пример закомментированный в jobs/40_dm_refresh.sql -- - См. пример закомментированный в sql/dm/40_dds_to_dm.sql
-- ============================================================================ -- ============================================================================
-- ---------------------------------------------------------------------------- -- ----------------------------------------------------------------------------