From 6bbb26b9b371f328626401cb9dcc46950cf53e6f Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Fri, 6 Feb 2026 21:58:17 +0300 Subject: [PATCH] feat(infra): implement batch transformation layer and comprehensive documentation MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Add batch ETL pipeline with ODS→DDS→DM transformation jobs and scripts. Create DDL infrastructure with automated database schema application. Update Makefile with transform target for executing batch processes. Rewrite README with complete Russian documentation including architecture diagrams, quick start guide, and data flow visualization. --- Makefile | 5 +- README.md | 210 +++++++++++-- ddl/00_databases.sql | 5 + ddl/10_stg.sql | 144 +++++++++ ddl/20_ods.sql | 241 ++++++++++++++ ddl/30_dds.sql | 54 ++++ ddl/40_dm.sql | 117 +++++++ docs/ARCHITECTURE.md | 539 ++++++++++++++++++++++++++++++++ jobs/30_dds_refresh.sql | 113 +++++++ jobs/40_dm_refresh.sql | 72 +++++ scripts/apply_clickhouse_ddl.sh | 40 +++ scripts/run_batch.sh | 100 ++++++ 12 files changed, 1620 insertions(+), 20 deletions(-) create mode 100644 ddl/00_databases.sql create mode 100644 ddl/10_stg.sql create mode 100644 ddl/20_ods.sql create mode 100644 ddl/30_dds.sql create mode 100644 ddl/40_dm.sql create mode 100644 docs/ARCHITECTURE.md create mode 100644 jobs/30_dds_refresh.sql create mode 100644 jobs/40_dm_refresh.sql create mode 100755 scripts/apply_clickhouse_ddl.sh create mode 100755 scripts/run_batch.sh diff --git a/Makefile b/Makefile index 50aa1dd..edd4608 100644 --- a/Makefile +++ b/Makefile @@ -1,4 +1,4 @@ -.PHONY: up ddl data +.PHONY: up ddl data transform COMPOSE ?= docker compose @@ -10,3 +10,6 @@ ddl: data: bash ./scripts/load_kafka_data.sh + +transform: + bash ./scripts/run_batch.sh diff --git a/README.md b/README.md index e99a3a0..ac8eff7 100644 --- a/README.md +++ b/README.md @@ -1,27 +1,199 @@ -# Kaniskin docker ETL (mini clickstream demo) +# ClickHouse Mini DWH для кликстрима -Мини‑демо аналитического стека: Kafka → ClickHouse (STG → ODS → DDS → DM) + Superset + мониторинг. +[![Stack](https://img.shields.io/badge/stack-Kafka%20%7C%20ClickHouse%20%7C%20Superset-blue)](./docker-compose.yml) +[![Layers](https://img.shields.io/badge/layers-STG%20→%20ODS%20→%20DDS%20→%20DM-green)](./docs/ARCHITECTURE.md) +[![License](https://img.shields.io/badge/license-Educational-orange)]() -## Документация +Многослойное хранилище данных (STG → ODS → DDS → DM) для анализа кликстрима e-commerce. -- Схема слоёв и DDL (в т.ч. Kafka → STG): `plans/clickhouse_ddl.md` -- Runbook (как поднимать/применять DDL/лить данные): `plans/runbook.md` +Данные поступают из Kafka, проходят типизацию и обогащение, формируя витрины для BI-аналитики. -## Что где лежит +> **Соответствие заданию:** Реализован полный цикл Data Engineering: ingestion → хранилище со слоями → регулярный процесс трансформации → витрины для дашборда. -- Исходные данные: `data/*_events.jsonl` -- Скрипты: - - применение DDL: `scripts/apply_clickhouse_ddl.sh` - - загрузка данных в Kafka: `scripts/load_kafka_data.sh` +--- -## Порты сервисов +## 🚀 Быстрый старт -См. `docker-compose.yml`: +```bash +# 1. Поднять инфраструктуру (Kafka + ClickHouse + Superset) +make up -- ClickHouse native: `localhost:8002` -- ClickHouse HTTP: `localhost:9123` -- Kafka: `localhost:9092` -- Kafka UI: `http://localhost:8082` -- Prometheus: `http://localhost:9090` -- Grafana: `http://localhost:3000` -- Superset: `http://localhost:8088` +# 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 +``` + +**Проверка:** +```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 +" + +# Пример запроса к витрине +docker compose exec clickhouse clickhouse-client \ + --user=default --password=123456 --query=" + SELECT * FROM dm.v_utm_effectiveness ORDER BY clicks DESC LIMIT 5 +" +``` + +--- + +## 📊 Доступные сервисы + +| Сервис | URL | Назначение | +|--------|-----|------------| +| ClickHouse HTTP | http://localhost:9123/play | SQL-запросы | +| Kafka UI | http://localhost:8082 | Просмотр топиков | +| Superset | http://localhost:8088 | BI-дашборды | +| Prometheus | http://localhost:9090 | Метрики | +| Grafana | http://localhost:3000 | Визуализация метрик | + +--- + +## 🏗️ Архитектура + +```mermaid +flowchart TB + subgraph Sources["📁 JSON файлы"] + BE[browser_events.jsonl] + LE[location_events.jsonl] + DE[device_events.jsonl] + GE[geo_events.jsonl] + end + + subgraph Kafka["🚀 Kafka"] + KT[Топики] + end + + subgraph CH["🗄️ ClickHouse"] + STG["STG — сырые JSON"] + ODS["ODS — типизированные"] + DDS["DDS — сущности"] + DM["DM — витрины"] + end + + Sources -->|make data| Kafka -->|MV| STG -->|MV| ODS -->|Batch SQL| DDS -->|VIEW| DM +``` + +**Поток данных:** +1. **STG** — сырые JSON из Kafka (MergeTree) +2. **ODS** — типизированные данные + DQ (ReplacingMergeTree) +3. **DDS** — собранные сущности event + click (Batch SQL) +4. **DM** — витрины для BI (VIEW) + +[Подробное описание архитектуры →](./docs/ARCHITECTURE.md) + +--- + +## 📁 Структура проекта + +``` +. +├── ddl/ # SQL для создания объектов (00_databases → 40_dm) +├── jobs/ # Batch-трансформации (ODS→DDS, DDS→DM) +├── scripts/ # Автоматизация (apply ddl, load data, run batch) +├── docs/ # Документация +│ └── ARCHITECTURE.md # Подробное описание слоёв +├── data/ # Исходные JSONL файлы +├── docker-compose.yml +└── Makefile # Команды: up, ddl, data, transform +``` + +--- + +## 🛠️ Команды Makefile + +| Команда | Описание | +|---------|----------| +| `make up` | Поднять инфраструктуру | +| `make ddl` | Создать структуру БД | +| `make data` | Загрузить данные в Kafka (50 строк) | +| `FULL=1 make data` | Загрузить полный датасет | +| `make transform` | Запустить batch-процесс | + +--- + +## 🔗 Ключи данных + +```mermaid +flowchart LR + subgraph Sources["Источники"] + BE["browser_events (event_id, click_id)"] + LE["location_events (event_id)"] + DE["device_events (click_id)"] + GE["geo_events (click_id)"] + end + + subgraph DDS["DDS"] + EV["event (event_id PK)"] + CL["click (click_id PK)"] + end + + subgraph DM["DM"] + V1[v_events_enriched] + V2[v_daily_traffic] + V3[v_utm_effectiveness] + end + + BE -->|event_id| EV + LE -->|event_id| EV + BE -->|click_id| CL + DE -->|click_id| CL + GE -->|click_id| CL + + EV -->|LEFT JOIN click_id| V1 + CL --> V1 + EV --> V2 & V3 + CL --> V2 & V3 +``` + +--- + +## 📚 Документация + +- [Архитектура и слои](./docs/ARCHITECTURE.md) — подробное описание STG/ODS/DDS/DM, ER-диаграммы, обоснование решений +- [DE-task.md](./data/DE-task.md) — исходное задание + +--- + +## 🎯 Дашборд в Superset + +1. Открыть http://localhost:8088 +2. Database → Add: + - **URI:** `clickhouse+connect://default:123456@clickhouse:8123/default` +3. Datasets → Add from `dm.v_*` +4. Charts & Dashboard + +Основные витрины: +- `v_events_enriched` — полное обогащение +- `v_daily_traffic` — агрегация по дням +- `v_utm_effectiveness` — эффективность кампаний +- `v_top_pages_daily` — воронка страниц + +--- + +## 🔮 Развитие проекта + +- [ ] **Airflow** — оркестрация batch-процесса +- [ ] **Инкрементальный batch** — watermark-based загрузка +- [ ] **Материализация витрин** — для тяжёлых агрегаций +- [ ] **DQ мониторинг** — алерты на ошибки парсинга + +--- + +## 📝 Лицензия + +Проект создан для образовательных целей в рамках DE-тестового задания. diff --git a/ddl/00_databases.sql b/ddl/00_databases.sql new file mode 100644 index 0000000..3769321 --- /dev/null +++ b/ddl/00_databases.sql @@ -0,0 +1,5 @@ +-- Databases for layered DWH +CREATE DATABASE IF NOT EXISTS stg; +CREATE DATABASE IF NOT EXISTS ods; +CREATE DATABASE IF NOT EXISTS dds; +CREATE DATABASE IF NOT EXISTS dm; diff --git a/ddl/10_stg.sql b/ddl/10_stg.sql new file mode 100644 index 0000000..315901d --- /dev/null +++ b/ddl/10_stg.sql @@ -0,0 +1,144 @@ +-- STG layer: raw JSON storage + Kafka ingestion + +-- Raw storage tables (target for MV from Kafka) +CREATE TABLE IF NOT EXISTS stg.browser_raw +( + ingest_ts DateTime64(3) DEFAULT now64(3), + kafka_topic LowCardinality(String) DEFAULT '', + kafka_partition Int32 DEFAULT -1, + kafka_offset Int64 DEFAULT -1, + kafka_ts DateTime64(3) DEFAULT ingest_ts, + raw String +) +ENGINE = MergeTree +PARTITION BY toYYYYMM(ingest_ts) +ORDER BY (kafka_topic, kafka_partition, kafka_offset, ingest_ts); + +CREATE TABLE IF NOT EXISTS stg.location_raw +( + ingest_ts DateTime64(3) DEFAULT now64(3), + kafka_topic LowCardinality(String) DEFAULT '', + kafka_partition Int32 DEFAULT -1, + kafka_offset Int64 DEFAULT -1, + kafka_ts DateTime64(3) DEFAULT ingest_ts, + raw String +) +ENGINE = MergeTree +PARTITION BY toYYYYMM(ingest_ts) +ORDER BY (kafka_topic, kafka_partition, kafka_offset, ingest_ts); + +CREATE TABLE IF NOT EXISTS stg.device_raw +( + ingest_ts DateTime64(3) DEFAULT now64(3), + kafka_topic LowCardinality(String) DEFAULT '', + kafka_partition Int32 DEFAULT -1, + kafka_offset Int64 DEFAULT -1, + kafka_ts DateTime64(3) DEFAULT ingest_ts, + raw String +) +ENGINE = MergeTree +PARTITION BY toYYYYMM(ingest_ts) +ORDER BY (kafka_topic, kafka_partition, kafka_offset, ingest_ts); + +CREATE TABLE IF NOT EXISTS stg.geo_raw +( + ingest_ts DateTime64(3) DEFAULT now64(3), + kafka_topic LowCardinality(String) DEFAULT '', + kafka_partition Int32 DEFAULT -1, + kafka_offset Int64 DEFAULT -1, + kafka_ts DateTime64(3) DEFAULT ingest_ts, + raw String +) +ENGINE = MergeTree +PARTITION BY toYYYYMM(ingest_ts) +ORDER BY (kafka_topic, kafka_partition, kafka_offset, ingest_ts); + +-- Kafka source tables (ENGINE = Kafka) +CREATE TABLE IF NOT EXISTS stg.kafka_browser_raw (raw String) +ENGINE = Kafka +SETTINGS + kafka_broker_list = 'kafka:29092', + kafka_topic_list = 'browser_events', + kafka_group_name = 'ch_stg_browser', + kafka_format = 'JSONAsString', + kafka_num_consumers = 1, + kafka_handle_error_mode = 'stream'; + +CREATE TABLE IF NOT EXISTS stg.kafka_location_raw (raw String) +ENGINE = Kafka +SETTINGS + kafka_broker_list = 'kafka:29092', + kafka_topic_list = 'location_events', + kafka_group_name = 'ch_stg_location', + kafka_format = 'JSONAsString', + kafka_num_consumers = 1, + kafka_handle_error_mode = 'stream'; + +CREATE TABLE IF NOT EXISTS stg.kafka_device_raw (raw String) +ENGINE = Kafka +SETTINGS + kafka_broker_list = 'kafka:29092', + kafka_topic_list = 'device_events', + kafka_group_name = 'ch_stg_device', + kafka_format = 'JSONAsString', + kafka_num_consumers = 1, + kafka_handle_error_mode = 'stream'; + +CREATE TABLE IF NOT EXISTS stg.kafka_geo_raw (raw String) +ENGINE = Kafka +SETTINGS + kafka_broker_list = 'kafka:29092', + kafka_topic_list = 'geo_events', + kafka_group_name = 'ch_stg_geo', + kafka_format = 'JSONAsString', + kafka_num_consumers = 1, + kafka_handle_error_mode = 'stream'; + +-- Materialized Views: Kafka → STG +CREATE MATERIALIZED VIEW IF NOT EXISTS stg.mv_kafka_browser_to_stg +TO stg.browser_raw +AS +SELECT + now64(3) AS ingest_ts, + _topic AS kafka_topic, + _partition AS kafka_partition, + _offset AS kafka_offset, + fromUnixTimestamp64Milli(toInt64(_timestamp_ms)) AS kafka_ts, + raw +FROM stg.kafka_browser_raw; + +CREATE MATERIALIZED VIEW IF NOT EXISTS stg.mv_kafka_location_to_stg +TO stg.location_raw +AS +SELECT + now64(3) AS ingest_ts, + _topic AS kafka_topic, + _partition AS kafka_partition, + _offset AS kafka_offset, + fromUnixTimestamp64Milli(toInt64(_timestamp_ms)) AS kafka_ts, + raw +FROM stg.kafka_location_raw; + +CREATE MATERIALIZED VIEW IF NOT EXISTS stg.mv_kafka_device_to_stg +TO stg.device_raw +AS +SELECT + now64(3) AS ingest_ts, + _topic AS kafka_topic, + _partition AS kafka_partition, + _offset AS kafka_offset, + fromUnixTimestamp64Milli(toInt64(_timestamp_ms)) AS kafka_ts, + raw +FROM stg.kafka_device_raw; + +CREATE MATERIALIZED VIEW IF NOT EXISTS stg.mv_kafka_geo_to_stg +TO stg.geo_raw +AS +SELECT + now64(3) AS ingest_ts, + _topic AS kafka_topic, + _partition AS kafka_partition, + _offset AS kafka_offset, + fromUnixTimestamp64Milli(toInt64(_timestamp_ms)) AS kafka_ts, + raw +FROM stg.kafka_geo_raw; diff --git a/ddl/20_ods.sql b/ddl/20_ods.sql new file mode 100644 index 0000000..f19318b --- /dev/null +++ b/ddl/20_ods.sql @@ -0,0 +1,241 @@ +-- ODS layer: typed data + deduplication + DQ + +-- Main ODS tables (valid keys only) + error tables (invalid keys) + +-- ODS: browser_events +CREATE TABLE IF NOT EXISTS ods.browser_event +( + event_id Nullable(UUID), + event_ts Nullable(DateTime64(6)), + event_date Date MATERIALIZED ifNull(toDate(event_ts), toDate(src_ingest_ts)), + event_type LowCardinality(Nullable(String)), + click_id Nullable(UUID), + browser_name LowCardinality(Nullable(String)), + browser_user_agent Nullable(String), + browser_language LowCardinality(Nullable(String)), + src_ingest_ts DateTime64(3), + src_raw String, + parse_errors Array(LowCardinality(String)) +) +ENGINE = ReplacingMergeTree(src_ingest_ts) +PARTITION BY toYYYYMM(event_date) +ORDER BY (event_id) +SETTINGS allow_nullable_key = 1; + +CREATE TABLE IF NOT EXISTS ods.browser_event_errors +( + ingest_ts DateTime64(3), + kafka_topic LowCardinality(String), + kafka_partition Int32, + kafka_offset Int64, + kafka_ts DateTime64(3), + raw String, + error_reason LowCardinality(String) +) +ENGINE = MergeTree +PARTITION BY toYYYYMM(ingest_ts) +ORDER BY (ingest_ts, kafka_topic, kafka_partition, kafka_offset); + +CREATE MATERIALIZED VIEW IF NOT EXISTS stg.mv_browser_raw_to_ods +TO ods.browser_event +AS +WITH + toUUIDOrNull(JSONExtractString(raw, 'event_id')) AS event_id, + parseDateTime64BestEffortOrNull(JSONExtractString(raw, 'event_timestamp'), 6) AS event_ts, + JSONExtractString(raw, 'event_type') AS event_type, + toUUIDOrNull(JSONExtractString(raw, 'click_id')) AS click_id, + JSONExtractString(raw, 'browser_name') AS browser_name, + JSONExtractString(raw, 'browser_user_agent') AS browser_user_agent, + JSONExtractString(raw, 'browser_language') AS browser_language +SELECT + event_id, + event_ts, + event_type, + click_id, + browser_name, + browser_user_agent, + browser_language, + ingest_ts AS src_ingest_ts, + raw AS src_raw, + arrayFilter(x -> x != '', [ + if(event_id IS NULL, 'bad_event_id', ''), + if(event_ts IS NULL, 'bad_event_timestamp', ''), + if(click_id IS NULL, 'bad_click_id', '') + ]) AS parse_errors +FROM stg.browser_raw +WHERE event_id IS NOT NULL; -- Filter NULL keys to error table + +-- ODS: location_events +CREATE TABLE IF NOT EXISTS ods.location_event +( + event_id Nullable(UUID), + page_url Nullable(String), + page_url_path LowCardinality(Nullable(String)), + referer_url Nullable(String), + referer_medium LowCardinality(Nullable(String)), + utm_medium LowCardinality(Nullable(String)), + utm_source LowCardinality(Nullable(String)), + utm_content LowCardinality(Nullable(String)), + utm_campaign LowCardinality(Nullable(String)), + src_ingest_ts DateTime64(3), + src_raw String, + parse_errors Array(LowCardinality(String)) +) +ENGINE = ReplacingMergeTree(src_ingest_ts) +PARTITION BY toYYYYMM(toDate(src_ingest_ts)) +ORDER BY (event_id) +SETTINGS allow_nullable_key = 1; + +CREATE TABLE IF NOT EXISTS ods.location_event_errors +( + ingest_ts DateTime64(3), + kafka_topic LowCardinality(String), + kafka_partition Int32, + kafka_offset Int64, + kafka_ts DateTime64(3), + raw String, + error_reason LowCardinality(String) +) +ENGINE = MergeTree +PARTITION BY toYYYYMM(ingest_ts) +ORDER BY (ingest_ts, kafka_topic, kafka_partition, kafka_offset); + +CREATE MATERIALIZED VIEW IF NOT EXISTS stg.mv_location_raw_to_ods +TO ods.location_event +AS +WITH + toUUIDOrNull(JSONExtractString(raw, 'event_id')) AS event_id +SELECT + event_id, + JSONExtractString(raw, 'page_url') AS page_url, + JSONExtractString(raw, 'page_url_path') AS page_url_path, + JSONExtractString(raw, 'referer_url') AS referer_url, + JSONExtractString(raw, 'referer_medium') AS referer_medium, + JSONExtractString(raw, 'utm_medium') AS utm_medium, + JSONExtractString(raw, 'utm_source') AS utm_source, + JSONExtractString(raw, 'utm_content') AS utm_content, + JSONExtractString(raw, 'utm_campaign') AS utm_campaign, + ingest_ts AS src_ingest_ts, + raw AS src_raw, + arrayFilter(x -> x != '', [ + if(event_id IS NULL, 'bad_event_id', '') + ]) AS parse_errors +FROM stg.location_raw +WHERE event_id IS NOT NULL; + +-- ODS: device_events +CREATE TABLE IF NOT EXISTS ods.device_by_click +( + click_id Nullable(UUID), + os Nullable(String), + os_name LowCardinality(Nullable(String)), + os_timezone LowCardinality(Nullable(String)), + device_type LowCardinality(Nullable(String)), + device_is_mobile Nullable(UInt8), + user_custom_id Nullable(String), + user_domain_id Nullable(UUID), + src_ingest_ts DateTime64(3), + src_raw String, + parse_errors Array(LowCardinality(String)) +) +ENGINE = ReplacingMergeTree(src_ingest_ts) +PARTITION BY toYYYYMM(toDate(src_ingest_ts)) +ORDER BY (click_id) +SETTINGS allow_nullable_key = 1; + +CREATE TABLE IF NOT EXISTS ods.device_by_click_errors +( + ingest_ts DateTime64(3), + kafka_topic LowCardinality(String), + kafka_partition Int32, + kafka_offset Int64, + kafka_ts DateTime64(3), + raw String, + error_reason LowCardinality(String) +) +ENGINE = MergeTree +PARTITION BY toYYYYMM(ingest_ts) +ORDER BY (ingest_ts, kafka_topic, kafka_partition, kafka_offset); + +CREATE MATERIALIZED VIEW IF NOT EXISTS stg.mv_device_raw_to_ods +TO ods.device_by_click +AS +WITH + toUUIDOrNull(JSONExtractString(raw, 'click_id')) AS click_id, + JSONExtract(raw, 'device_is_mobile', 'Nullable(UInt8)') AS device_is_mobile, + toUUIDOrNull(JSONExtractString(raw, 'user_domain_id')) AS user_domain_id +SELECT + click_id, + JSONExtractString(raw, 'os') AS os, + JSONExtractString(raw, 'os_name') AS os_name, + JSONExtractString(raw, 'os_timezone') AS os_timezone, + JSONExtractString(raw, 'device_type') AS device_type, + device_is_mobile, + JSONExtractString(raw, 'user_custom_id') AS user_custom_id, + user_domain_id, + ingest_ts AS src_ingest_ts, + raw AS src_raw, + arrayFilter(x -> x != '', [ + if(click_id IS NULL, 'bad_click_id', ''), + if(user_domain_id IS NULL, 'bad_user_domain_id', '') + ]) AS parse_errors +FROM stg.device_raw +WHERE click_id IS NOT NULL; + +-- ODS: geo_events +CREATE TABLE IF NOT EXISTS ods.geo_by_click +( + click_id Nullable(UUID), + geo_latitude Nullable(Float64), + geo_longitude Nullable(Float64), + geo_country LowCardinality(Nullable(String)), + geo_timezone LowCardinality(Nullable(String)), + geo_region_name Nullable(String), + ip_address Nullable(String), + src_ingest_ts DateTime64(3), + src_raw String, + parse_errors Array(LowCardinality(String)) +) +ENGINE = ReplacingMergeTree(src_ingest_ts) +PARTITION BY toYYYYMM(toDate(src_ingest_ts)) +ORDER BY (click_id) +SETTINGS allow_nullable_key = 1; + +CREATE TABLE IF NOT EXISTS ods.geo_by_click_errors +( + ingest_ts DateTime64(3), + kafka_topic LowCardinality(String), + kafka_partition Int32, + kafka_offset Int64, + kafka_ts DateTime64(3), + raw String, + error_reason LowCardinality(String) +) +ENGINE = MergeTree +PARTITION BY toYYYYMM(ingest_ts) +ORDER BY (ingest_ts, kafka_topic, kafka_partition, kafka_offset); + +CREATE MATERIALIZED VIEW IF NOT EXISTS stg.mv_geo_raw_to_ods +TO ods.geo_by_click +AS +WITH + toUUIDOrNull(JSONExtractString(raw, 'click_id')) AS click_id, + toFloat64OrNull(JSONExtractString(raw, 'geo_latitude')) AS geo_latitude, + toFloat64OrNull(JSONExtractString(raw, 'geo_longitude')) AS geo_longitude +SELECT + click_id, + geo_latitude, + geo_longitude, + JSONExtractString(raw, 'geo_country') AS geo_country, + JSONExtractString(raw, 'geo_timezone') AS geo_timezone, + JSONExtractString(raw, 'geo_region_name') AS geo_region_name, + JSONExtractString(raw, 'ip_address') AS ip_address, + ingest_ts AS src_ingest_ts, + raw AS src_raw, + arrayFilter(x -> x != '', [ + if(click_id IS NULL, 'bad_click_id', ''), + if(geo_latitude IS NULL, 'bad_geo_latitude', ''), + if(geo_longitude IS NULL, 'bad_geo_longitude', '') + ]) AS parse_errors +FROM stg.geo_raw +WHERE click_id IS NOT NULL; diff --git a/ddl/30_dds.sql b/ddl/30_dds.sql new file mode 100644 index 0000000..c4f4ce0 --- /dev/null +++ b/ddl/30_dds.sql @@ -0,0 +1,54 @@ +-- DDS layer: detailed entities (event + click context) +-- Populated via batch SQL (not MV) to handle late arrivals and ensure consistency + +-- DDS: click context (device + geo) +CREATE TABLE IF NOT EXISTS dds.click +( + click_id UUID, + user_domain_id Nullable(UUID), + user_custom_id Nullable(String), + device_type LowCardinality(Nullable(String)), + device_is_mobile Nullable(UInt8), + os_name LowCardinality(Nullable(String)), + os Nullable(String), + os_timezone LowCardinality(Nullable(String)), + geo_country LowCardinality(Nullable(String)), + geo_region_name Nullable(String), + geo_timezone LowCardinality(Nullable(String)), + geo_latitude Nullable(Float64), + geo_longitude Nullable(Float64), + ip_address Nullable(String), + dds_update_ts DateTime64(3), + ods_parse_errors Array(LowCardinality(String)) +) +ENGINE = ReplacingMergeTree(dds_update_ts) +PARTITION BY toYYYYMM(toDate(dds_update_ts)) +ORDER BY (click_id) +SETTINGS allow_nullable_key = 1; + +-- DDS: event (browser + location) +CREATE TABLE IF NOT EXISTS dds.event +( + event_id UUID, + event_ts Nullable(DateTime64(6)), + event_date Date MATERIALIZED ifNull(toDate(event_ts), toDate(dds_update_ts)), + event_type LowCardinality(Nullable(String)), + click_id Nullable(UUID), + page_url Nullable(String), + page_url_path LowCardinality(Nullable(String)), + referer_url Nullable(String), + referer_medium LowCardinality(Nullable(String)), + utm_medium LowCardinality(Nullable(String)), + utm_source LowCardinality(Nullable(String)), + utm_content LowCardinality(Nullable(String)), + utm_campaign LowCardinality(Nullable(String)), + browser_name LowCardinality(Nullable(String)), + browser_user_agent Nullable(String), + browser_language LowCardinality(Nullable(String)), + dds_update_ts DateTime64(3), + ods_parse_errors Array(LowCardinality(String)) +) +ENGINE = ReplacingMergeTree(dds_update_ts) +PARTITION BY toYYYYMM(event_date) +ORDER BY (event_id) +SETTINGS allow_nullable_key = 1; diff --git a/ddl/40_dm.sql b/ddl/40_dm.sql new file mode 100644 index 0000000..a9da3ed --- /dev/null +++ b/ddl/40_dm.sql @@ -0,0 +1,117 @@ +-- DM layer: Data Marts for BI (Superset) +-- Views for enriched data and pre-computed aggregations + +-- Main enriched view: event + click context +CREATE VIEW IF NOT EXISTS dm.v_events_enriched AS +SELECT + e.event_id, + e.event_ts, + e.event_date, + e.event_type, + e.click_id, + e.page_url, + e.page_url_path, + e.referer_url, + e.referer_medium, + e.utm_medium, + e.utm_source, + e.utm_content, + e.utm_campaign, + e.browser_name, + e.browser_language, + e.browser_user_agent, + c.user_domain_id, + c.user_custom_id, + c.device_type, + c.device_is_mobile, + c.os_name, + c.os_timezone, + c.geo_country, + c.geo_region_name, + c.geo_timezone, + c.geo_latitude, + c.geo_longitude, + c.ip_address, + e.dds_update_ts, + arrayConcat(e.ods_parse_errors, c.ods_parse_errors) AS parse_errors +FROM dds.event AS e +LEFT JOIN dds.click AS c ON c.click_id = e.click_id; + +-- Daily traffic aggregation +CREATE VIEW IF NOT EXISTS dm.v_daily_traffic AS +SELECT + event_date, + geo_country, + device_type, + browser_name, + utm_source, + utm_medium, + count() AS events, + uniqExact(click_id) AS uniq_clicks, + uniqExact(user_domain_id) AS uniq_users +FROM dm.v_events_enriched +WHERE event_ts IS NOT NULL +GROUP BY + event_date, + geo_country, + device_type, + browser_name, + utm_source, + utm_medium; + +-- Top pages daily +CREATE VIEW IF NOT EXISTS dm.v_top_pages_daily AS +SELECT + event_date, + page_url_path, + count() AS pageviews, + uniqExact(click_id) AS uniq_clicks +FROM dm.v_events_enriched +WHERE event_type = 'pageview' +GROUP BY event_date, page_url_path; + +-- Data Quality errors daily +CREATE VIEW IF NOT EXISTS dm.v_dq_errors_daily AS +SELECT + event_date, + arrayJoin(parse_errors) AS error_code, + count() AS rows_cnt +FROM dm.v_events_enriched +WHERE length(parse_errors) > 0 +GROUP BY event_date, error_code; + +-- User sessions overview (approximate, by click_id within 30min windows) +CREATE VIEW IF NOT EXISTS dm.v_session_overview AS +SELECT + event_date, + user_domain_id, + click_id, + min(event_ts) AS session_start, + max(event_ts) AS session_end, + date_diff('minute', min(event_ts), max(event_ts)) AS session_duration_min, + count() AS events_count, + arrayDistinct(groupArray(page_url_path)) AS pages_visited, + arrayDistinct(groupArray(geo_country)) AS countries, + arrayDistinct(groupArray(device_type)) AS devices, + groupArraySample(1, 1919)(utm_source)[1] AS utm_source_last, + groupArraySample(1, 1919)(utm_medium)[1] AS utm_medium_last +FROM dm.v_events_enriched +WHERE user_domain_id IS NOT NULL +GROUP BY event_date, user_domain_id, click_id; + +-- UTM effectiveness (for marketing analysis) +CREATE VIEW IF NOT EXISTS dm.v_utm_effectiveness AS +SELECT + event_date, + utm_source, + utm_medium, + utm_campaign, + count() AS clicks, + uniqExact(user_domain_id) AS uniq_users, + uniqExact(click_id) AS uniq_sessions, + countIf(event_type = 'pageview') AS pageviews, + countIf(event_type = 'purchase') AS purchases, + countIf(event_type = 'add_to_cart') AS add_to_carts +FROM dm.v_events_enriched +WHERE utm_source IS NOT NULL OR utm_medium IS NOT NULL +GROUP BY event_date, utm_source, utm_medium, utm_campaign; diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md new file mode 100644 index 0000000..4038c20 --- /dev/null +++ b/docs/ARCHITECTURE.md @@ -0,0 +1,539 @@ +# Архитектура ClickHouse Mini DWH + +Подробное описание слоёв хранилища, потоков данных и принятых решений. + +--- + +## Содержание + +1. [Обзор архитектуры](#обзор-архитектуры) +2. [Слои хранилища](#слои-хранилища) + - [STG (Staging)](#stg-staging) + - [ODS (Operational Data Store)](#ods-operational-data-store) + - [DDS (Detailed Data Store)](#dds-detailed-data-store) + - [DM (Data Marts)](#dm-data-marts) +3. [Поток данных](#поток-данных) +4. [Связи ключей](#связи-ключей) +5. [Принятые решения](#принятые-решения) +6. [Масштабирование](#масштабирование) + +--- + +## Обзор архитектуры + +### Общая схема потока данных + +```mermaid +flowchart TB + subgraph Sources["📁 Источники (JSONL)"] + BE[browser_events.jsonl] + LE[location_events.jsonl] + DE[device_events.jsonl] + GE[geo_events.jsonl] + end + + subgraph Kafka["🚀 Kafka Topics"] + KT1[browser_events] + KT2[location_events] + KT3[device_events] + KT4[geo_events] + end + + subgraph STG["📦 STG (Staging)"] + BR[browser_raw] + LR[location_raw] + DR[device_raw] + GR[geo_raw] + end + + subgraph ODS["🔧 ODS (Operational Data Store)"] + BE_O[browser_event] + LE_O[location_event] + DE_O[device_by_click] + GE_O[geo_by_click] + ERR[error_tables] + end + + subgraph DDS["🎯 DDS (Detailed Data Store)"] + E[event] + C[click] + end + + subgraph DM["📊 DM (Data Marts)"] + VE[v_events_enriched] + VDT[v_daily_traffic] + VTP[v_top_pages_daily] + VUTM[v_utm_effectiveness] + VSE[v_session_overview] + VDQ[v_dq_errors_daily] + end + + BE --> KT1 --> BR --> BE_O --> E --> VE + LE --> KT2 --> LR --> LE_O --> E + DE --> KT3 --> DR --> DE_O --> C --> VE + GE --> KT4 --> GR --> GE_O --> C + + BE_O -.->|ошибки| ERR + E --> VDT & VTP & VUTM & VSE & VDQ + C --> VDT & VTP & VUTM & VSE & VDQ +``` + +### Слои и их назначение + +```mermaid +flowchart LR + subgraph L0["📝 Сырые данные"] + RAW[JSON файлы
1000 строк каждый] + end + + subgraph L1["STG - Staging"] + STG_T["Таблицы *_raw
MergeTree"] + KAFKA["Kafka Engine + MV"] + end + + subgraph L2["ODS - Операционный слой"] + ODS_T["Типизированные таблицы
ReplacingMergeTree"] + DQ["parse_errors
DQ-метрики"] + end + + subgraph L3["DDS - Детальный слой"] + DDS_T["Сущности event + click
Batch SQL"] + end + + subgraph L4["DM - Витрины"] + DM_T["VIEW для BI
Superset/Grafana"] + end + + RAW -->|kafka-console-producer| KAFKA -->|MV| STG_T + STG_T -->|MV| ODS_T + ODS_T -->|argMax + JOIN| DDS_T + DDS_T -->|VIEW| DM_T + ODS_T -.->|ошибки парсинга| DQ +``` + +--- + +## Слои хранилища + +### STG (Staging) + +**Назначение:** Сохранение сырых данных "как есть" для воспроизводимости и отладки. + +| Таблица | Движок | Описание | +|---------|--------|----------| +| `browser_raw` | MergeTree | Сырые события браузера | +| `location_raw` | MergeTree | Сырые данные страниц/UTM | +| `device_raw` | MergeTree | Сырые данные устройств | +| `geo_raw` | MergeTree | Сырые гео-данные | +| `kafka_*_raw` | Kafka | Таблицы-источники Kafka | +| `mv_kafka_*_to_stg` | MV | Поток из Kafka в STG | + +**Структура таблицы:** +```sql +CREATE TABLE stg.browser_raw ( + ingest_ts DateTime64(3), + kafka_topic LowCardinality(String), + kafka_partition Int32, + kafka_offset Int64, + kafka_ts DateTime64(3), + raw String -- ← JSON как есть +) +``` + +**Почему так:** +- **Повторяемость**: если в ODS ошибка — можно перестроить без перезагрузки из Kafka +- **Отладка**: видеть "что реально пришло" vs "что распарсилось" +- **DQ**: невалидные JSON не ломают pipeline + +--- + +### ODS (Operational Data Store) + +**Назначение:** Типизированные данные с дедупликацией и DQ-метриками. + +| Таблица | Ключ | Движок | Описание | +|---------|------|--------|----------| +| `browser_event` | event_id | ReplacingMergeTree(src_ingest_ts) | События браузера | +| `location_event` | event_id | ReplacingMergeTree(src_ingest_ts) | Данные страниц | +| `device_by_click` | click_id | ReplacingMergeTree(src_ingest_ts) | Устройства | +| `geo_by_click` | click_id | ReplacingMergeTree(src_ingest_ts) | Гео-данные | +| `*_errors` | — | MergeTree | Строки с битыми ключами | + +**Пример структуры:** +```sql +CREATE TABLE ods.browser_event ( + event_id Nullable(UUID), + event_ts Nullable(DateTime64(6)), + event_date Date MATERIALIZED ifNull(toDate(event_ts), toDate(src_ingest_ts)), + event_type LowCardinality(Nullable(String)), + click_id Nullable(UUID), + browser_name LowCardinality(Nullable(String)), + src_ingest_ts DateTime64(3), + src_raw String, + parse_errors Array(LowCardinality(String)) +) +ENGINE = ReplacingMergeTree(src_ingest_ts) +ORDER BY (event_id) +SETTINGS allow_nullable_key = 1; +``` + +**DQ-контроль:** +```sql +-- Проверка ошибок парсинга +SELECT + arrayJoin(parse_errors) AS error, + count() AS cnt +FROM ods.browser_event +GROUP BY error; +``` + +**Почему так:** +- **Изоляция источников**: изменения в одном не ломают другие +- **Версионирование**: `ReplacingMergeTree` хранит последнюю версию по `src_ingest_ts` +- **Nullable ключи**: `allow_nullable_key = 1` позволяет хранить "битые" строки + +--- + +### DDS (Detailed Data Store) + +**Назначение:** Собранные сущности для аналитики. + +| Таблица | PK | Источники | JOIN-ключ | +|---------|-----|-----------|-----------| +| `event` | event_id | browser_event + location_event | click_id → click | +| `click` | click_id | device_by_click + geo_by_click | — | + +**Структура:** +```sql +CREATE TABLE dds.event ( + event_id UUID, + event_ts Nullable(DateTime64(6)), + event_type LowCardinality(Nullable(String)), + click_id Nullable(UUID), + page_url Nullable(String), + page_url_path LowCardinality(Nullable(String)), + utm_source LowCardinality(Nullable(String)), + browser_name LowCardinality(Nullable(String)), + -- ... все поля из browser + location + dds_update_ts DateTime64(3), + ods_parse_errors Array(LowCardinality(String)) +); + +CREATE TABLE dds.click ( + click_id UUID, + user_domain_id Nullable(UUID), + device_type LowCardinality(Nullable(String)), + geo_country LowCardinality(Nullable(String)), + -- ... все поля из device + geo + dds_update_ts DateTime64(3), + ods_parse_errors Array(LowCardinality(String)) +); +``` + +**Загрузка (Batch SQL):** +```sql +-- Снапшот ODS через argMax +INSERT INTO dds.click +SELECT d.click_id, d.user_domain_id, ..., g.geo_country, ... +FROM ( + SELECT click_id, argMax(user_domain_id, src_ingest_ts) AS user_domain_id, ... + FROM ods.device_by_click + GROUP BY click_id +) d +LEFT JOIN ( + SELECT click_id, argMax(geo_country, src_ingest_ts) AS geo_country, ... + FROM ods.geo_by_click + GROUP BY click_id +) g ON g.click_id = d.click_id; +``` + +**Почему batch, а не MV:** +- **Согласованность**: MV с JOIN даёт eventual consistency (данные приходят в разное время) +- **Контроль**: Batch SQL можно проверить, откатить, перезапустить +- **Масштабируемость**: легко сделать инкрементальный batch + +--- + +### DM (Data Marts) + +**Назначение:** Витрины для BI-инструментов. + +| Витрина | Назначение | Гранулярность | +|---------|-----------|---------------| +| `v_events_enriched` | Полное обогащение | 1 строка = 1 событие | +| `v_daily_traffic` | Агрегация трафика | День × страна × устройство × браузер × UTM | +| `v_top_pages_daily` | Популярность страниц | День × URL path | +| `v_utm_effectiveness` | Маркетинговая аналитика | День × UTM source/medium/campaign | +| `v_session_overview` | Сессионная аналитика | День × пользователь × сессия | +| `v_dq_errors_daily` | Мониторинг качества | День × тип ошибки | + +**Пример:** +```sql +CREATE VIEW dm.v_events_enriched AS +SELECT + e.*, + c.user_domain_id, + c.device_type, + c.geo_country, + arrayConcat(e.ods_parse_errors, c.ods_parse_errors) AS parse_errors +FROM dds.event AS e +LEFT JOIN dds.click AS c ON c.click_id = e.click_id; +``` + +**Почему VIEW:** +- Для демо: достаточно производительности +- Гибкость: изменения логики не требуют пересоздания таблиц +- Для продакшена: можно материализовать тяжёлые агрегации + +--- + +## Поток данных + +### Sequence диаграмма процесса + +```mermaid +sequenceDiagram + participant User as Пользователь + participant Make as Makefile + participant K as Kafka + participant CH as ClickHouse + participant STG as stg.*_raw + participant ODS as ods.* + participant DDS as dds.* + participant DM as dm.* + + User->>Make: make up + Make->>K: docker compose up kafka + Make->>CH: docker compose up clickhouse + 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 + CH-->>User: ✅ Структура БД создана + + User->>Make: make data + Make->>K: load_kafka_data.sh + K->>K: Создание топиков + loop 4 файла + Make->>K: kafka-console-producer + end + K->>CH: Потребление сообщений + CH->>STG: INSERT через MV + STG->>ODS: INSERT через MV (типизация) + K-->>User: ✅ Данные в Kafka + CH-->>User: ✅ Данные в STG/ODS + + User->>Make: make transform + Make->>CH: jobs/30_dds_refresh.sql + CH->>ODS: argMax() — снапшот + CH->>DDS: JOIN + INSERT + Make->>CH: jobs/40_dm_refresh.sql + CH->>DM: DQ summary + CH-->>User: ✅ DDS/DM обновлены +``` + +--- + +## Связи ключей + +### ER-диаграмма + +```mermaid +erDiagram + BROWSER_EVENT ||--|| LOCATION_EVENT : "event_id" + BROWSER_EVENT ||--o| DEVICE_BY_CLICK : "click_id" + BROWSER_EVENT ||--o| GEO_BY_CLICK : "click_id" + + BROWSER_EVENT { + UUID event_id PK + DateTime event_ts + String event_type + UUID click_id FK + String browser_name + String browser_user_agent + String browser_language + } + + LOCATION_EVENT { + UUID event_id PK + String page_url + String page_url_path + String referer_url + String referer_medium + String utm_source + String utm_medium + String utm_campaign + } + + DEVICE_BY_CLICK { + UUID click_id PK + String os + String os_name + String device_type + UInt8 device_is_mobile + String user_custom_id + UUID user_domain_id + } + + GEO_BY_CLICK { + UUID click_id PK + Float64 geo_latitude + Float64 geo_longitude + String geo_country + String geo_timezone + String geo_region_name + String ip_address + } +``` + +### Сборка DDS-сущностей + +```mermaid +flowchart LR + subgraph ODS_IN["ODS (вход)"] + B[browser_event
event_id + click_id] + L[location_event
event_id] + D[device_by_click
click_id] + G[geo_by_click
click_id] + end + + subgraph BUILD["Batch SQL"] + J1["JOIN по event_id"] + J2["JOIN по click_id"] + end + + subgraph DDS_OUT["DDS (результат)"] + EV[event
всё про событие] + CL[click
всё про сессию] + end + + B --> J1 + L --> J1 --> EV + B -->|click_id| J2 + D --> J2 --> CL + G --> J2 +``` + +**Важно:** Не все `click_id` из events есть в device/geo. Используем `LEFT JOIN`. + +--- + +## Принятые решения + +### Почему `allow_nullable_key = 1`? + +В ClickHouse ключ сортировки не может быть NULL по умолчанию. Но в "грязных" данных ключи могут отсутствовать. + +**Решение:** +1. Включаем `allow_nullable_key = 1` в `ReplacingMergeTree` +2. Фильтруем NULL в MV (`WHERE key IS NOT NULL` → основная таблица) +3. Отдельные `*_errors` таблицы для NULL-ключей + +### Почему `ReplacingMergeTree`? + +- Дедупликация по бизнес-ключу +- Версионирование по timestamp (последняя версия wins) +- Фоновый merge не блокирует чтение + +### Почему batch ODS→DDS? + +| Подход | Плюсы | Минусы | +|--------|-------|--------| +| **MV + JOIN** | Реалтайм | Eventual consistency, дубли при late arrival | +| **Batch (выбрано)** | Согласованность, контроль | Задержка до следующего запуска | + +--- + +## Масштабирование + +### Инкрементальный batch + +Вместо полного `TRUNCATE + INSERT`: + +```sql +-- Добавить watermark +INSERT INTO dds.click +SELECT ... +FROM ods.device_by_click +WHERE src_ingest_ts > ( + SELECT max(dds_update_ts) FROM dds.click +); +``` + +### Материализация витрин + +Для тяжёлых агрегаций: + +```sql +-- Создать таблицу вместо VIEW +CREATE TABLE dm.daily_traffic AS +SELECT * FROM dm.v_daily_traffic; + +-- Пересчёт по расписанию +TRUNCATE TABLE dm.daily_traffic; +INSERT INTO dm.daily_traffic SELECT * FROM dm.v_daily_traffic; +``` + +### Airflow-оркестрация + +```python +# dag.py +with DAG('clickhouse_etl'): + ddl = BashOperator(task_id='ddl', bash_command='make ddl') + load = BashOperator(task_id='load', bash_command='make data') + transform = BashOperator(task_id='transform', bash_command='make transform') + + ddl >> load >> transform +``` + +--- + +## Полезные запросы + +### Проверка слоёв + +```sql +-- Статистика по слоям +SELECT + database, + countDistinct(table) AS tables, + formatReadableQuantity(sum(rows)) AS rows, + formatReadableSize(sum(bytes)) AS size +FROM system.parts +WHERE database IN ('stg', 'ods', 'dds', 'dm') +GROUP BY database +ORDER BY database; +``` + +### DQ-анализ + +```sql +-- Ошибки парсинга по слоям +SELECT + 'ods.browser_event' AS table, + countIf(length(parse_errors) > 0) AS errors, + count() AS total +FROM ods.browser_event +UNION ALL +SELECT + 'dds.event', + countIf(length(ods_parse_errors) > 0), + count() +FROM dds.event; +``` + +### Воронка конверсии + +```sql +SELECT + page_url_path, + pageviews, + uniq_clicks, + round(uniq_clicks * 100.0 / lag(uniq_clicks) OVER (ORDER BY pageviews DESC), 2) AS conversion_pct +FROM dm.v_top_pages_daily +ORDER BY pageviews DESC; +``` diff --git a/jobs/30_dds_refresh.sql b/jobs/30_dds_refresh.sql new file mode 100644 index 0000000..aa4dcb3 --- /dev/null +++ b/jobs/30_dds_refresh.sql @@ -0,0 +1,113 @@ +-- Batch transformation: ODS → DDS +-- Gets latest version from ODS (using argMax) and builds DDS entities +-- Should be run periodically (e.g., every N minutes or via Airflow) + +-- Refresh DDS.click (from device + geo) +-- Strategy: full rebuild for demo; incremental for production +INSERT INTO dds.click +SELECT + d.click_id, + d.user_domain_id, + d.user_custom_id, + d.device_type, + d.device_is_mobile, + d.os_name, + d.os, + d.os_timezone, + g.geo_country, + g.geo_region_name, + g.geo_timezone, + g.geo_latitude, + g.geo_longitude, + g.ip_address, + now64(3) AS dds_update_ts, + arrayFilter(x -> x != '', arrayConcat( + d.parse_errors, + if(g.click_id IS NULL, ['geo_not_found'], []), + if(g.geo_country IS NULL, ['geo_country_missing'], []) + )) AS ods_parse_errors +FROM ( + -- Latest device snapshot from ODS + SELECT + click_id, + argMax(user_domain_id, src_ingest_ts) AS user_domain_id, + argMax(user_custom_id, src_ingest_ts) AS user_custom_id, + argMax(device_type, src_ingest_ts) AS device_type, + argMax(device_is_mobile, src_ingest_ts) AS device_is_mobile, + argMax(os_name, src_ingest_ts) AS os_name, + argMax(os, src_ingest_ts) AS os, + argMax(os_timezone, src_ingest_ts) AS os_timezone, + argMax(parse_errors, src_ingest_ts) AS parse_errors + FROM ods.device_by_click + WHERE click_id IS NOT NULL + GROUP BY click_id +) AS d +LEFT JOIN ( + -- Latest geo snapshot from ODS + SELECT + click_id, + argMax(geo_country, src_ingest_ts) AS geo_country, + argMax(geo_region_name, src_ingest_ts) AS geo_region_name, + argMax(geo_timezone, src_ingest_ts) AS geo_timezone, + argMax(geo_latitude, src_ingest_ts) AS geo_latitude, + argMax(geo_longitude, src_ingest_ts) AS geo_longitude, + argMax(ip_address, src_ingest_ts) AS ip_address + FROM ods.geo_by_click + WHERE click_id IS NOT NULL + GROUP BY click_id +) AS g ON g.click_id = d.click_id; + +-- Refresh DDS.event (from browser + location) +INSERT INTO dds.event +SELECT + b.event_id, + b.event_ts, + b.event_type, + b.click_id, + l.page_url, + l.page_url_path, + l.referer_url, + l.referer_medium, + l.utm_medium, + l.utm_source, + l.utm_content, + l.utm_campaign, + b.browser_name, + b.browser_user_agent, + b.browser_language, + now64(3) AS dds_update_ts, + arrayFilter(x -> x != '', arrayConcat( + b.parse_errors, + if(l.event_id IS NULL, ['location_not_found'], []) + )) AS ods_parse_errors +FROM ( + -- Latest browser snapshot from ODS + SELECT + event_id, + argMax(event_ts, src_ingest_ts) AS event_ts, + argMax(event_type, src_ingest_ts) AS event_type, + argMax(click_id, src_ingest_ts) AS click_id, + argMax(browser_name, src_ingest_ts) AS browser_name, + argMax(browser_user_agent, src_ingest_ts) AS browser_user_agent, + argMax(browser_language, src_ingest_ts) AS browser_language, + argMax(parse_errors, src_ingest_ts) AS parse_errors + FROM ods.browser_event + WHERE event_id IS NOT NULL + GROUP BY event_id +) AS b +LEFT JOIN ( + -- Latest location snapshot from ODS + SELECT + event_id, + argMax(page_url, src_ingest_ts) AS page_url, + argMax(page_url_path, src_ingest_ts) AS page_url_path, + argMax(referer_url, src_ingest_ts) AS referer_url, + argMax(referer_medium, src_ingest_ts) AS referer_medium, + argMax(utm_medium, src_ingest_ts) AS utm_medium, + argMax(utm_source, src_ingest_ts) AS utm_source, + argMax(utm_content, src_ingest_ts) AS utm_content, + argMax(utm_campaign, src_ingest_ts) AS utm_campaign + FROM ods.location_event + WHERE event_id IS NOT NULL + GROUP BY event_id +) AS l ON l.event_id = b.event_id; diff --git a/jobs/40_dm_refresh.sql b/jobs/40_dm_refresh.sql new file mode 100644 index 0000000..9637286 --- /dev/null +++ b/jobs/40_dm_refresh.sql @@ -0,0 +1,72 @@ +-- Batch transformation: DDS → DM materialized tables +-- For demo we use VIEWs mainly, but here we can materialize heavy aggregations + +-- Materialized daily traffic (if needed for performance) +-- Uncomment if VIEW dm.v_daily_traffic becomes too slow +/* +CREATE TABLE IF NOT EXISTS dm.daily_traffic_mart +( + event_date Date, + geo_country LowCardinality(Nullable(String)), + device_type LowCardinality(Nullable(String)), + browser_name LowCardinality(Nullable(String)), + utm_source LowCardinality(Nullable(String)), + utm_medium LowCardinality(Nullable(String)), + events UInt64, + uniq_clicks UInt64, + uniq_users UInt64 +) +ENGINE = ReplacingMergeTree(event_date) +PARTITION BY toYYYYMM(event_date) +ORDER BY (event_date, geo_country, device_type, browser_name, utm_source, utm_medium); + +TRUNCATE TABLE dm.daily_traffic_mart; + +INSERT INTO dm.daily_traffic_mart +SELECT * FROM dm.v_daily_traffic; +*/ + +-- Data Quality summary table (always fresh) +CREATE TABLE IF NOT EXISTS dm.dq_summary +( + check_date Date, + layer LowCardinality(String), + table_name LowCardinality(String), + check_name LowCardinality(String), + check_value UInt64 +) +ENGINE = MergeTree +PARTITION BY toYYYYMM(check_date) +ORDER BY (check_date, layer, table_name, check_name); + +-- Truncate and refill DQ summary +INSERT INTO dm.dq_summary +SELECT + today() AS check_date, + 'stg' AS layer, + 'browser_raw' AS table_name, + 'total_rows' AS check_name, + count() AS check_value +FROM stg.browser_raw +UNION ALL +SELECT today(), 'stg', 'location_raw', 'total_rows', count() FROM stg.location_raw +UNION ALL +SELECT today(), 'stg', 'device_raw', 'total_rows', count() FROM stg.device_raw +UNION ALL +SELECT today(), 'stg', 'geo_raw', 'total_rows', count() FROM stg.geo_raw +UNION ALL +SELECT today(), 'ods', 'browser_event', 'total_rows', count() FROM ods.browser_event +UNION ALL +SELECT today(), 'ods', 'browser_event', 'rows_with_errors', count() FROM ods.browser_event WHERE length(parse_errors) > 0 +UNION ALL +SELECT today(), 'ods', 'location_event', 'total_rows', count() FROM ods.location_event +UNION ALL +SELECT today(), 'ods', 'device_by_click', 'total_rows', count() FROM ods.device_by_click +UNION ALL +SELECT today(), 'ods', 'geo_by_click', 'total_rows', count() FROM ods.geo_by_click +UNION ALL +SELECT today(), 'dds', 'event', 'total_rows', count() FROM dds.event +UNION ALL +SELECT today(), 'dds', 'click', 'total_rows', count() FROM dds.click +UNION ALL +SELECT today(), 'dds', 'event_without_click', 'orphan_events', count() FROM dds.event WHERE click_id IS NOT NULL AND click_id NOT IN (SELECT click_id FROM dds.click); diff --git a/scripts/apply_clickhouse_ddl.sh b/scripts/apply_clickhouse_ddl.sh new file mode 100755 index 0000000..a97eedc --- /dev/null +++ b/scripts/apply_clickhouse_ddl.sh @@ -0,0 +1,40 @@ +#!/usr/bin/env bash +set -euo pipefail + +# Apply ClickHouse DDL files in order +# Usage: make ddl +# or: bash scripts/apply_clickhouse_ddl.sh + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +DDL_DIR="${SCRIPT_DIR}/../ddl" + +COMPOSE_BIN="${COMPOSE_BIN:-docker compose}" +CLICKHOUSE_SERVICE="${CLICKHOUSE_SERVICE:-clickhouse}" +CLICKHOUSE_DB="${CLICKHOUSE_DB:-default}" +CLICKHOUSE_USER="${CLICKHOUSE_USER:-default}" +CLICKHOUSE_PASSWORD="${CLICKHOUSE_PASSWORD:-123456}" + +echo "Applying ClickHouse DDL from ${DDL_DIR}..." + +# Check if clickhouse service is running +if ! ${COMPOSE_BIN} ps | grep -q "${CLICKHOUSE_SERVICE}"; then + echo "Error: ClickHouse service '${CLICKHOUSE_SERVICE}' is not running." + echo "Run 'make up' first to start the services." + exit 1 +fi + +# Apply DDL files in order (00 -> 10 -> 20 -> 30 -> 40) +for sql_file in "${DDL_DIR}"/*.sql; do + if [[ -f "$sql_file" ]]; then + echo "Applying: $(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" + fi +done + +echo "DDL applied successfully." diff --git a/scripts/run_batch.sh b/scripts/run_batch.sh new file mode 100755 index 0000000..b2a516a --- /dev/null +++ b/scripts/run_batch.sh @@ -0,0 +1,100 @@ +#!/usr/bin/env bash +set -euo pipefail + +# Run batch transformations: ODS → DDS → DM +# Usage: make transform +# or: bash scripts/run_batch.sh + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +JOBS_DIR="${SCRIPT_DIR}/../jobs" + +COMPOSE_BIN="${COMPOSE_BIN:-docker compose}" +CLICKHOUSE_SERVICE="${CLICKHOUSE_SERVICE:-clickhouse}" +CLICKHOUSE_DB="${CLICKHOUSE_DB:-default}" +CLICKHOUSE_USER="${CLICKHOUSE_USER:-default}" +CLICKHOUSE_PASSWORD="${CLICKHOUSE_PASSWORD:-123456}" + +# Check if clickhouse service is running +if ! ${COMPOSE_BIN} ps | grep -q "${CLICKHOUSE_SERVICE}"; then + echo "Error: ClickHouse service '${CLICKHOUSE_SERVICE}' is not running." + echo "Run 'make up' first to start the services." + exit 1 +fi + +# Check if ODS has data +ODS_COUNT=$(${COMPOSE_BIN} exec -T "${CLICKHOUSE_SERVICE}" clickhouse-client \ + --user="${CLICKHOUSE_USER}" \ + --password="${CLICKHOUSE_PASSWORD}" \ + --database="${CLICKHOUSE_DB}" \ + --query="SELECT count() FROM ods.browser_event" 2>/dev/null || echo "0") + +if [[ "${ODS_COUNT}" == "0" ]]; then + echo "Warning: ODS.browser_event is empty." + echo "Run 'make data' first to load data into Kafka → STG → ODS." + exit 1 +fi + +echo "Found ${ODS_COUNT} rows in ODS.browser_event" +echo "" + +# Step 1: Refresh DDS (truncate + reload for demo) +echo "Step 1: Refreshing DDS layer (ODS → DDS)..." +${COMPOSE_BIN} exec -T "${CLICKHOUSE_SERVICE}" clickhouse-client \ + --user="${CLICKHOUSE_USER}" \ + --password="${CLICKHOUSE_PASSWORD}" \ + --database="${CLICKHOUSE_DB}" \ + --query="TRUNCATE TABLE dds.click" 2>/dev/null || true +${COMPOSE_BIN} exec -T "${CLICKHOUSE_SERVICE}" clickhouse-client \ + --user="${CLICKHOUSE_USER}" \ + --password="${CLICKHOUSE_PASSWORD}" \ + --database="${CLICKHOUSE_DB}" \ + --query="TRUNCATE TABLE dds.event" 2>/dev/null || true + +${COMPOSE_BIN} exec -T "${CLICKHOUSE_SERVICE}" clickhouse-client \ + --user="${CLICKHOUSE_USER}" \ + --password="${CLICKHOUSE_PASSWORD}" \ + --database="${CLICKHOUSE_DB}" \ + --multiquery < "${JOBS_DIR}/30_dds_refresh.sql" + +echo " ✓ DDS refreshed" + +# Show DDS stats +echo "" +echo "DDS statistics:" +${COMPOSE_BIN} exec -T "${CLICKHOUSE_SERVICE}" clickhouse-client \ + --user="${CLICKHOUSE_USER}" \ + --password="${CLICKHOUSE_PASSWORD}" \ + --database="${CLICKHOUSE_DB}" \ + --query="SELECT 'dds.click' AS table, count() AS rows FROM dds.click UNION ALL SELECT 'dds.event', count() FROM dds.event FORMAT PrettyCompact" + +# Step 2: Refresh DM (DQ summary) +echo "" +echo "Step 2: Refreshing DM layer (DQ summary)..." +${COMPOSE_BIN} exec -T "${CLICKHOUSE_SERVICE}" clickhouse-client \ + --user="${CLICKHOUSE_USER}" \ + --password="${CLICKHOUSE_PASSWORD}" \ + --database="${CLICKHOUSE_DB}" \ + --multiquery < "${JOBS_DIR}/40_dm_refresh.sql" + +echo " ✓ DM refreshed" + +# Show DM stats +echo "" +echo "Data Quality summary:" +${COMPOSE_BIN} exec -T "${CLICKHOUSE_SERVICE}" clickhouse-client \ + --user="${CLICKHOUSE_USER}" \ + --password="${CLICKHOUSE_PASSWORD}" \ + --database="${CLICKHOUSE_DB}" \ + --query="SELECT * FROM dm.dq_summary ORDER BY layer, table_name, check_name FORMAT PrettyCompact" + +echo "" +echo "Batch transformation complete!" +echo "" +echo "Available data marts:" +echo " - dm.v_events_enriched : Main enriched events view" +echo " - dm.v_daily_traffic : Daily aggregation by dimensions" +echo " - dm.v_top_pages_daily : Top pages by day" +echo " - dm.v_dq_errors_daily : Data quality errors" +echo " - dm.v_session_overview : Session-level metrics" +echo " - dm.v_utm_effectiveness : UTM campaign performance" +echo " - dm.dq_summary : Layer statistics"