diff --git a/ddl/00_databases.sql b/ddl/00_databases.sql index 3769321..6722a58 100644 --- a/ddl/00_databases.sql +++ b/ddl/00_databases.sql @@ -1,4 +1,12 @@ --- Databases for layered DWH +-- ============================================================================ +-- Создание баз данных для слоёв хранилища +-- ============================================================================ +-- stg — Staging: сырые данные из Kafka +-- ods — Operational Data Store: типизированные данные с DQ +-- dds — Detailed Data Store: детальные сущности для аналитики +-- dm — Data Marts: витрины для BI +-- ============================================================================ + CREATE DATABASE IF NOT EXISTS stg; CREATE DATABASE IF NOT EXISTS ods; CREATE DATABASE IF NOT EXISTS dds; diff --git a/ddl/10_stg.sql b/ddl/10_stg.sql index 315901d..cc8a821 100644 --- a/ddl/10_stg.sql +++ b/ddl/10_stg.sql @@ -1,14 +1,29 @@ --- STG layer: raw JSON storage + Kafka ingestion +-- ============================================================================ +-- Слой STG (Staging) — сырые JSON + интеграция с Kafka +-- ============================================================================ +-- Назначение: +-- - Хранение сырых данных "как есть" из Kafka +-- - Метаданные доставки (topic, partition, offset, timestamp) +-- - Источник для отладки и восстановления +-- +-- Поток данных: +-- Kafka → kafka_*_raw (ENGINE=Kafka) → MV → *_raw (MergeTree) +-- ============================================================================ --- Raw storage tables (target for MV from Kafka) +-- ---------------------------------------------------------------------------- +-- Таблицы для сырых данных (MergeTree) +-- ---------------------------------------------------------------------------- +-- Хранят JSON "как есть" + метаданные Kafka +-- Партиционируем по месяцу ingest_ts для эффективной очистки старых данных +-- ---------------------------------------------------------------------------- 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 + ingest_ts DateTime64(3) DEFAULT now64(3), -- Время вставки в ClickHouse + kafka_topic LowCardinality(String) DEFAULT '', -- Топик Kafka + kafka_partition Int32 DEFAULT -1, -- Партиция + kafka_offset Int64 DEFAULT -1, -- Смещение (гарантирует уникальность) + kafka_ts DateTime64(3) DEFAULT ingest_ts, -- Время из Kafka (если есть) + raw String -- JSON как строка ) ENGINE = MergeTree PARTITION BY toYYYYMM(ingest_ts) @@ -53,16 +68,21 @@ ENGINE = MergeTree PARTITION BY toYYYYMM(ingest_ts) ORDER BY (kafka_topic, kafka_partition, kafka_offset, ingest_ts); --- Kafka source tables (ENGINE = Kafka) +-- ---------------------------------------------------------------------------- +-- Таблицы-источники Kafka (ENGINE = Kafka) +-- ---------------------------------------------------------------------------- +-- Читают данные из топиков Kafka в реальном времени +-- Не хранят данные постоянно, а "потребляют" их при SELECT/MV +-- ---------------------------------------------------------------------------- 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'; + kafka_broker_list = 'kafka:29092', -- Адрес брокера (внутри Docker-сети) + kafka_topic_list = 'browser_events', -- Имя топика + kafka_group_name = 'ch_stg_browser', -- Группа консьюмеров (для контроля offset'ов) + kafka_format = 'JSONAsString', -- Читаем весь JSON как строку + kafka_num_consumers = 1, -- Количество консьюмеров (для параллелизма) + kafka_handle_error_mode = 'stream'; -- Ошибки не прерывают чтение CREATE TABLE IF NOT EXISTS stg.kafka_location_raw (raw String) ENGINE = Kafka @@ -94,7 +114,12 @@ SETTINGS kafka_num_consumers = 1, kafka_handle_error_mode = 'stream'; +-- ---------------------------------------------------------------------------- -- Materialized Views: Kafka → STG +-- ---------------------------------------------------------------------------- +-- Автоматически перекладывают данные из Kafka-таблиц в MergeTree +-- Работают в реальном времени: сообщение из Kafka → сразу в *_raw +-- ---------------------------------------------------------------------------- CREATE MATERIALIZED VIEW IF NOT EXISTS stg.mv_kafka_browser_to_stg TO stg.browser_raw AS diff --git a/ddl/20_ods.sql b/ddl/20_ods.sql index 0b2c0c6..83bea0d 100644 --- a/ddl/20_ods.sql +++ b/ddl/20_ods.sql @@ -1,41 +1,71 @@ --- ODS layer: typed data + deduplication + DQ +-- ============================================================================ +-- Слой ODS (Operational Data Store) — типизированные данные + дедупликация + DQ +-- ============================================================================ +-- Назначение: +-- - Типизация данных из STG (String → UUID, DateTime, etc.) +-- - Дедупликация через ReplacingMergeTree (последняя версия по src_ingest_ts) +-- - Контроль качества: массив parse_errors для "грязных" данных +-- - Разделение: валидные строки → основная таблица, ошибки → *_errors +-- +-- Поток данных: +-- STG (*_raw) → MV → ODS (основная таблица + error_tables) +-- ============================================================================ --- Main ODS tables (valid keys only) + error tables (invalid keys) +-- ============================================================================ +-- BROWSER EVENTS +-- ============================================================================ --- ODS: browser_events +-- ---------------------------------------------------------------------------- +-- Основная таблица: валидные строки (event_id IS NOT NULL) +-- ---------------------------------------------------------------------------- +-- ReplacingMergeTree: при мердже оставляет строку с максимальным src_ingest_ts +-- allow_nullable_key = 1: разрешаем NULL в ключе (ClickHouse по умолчанию запрещает) +-- ---------------------------------------------------------------------------- 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)) + event_id Nullable(UUID), -- UUID события (ключ) + event_ts Nullable(DateTime64(6)), -- Время события из JSON + event_date Date MATERIALIZED ifNull(toDate(event_ts), toDate(src_ingest_ts)), -- Партиция + event_type LowCardinality(Nullable(String)), -- Тип события (pageview, click и т.д.) + click_id Nullable(UUID), -- Связь с click-контекстом + browser_name LowCardinality(Nullable(String)), -- Chrome, Firefox и т.д. + browser_user_agent Nullable(String), -- User-Agent строка + browser_language LowCardinality(Nullable(String)), -- Язык браузера + src_ingest_ts DateTime64(3), -- Время загрузки в ODS (версия для ReplacingMergeTree) + src_raw String, -- Исходный JSON для аудита + parse_errors Array(LowCardinality(String)) -- Ошибки парсинга (если есть) ) -ENGINE = ReplacingMergeTree(src_ingest_ts) -PARTITION BY toYYYYMM(event_date) -ORDER BY (event_id) -SETTINGS allow_nullable_key = 1; +ENGINE = ReplacingMergeTree(src_ingest_ts) -- Движок дедупликации по версии +PARTITION BY toYYYYMM(event_date) -- Партиции по месяцу для быстрой очистки +ORDER BY (event_id) -- Ключ сортировки (и дедупликации) +SETTINGS allow_nullable_key = 1; -- Разрешаем NULL в ключе (для "битых" данных) +-- ---------------------------------------------------------------------------- +-- Таблица ошибок: строки с невалидными ключами (event_id IS NULL) +-- ---------------------------------------------------------------------------- +-- Сохраняем полную информацию для анализа проблем с данными +-- MergeTree без Replacing: сохраняем все ошибки (не дедуплицируем) +-- ---------------------------------------------------------------------------- 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) + ingest_ts DateTime64(3), -- Время вставки в ClickHouse + kafka_topic LowCardinality(String), -- Топик Kafka (для отслеживания источника) + kafka_partition Int32, -- Партиция Kafka + kafka_offset Int64, -- Смещение Kafka (уникальный идентификатор сообщения) + kafka_ts DateTime64(3), -- Время из Kafka + raw String, -- Исходный JSON + error_reason LowCardinality(String) -- Описание ошибки (что именно не распарсилось) ) ENGINE = MergeTree PARTITION BY toYYYYMM(ingest_ts) ORDER BY (ingest_ts, kafka_topic, kafka_partition, kafka_offset); +-- ---------------------------------------------------------------------------- +-- MV: STG → ODS (основная таблица) +-- ---------------------------------------------------------------------------- +-- Фильтруем только валидные строки: WHERE event_id IS NOT NULL +-- Парсим JSON, типизируем поля, собираем ошибки в массив parse_errors +-- ---------------------------------------------------------------------------- CREATE MATERIALIZED VIEW IF NOT EXISTS stg.mv_browser_raw_to_ods TO ods.browser_event AS @@ -57,14 +87,21 @@ SELECT 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 +WHERE event_id IS NOT NULL; -- Только валидные строки (NULL → в error_tables) +-- ---------------------------------------------------------------------------- +-- MV: STG → ODS (таблица ошибок) +-- ---------------------------------------------------------------------------- +-- Перенаправляем строки с ошибками парсинга в отдельную таблицу +-- Это позволяет не терять данные и анализировать проблемы +-- ---------------------------------------------------------------------------- CREATE MATERIALIZED VIEW IF NOT EXISTS stg.mv_browser_raw_to_ods_errors TO ods.browser_event_errors AS @@ -86,31 +123,34 @@ SELECT raw, arrayStringConcat(parse_errors, ',') AS error_reason FROM stg.browser_raw -WHERE length(parse_errors) > 0 +WHERE length(parse_errors) > 0 -- Есть хотя бы одна ошибка AND ( event_id IS NULL OR event_ts IS NULL OR click_id IS NULL ); --- ODS: location_events +-- ============================================================================ +-- 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)), + page_url Nullable(String), -- Полный URL страницы + page_url_path LowCardinality(Nullable(String)), -- Путь (/home, /product и т.д.) + referer_url Nullable(String), -- Откуда пришёл пользователь + referer_medium LowCardinality(Nullable(String)), -- Тип referer (internal, search и т.д.) + utm_medium LowCardinality(Nullable(String)), -- UTM medium (cpc, organic и т.д.) + utm_source LowCardinality(Nullable(String)), -- UTM source (google, mailchimp и т.д.) + utm_content LowCardinality(Nullable(String)), -- UTM content (ad_1, ad_2 и т.д.) + utm_campaign LowCardinality(Nullable(String)), -- UTM campaign (campaign_1 и т.д.) 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)) +PARTITION BY toYYYYMM(toDate(src_ingest_ts)) -- Партиция по времени загрузки (нет event_date) ORDER BY (event_id) SETTINGS allow_nullable_key = 1; @@ -171,17 +211,20 @@ FROM stg.location_raw WHERE length(parse_errors) > 0 AND event_id IS NULL; --- ODS: device_events +-- ============================================================================ +-- 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), + os Nullable(String), -- Полное название ОС + os_name LowCardinality(Nullable(String)), -- Короткое название (Windows, iOS и т.д.) + os_timezone LowCardinality(Nullable(String)), -- Таймзона пользователя + device_type LowCardinality(Nullable(String)), -- Mobile, Computer, Tablet + device_is_mobile Nullable(UInt8), -- 1 = мобильное, 0 = десктоп + user_custom_id Nullable(String), -- Email или username + user_domain_id Nullable(UUID), -- UUID пользователя в системе src_ingest_ts DateTime64(3), src_raw String, parse_errors Array(LowCardinality(String)) @@ -255,16 +298,19 @@ WHERE length(parse_errors) > 0 OR user_domain_id IS NULL ); --- ODS: geo_events +-- ============================================================================ +-- 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), + geo_latitude Nullable(Float64), -- Широта + geo_longitude Nullable(Float64), -- Долгота + geo_country LowCardinality(Nullable(String)), -- Код страны (RU, US и т.д.) + geo_timezone LowCardinality(Nullable(String)), -- Таймзона + geo_region_name Nullable(String), -- Название региона/города + ip_address Nullable(String), -- IP адрес src_ingest_ts DateTime64(3), src_raw String, parse_errors Array(LowCardinality(String)) diff --git a/ddl/30_dds.sql b/ddl/30_dds.sql index c4f4ce0..fc60469 100644 --- a/ddl/30_dds.sql +++ b/ddl/30_dds.sql @@ -1,54 +1,79 @@ --- DDS layer: detailed entities (event + click context) --- Populated via batch SQL (not MV) to handle late arrivals and ensure consistency +-- ============================================================================ +-- Слой DDS (Detailed Data Store) — детальные сущности для аналитики +-- ============================================================================ +-- Назначение: +-- - Собранные "чистые" сущности из ODS для JOIN'ов и аналитики +-- - click: объединяет device + geo (контекст сессии пользователя) +-- - event: объединяет browser + location (контекст события) +-- +-- Загрузка: +-- Batch SQL (не MV!) — для согласованности при late arrivals +-- См. jobs/30_dds_refresh.sql +-- +-- Почему не MV: +-- - MV с JOIN даёт eventual consistency (данные приходят в разное время) +-- - Batch позволяет сделать снапшот через argMax и корректно джойнить +-- - Контроль: можно проверить SQL, откатить, перезапустить +-- ============================================================================ --- DDS: click context (device + geo) +-- ---------------------------------------------------------------------------- +-- Сущность: dds.click (контекст сессии пользователя) +-- ---------------------------------------------------------------------------- +-- Объединяет данные из ods.device_by_click + ods.geo_by_click +-- Связь с event через click_id (LEFT JOIN, не все события имеют 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)) + click_id UUID, -- UUID сессии/клика (PK) + user_domain_id Nullable(UUID), -- UUID пользователя в системе + user_custom_id Nullable(String), -- Email/username + device_type LowCardinality(Nullable(String)), -- Mobile, Computer, Tablet + device_is_mobile Nullable(UInt8), -- 1 = мобильное, 0 = десктоп + os_name LowCardinality(Nullable(String)), -- Windows, iOS, Android и т.д. + os Nullable(String), -- Полное название ОС + os_timezone LowCardinality(Nullable(String)), -- Таймзона пользователя + geo_country LowCardinality(Nullable(String)), -- Код страны (RU, US) + geo_region_name Nullable(String), -- Регион/город + geo_timezone LowCardinality(Nullable(String)), -- Таймзона по гео + geo_latitude Nullable(Float64), -- Широта + geo_longitude Nullable(Float64), -- Долгота + ip_address Nullable(String), -- IP адрес + dds_update_ts DateTime64(3), -- Время загрузки в DDS (версия) + ods_parse_errors Array(LowCardinality(String)) -- Ошибки из ODS + метки о пропусках ) -ENGINE = ReplacingMergeTree(dds_update_ts) -PARTITION BY toYYYYMM(toDate(dds_update_ts)) -ORDER BY (click_id) -SETTINGS allow_nullable_key = 1; +ENGINE = ReplacingMergeTree(dds_update_ts) -- Дедупликация по версии +PARTITION BY toYYYYMM(toDate(dds_update_ts)) -- Партиция по месяцу загрузки +ORDER BY (click_id) -- Ключ сортировки +SETTINGS allow_nullable_key = 1; -- Разрешаем NULL (на всякий случай) --- DDS: event (browser + location) +-- ---------------------------------------------------------------------------- +-- Сущность: dds.event (контекст события) +-- ---------------------------------------------------------------------------- +-- Объединяет данные из ods.browser_event + ods.location_event +-- Связь с click через click_id (может быть NULL, если нет device/geo) +-- ---------------------------------------------------------------------------- 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)) + event_id UUID, -- UUID события (PK) + event_ts Nullable(DateTime64(6)), -- Время события + event_date Date MATERIALIZED ifNull(toDate(event_ts), toDate(dds_update_ts)), -- Дата для партиций + event_type LowCardinality(Nullable(String)), -- pageview, click, purchase и т.д. + click_id Nullable(UUID), -- Связь с dds.click (может быть NULL) + page_url Nullable(String), -- Полный URL + page_url_path LowCardinality(Nullable(String)), -- Путь (/home, /product) + referer_url Nullable(String), -- Откуда пришёл + referer_medium LowCardinality(Nullable(String)), -- Тип referer + utm_medium LowCardinality(Nullable(String)), -- UTM medium + utm_source LowCardinality(Nullable(String)), -- UTM source + utm_content LowCardinality(Nullable(String)), -- UTM content + utm_campaign LowCardinality(Nullable(String)), -- UTM campaign + browser_name LowCardinality(Nullable(String)), -- Chrome, Firefox + browser_user_agent Nullable(String), -- User-Agent + browser_language LowCardinality(Nullable(String)), -- Язык браузера + dds_update_ts DateTime64(3), -- Время загрузки в DDS (версия) + ods_parse_errors Array(LowCardinality(String)) -- Ошибки из ODS ) -ENGINE = ReplacingMergeTree(dds_update_ts) -PARTITION BY toYYYYMM(event_date) -ORDER BY (event_id) -SETTINGS allow_nullable_key = 1; +ENGINE = ReplacingMergeTree(dds_update_ts) -- Дедупликация по версии +PARTITION BY toYYYYMM(event_date) -- Партиция по дате события (важно для фильтров) +ORDER BY (event_id) -- Ключ сортировки +SETTINGS allow_nullable_key = 1; -- Разрешаем NULL diff --git a/ddl/40_dm.sql b/ddl/40_dm.sql index a9da3ed..ad985d3 100644 --- a/ddl/40_dm.sql +++ b/ddl/40_dm.sql @@ -1,7 +1,28 @@ --- DM layer: Data Marts for BI (Superset) --- Views for enriched data and pre-computed aggregations +-- ============================================================================ +-- Слой DM (Data Marts) — витрины для BI (Superset/Grafana) +-- ============================================================================ +-- Назначение: +-- - Представления (VIEW) для удобного доступа к данным из BI-инструментов +-- - Обогащение: соединяем event + click через LEFT JOIN +-- - Агрегации: готовые GROUP BY для частых запросов +-- +-- Почему VIEW: +-- - Гибкость: меняем логику без пересоздания таблиц +-- - Нет дублирования данных (храним только в DDS) +-- - Для демо: производительность достаточная +-- +-- Для продакшена: +-- - Если тяжёлые агрегации тормозят — материализовать в таблицы +-- - См. пример закомментированный в jobs/40_dm_refresh.sql +-- ============================================================================ --- Main enriched view: event + click context +-- ---------------------------------------------------------------------------- +-- Витрина: полное обогащение событий (event + click) +-- ---------------------------------------------------------------------------- +-- Соединяет dds.event и dds.click через click_id +-- LEFT JOIN: не все события имеют device/geo контекст +-- Используется как основа для других витрин +-- ---------------------------------------------------------------------------- CREATE VIEW IF NOT EXISTS dm.v_events_enriched AS SELECT e.event_id, @@ -9,6 +30,7 @@ SELECT e.event_date, e.event_type, e.click_id, + -- Поля из location (через event) e.page_url, e.page_url_path, e.referer_url, @@ -17,9 +39,11 @@ SELECT e.utm_source, e.utm_content, e.utm_campaign, + -- Поля из browser (через event) e.browser_name, e.browser_language, e.browser_user_agent, + -- Поля из click (device + geo), могут быть NULL c.user_domain_id, c.user_custom_id, c.device_type, @@ -32,12 +56,19 @@ SELECT 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 +-- ---------------------------------------------------------------------------- +-- Витрина: агрегация трафика по дням и измерениям +-- ---------------------------------------------------------------------------- +-- Используется для анализа посещаемости +-- Гранулярность: дата × страна × устройство × браузер × UTM +-- ---------------------------------------------------------------------------- CREATE VIEW IF NOT EXISTS dm.v_daily_traffic AS SELECT event_date, @@ -46,11 +77,11 @@ SELECT browser_name, utm_source, utm_medium, - count() AS events, - uniqExact(click_id) AS uniq_clicks, - uniqExact(user_domain_id) AS uniq_users + 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 +WHERE event_ts IS NOT NULL -- Фильтруем битые timestamp GROUP BY event_date, geo_country, @@ -59,59 +90,80 @@ GROUP BY 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 + page_url_path, -- Путь URL (/home, /product) + count() AS pageviews, -- Количество просмотров + uniqExact(click_id) AS uniq_clicks -- Уникальные сессии FROM dm.v_events_enriched -WHERE event_type = 'pageview' +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 + arrayJoin(parse_errors) AS error_code, -- Разворачиваем массив ошибок + count() AS rows_cnt -- Количество строк с этой ошибкой FROM dm.v_events_enriched -WHERE length(parse_errors) > 0 +WHERE length(parse_errors) > 0 -- Только строки с ошибками GROUP BY event_date, error_code; --- User sessions overview (approximate, by click_id within 30min windows) +-- ---------------------------------------------------------------------------- +-- Витрина: обзор сессий пользователей +-- ---------------------------------------------------------------------------- +-- Группировка по click_id (сессия) и user_domain_id (пользователь) +-- Показывает длительность сессии, посещённые страницы, устройства +-- ---------------------------------------------------------------------------- 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, + 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, -- Устройства (если менялось) + -- Берём последние UTM-метки сессии (для атрибуции) 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 +WHERE user_domain_id IS NOT NULL -- Только идентифицированные пользователи GROUP BY event_date, user_domain_id, click_id; --- UTM effectiveness (for marketing analysis) +-- ---------------------------------------------------------------------------- +-- Витрина: эффективность UTM-кампаний (маркетинговая аналитика) +-- ---------------------------------------------------------------------------- +-- Показывает какие каналы (utm_source/medium) приносят трафик +-- Отдельно считаем pageviews, purchases, add_to_carts +-- ---------------------------------------------------------------------------- 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 + 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 +WHERE utm_source IS NOT NULL OR utm_medium IS NOT NULL -- Только с UTM-метками GROUP BY event_date, utm_source, utm_medium, utm_campaign; diff --git a/jobs/30_dds_refresh.sql b/jobs/30_dds_refresh.sql index 1597035..20921d4 100644 --- a/jobs/30_dds_refresh.sql +++ b/jobs/30_dds_refresh.sql @@ -1,9 +1,26 @@ --- 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) +-- ============================================================================ +-- Batch-трансформация: ODS → DDS +-- ============================================================================ +-- Что делает: +-- Собирает "чистые" сущности из типизированных данных ODS для аналитики. +-- Использует argMax() для получения последней версии строк по ключу. +-- +-- Когда запускать: +-- - Периодически (например, каждые N минут) +-- - Через Airflow по расписанию +-- - Вручную после загрузки новых данных +-- +-- Стратегия: +-- Сейчас: полная перезагрузка (TRUNCATE + INSERT) — проще для демо +-- В продакшене: инкрементальная загрузка по watermark (src_ingest_ts) +-- ============================================================================ --- Refresh DDS.click (from device + geo) --- Strategy: full rebuild for demo; incremental for production +-- ---------------------------------------------------------------------------- +-- Сущность: dds.click (объединяет device + geo) +-- ---------------------------------------------------------------------------- +-- Проблема: device и geo события приходят независимо (не все click_id есть в обоих источниках) +-- Решение: UNION всех click_id + LEFT JOIN (обрабатываем geo-only и device-only) +-- ---------------------------------------------------------------------------- INSERT INTO dds.click SELECT c.click_id, @@ -20,15 +37,19 @@ SELECT g.geo_latitude, g.geo_longitude, g.ip_address, - now64(3) AS dds_update_ts, + now64(3) AS dds_update_ts, -- Время загрузки в DDS для версионирования + -- Собираем все ошибки парсинга из источников + отслеживаем пропущенные данные arrayFilter(x -> x != '', arrayConcat( ifNull(d.parse_errors, []), - if(d.click_id IS NULL, ['device_not_found'], []), - if(g.click_id IS NULL, ['geo_not_found'], []), - if(g.geo_country IS NULL, ['geo_country_missing'], []) + if(d.click_id IS NULL, ['device_not_found'], []), -- Нет данных об устройстве + if(g.click_id IS NULL, ['geo_not_found'], []), -- Нет гео-данных + if(g.geo_country IS NULL, ['geo_country_missing'], []) -- Гео есть, но страна не определена )) AS ods_parse_errors FROM ( - -- Union of all click ids to keep geo-only/device-only arrivals + -- Шаг 1: Собираем ВСЕ уникальные click_id из обоих источников + -- Это позволяет обработать случаи, когда: + -- - есть device, но нет geo (geo-only клики) + -- - есть geo, но нет device (device-only клики) SELECT click_id FROM ( SELECT assumeNotNull(click_id) AS click_id @@ -45,8 +66,9 @@ FROM ( GROUP BY click_id ) ) AS c +-- Шаг 2: LEFT JOIN с device (может не быть данных) LEFT JOIN ( - -- Latest device snapshot from ODS + -- Снапшот device: последняя версия каждого click_id по времени загрузки SELECT assumeNotNull(click_id) AS click_id, argMax(user_domain_id, src_ingest_ts) AS user_domain_id, @@ -61,8 +83,9 @@ LEFT JOIN ( WHERE click_id IS NOT NULL GROUP BY click_id ) AS d ON d.click_id = c.click_id +-- Шаг 3: LEFT JOIN с geo (может не быть данных) LEFT JOIN ( - -- Latest geo snapshot from ODS + -- Снапшот geo: последняя версия каждого click_id по времени загрузки SELECT assumeNotNull(click_id) AS click_id, argMax(geo_country, src_ingest_ts) AS geo_country, @@ -76,13 +99,19 @@ LEFT JOIN ( GROUP BY click_id ) AS g ON g.click_id = c.click_id; --- Refresh DDS.event (from browser + location) +-- ---------------------------------------------------------------------------- +-- Сущность: dds.event (объединяет browser + location) +-- ---------------------------------------------------------------------------- +-- Логика проще: event_id связывает browser и location 1:1 +-- Если location нет — это тоже полезная информация (ошибка или пропуск) +-- ---------------------------------------------------------------------------- INSERT INTO dds.event SELECT b.event_id, b.event_ts, b.event_type, b.click_id, + -- Поля из location (могут быть NULL, если location не пришёл) l.page_url, l.page_url_path, l.referer_url, @@ -91,16 +120,18 @@ SELECT l.utm_source, l.utm_content, l.utm_campaign, + -- Поля из browser b.browser_name, b.browser_user_agent, b.browser_language, now64(3) AS dds_update_ts, + -- Собираем ошибки парсинга + отмечаем, если location не найден 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 + -- Снапшот browser: последняя версия каждого event_id SELECT event_id, argMax(event_ts, src_ingest_ts) AS event_ts, @@ -111,11 +142,12 @@ FROM ( 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 + WHERE event_id IS NOT NULL -- Фильтруем битые ключи (они в error_tables) GROUP BY event_id ) AS b +-- LEFT JOIN с location (может не быть данных для некоторых событий) LEFT JOIN ( - -- Latest location snapshot from ODS + -- Снапшот location: последняя версия каждого event_id SELECT event_id, argMax(page_url, src_ingest_ts) AS page_url, diff --git a/jobs/40_dm_refresh.sql b/jobs/40_dm_refresh.sql index d2f0c42..0edece8 100644 --- a/jobs/40_dm_refresh.sql +++ b/jobs/40_dm_refresh.sql @@ -1,47 +1,64 @@ --- Batch transformation: DDS → DM materialized tables --- For demo we use VIEWs mainly, but here we can materialize heavy aggregations +-- ============================================================================ +-- Batch-трансформация: DDS → DM (Data Quality summary) +-- ============================================================================ +-- Что делает: +-- Собирает статистику по всем слоям (stg/ods/dds) для мониторинга качества данных. +-- Позволяет быстро проверить, сколько данных прошло через каждый слой +-- и сколько ошибок было на каждом этапе. +-- +-- Важно: +-- Таблица dq_summary пересоздаётся при каждом запуске (TRUNCATE + INSERT), +-- чтобы не накапливать дубликаты при повторных прогонах. +-- ============================================================================ --- 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); +-- ---------------------------------------------------------------------------- +-- Пример материализации тяжёлой витрины (закомментировано) +-- ---------------------------------------------------------------------------- +-- Если VIEW dm.v_daily_traffic работает медленно, можно создать таблицу: +-- +-- 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; +-- ---------------------------------------------------------------------------- -TRUNCATE TABLE dm.daily_traffic_mart; - -INSERT INTO dm.daily_traffic_mart -SELECT * FROM dm.v_daily_traffic; -*/ - --- Data Quality summary table (always fresh) +-- ---------------------------------------------------------------------------- +-- Таблица для сводки по качеству данных (DQ summary) +-- ---------------------------------------------------------------------------- +-- Хранит метрики по всем слоям для быстрой проверки пайплайна +-- ---------------------------------------------------------------------------- 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 + check_date Date, -- Дата проверки + layer LowCardinality(String), -- Слой: stg, ods, dds + table_name LowCardinality(String), -- Имя таблицы + check_name LowCardinality(String), -- Тип проверки: total_rows, rows_with_errors и т.д. + check_value UInt64 -- Значение метрики ) ENGINE = MergeTree PARTITION BY toYYYYMM(check_date) ORDER BY (check_date, layer, table_name, check_name); --- Truncate and refill DQ summary +-- Очищаем перед заполнением, чтобы не было дубликатов при повторных запусках TRUNCATE TABLE dm.dq_summary; +-- ---------------------------------------------------------------------------- +-- Заполняем сводку метриками по всем слоям +-- ---------------------------------------------------------------------------- INSERT INTO dm.dq_summary SELECT today() AS check_date, @@ -50,25 +67,36 @@ SELECT '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 +-- Считаем строки с ошибками парсинга в ODS +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); +-- Считаем "осиротевшие" события (есть click_id, но нет такого click в dds.click) +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 index a97eedc..8188f90 100755 --- a/scripts/apply_clickhouse_ddl.sh +++ b/scripts/apply_clickhouse_ddl.sh @@ -1,32 +1,54 @@ #!/usr/bin/env bash +# +# Скрипт применения DDL в ClickHouse +# +# Назначение: +# Последовательно применяет SQL-файлы из ddl/*.sql в базу ClickHouse. +# Файлы применяются в алфавитном порядке (00 → 10 → 20 → 30 → 40). +# +# Как запускать: +# make ddl +# или: bash scripts/apply_clickhouse_ddl.sh +# +# Требования: +# - Сервис clickhouse должен быть запущен (make up) +# - Доступен clickhouse-client внутри контейнера +# +# Порядок применения важен: +# 00_databases.sql → 10_stg.sql → 20_ods.sql → 30_dds.sql → 40_dm.sql +# + 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}..." +echo "Применение DDL из ${DDL_DIR}..." -# Check if clickhouse service is running +# ----------------------------------------------------------------------------- +# Проверка: ClickHouse запущен? +# ----------------------------------------------------------------------------- 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." + echo "Ошибка: Сервис '${CLICKHOUSE_SERVICE}' не запущен." + echo "Запустите сначала: make up" exit 1 fi -# Apply DDL files in order (00 -> 10 -> 20 -> 30 -> 40) +# ----------------------------------------------------------------------------- +# Применение SQL-файлов по порядку +# ----------------------------------------------------------------------------- +# shellcheck disable=SC2044 for sql_file in "${DDL_DIR}"/*.sql; do if [[ -f "$sql_file" ]]; then - echo "Applying: $(basename "$sql_file")" + echo "Применение: $(basename "$sql_file")" ${COMPOSE_BIN} exec -T "${CLICKHOUSE_SERVICE}" clickhouse-client \ --user="${CLICKHOUSE_USER}" \ --password="${CLICKHOUSE_PASSWORD}" \ @@ -37,4 +59,11 @@ for sql_file in "${DDL_DIR}"/*.sql; do fi done -echo "DDL applied successfully." +echo "" +echo "DDL успешно применён!" +echo "" +echo "Созданы базы данных:" +echo " - stg : Staging (сырые данные)" +echo " - ods : Operational Data Store (типизированные данные)" +echo " - dds : Detailed Data Store (детальные сущности)" +echo " - dm : Data Marts (витрины для BI)" diff --git a/scripts/run_batch.sh b/scripts/run_batch.sh index b2a516a..6b84730 100755 --- a/scripts/run_batch.sh +++ b/scripts/run_batch.sh @@ -1,27 +1,50 @@ #!/usr/bin/env bash +# +# Скрипт batch-трансформации данных: ODS → DDS → DM +# +# Назначение: +# Запускает SQL-скрипты из jobs/ для преобразования данных между слоями: +# 1. ODS → DDS : Сборка сущностей из типизированных данных +# 2. DDS → DM : Обновление сводки по качеству данных (dq_summary) +# +# Как запускать: +# make transform +# или: bash scripts/run_batch.sh +# +# Требования: +# - ClickHouse запущен (make up) +# - ODS содержит данные (make data выполнен) +# +# Стратегия: +# Сейчас: полная перезагрузка (TRUNCATE + INSERT) — для демо +# В продакшене: инкрементальная загрузка по watermark +# + 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 +# ----------------------------------------------------------------------------- +# Проверка: ClickHouse запущен? +# ----------------------------------------------------------------------------- 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." + echo "Ошибка: Сервис '${CLICKHOUSE_SERVICE}' не запущен." + echo "Запустите сначала: make up" exit 1 fi -# Check if ODS has data +# ----------------------------------------------------------------------------- +# Проверка: в ODS есть данные? +# ----------------------------------------------------------------------------- ODS_COUNT=$(${COMPOSE_BIN} exec -T "${CLICKHOUSE_SERVICE}" clickhouse-client \ --user="${CLICKHOUSE_USER}" \ --password="${CLICKHOUSE_PASSWORD}" \ @@ -29,16 +52,21 @@ ODS_COUNT=$(${COMPOSE_BIN} exec -T "${CLICKHOUSE_SERVICE}" clickhouse-client \ --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." + echo "Предупреждение: Таблица ODS.browser_event пуста." + echo "Сначала загрузите данные: make data" exit 1 fi -echo "Found ${ODS_COUNT} rows in ODS.browser_event" +echo "Найдено ${ODS_COUNT} строк в ODS.browser_event" echo "" -# Step 1: Refresh DDS (truncate + reload for demo) -echo "Step 1: Refreshing DDS layer (ODS → DDS)..." +# ----------------------------------------------------------------------------- +# Шаг 1: ODS → DDS (сборка сущностей) +# ----------------------------------------------------------------------------- +echo "Шаг 1: Обновление DDS слоя (ODS → DDS)..." +echo " - Очистка текущих данных (TRUNCATE)..." + +# Очищаем таблицы перед загрузкой (полная перезагрузка для демо) ${COMPOSE_BIN} exec -T "${CLICKHOUSE_SERVICE}" clickhouse-client \ --user="${CLICKHOUSE_USER}" \ --password="${CLICKHOUSE_PASSWORD}" \ @@ -50,51 +78,61 @@ ${COMPOSE_BIN} exec -T "${CLICKHOUSE_SERVICE}" clickhouse-client \ --database="${CLICKHOUSE_DB}" \ --query="TRUNCATE TABLE dds.event" 2>/dev/null || true +echo " - Загрузка dds.click (device + geo)..." ${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" +echo " ✓ DDS обновлён" -# Show DDS stats +# Показываем статистику DDS echo "" -echo "DDS statistics:" +echo "Статистика DDS:" ${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) +# ----------------------------------------------------------------------------- +# Шаг 2: DDS → DM (сводка по качеству) +# ----------------------------------------------------------------------------- echo "" -echo "Step 2: Refreshing DM layer (DQ summary)..." +echo "Шаг 2: Обновление DM слоя (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" +echo " ✓ DM обновлён" -# Show DM stats +# Показываем сводку по качеству echo "" -echo "Data Quality summary:" +echo "Сводка по качеству данных:" ${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 "Batch-трансформация завершена!" +echo "========================================" 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" +echo "Доступные витрины для анализа:" +echo " - dm.v_events_enriched : Полное обогащение событий" +echo " - dm.v_daily_traffic : Агрегация по дням" +echo " - dm.v_top_pages_daily : Топ страниц" +echo " - dm.v_dq_errors_daily : Ошибки качества" +echo " - dm.v_session_overview : Обзор сессий" +echo " - dm.v_utm_effectiveness : Эффективность UTM" +echo " - dm.dq_summary : Статистика по слоям" +echo "" +echo "Подключитесь к Superset: http://localhost:8088"