- Зачем:
- чарт «Data Quality Summary» суммировал total_rows по всем таблицам слоя,
складывал таблицы разного зерна (события 1000 + визиты 99 + error-таблицы 0)
и рисовал убывающую «воронку потерь» (stg≈4250→ods≈2198→dds≈1099), которой
в данных нет. На учебном стенде это активно вводит в заблуждение.
- Что:
- чарт переделан в row-lineage одного event-зерна и переименован в
«🧱 Rows by Layer (event)»; rename идемпотентный через previous_slice_names.
- чарт берёт по одной канонической таблице на слой
(browser_raw→browser_event→event→v_events_enriched), порядок слоёв задан
числовым префиксом в groupby + order_bars.
- в dm.dq_summary добавлена строка total_rows для слоя dm, чтобы цепочка
замыкалась до витрины.
- описание дашборда обновлено под новый смысл.
- Проверка:
- python3 -m py_compile superset/create_dashboard.py.
- make transform / прогон sql/dm/40_dds_to_dm.sql; в dq_summary есть строка dm.
- make superset-dashboard (идемпотентно, 10 чартов, дублей нет).
- визуально через playwright-cli: 4 столбца 1·stg→2·ods→3·dds→4·dm,
видимый шаг дедупликации 1050→1000, консоль без ошибок.
116 lines
5.7 KiB
SQL
116 lines
5.7 KiB
SQL
-- ============================================================================
|
|
-- Batch-трансформация: DDS → DM (Data Quality summary)
|
|
-- ============================================================================
|
|
-- Поток данных:
|
|
-- stg.* + ods.* + dds.* + dm.v_events_enriched → dm.dq_summary (сводка по слоям)
|
|
-- Сами витрины (dm.v_*) — это VIEW поверх DDS, создаются в sql/ddl/dm/40_dm.sql.
|
|
--
|
|
-- Что делает:
|
|
-- Собирает статистику по всем слоям (stg/ods/dds) для мониторинга качества данных.
|
|
-- Позволяет быстро проверить, сколько данных прошло через каждый слой
|
|
-- и сколько ошибок было на каждом этапе.
|
|
--
|
|
-- Важно:
|
|
-- Таблица dq_summary пересоздаётся при каждом запуске (TRUNCATE + INSERT),
|
|
-- чтобы не накапливать дубликаты при повторных прогонах.
|
|
--
|
|
-- Витрины DM сейчас — это VIEW (логика без копии данных). Если тяжёлая агрегация
|
|
-- начнёт тормозить, её материализуют в таблицу — пример в docs/ARCHITECTURE.md,
|
|
-- раздел «Материализация витрин».
|
|
-- ============================================================================
|
|
|
|
-- ----------------------------------------------------------------------------
|
|
-- Таблица для сводки по качеству данных (DQ summary)
|
|
-- ----------------------------------------------------------------------------
|
|
-- Хранит метрики по всем слоям для быстрой проверки пайплайна
|
|
-- ----------------------------------------------------------------------------
|
|
CREATE TABLE IF NOT EXISTS dm.dq_summary
|
|
(
|
|
check_date Date, -- Дата проверки
|
|
layer LowCardinality(String), -- Слой: stg, ods, dds, dm
|
|
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 TABLE dm.dq_summary;
|
|
|
|
-- ----------------------------------------------------------------------------
|
|
-- Заполняем сводку метриками по всем слоям
|
|
-- ----------------------------------------------------------------------------
|
|
INSERT INTO dm.dq_summary
|
|
SELECT
|
|
today() AS check_date,
|
|
'stg' AS layer,
|
|
'browser_raw' AS table_name,
|
|
'total_rows' AS check_name,
|
|
count() AS check_value
|
|
FROM stg.browser_raw
|
|
|
|
UNION ALL
|
|
SELECT today(), 'stg', 'location_raw', 'total_rows', count() FROM stg.location_raw
|
|
UNION ALL
|
|
SELECT today(), 'stg', 'device_raw', 'total_rows', count() FROM stg.device_raw
|
|
UNION ALL
|
|
SELECT today(), 'stg', 'geo_raw', 'total_rows', count() FROM stg.geo_raw
|
|
|
|
UNION ALL
|
|
SELECT today(), 'ods', 'browser_event', 'total_rows', count() FROM ods.browser_event
|
|
UNION ALL
|
|
-- Считаем строки с ошибками парсинга в 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', 'location_event', 'rows_with_errors', count()
|
|
FROM ods.location_event
|
|
WHERE length(parse_errors) > 0
|
|
|
|
UNION ALL
|
|
SELECT today(), 'ods', 'device_by_click', 'total_rows', count() FROM ods.device_by_click
|
|
UNION ALL
|
|
SELECT today(), 'ods', 'device_by_click', 'rows_with_errors', count()
|
|
FROM ods.device_by_click
|
|
WHERE length(parse_errors) > 0
|
|
|
|
UNION ALL
|
|
SELECT today(), 'ods', 'geo_by_click', 'total_rows', count() FROM ods.geo_by_click
|
|
UNION ALL
|
|
SELECT today(), 'ods', 'geo_by_click', 'rows_with_errors', count()
|
|
FROM ods.geo_by_click
|
|
WHERE length(parse_errors) > 0
|
|
|
|
UNION ALL
|
|
SELECT today(), 'ods', 'browser_event_errors', 'total_rows', count() FROM ods.browser_event_errors
|
|
UNION ALL
|
|
SELECT today(), 'ods', 'location_event_errors', 'total_rows', count() FROM ods.location_event_errors
|
|
UNION ALL
|
|
SELECT today(), 'ods', 'device_by_click_errors', 'total_rows', count() FROM ods.device_by_click_errors
|
|
UNION ALL
|
|
SELECT today(), 'ods', 'geo_by_click_errors', 'total_rows', count() FROM ods.geo_by_click_errors
|
|
|
|
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
|
|
-- Считаем "осиротевшие" события (есть 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)
|
|
|
|
UNION ALL
|
|
-- DM-слой: финальная витрина событий (VIEW поверх dds.event).
|
|
-- Нужна, чтобы lineage-чарт замыкал цепочку stg→ods→dds→dm на одном (event) зерне.
|
|
-- VIEW создаётся в DDL (sql/ddl/dm/40_dm.sql) до трансформаций, а этот шаг идёт
|
|
-- после наполнения dds.event — поэтому count() здесь корректен.
|
|
SELECT today(), 'dm', 'v_events_enriched', 'total_rows', count() FROM dm.v_events_enriched;
|