From cbf5f220645e25c2855a1782666430f430255609 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Fri, 6 Feb 2026 22:21:16 +0300 Subject: [PATCH] docs(architecture): update ODS error handling and DDS partial data support Refine data flow diagrams and documentation to clarify error handling in the ODS layer and partial data processing in the DDS layer. Add detailed explanations for materialized views, batch SQL transformations, and data quality metrics. Split DDS entity assembly diagrams for better readability of event and click processing pipelines. --- docs/ARCHITECTURE.md | 268 ++++++++++++++++++++++++++++++------------- 1 file changed, 187 insertions(+), 81 deletions(-) diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 4038c20..5400f5c 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -24,91 +24,91 @@ ### Общая схема потока данных ```mermaid -flowchart TB +flowchart LR subgraph Sources["📁 Источники (JSONL)"] - BE[browser_events.jsonl] - LE[location_events.jsonl] - DE[device_events.jsonl] - GE[geo_events.jsonl] + BE[browser_events] + LE[location_events] + DE[device_events] + GE[geo_events] end - subgraph Kafka["🚀 Kafka Topics"] - KT1[browser_events] - KT2[location_events] - KT3[device_events] - KT4[geo_events] + subgraph Kafka["🚀 Kafka"] + K1[browser_events] + K2[location_events] + K3[device_events] + K4[geo_events] end - subgraph STG["📦 STG (Staging)"] - BR[browser_raw] - LR[location_raw] - DR[device_raw] - GR[geo_raw] + subgraph STG["📦 STG"] + S1[browser_raw] + S2[location_raw] + S3[device_raw] + S4[geo_raw] + MV1[mv_*_to_ods] + MV2[mv_*_to_errors] end - subgraph ODS["🔧 ODS (Operational Data Store)"] - BE_O[browser_event] - LE_O[location_event] - DE_O[device_by_click] - GE_O[geo_by_click] - ERR[error_tables] + subgraph ODS["🔧 ODS"] + O1[browser_event] + O2[location_event] + O3[device_by_click] + O4[geo_by_click] + OE[error_tables] end - subgraph DDS["🎯 DDS (Detailed Data Store)"] - E[event] - C[click] + subgraph DDS["🎯 DDS"] + DE1[event] + DC1[click] end - subgraph DM["📊 DM (Data Marts)"] - VE[v_events_enriched] - VDT[v_daily_traffic] - VTP[v_top_pages_daily] - VUTM[v_utm_effectiveness] - VSE[v_session_overview] - VDQ[v_dq_errors_daily] + subgraph DM["📊 DM"] + DM1[v_events_enriched] + DM2[v_daily_traffic] + DM3[v_utm_effectiveness] + DM4[v_top_pages] end - BE --> KT1 --> BR --> BE_O --> E --> VE - LE --> KT2 --> LR --> LE_O --> E - DE --> KT3 --> DR --> DE_O --> C --> VE - GE --> KT4 --> GR --> GE_O --> C - - BE_O -.->|ошибки| ERR - E --> VDT & VTP & VUTM & VSE & VDQ - C --> VDT & VTP & VUTM & VSE & VDQ + BE --> K1 --> S1 --> MV1 --> O1 --> DE1 --> DM1 + LE --> K2 --> S2 --> MV1 --> O2 --> DE1 + DE --> K3 --> S3 --> MV1 --> O3 --> DC1 --> DM1 + GE --> K4 --> S4 --> MV1 --> O4 --> DC1 + + S1 & S2 & S3 & S4 --> MV2 -.-> OE + DE1 --> DM2 & DM3 & DM4 + DC1 --> DM2 & DM3 & DM4 ``` ### Слои и их назначение ```mermaid -flowchart LR - subgraph L0["📝 Сырые данные"] - RAW[JSON файлы
1000 строк каждый] +flowchart TB + subgraph L0["📝 Источники"] + RAW["JSON файлы (1000 строк)"] end - subgraph L1["STG - Staging"] - STG_T["Таблицы *_raw
MergeTree"] - KAFKA["Kafka Engine + MV"] + subgraph L1["📦 STG - Staging"] + direction LR + KAFKA["Kafka Engine"] + STG_T["*_raw таблицы
(MergeTree)"] end - subgraph L2["ODS - Операционный слой"] - ODS_T["Типизированные таблицы
ReplacingMergeTree"] + subgraph L2["🔧 ODS - Операционный слой"] + direction LR + ODS_T["Типизированные таблицы
(ReplacingMergeTree)"] DQ["parse_errors
DQ-метрики"] end - subgraph L3["DDS - Детальный слой"] - DDS_T["Сущности event + click
Batch SQL"] + subgraph L3["🎯 DDS - Детальный слой"] + DDS_T["event + click
(Batch SQL)"] end - subgraph L4["DM - Витрины"] - DM_T["VIEW для BI
Superset/Grafana"] + subgraph L4["📊 DM - Витрины"] + DM_T["VIEW для BI
(Superset/Grafana)"] end - RAW -->|kafka-console-producer| KAFKA -->|MV| STG_T - STG_T -->|MV| ODS_T - ODS_T -->|argMax + JOIN| DDS_T - DDS_T -->|VIEW| DM_T - ODS_T -.->|ошибки парсинга| DQ + RAW -->|kafka-console-producer| KAFKA -->|MV| STG_T -->|MV| ODS_T + ODS_T -->|argMax + JOIN| DDS_T -->|VIEW| DM_T + ODS_T -.->|ошибки| DQ ``` --- @@ -159,6 +159,19 @@ CREATE TABLE stg.browser_raw ( | `geo_by_click` | click_id | ReplacingMergeTree(src_ingest_ts) | Гео-данные | | `*_errors` | — | MergeTree | Строки с битыми ключами | +**Materialized Views для обработки ошибок:** + +| MV | Назначение | +|----|-----------| +| `mv_browser_raw_to_ods_errors` | Переносит строки с ошибками в `browser_event_errors` | +| `mv_location_raw_to_ods_errors` | Переносит строки с ошибками в `location_event_errors` | +| `mv_device_raw_to_ods_errors` | Переносит строки с ошибками в `device_by_click_errors` | +| `mv_geo_raw_to_ods_errors` | Переносит строки с ошибками в `geo_by_click_errors` | + +**Логика разделения:** +- **Основная таблица**: строки с валидными ключами (`WHERE key IS NOT NULL`) +- **Таблица ошибок**: строки с невалидными ключами (`WHERE key IS NULL`) + **Пример структуры:** ```sql CREATE TABLE ods.browser_event ( @@ -231,22 +244,57 @@ CREATE TABLE dds.click ( ``` **Загрузка (Batch SQL):** + +Загрузка `dds.click` с поддержкой partial data (когда device и geo приходят независимо): + ```sql --- Снапшот ODS через argMax +-- UNION всех click_id из device и geo INSERT INTO dds.click -SELECT d.click_id, d.user_domain_id, ..., g.geo_country, ... +SELECT + c.click_id, + d.user_domain_id, + d.device_type, + g.geo_country, + g.geo_latitude, + -- ... остальные поля + now64(3) AS dds_update_ts, + 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'], []) + )) AS ods_parse_errors FROM ( - SELECT click_id, argMax(user_domain_id, src_ingest_ts) AS user_domain_id, ... - FROM ods.device_by_click - GROUP BY click_id -) d + -- Union всех click_id для обработки geo-only и device-only + SELECT click_id FROM ( + SELECT assumeNotNull(click_id) AS click_id + FROM ods.device_by_click WHERE click_id IS NOT NULL + GROUP BY click_id + ) + UNION DISTINCT + SELECT click_id FROM ( + SELECT assumeNotNull(click_id) AS click_id + FROM ods.geo_by_click WHERE click_id IS NOT NULL + GROUP BY click_id + ) +) AS c LEFT JOIN ( - SELECT click_id, argMax(geo_country, src_ingest_ts) AS geo_country, ... - FROM ods.geo_by_click - GROUP BY click_id -) g ON g.click_id = d.click_id; + -- Снапшот device + SELECT assumeNotNull(click_id) AS click_id, ... + FROM ods.device_by_click GROUP BY click_id +) AS d ON d.click_id = c.click_id +LEFT JOIN ( + -- Снапшот geo + SELECT assumeNotNull(click_id) AS click_id, ... + FROM ods.geo_by_click GROUP BY click_id +) AS g ON g.click_id = c.click_id; ``` +**Ключевые особенности:** +- **UNION click_id**: собираем все уникальные click_id из обоих источников +- **LEFT JOIN**: обрабатываем случаи когда есть только device или только geo +- **`assumeNotNull`**: типобезопасное преобразование после фильтрации NULL +- **DQ-метрики**: маркируем отсутствующие данные (`device_not_found`, `geo_not_found`) + **Почему batch, а не MV:** - **Согласованность**: MV с JOIN даёт eventual consistency (данные приходят в разное время) - **Контроль**: Batch SQL можно проверить, откатить, перезапустить @@ -280,6 +328,19 @@ FROM dds.event AS e LEFT JOIN dds.click AS c ON c.click_id = e.click_id; ``` +**Материализованная таблица DQ:** + +```sql +-- Таблица для мониторинга качества (пересоздаётся при каждом batch) +TRUNCATE TABLE dm.dq_summary; +INSERT INTO dm.dq_summary +SELECT today() AS check_date, 'stg' AS layer, ... +FROM ... +``` + +- `TRUNCATE` предотвращает накопление дубликатов при повторных запусках +- Хранит статистику по всем слоям (stg/ods/dds) для быстрой проверки + **Почему VIEW:** - Для демо: достаточно производительности - Гибкость: изменения логики не требуют пересоздания таблиц @@ -392,30 +453,45 @@ erDiagram ### Сборка DDS-сущностей +**event** (browser + location): ```mermaid flowchart LR - subgraph ODS_IN["ODS (вход)"] - B[browser_event
event_id + click_id] - L[location_event
event_id] - D[device_by_click
click_id] - G[geo_by_click
click_id] + subgraph ODS["ODS"] + B["browser_event"] + L["location_event"] + end + + subgraph DDS["DDS"] + EV["event"] + end + + B -->|JOIN по event_id| EV + L -->|JOIN по event_id| EV +``` + +**click** (device + geo) с поддержкой partial data: +```mermaid +flowchart LR + subgraph ODS["ODS"] + D["device_by_click"] + G["geo_by_click"] end subgraph BUILD["Batch SQL"] - J1["JOIN по event_id"] - J2["JOIN по click_id"] + U["UNION DISTINCT
click_id"] + J["LEFT JOIN"] end - subgraph DDS_OUT["DDS (результат)"] - EV[event
всё про событие] - CL[click
всё про сессию] + subgraph DDS["DDS"] + CL["click"] end - B --> J1 - L --> J1 --> EV - B -->|click_id| J2 - D --> J2 --> CL - G --> J2 + D -->|все click_id| U + G -->|все click_id| U + U --> J + D -->|данные| J + G -->|данные| J + J --> CL ``` **Важно:** Не все `click_id` из events есть в device/geo. Используем `LEFT JOIN`. @@ -446,6 +522,36 @@ flowchart LR | **MV + JOIN** | Реалтайм | Eventual consistency, дубли при late arrival | | **Batch (выбрано)** | Согласованность, контроль | Задержка до следующего запуска | +### Обработка ошибок в ODS + +**Проблема:** Грязные данные с невалидными ключами (NULL event_id/click_id). + +**Решение:** +1. **Основная таблица**: только валидные строки (`WHERE key IS NOT NULL`) +2. **Таблица ошибок**: строки с невалидными ключами через отдельные MV +3. **DQ-метрики**: массив `parse_errors` для аудита + +```sql +-- Основная таблица +CREATE MV mv_browser_raw_to_ods_browser_event +TO ods.browser_event +SELECT ... FROM stg.browser_raw WHERE event_id IS NOT NULL; + +-- Таблица ошибок +CREATE MV mv_browser_raw_to_ods_errors +TO ods.browser_event_errors +SELECT ... FROM stg.browser_raw WHERE event_id IS NULL; +``` + +### Partial data в DDS + +**Проблема:** Device и geo события приходят независимо (не все click_id есть в обоих источниках). + +**Решение:** +1. **UNION DISTINCT** всех click_id из обоих источников +2. **LEFT JOIN** для получения данных (обрабатываем device-only и geo-only) +3. **DQ-маркеры**: `device_not_found`, `geo_not_found` в `parse_errors` + --- ## Масштабирование