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"