docs(plans): update ddl architecture to use batch transforms for ods to dds
Replace materialized view joins with batch SQL transformations to avoid consistency issues with out-of-order data. Document the reasoning for using batch processing for ODS to DDS layer, including handling of eventual consistency and versioning in ReplacingMergeTree. Update data flow diagrams and remove MV creation DDL for DDS tables. Add documentation for error handling tables and batch transformation jobs.
This commit is contained in:
+52
-118
@@ -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/<step>.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;
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
Reference in New Issue
Block a user