docs: add russian comments to pipeline files
This commit is contained in:
+48
-16
@@ -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,
|
||||
|
||||
+63
-35
@@ -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);
|
||||
|
||||
Reference in New Issue
Block a user