diff --git a/plans/clickhouse_ddl.md b/plans/clickhouse_ddl.md index b9700b9..685682c 100644 --- a/plans/clickhouse_ddl.md +++ b/plans/clickhouse_ddl.md @@ -31,18 +31,18 @@ Эта схема укладывается в задание так: -- **Kafka → STG → ODS → DDS** можно сделать полностью внутри ClickHouse через `ENGINE = Kafka` + MV (стриминг, 1 json = 1 row). -- “**регулярный процесс**” в MVP заменяется встроенной реактивностью MV (это тоже регулярная обработка, просто без внешнего оркестратора). -- Если хочется буквально следовать формулировке “в Airflow преобразовать и выгрузить агрегаты”, то Airflow можно использовать как **тонкий оркестратор**: - - либо для `ODS → DDS/DM` батчами (`INSERT INTO … SELECT …`) по расписанию, - - либо для материализации/пересчёта DM-агрегатов (daily/topN) под Superset. +- **Kafka → STG → ODS** можно сделать полностью внутри ClickHouse через `ENGINE = Kafka` + MV (стриминг, 1 json = 1 row). +- “**регулярный процесс**” — батч‑трансформации `ODS → DDS (→ DM)` в виде SQL (`INSERT INTO … SELECT …`) по расписанию (в будущем Airflow; пока — ручной запуск). +- Airflow (когда появится) можно использовать как **тонкий оркестратор**: применить DDL и запускать batch‑SQL по расписанию. ### Выбранное решение (MVP) Чтобы сделать “хорошо, но без оверинжиниринга”, фиксируем такое MVP: -- **STG → ODS → DDS** — инкрементально через MV (всё рядом с данными, минимум внешних компонентов). -- **DM** — `VIEW` (витрины “на чтении”), чтобы не плодить лишние таблицы и джобы под демо. +- **STG → ODS** — инкрементально через MV (парсинг/типизация рядом с ingest). +- **DDS** — батч‑сборка из ODS через SQL (`INSERT INTO … SELECT …`) как “регулярный процесс”. + - Причина: `MV + JOIN` в `ODS → DDS` плохо переносит произвольный порядок прихода данных и может давать некорректные результаты (eventual consistency ODS, версии в разных партициях и т.п.). +- **DM** — `VIEW` (витрины “на чтении”) поверх DDS, чтобы не плодить лишние таблицы и джобы под демо. - На стороне BI считаем, что запросы всегда идут с фильтрами по времени (`event_date`/`event_ts`) и не сканируют всю историю. Диаграмма витрин и потоков данных: @@ -80,11 +80,11 @@ flowchart LR stg_device -->|MV parse| ods_device stg_geo -->|MV parse| ods_geo - ods_browser -->|MV join| dds_event - ods_location -->|MV join| dds_event + ods_browser -->|Batch SQL| dds_event + ods_location -->|Batch SQL| dds_event - ods_device -->|MV join| dds_click - ods_geo -->|MV join| dds_click + ods_device -->|Batch SQL| dds_click + ods_geo -->|Batch SQL| dds_click dds_event -->|VIEW join| v_enriched dds_click -->|VIEW join| v_enriched @@ -94,13 +94,21 @@ flowchart LR v_enriched --> v_dq ``` -### Эволюция при росте (scale-out story) +### Почему DDS батчами (а не MV join) -Если система растёт и MV-джойны в `ODS → DDS` начинают ограничивать ingest, план такой: +Мы сознательно уходим от `MV + JOIN` в `ODS → DDS`, потому что это решение: -- перевод `ODS → DDS` на батчи (по watermark/окнам) через `INSERT INTO … SELECT …` по расписанию; -- материализация самых популярных/тяжёлых витрин в `DM` (daily/topN), чтобы BI не джойнил “деталь” на лету; -- ресурсные лимиты для BI-пользователя ClickHouse (time/memory/rows), чтобы Superset не “утопил” БД. +- чувствительно к произвольному порядку прихода сообщений между топиками; +- может давать неконсистентные “снимки” из‑за версионирования в `ReplacingMergeTree` и отсутствия гарантий “последней версии” в момент выполнения MV; +- может порождать дубли, если одна и та же сущность попадает в разные партиции (например, когда часть полей для партиционирования появляется “позже”). + +Поэтому DDS считаем батчами: сначала получаем “current snapshot” ODS (например, через `argMax(..., src_ingest_ts)` по ключу), потом делаем join и грузим результат в DDS. + +Дальше, при росте нагрузки: + +- делаем инкрементальные батчи по watermark/окнам (а не full rebuild); +- материализуем самые тяжёлые витрины в `DM` (daily/topN), чтобы BI не джойнил “деталь” на лету; +- вводим ресурсные лимиты для BI-пользователя ClickHouse (time/memory/rows), чтобы Superset не “утопил” БД. --- @@ -117,12 +125,18 @@ flowchart LR - `ddl/00_databases.sql` — базы `stg/ods/dds/dm`. - `ddl/10_stg.sql` — STG raw (`stg.*_raw`) + Kafka source tables (`ENGINE = Kafka`) + MV `Kafka → STG`. -- `ddl/20_ods.sql` — ODS таблицы типизации + DQ (`parse_errors`) + MV `STG → ODS`. -- `ddl/30_dds.sql` — DDS сущности (`dds.event`, `dds.click`) + MV `ODS → DDS`. +- `ddl/20_ods.sql` — ODS таблицы типизации + DQ (`parse_errors`) + MV `STG → ODS` + таблицы `ods_*_errors` для строк с битыми ключами. +- `ddl/30_dds.sql` — DDS таблицы (`dds.event`, `dds.click`) **без MV** (только `CREATE TABLE`). - `ddl/40_dm.sql` — витрины `VIEW` для Superset (`dm.v_*`). BI-ограничения (ресурсы/пользователь) **не выносим в `ddl/*.sql`**: оставляем это только как текст/пример в этом плане, чтобы не смешивать инфраструктуру доступа с DDL витрин. +### Артефакты batch-трансформаций (планируемые файлы) + +- `jobs/30_dds_refresh.sql` — регулярная батч‑сборка DDS из ODS: + - получить “последнюю версию” строк по ключам (`event_id`/`click_id`) через `argMax(..., src_ingest_ts)` (или эквивалент); + - выполнить join snapshot’ов и загрузить в `dds.event`/`dds.click` (для демо возможно “full rebuild”; позже — инкрементально). + ### Исполнение DDL (make сейчас / Airflow потом) Требования к файлам `ddl/*.sql`: @@ -131,12 +145,30 @@ BI-ограничения (ресурсы/пользователь) **не вы - строгий порядок исполнения: `00 → 10 → 20 → 30 → 40` (из‑за зависимостей MV); - единые имена топиков Kafka: `browser_events`, `location_events`, `device_events`, `geo_events` (их создаёт `make data`). +### Дедупликация и обработка “битых” ключей (ODS) + +Проблема: данные “грязные”, а ключи стыковки (`event_id`, `click_id`) могут быть `NULL`/невалидными. Если хранить такие строки в основной ODS‑таблице на `ReplacingMergeTree` с `ORDER BY (event_id/click_id)`, то строки с `NULL` ключом могут схлопываться друг с другом на мерджах, и мы потеряем часть ошибок. + +Решение в target state: + +- основная ODS (`ods.browser_event`, `ods.location_event`, `ods.device_by_click`, `ods.geo_by_click`) хранит только строки с валидными ключами (key `IS NOT NULL`) и подходит для join’ов/сборки DDS; +- отдельные таблицы `ods.*_errors` хранят строки с битыми ключами (и/или критичными ошибками парсинга) для DQ‑аналитики и дебага; в них важно сохранять “уникальность строки” через Kafka‑метаданные (`kafka_topic/partition/offset`, `kafka_ts`) + `src_ingest_ts` + `raw`. + +Про дедуп: + +- STG хранит все сообщения как есть; при чтении из Kafka уникальность сообщения определяется `(kafka_topic, kafka_partition, kafka_offset)`. +- В основной ODS дедупликация — по бизнес‑ключу (`event_id` / `click_id`) с версией `src_ingest_ts` (ReplacingMergeTree). +- Для `ods.*_errors` дедуп/уникальность (если потребуется) делаем по Kafka‑метаданным; но в демо допустимо хранить “как пришло” без схлопывания. + Текущее “как запускаем” (целевое, для реализации следующим шагом): - `make ddl` вызывает `scripts/apply_clickhouse_ddl.sh`; - скрипт прогоняет `ddl/*.sql` по порядку через `clickhouse-client --multiquery` внутри контейнера ClickHouse. -Airflow-версия (целевое): тот же порядок, но каждый `ddl/.sql` — отдельный таск, зависимости линейные. +Batch‑трансформации (целевое, для реализации следующим шагом): + +- `make transform` (или аналогичная команда) запускает `jobs/30_dds_refresh.sql` через `clickhouse-client`; +- в будущем Airflow будет делать то же самое по расписанию (один job‑SQL = один task). ### Параметры окружения (docker compose) @@ -593,7 +625,7 @@ DDS хранит **минимально необходимую детализа - `dds.event` — 1 строка на `event_id` (browser + location). - `dds.click` — 1 строка на `click_id` (device + geo + user). -Ключевой момент для стриминга: части могут приходить **в разное время**. Поэтому DDS делаем на `ReplacingMergeTree` и создаём **две MV на сущность** — по одной на каждую сторону джойна, чтобы поздние данные “перезатирали” строку новой версией. +Важно: MV с `JOIN` для `ODS → DDS` мы в target state **не используем** (см. обоснование выше). DDS собирается батчами из “снапшота” ODS. Поэтому в приложении ниже оставляем только DDL таблиц DDS. ### 3.1 DDS: click @@ -624,54 +656,6 @@ CREATE TABLE IF NOT EXISTS dds.click ENGINE = ReplacingMergeTree(dds_update_ts) PARTITION BY toYYYYMM(toDate(dds_update_ts)) ORDER BY (click_id); - -CREATE MATERIALIZED VIEW IF NOT EXISTS ods.mv_device_to_dds_click -TO dds.click -AS -SELECT - d.click_id, - d.user_domain_id, - d.user_custom_id, - d.device_type, - d.device_is_mobile, - d.os_name, - d.os, - d.os_timezone, - g.geo_country, - g.geo_region_name, - g.geo_timezone, - g.geo_latitude, - g.geo_longitude, - g.ip_address, - greatest(d.src_ingest_ts, ifNull(g.src_ingest_ts, toDateTime64(0, 3))) AS dds_update_ts, - arrayConcat(d.parse_errors, ifNull(g.parse_errors, [])) AS parse_errors -FROM ods.device_by_click AS d -LEFT JOIN ods.geo_by_click AS g - ON g.click_id = d.click_id; - -CREATE MATERIALIZED VIEW IF NOT EXISTS ods.mv_geo_to_dds_click -TO dds.click -AS -SELECT - g.click_id, - d.user_domain_id, - d.user_custom_id, - d.device_type, - d.device_is_mobile, - d.os_name, - d.os, - d.os_timezone, - g.geo_country, - g.geo_region_name, - g.geo_timezone, - g.geo_latitude, - g.geo_longitude, - g.ip_address, - greatest(g.src_ingest_ts, ifNull(d.src_ingest_ts, toDateTime64(0, 3))) AS dds_update_ts, - arrayConcat(g.parse_errors, ifNull(d.parse_errors, [])) AS parse_errors -FROM ods.geo_by_click AS g -LEFT JOIN ods.device_by_click AS d - ON d.click_id = g.click_id; ``` ### 3.2 DDS: event @@ -704,56 +688,6 @@ CREATE TABLE IF NOT EXISTS dds.event ENGINE = ReplacingMergeTree(dds_update_ts) PARTITION BY toYYYYMM(event_date) ORDER BY (event_id); - -CREATE MATERIALIZED VIEW IF NOT EXISTS ods.mv_browser_to_dds_event -TO dds.event -AS -SELECT - b.event_id, - b.event_ts, - b.event_type, - b.click_id, - l.page_url, - l.page_url_path, - l.referer_url, - l.referer_medium, - l.utm_medium, - l.utm_source, - l.utm_content, - l.utm_campaign, - b.browser_name, - b.browser_user_agent, - b.browser_language, - greatest(b.src_ingest_ts, ifNull(l.src_ingest_ts, toDateTime64(0, 3))) AS dds_update_ts, - arrayConcat(b.parse_errors, ifNull(l.parse_errors, [])) AS parse_errors -FROM ods.browser_event AS b -LEFT JOIN ods.location_event AS l - ON l.event_id = b.event_id; - -CREATE MATERIALIZED VIEW IF NOT EXISTS ods.mv_location_to_dds_event -TO dds.event -AS -SELECT - l.event_id, - b.event_ts, - b.event_type, - b.click_id, - l.page_url, - l.page_url_path, - l.referer_url, - l.referer_medium, - l.utm_medium, - l.utm_source, - l.utm_content, - l.utm_campaign, - b.browser_name, - b.browser_user_agent, - b.browser_language, - greatest(l.src_ingest_ts, ifNull(b.src_ingest_ts, toDateTime64(0, 3))) AS dds_update_ts, - arrayConcat(l.parse_errors, ifNull(b.parse_errors, [])) AS parse_errors -FROM ods.location_event AS l -LEFT JOIN ods.browser_event AS b - ON b.event_id = l.event_id; ``` ---