refactor(sql): reorganize sql files into structured directory hierarchy

Move DDL files from flat ddl/ directory to sql/ddl/ with layer-based
subdirectories (stg, ods, dds, dm). Move batch transformation SQL from
jobs/ to sql/ layer directories. Update scripts and documentation to
reflect new paths for improved organization and Airflow integration.
This commit is contained in:
2026-02-07 20:51:51 +03:00
parent 6466921bda
commit 9e340bb729
14 changed files with 95 additions and 64 deletions
+13
View File
@@ -0,0 +1,13 @@
-- ============================================================================
-- Создание баз данных для слоёв хранилища
-- ============================================================================
-- 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;
CREATE DATABASE IF NOT EXISTS dm;
+79
View File
@@ -0,0 +1,79 @@
-- ============================================================================
-- Слой DDS (Detailed Data Store) — детальные сущности для аналитики
-- ============================================================================
-- Назначение:
-- - Собранные "чистые" сущности из ODS для JOIN'ов и аналитики
-- - click: объединяет device + geo (контекст сессии пользователя)
-- - event: объединяет browser + location (контекст события)
--
-- Загрузка:
-- Batch SQL (не MV!) — для согласованности при late arrivals
-- См. sql/dds/30_ods_to_dds.sql
--
-- Почему не MV:
-- - MV с JOIN даёт eventual consistency (данные приходят в разное время)
-- - Batch позволяет сделать снапшот через argMax и корректно джойнить
-- - Контроль: можно проверить SQL, откатить, перезапустить
-- ============================================================================
-- ----------------------------------------------------------------------------
-- Сущность: 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, -- 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; -- Разрешаем NULL (на всякий случай)
-- ----------------------------------------------------------------------------
-- Сущность: dds.event (контекст события)
-- ----------------------------------------------------------------------------
-- Объединяет данные из ods.browser_event + ods.location_event
-- Связь с click через click_id (может быть NULL, если нет device/geo)
-- ----------------------------------------------------------------------------
CREATE TABLE IF NOT EXISTS dds.event
(
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; -- Разрешаем NULL
+169
View File
@@ -0,0 +1,169 @@
-- ============================================================================
-- Слой DM (Data Marts) — витрины для BI (Superset/Grafana)
-- ============================================================================
-- Назначение:
-- - Представления (VIEW) для удобного доступа к данным из BI-инструментов
-- - Обогащение: соединяем event + click через LEFT JOIN
-- - Агрегации: готовые GROUP BY для частых запросов
--
-- Почему VIEW:
-- - Гибкость: меняем логику без пересоздания таблиц
-- - Нет дублирования данных (храним только в DDS)
-- - Для демо: производительность достаточная
--
-- Для продакшена:
-- - Если тяжёлые агрегации тормозят — материализовать в таблицы
-- - См. пример закомментированный в sql/dm/40_dds_to_dm.sql
-- ============================================================================
-- ----------------------------------------------------------------------------
-- Витрина: полное обогащение событий (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,
e.event_ts,
e.event_date,
e.event_type,
e.click_id,
-- Поля из location (через event)
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,
-- Поля из 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,
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;
-- ----------------------------------------------------------------------------
-- Витрина: агрегация трафика по дням и измерениям
-- ----------------------------------------------------------------------------
-- Используется для анализа посещаемости
-- Гранулярность: дата × страна × устройство × браузер × UTM
-- ----------------------------------------------------------------------------
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 -- Фильтруем битые timestamp
GROUP BY
event_date,
geo_country,
device_type,
browser_name,
utm_source,
utm_medium;
-- ----------------------------------------------------------------------------
-- Витрина: популярность страниц (воронка)
-- ----------------------------------------------------------------------------
-- Показывает какие страницы чаще всего просматривают
-- Используется для анализа воронки конверсии
-- ----------------------------------------------------------------------------
CREATE VIEW IF NOT EXISTS dm.v_top_pages_daily AS
SELECT
event_date,
page_url_path, -- Путь URL (/home, /product)
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;
-- ----------------------------------------------------------------------------
-- Витрина: ошибки парсинга по дням
-- ----------------------------------------------------------------------------
-- Для мониторинга качества данных
-- Показывает сколько строк с какими ошибками за каждый день
-- ----------------------------------------------------------------------------
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;
-- ----------------------------------------------------------------------------
-- Витрина: обзор сессий пользователей
-- ----------------------------------------------------------------------------
-- Группировка по 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, -- Устройства (если менялось)
-- Берём последние 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 -- Только идентифицированные пользователи
GROUP BY event_date, user_domain_id, click_id;
-- ----------------------------------------------------------------------------
-- Витрина: эффективность 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 -- Добавления в корзину
FROM dm.v_events_enriched
WHERE utm_source IS NOT NULL OR utm_medium IS NOT NULL -- Только с UTM-метками
GROUP BY event_date, utm_source, utm_medium, utm_campaign;
+388
View File
@@ -0,0 +1,388 @@
-- ============================================================================
-- Слой ODS (Operational Data Store) — типизированные данные + дедупликация + DQ
-- ============================================================================
-- Назначение:
-- - Типизация данных из STG (String → UUID, DateTime, etc.)
-- - Дедупликация через ReplacingMergeTree (последняя версия по src_ingest_ts)
-- - Контроль качества: массив parse_errors для "грязных" данных
-- - Разделение: валидные строки → основная таблица, ошибки → *_errors
--
-- Поток данных:
-- STG (*_raw) → MV → ODS (основная таблица + error_tables)
-- ============================================================================
-- ============================================================================
-- 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), -- 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; -- Разрешаем NULL в ключе (для "битых" данных)
-- ----------------------------------------------------------------------------
-- Таблица ошибок: строки с невалидными ключами (event_id IS NULL)
-- ----------------------------------------------------------------------------
-- Сохраняем полную информацию для анализа проблем с данными
-- MergeTree без Replacing: сохраняем все ошибки (не дедуплицируем)
-- ----------------------------------------------------------------------------
CREATE TABLE IF NOT EXISTS ods.browser_event_errors
(
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
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; -- Только валидные строки (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
WITH
toUUIDOrNull(JSONExtractString(raw, 'event_id')) AS event_id,
parseDateTime64BestEffortOrNull(JSONExtractString(raw, 'event_timestamp'), 6) AS event_ts,
toUUIDOrNull(JSONExtractString(raw, 'click_id')) AS click_id,
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
SELECT
ingest_ts,
kafka_topic,
kafka_partition,
kafka_offset,
kafka_ts,
raw,
arrayStringConcat(parse_errors, ',') AS error_reason
FROM stg.browser_raw
WHERE length(parse_errors) > 0 -- Есть хотя бы одна ошибка
AND (
event_id IS NULL
OR event_ts IS NULL
OR click_id IS NULL
);
-- ============================================================================
-- LOCATION EVENTS
-- ============================================================================
CREATE TABLE IF NOT EXISTS ods.location_event
(
event_id Nullable(UUID),
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)) -- Партиция по времени загрузки (нет event_date)
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;
CREATE MATERIALIZED VIEW IF NOT EXISTS stg.mv_location_raw_to_ods_errors
TO ods.location_event_errors
AS
WITH
toUUIDOrNull(JSONExtractString(raw, 'event_id')) AS event_id,
arrayFilter(x -> x != '', [
if(event_id IS NULL, 'bad_event_id', '')
]) AS parse_errors
SELECT
ingest_ts,
kafka_topic,
kafka_partition,
kafka_offset,
kafka_ts,
raw,
arrayStringConcat(parse_errors, ',') AS error_reason
FROM stg.location_raw
WHERE length(parse_errors) > 0
AND event_id IS NULL;
-- ============================================================================
-- DEVICE EVENTS
-- ============================================================================
CREATE TABLE IF NOT EXISTS ods.device_by_click
(
click_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))
)
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;
CREATE MATERIALIZED VIEW IF NOT EXISTS stg.mv_device_raw_to_ods_errors
TO ods.device_by_click_errors
AS
WITH
toUUIDOrNull(JSONExtractString(raw, 'click_id')) AS click_id,
toUUIDOrNull(JSONExtractString(raw, 'user_domain_id')) AS user_domain_id,
arrayFilter(x -> x != '', [
if(click_id IS NULL, 'bad_click_id', ''),
if(user_domain_id IS NULL, 'bad_user_domain_id', '')
]) AS parse_errors
SELECT
ingest_ts,
kafka_topic,
kafka_partition,
kafka_offset,
kafka_ts,
raw,
arrayStringConcat(parse_errors, ',') AS error_reason
FROM stg.device_raw
WHERE length(parse_errors) > 0
AND (
click_id IS NULL
OR user_domain_id IS NULL
);
-- ============================================================================
-- 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)), -- Код страны (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))
)
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;
CREATE MATERIALIZED VIEW IF NOT EXISTS stg.mv_geo_raw_to_ods_errors
TO ods.geo_by_click_errors
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,
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
SELECT
ingest_ts,
kafka_topic,
kafka_partition,
kafka_offset,
kafka_ts,
raw,
arrayStringConcat(parse_errors, ',') AS error_reason
FROM stg.geo_raw
WHERE length(parse_errors) > 0
AND (
click_id IS NULL
OR geo_latitude IS NULL
OR geo_longitude IS NULL
);
+169
View File
@@ -0,0 +1,169 @@
-- ============================================================================
-- Слой STG (Staging) — сырые JSON + интеграция с Kafka
-- ============================================================================
-- Назначение:
-- - Хранение сырых данных "как есть" из Kafka
-- - Метаданные доставки (topic, partition, offset, timestamp)
-- - Источник для отладки и восстановления
--
-- Поток данных:
-- Kafka → kafka_*_raw (ENGINE=Kafka) → MV → *_raw (MergeTree)
-- ============================================================================
-- ----------------------------------------------------------------------------
-- Таблицы для сырых данных (MergeTree)
-- ----------------------------------------------------------------------------
-- Хранят JSON "как есть" + метаданные Kafka
-- Партиционируем по месяцу ingest_ts для эффективной очистки старых данных
-- ----------------------------------------------------------------------------
CREATE TABLE IF NOT EXISTS stg.browser_raw
(
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)
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 (ENGINE = Kafka)
-- ----------------------------------------------------------------------------
-- Читают данные из топиков Kafka в реальном времени
-- Не хранят данные постоянно, а "потребляют" их при SELECT/MV
-- ----------------------------------------------------------------------------
CREATE TABLE IF NOT EXISTS stg.kafka_browser_raw (raw String)
ENGINE = Kafka
SETTINGS
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
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
-- ----------------------------------------------------------------------------
-- Автоматически перекладывают данные из Kafka-таблиц в MergeTree
-- Работают в реальном времени: сообщение из Kafka → сразу в *_raw
-- ----------------------------------------------------------------------------
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;