From 284dc3dc1f1b8bf1924d6f331afd3f1c5dd2078d Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sun, 8 Feb 2026 16:36:01 +0300 Subject: [PATCH] Move STG->ODS to Airflow batch and align monitoring --- README.md | 13 +- dags/etl_pipeline_dag.py | 111 ++++++++++++--- docs/ARCHITECTURE.md | 69 +++++---- plans/airflow_dags_plan.md | 13 +- plans/clickhouse_ddl.md | 30 ++-- plans/runbook.md | 5 +- scripts/run_batch.sh | 75 ++++++++-- sql/ddl/ods/20_ods.sql | 281 +++++-------------------------------- sql/dm/40_dds_to_dm.sql | 23 +++ sql/ods/20_stg_to_ods.sql | 251 +++++++++++++++++++++++++++++++++ 10 files changed, 536 insertions(+), 335 deletions(-) create mode 100644 sql/ods/20_stg_to_ods.sql diff --git a/README.md b/README.md index 02867ff..af8e27e 100644 --- a/README.md +++ b/README.md @@ -9,7 +9,7 @@ Фокус проекта: быстро показать работающий end-to-end сценарий и понятным языком объяснить, как устроены слои и почему пайплайн не падает на "грязных" данных. Коротко про поток: -`data/*.jsonl` -> Kafka (1 строка = 1 сообщение) -> ClickHouse `stg` (сырые JSON) -> `ods` (типизация + DQ) -> Airflow batch -> `dds` (сущности) -> `dm` (витрины VIEW) -> Superset. +`data/*.jsonl` -> Kafka (1 строка = 1 сообщение) -> ClickHouse `stg` (сырые JSON) -> Airflow batch `stg -> ods -> dds -> dm` -> Superset. --- @@ -39,7 +39,7 @@ make data # по умолчанию первые 50 строк # или: FULL=1 make data # полный датасет (1000 строк) ``` -Запуск batch-трансформации (ODS -> DDS -> DM) в Airflow (если DAG выключен, сначала unpause): +Запуск batch-трансформации (STG -> ODS -> DDS -> DM) в Airflow (если DAG выключен, сначала unpause): ```bash docker compose exec -T airflow-webserver airflow dags trigger etl_pipeline \ --conf '{"full_refresh": true}' @@ -98,8 +98,8 @@ flowchart TB DAG[DAG: ddl_init / etl_pipeline] end - Sources -->|make data| Kafka -->|MV| STG -->|MV| ODS -->|Batch SQL| DDS -->|VIEW| DM - DAG -.->|оркестрация| ODS & DDS & DM + Sources -->|make data| Kafka -->|MV| STG -->|Batch SQL| ODS -->|Batch SQL| DDS -->|VIEW| DM + DAG -.->|оркестрация| STG & ODS & DDS & DM ``` Особенность задания про "грязные данные": парсинг не валит pipeline, ошибки фиксируются в `ods.*_errors` и в поле `parse_errors`. @@ -120,6 +120,7 @@ flowchart TB │ │ ├── ods/20_ods.sql │ │ ├── dds/30_dds.sql │ │ └── dm/40_dm.sql +│ ├── ods/ # Batch SQL: STG -> ODS │ ├── dds/ # Batch SQL: ODS -> DDS │ └── dm/ # Batch SQL: DDS -> DM ├── scripts/ # Автоматизация (apply ddl, load data, run batch) @@ -142,7 +143,7 @@ flowchart TB | `make ddl` | Применить DDL в ClickHouse (вне Airflow) | | `make data` | Загрузить данные в Kafka (50 строк) | | `FULL=1 make data` | Загрузить полный датасет | -| `make transform` | Запустить batch-процесс (вне Airflow) | +| `make transform` | Запустить batch-процесс `STG -> ODS -> DDS -> DM` (вне Airflow) | Примечания про сохранность данных: - Данные ClickHouse сохраняются в Docker volume `clickhouse-data`. @@ -216,7 +217,7 @@ flowchart LR Реализовано (Этап 1): - DAG `ddl_init`: последовательное применение DDL + проверка схемы. -- DAG `etl_pipeline`: precheck, ожидание данных в ODS, пересчёт DDS/DM, базовые проверки. +- DAG `etl_pipeline`: precheck, ожидание данных в STG, batch-пересчёт ODS/DDS/DM, базовые проверки. - Устойчивость к "грязным" данным: ошибки парсинга сохраняются в ODS, а не валят ingest. В планах (не требуется для MVP задания): diff --git a/dags/etl_pipeline_dag.py b/dags/etl_pipeline_dag.py index 586f0dc..eaff6c2 100644 --- a/dags/etl_pipeline_dag.py +++ b/dags/etl_pipeline_dag.py @@ -1,5 +1,5 @@ """ -DAG ETL-процесса ODS -> DDS -> DM для учебного проекта. +DAG ETL-процесса STG -> ODS -> DDS -> DM для учебного проекта. Принципы реализации: - SQL выполняется явными task на ClickHouseOperator; @@ -80,10 +80,61 @@ SELECT SQL_CHECK_ODS_QUALITY = """ SELECT - count() AS total_rows, - countIf(length(parse_errors) > 0) AS rows_with_errors, - round(if(count() = 0, 0, countIf(length(parse_errors) > 0) / count() * 100), 2) AS error_pct -FROM ods.browser_event + table_name, + total_rows, + rows_with_errors, + round(if(total_rows = 0, 0, rows_with_errors / total_rows * 100), 2) AS error_pct +FROM +( + SELECT + 'browser_event' AS table_name, + toFloat64(count()) AS total_rows, + toFloat64(countIf(length(parse_errors) > 0)) AS rows_with_errors + FROM ods.browser_event + UNION ALL + SELECT + 'location_event', + toFloat64(count()), + toFloat64(countIf(length(parse_errors) > 0)) + FROM ods.location_event + UNION ALL + SELECT + 'device_by_click', + toFloat64(count()), + toFloat64(countIf(length(parse_errors) > 0)) + FROM ods.device_by_click + UNION ALL + SELECT + 'geo_by_click', + toFloat64(count()), + toFloat64(countIf(length(parse_errors) > 0)) + FROM ods.geo_by_click + UNION ALL + SELECT + 'browser_event_errors', + toFloat64(count()), + toFloat64(count()) + FROM ods.browser_event_errors + UNION ALL + SELECT + 'location_event_errors', + toFloat64(count()), + toFloat64(count()) + FROM ods.location_event_errors + UNION ALL + SELECT + 'device_by_click_errors', + toFloat64(count()), + toFloat64(count()) + FROM ods.device_by_click_errors + UNION ALL + SELECT + 'geo_by_click_errors', + toFloat64(count()), + toFloat64(count()) + FROM ods.geo_by_click_errors +) +ORDER BY table_name """ SQL_TRUNCATE_DDS_CLICK = "TRUNCATE TABLE dds.click" @@ -115,21 +166,38 @@ def assert_schema_ready(**context) -> None: ) -def wait_for_ods_data(**context) -> None: +def wait_for_stg_data(**context) -> None: """ - Ожидает появления строк в ods.browser_event до заданного таймаута. - Таймаут берётся из dag_run.conf.wait_ods_timeout_sec или из params. + Ожидает появления строк в STG до заданного таймаута. + Таймаут берётся из dag_run.conf.wait_stg_timeout_sec (или legacy wait_ods_timeout_sec) + либо из params. """ dag_run = context.get("dag_run") conf = dag_run.conf if dag_run else {} - timeout_sec = int(conf.get("wait_ods_timeout_sec", context["params"]["wait_ods_timeout_sec"])) + timeout_sec = int( + conf.get( + "wait_stg_timeout_sec", + conf.get( + "wait_ods_timeout_sec", + context["params"]["wait_stg_timeout_sec"], + ), + ) + ) poll_interval_sec = 10 hook = ClickHouseHook(clickhouse_conn_id="clickhouse_default", database="default") started = time.monotonic() while True: - rows = hook.execute("SELECT count() FROM ods.browser_event") + rows = hook.execute( + """ + SELECT + (SELECT count() FROM stg.browser_raw) + + (SELECT count() FROM stg.location_raw) + + (SELECT count() FROM stg.device_raw) + + (SELECT count() FROM stg.geo_raw) AS stg_rows_total + """ + ) count_rows = int(rows[0][0]) if rows else 0 if count_rows > 0: return @@ -137,8 +205,8 @@ def wait_for_ods_data(**context) -> None: elapsed = int(time.monotonic() - started) if elapsed >= timeout_sec: raise AirflowException( - f"Таймаут ожидания ODS истёк ({timeout_sec} сек). " - "Таблица ods.browser_event всё ещё пуста." + f"Таймаут ожидания STG истёк ({timeout_sec} сек). " + "Таблицы stg.*_raw всё ещё пусты." ) time.sleep(poll_interval_sec) @@ -167,7 +235,7 @@ def assert_dm_summary_not_empty(**context) -> None: with DAG( dag_id="etl_pipeline", - description="ETL ODS -> DDS -> DM для demo-проекта", + description="ETL STG -> ODS -> DDS -> DM для demo-проекта", default_args=default_args, schedule=None, start_date=datetime(2024, 1, 1), @@ -177,7 +245,7 @@ with DAG( tags=["etl", "clickhouse", "demo"], params={ "full_refresh": Param(True, type="boolean"), - "wait_ods_timeout_sec": Param(600, type="integer", minimum=30), + "wait_stg_timeout_sec": Param(600, type="integer", minimum=30), }, ) as dag: with TaskGroup(group_id="precheck") as precheck: @@ -203,9 +271,16 @@ with DAG( check_clickhouse >> check_schema_ready_sql >> check_schema_ready with TaskGroup(group_id="transform") as transform: - wait_for_ods_data_task = PythonOperator( - task_id="wait_for_ods_data", - python_callable=wait_for_ods_data, + wait_for_stg_data_task = PythonOperator( + task_id="wait_for_stg_data", + python_callable=wait_for_stg_data, + ) + + load_ods = ClickHouseOperator( + task_id="load_ods", + sql=load_sql_statements("ods/20_stg_to_ods.sql"), + clickhouse_conn_id="clickhouse_default", + database="default", ) check_ods_quality = ClickHouseOperator( @@ -274,7 +349,7 @@ with DAG( python_callable=assert_dm_summary_not_empty, ) - wait_for_ods_data_task >> check_ods_quality >> choose_refresh_mode + wait_for_stg_data_task >> load_ods >> check_ods_quality >> choose_refresh_mode choose_refresh_mode >> truncate_dds_click >> truncate_dds_event >> truncate_complete choose_refresh_mode >> skip_truncate >> truncate_complete truncate_complete >> load_dds >> check_dds_integrity >> load_dm_summary >> validate_dm_summary_sql >> validate_dm_summary diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 6e5c37d..e652521 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -44,8 +44,10 @@ flowchart LR S2[location_raw] S3[device_raw] S4[geo_raw] - MV1[mv_*_to_ods] - MV2[mv_*_to_errors] + end + + subgraph AF["⚙️ Airflow"] + B1[load_ods
sql/ods/20_stg_to_ods.sql] end subgraph ODS["🔧 ODS"] @@ -68,12 +70,20 @@ flowchart LR DM4[v_top_pages] end - 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 + BE --> K1 --> S1 + LE --> K2 --> S2 + DE --> K3 --> S3 + GE --> K4 --> S4 + + S1 & S2 & S3 & S4 -->|Batch SQL| B1 + B1 --> O1 & O2 & O3 & O4 + B1 -.-> OE + O1 --> DE1 + O2 --> DE1 + O3 --> DC1 + O4 --> DC1 + DE1 --> DM1 + DC1 --> DM1 DE1 --> DM2 & DM3 & DM4 DC1 --> DM2 & DM3 & DM4 ``` @@ -106,7 +116,7 @@ flowchart TB DM_T["VIEW для BI
(Superset/Grafana)"] end - RAW -->|kafka-console-producer| KAFKA -->|MV| STG_T -->|MV| ODS_T + RAW -->|kafka-console-producer| KAFKA -->|MV| STG_T -->|Batch SQL (Airflow)| ODS_T ODS_T -->|argMax + JOIN| DDS_T -->|VIEW| DM_T ODS_T -.->|ошибки| DQ ``` @@ -159,14 +169,14 @@ CREATE TABLE stg.browser_raw ( | `geo_by_click` | click_id | ReplacingMergeTree(src_ingest_ts) | Гео-данные | | `*_errors` | — | MergeTree | Строки с битыми ключами | -**Materialized Views для обработки ошибок:** +**Batch шаг наполнения ODS:** -| 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` | +| Шаг | Назначение | +|-----|-----------| +| `load_ods` (`sql/ods/20_stg_to_ods.sql`) | Полная пересборка ODS из STG в рамках DAG `etl_pipeline` | +| `TRUNCATE ods.*` | Очистка перед пересборкой для детерминированного результата | +| `INSERT ... SELECT` | Типизация валидных строк в основные ODS таблицы | +| `INSERT ... SELECT` в `*_errors` | Сохранение строк с критичными ошибками парсинга | **Логика разделения:** - **Основная таблица**: строки с валидными ключами (`WHERE key IS NOT NULL`) @@ -203,7 +213,7 @@ GROUP BY error; **Почему так:** - **Изоляция источников**: изменения в одном не ломают другие - **Версионирование**: `ReplacingMergeTree` хранит последнюю версию по `src_ingest_ts` -- **Nullable ключи**: `allow_nullable_key = 1` позволяет хранить "битые" строки +- **Управляемость**: шаг `load_ods` виден в Airflow, есть task-level мониторинг и ретраи --- @@ -356,6 +366,7 @@ FROM ... sequenceDiagram participant User as Пользователь participant Make as Makefile + participant Airflow as Airflow participant K as Kafka participant CH as ClickHouse participant STG as stg.*_raw @@ -371,7 +382,7 @@ sequenceDiagram User->>Make: make ddl Make->>CH: sql/ddl/00_databases.sql Make->>CH: sql/ddl/stg/10_stg.sql (Kafka Engine) - Make->>CH: sql/ddl/ods/20_ods.sql (MV) + Make->>CH: sql/ddl/ods/20_ods.sql (таблицы ODS + drop legacy MV) Make->>CH: sql/ddl/dds/30_dds.sql Make->>CH: sql/ddl/dm/40_dm.sql CH-->>User: ✅ Структура БД создана @@ -384,17 +395,17 @@ sequenceDiagram end K->>CH: Потребление сообщений CH->>STG: INSERT через MV - STG->>ODS: INSERT через MV (типизация) K-->>User: ✅ Данные в Kafka - CH-->>User: ✅ Данные в STG/ODS + CH-->>User: ✅ Данные в STG - User->>Make: make transform - Make->>CH: sql/dds/30_ods_to_dds.sql + User->>Airflow: Trigger etl_pipeline + Airflow->>CH: sql/ods/20_stg_to_ods.sql + Airflow->>CH: sql/dds/30_ods_to_dds.sql CH->>ODS: argMax() — снапшот CH->>DDS: JOIN + INSERT - Make->>CH: sql/dm/40_dds_to_dm.sql + Airflow->>CH: sql/dm/40_dds_to_dm.sql CH->>DM: DQ summary - CH-->>User: ✅ DDS/DM обновлены + CH-->>User: ✅ ODS/DDS/DM обновлены ``` --- @@ -506,7 +517,7 @@ flowchart LR **Решение:** 1. Включаем `allow_nullable_key = 1` в `ReplacingMergeTree` -2. Фильтруем NULL в MV (`WHERE key IS NOT NULL` → основная таблица) +2. Фильтруем NULL в batch SQL (`WHERE key IS NOT NULL` → основная таблица) 3. Отдельные `*_errors` таблицы для NULL-ключей ### Почему `ReplacingMergeTree`? @@ -528,18 +539,16 @@ flowchart LR **Решение:** 1. **Основная таблица**: только валидные строки (`WHERE key IS NOT NULL`) -2. **Таблица ошибок**: строки с невалидными ключами через отдельные MV +2. **Таблица ошибок**: строки с невалидными ключами через отдельные `INSERT ... SELECT` 3. **DQ-метрики**: массив `parse_errors` для аудита ```sql -- Основная таблица -CREATE MV mv_browser_raw_to_ods_browser_event -TO ods.browser_event +INSERT INTO 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 +INSERT INTO ods.browser_event_errors SELECT ... FROM stg.browser_raw WHERE event_id IS NULL; ``` diff --git a/plans/airflow_dags_plan.md b/plans/airflow_dags_plan.md index c12a8eb..b1ecd10 100644 --- a/plans/airflow_dags_plan.md +++ b/plans/airflow_dags_plan.md @@ -46,7 +46,7 @@ | `check_clickhouse` | Проверка доступности CH (`SELECT 1`) | `PythonOperator` + `clickhouse-connect` | | `ddl_00_databases` | Создание БД | `sql/ddl/00_databases.sql` | | `ddl_10_stg` | STG + Kafka Engine + MV | `sql/ddl/stg/10_stg.sql` | -| `ddl_20_ods` | ODS + MV STG→ODS + *_errors | `sql/ddl/ods/20_ods.sql` | +| `ddl_20_ods` | ODS таблицы + drop legacy MV STG→ODS | `sql/ddl/ods/20_ods.sql` | | `ddl_30_dds` | Таблицы DDS | `sql/ddl/dds/30_dds.sql` | | `ddl_40_dm` | VIEW витрины DM | `sql/ddl/dm/40_dm.sql` | | `verify_schema` | Проверка ключевых таблиц/VIEW | SQL-check | @@ -98,7 +98,7 @@ precheck >> prepare_topics >> [load_browser_events, load_location_events, load_d ## Дизайн DAG `etl_pipeline` ### Params (через Trigger DAG with config) - `full_refresh`: bool, default `true`. -- `wait_ods_timeout_sec`: int, default `600`. +- `wait_stg_timeout_sec`: int, default `600`. ### TaskGroup `precheck` | Task ID | Что делает | Реализация | @@ -109,7 +109,8 @@ precheck >> prepare_topics >> [load_browser_events, load_location_events, load_d ### TaskGroup `transform` | Task ID | Что делает | Источник SQL | |---------|------------|--------------| -| `wait_for_ods_data` | Ожидание строк в `ods.browser_event` | SQL-check | +| `wait_for_stg_data` | Ожидание строк в `stg.*_raw` | SQL-check | +| `load_ods` | Batch STG → ODS (основные таблицы + *_errors) | `sql/ods/20_stg_to_ods.sql` | | `check_ods_quality` | Базовые DQ-метрики ODS (ошибки/total) | SQL-check | | `truncate_dds_click` | Очистка `dds.click` при `full_refresh=true` | inline SQL | | `truncate_dds_event` | Очистка `dds.event` при `full_refresh=true` | inline SQL | @@ -120,7 +121,7 @@ precheck >> prepare_topics >> [load_browser_events, load_location_events, load_d Зависимости: ```text -wait_for_ods_data >> check_ods_quality >> [truncate_dds_click, truncate_dds_event] >> load_dds >> check_dds_integrity >> load_dm_summary >> validate_dm_summary +wait_for_stg_data >> load_ods >> check_ods_quality >> [truncate_dds_click, truncate_dds_event] >> load_dds >> check_dds_integrity >> load_dm_summary >> validate_dm_summary ``` ### Итоговая цепочка `etl_pipeline` @@ -165,7 +166,7 @@ dags/ ├── __init__.py ├── ddl_init_dag.py # отдельный DAG для DDL (обязателен) ├── kafka_load_dag.py # отдельный DAG для ingest в Kafka (обязателен) -├── etl_pipeline_dag.py # основной DAG ODS -> DDS -> DM (обязателен) +├── etl_pipeline_dag.py # основной DAG STG -> ODS -> DDS -> DM (обязателен) ├── dq_monitor_dag.py # опциональный DAG мониторинга └── utils/ ├── __init__.py @@ -189,7 +190,7 @@ dags/ ## Критерии готовности - Этап 1: - В Airflow UI видны DAG `ddl_init` и `etl_pipeline`. - - `etl_pipeline` падает с понятной ошибкой, если схема не применена или ODS пуста. + - `etl_pipeline` падает с понятной ошибкой, если схема не применена или STG пуста. - После прогона `make data -> etl_pipeline`: - в `ods.browser_event` есть строки; - в `dds.click` и `dds.event` есть строки; diff --git a/plans/clickhouse_ddl.md b/plans/clickhouse_ddl.md index fdf80ff..e016977 100644 --- a/plans/clickhouse_ddl.md +++ b/plans/clickhouse_ddl.md @@ -31,16 +31,16 @@ Эта схема укладывается в задание так: -- **Kafka → STG → ODS** можно сделать полностью внутри ClickHouse через `ENGINE = Kafka` + MV (стриминг, 1 json = 1 row). -- “**регулярный процесс**” — батч‑трансформации `ODS → DDS (→ DM)` в виде SQL (`INSERT INTO … SELECT …`) по расписанию (в будущем Airflow; пока — ручной запуск). -- Airflow (когда появится) можно использовать как **тонкий оркестратор**: применить DDL и запускать batch‑SQL по расписанию. +- **Kafka → STG** реализуется внутри ClickHouse через `ENGINE = Kafka` + MV (стриминг, 1 json = 1 row). +- “**регулярный процесс**” — батч‑трансформации `STG → ODS → DDS (→ DM)` в виде SQL (`INSERT INTO … SELECT …`) по расписанию (Airflow). +- Airflow используется как основной оркестратор batch‑шагов и DQ‑проверок. ### Выбранное решение (MVP) Чтобы сделать “хорошо, но без оверинжиниринга”, фиксируем такое MVP: -- **STG → ODS** — инкрементально через MV (парсинг/типизация рядом с ingest). -- **DDS** — батч‑сборка из ODS через SQL (`INSERT INTO … SELECT …`) как “регулярный процесс”. +- **STG → ODS** — batch через SQL (`sql/ods/20_stg_to_ods.sql`) внутри `etl_pipeline`. +- **DDS** — batch‑сборка из ODS через SQL (`INSERT INTO … SELECT …`) как “регулярный процесс”. - Причина: `MV + JOIN` в `ODS → DDS` плохо переносит произвольный порядок прихода данных и может давать некорректные результаты (eventual consistency ODS, версии в разных партициях и т.п.). - **DM** — `VIEW` (витрины “на чтении”) поверх DDS, чтобы не плодить лишние таблицы и джобы под демо. - На стороне BI считаем, что запросы всегда идут с фильтрами по времени (`event_date`/`event_ts`) и не сканируют всю историю. @@ -75,10 +75,10 @@ flowchart LR v_dq[dm.v_dq_errors_daily] end - stg_browser -->|MV parse| ods_browser - stg_location -->|MV parse| ods_location - stg_device -->|MV parse| ods_device - stg_geo -->|MV parse| ods_geo + stg_browser -->|Batch SQL| ods_browser + stg_location -->|Batch SQL| ods_location + stg_device -->|Batch SQL| ods_device + stg_geo -->|Batch SQL| ods_geo ods_browser -->|Batch SQL| dds_event ods_location -->|Batch SQL| dds_event @@ -125,7 +125,7 @@ flowchart LR - `sql/ddl/00_databases.sql` — базы `stg/ods/dds/dm`. - `sql/ddl/stg/10_stg.sql` — STG raw (`stg.*_raw`) + Kafka source tables (`ENGINE = Kafka`) + MV `Kafka → STG`. -- `sql/ddl/ods/20_ods.sql` — ODS таблицы типизации + DQ (`parse_errors`) + MV `STG → ODS` + таблицы `ods_*_errors` для строк с битыми ключами. +- `sql/ddl/ods/20_ods.sql` — ODS таблицы типизации + DQ (`parse_errors`) + удаление legacy MV `STG → ODS`. - `sql/ddl/dds/30_dds.sql` — DDS таблицы (`dds.event`, `dds.click`) **без MV** (только `CREATE TABLE`). - `sql/ddl/dm/40_dm.sql` — витрины `VIEW` для Superset (`dm.v_*`). @@ -136,13 +136,17 @@ BI-ограничения (ресурсы/пользователь) **не вы - `sql/dds/30_ods_to_dds.sql` — регулярная батч‑сборка DDS из ODS: - получить “последнюю версию” строк по ключам (`event_id`/`click_id`) через `argMax(..., src_ingest_ts)` (или эквивалент); - выполнить join snapshot’ов и загрузить в `dds.event`/`dds.click` (для демо возможно “full rebuild”; позже — инкрементально). +- `sql/ods/20_stg_to_ods.sql` — регулярная батч‑сборка ODS из STG: + - очистить ODS (`TRUNCATE`) перед пересборкой; + - типизировать валидные строки в `ods.*`; + - сложить критичные ошибки парсинга в `ods.*_errors`. ### Исполнение DDL (make сейчас / Airflow потом) Требования к файлам `sql/*/*.sql`: - идемпотентность (`IF NOT EXISTS`), чтобы повторные прогоны были безопасны; -- строгий порядок исполнения: `00 → 10 → 20 → 30 → 40` (из‑за зависимостей MV); +- строгий порядок исполнения: `00 → 10 → 20 → 30 → 40` (из‑за зависимостей объектов); - единые имена топиков Kafka: `browser_events`, `location_events`, `device_events`, `geo_events` (их создаёт `make data`). ### Дедупликация и обработка “битых” ключей (ODS) @@ -167,8 +171,8 @@ BI-ограничения (ресурсы/пользователь) **не вы Batch‑трансформации (целевое, для реализации следующим шагом): -- `make transform` (или аналогичная команда) запускает `sql/dds/30_ods_to_dds.sql` через `clickhouse-client`; -- в будущем Airflow будет делать то же самое по расписанию (один job‑SQL = один task). +- `make transform` (или аналогичная команда) запускает `sql/ods/20_stg_to_ods.sql`, `sql/dds/30_ods_to_dds.sql` и `sql/dm/40_dds_to_dm.sql` через `clickhouse-client`; +- Airflow может выполнять те же шаги как отдельные task (с ретраями и мониторингом). ### Параметры окружения (docker compose) diff --git a/plans/runbook.md b/plans/runbook.md index bfe83ae..c191645 100644 --- a/plans/runbook.md +++ b/plans/runbook.md @@ -32,8 +32,9 @@ make ddl ## Make таргеты - `make up` — `docker compose up -d` (поднимает весь стек из `docker-compose.yml`). -- `make ddl` — применяет SQL из `plans/clickhouse_ddl.md` в контейнер ClickHouse (извлекает все блоки ```sql``` и исполняет их через `clickhouse-client`). +- `make ddl` — применяет исполняемые SQL-файлы из `sql/ddl/*` в контейнер ClickHouse через `clickhouse-client`. - `make data` — пересоздаёт топики (по умолчанию) и публикует события из `data/*.jsonl` в Kafka (1 строка = 1 Kafka message value). +- `make transform` — выполняет batch-процесс `STG -> ODS -> DDS -> DM` через `scripts/run_batch.sh`. План реализации механики заливки (дизайн/решения): `plans/kafka_ingest_plan.md`. @@ -75,7 +76,7 @@ RESET_TOPICS=0 make data ## Применение DDL в ClickHouse (`make ddl`) -Скрипт исполняет SQL из `plans/clickhouse_ddl.md`. Для `ENGINE = Kafka` важно, чтобы `kafka_broker_list` был доступен из контейнера ClickHouse. +Скрипт исполняет SQL-файлы из `sql/ddl/*` по фиксированному порядку. Для `ENGINE = Kafka` важно, чтобы `kafka_broker_list` был доступен из контейнера ClickHouse. В текущем compose: diff --git a/scripts/run_batch.sh b/scripts/run_batch.sh index 9c36b60..885e8da 100755 --- a/scripts/run_batch.sh +++ b/scripts/run_batch.sh @@ -1,11 +1,12 @@ #!/usr/bin/env bash # -# Скрипт batch-трансформации данных: ODS → DDS → DM +# Скрипт batch-трансформации данных: STG → ODS → DDS → DM # # Назначение: -# Запускает SQL-скрипты из sql/dds и sql/dm для преобразования данных между слоями: -# 1. ODS → DDS : Сборка сущностей из типизированных данных -# 2. DDS → DM : Обновление сводки по качеству данных (dq_summary) +# Запускает SQL-скрипты из sql/ods, sql/dds и sql/dm для преобразования данных между слоями: +# 1. STG → ODS : Типизация и перенос ошибок в *_errors +# 2. ODS → DDS : Сборка сущностей из типизированных данных +# 3. DDS → DM : Обновление сводки по качеству данных (dq_summary) # # Как запускать: # make transform @@ -13,7 +14,7 @@ # # Требования: # - ClickHouse запущен (make up) -# - ODS содержит данные (make data выполнен) +# - STG содержит данные (make data выполнен) # # Стратегия: # Сейчас: полная перезагрузка (TRUNCATE + INSERT) — для демо @@ -25,6 +26,7 @@ set -euo pipefail # Директория со скриптом SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" SQL_ROOT_DIR="${SCRIPT_DIR}/../sql" +ODS_TRANSFORM_SQL="${SQL_ROOT_DIR}/ods/20_stg_to_ods.sql" DDS_TRANSFORM_SQL="${SQL_ROOT_DIR}/dds/30_ods_to_dds.sql" DM_TRANSFORM_SQL="${SQL_ROOT_DIR}/dm/40_dds_to_dm.sql" @@ -47,6 +49,11 @@ fi # ----------------------------------------------------------------------------- # Проверка: SQL-файлы batch существуют? # ----------------------------------------------------------------------------- +if [[ ! -f "${ODS_TRANSFORM_SQL}" ]]; then + echo "Ошибка: Не найден SQL-файл: ${ODS_TRANSFORM_SQL}" + exit 1 +fi + if [[ ! -f "${DDS_TRANSFORM_SQL}" ]]; then echo "Ошибка: Не найден SQL-файл: ${DDS_TRANSFORM_SQL}" exit 1 @@ -58,27 +65,65 @@ if [[ ! -f "${DM_TRANSFORM_SQL}" ]]; then fi # ----------------------------------------------------------------------------- -# Проверка: в ODS есть данные? +# Проверка: в STG есть данные? # ----------------------------------------------------------------------------- -ODS_COUNT=$(${COMPOSE_BIN} exec -T "${CLICKHOUSE_SERVICE}" clickhouse-client \ +STG_COUNT=$(${COMPOSE_BIN} exec -T "${CLICKHOUSE_SERVICE}" clickhouse-client \ --user="${CLICKHOUSE_USER}" \ --password="${CLICKHOUSE_PASSWORD}" \ --database="${CLICKHOUSE_DB}" \ - --query="SELECT count() FROM ods.browser_event" 2>/dev/null || echo "0") + --query=" + SELECT + (SELECT count() FROM stg.browser_raw) + + (SELECT count() FROM stg.location_raw) + + (SELECT count() FROM stg.device_raw) + + (SELECT count() FROM stg.geo_raw) AS stg_rows_total + " 2>/dev/null || echo "0") -if [[ "${ODS_COUNT}" == "0" ]]; then - echo "Предупреждение: Таблица ODS.browser_event пуста." +if [[ "${STG_COUNT}" == "0" ]]; then + echo "Предупреждение: Таблицы STG пусты." echo "Сначала загрузите данные: make data" exit 1 fi -echo "Найдено ${ODS_COUNT} строк в ODS.browser_event" +echo "Найдено ${STG_COUNT} строк в STG (суммарно по 4 потокам)" echo "" # ----------------------------------------------------------------------------- -# Шаг 1: ODS → DDS (сборка сущностей) +# Шаг 1: STG → ODS (типизация и DQ) # ----------------------------------------------------------------------------- -echo "Шаг 1: Обновление DDS слоя (ODS → DDS)..." +echo "Шаг 1: Обновление ODS слоя (STG → ODS)..." +${COMPOSE_BIN} exec -T "${CLICKHOUSE_SERVICE}" clickhouse-client \ + --user="${CLICKHOUSE_USER}" \ + --password="${CLICKHOUSE_PASSWORD}" \ + --database="${CLICKHOUSE_DB}" \ + --multiquery < "${ODS_TRANSFORM_SQL}" + +echo " ✓ ODS обновлён" + +# Показываем статистику ODS +echo "" +echo "Статистика ODS:" +${COMPOSE_BIN} exec -T "${CLICKHOUSE_SERVICE}" clickhouse-client \ + --user="${CLICKHOUSE_USER}" \ + --password="${CLICKHOUSE_PASSWORD}" \ + --database="${CLICKHOUSE_DB}" \ + --query=" + SELECT 'ods.browser_event' AS table, count() AS rows FROM ods.browser_event + UNION ALL SELECT 'ods.location_event', count() FROM ods.location_event + UNION ALL SELECT 'ods.device_by_click', count() FROM ods.device_by_click + UNION ALL SELECT 'ods.geo_by_click', count() FROM ods.geo_by_click + UNION ALL SELECT 'ods.browser_event_errors', count() FROM ods.browser_event_errors + UNION ALL SELECT 'ods.location_event_errors', count() FROM ods.location_event_errors + UNION ALL SELECT 'ods.device_by_click_errors', count() FROM ods.device_by_click_errors + UNION ALL SELECT 'ods.geo_by_click_errors', count() FROM ods.geo_by_click_errors + FORMAT PrettyCompact + " + +# ----------------------------------------------------------------------------- +# Шаг 2: ODS → DDS (сборка сущностей) +# ----------------------------------------------------------------------------- +echo "" +echo "Шаг 2: Обновление DDS слоя (ODS → DDS)..." echo " - Очистка текущих данных (TRUNCATE)..." # Очищаем таблицы перед загрузкой (полная перезагрузка для демо) @@ -112,10 +157,10 @@ ${COMPOSE_BIN} exec -T "${CLICKHOUSE_SERVICE}" clickhouse-client \ --query="SELECT 'dds.click' AS table, count() AS rows FROM dds.click UNION ALL SELECT 'dds.event', count() FROM dds.event FORMAT PrettyCompact" # ----------------------------------------------------------------------------- -# Шаг 2: DDS → DM (сводка по качеству) +# Шаг 3: DDS → DM (сводка по качеству) # ----------------------------------------------------------------------------- echo "" -echo "Шаг 2: Обновление DM слоя (DQ summary)..." +echo "Шаг 3: Обновление DM слоя (DQ summary)..." ${COMPOSE_BIN} exec -T "${CLICKHOUSE_SERVICE}" clickhouse-client \ --user="${CLICKHOUSE_USER}" \ --password="${CLICKHOUSE_PASSWORD}" \ diff --git a/sql/ddl/ods/20_ods.sql b/sql/ddl/ods/20_ods.sql index 83bea0d..f9df01c 100644 --- a/sql/ddl/ods/20_ods.sql +++ b/sql/ddl/ods/20_ods.sql @@ -2,15 +2,28 @@ -- Слой ODS (Operational Data Store) — типизированные данные + дедупликация + DQ -- ============================================================================ -- Назначение: --- - Типизация данных из STG (String → UUID, DateTime, etc.) +-- - Хранение типизированных данных для дальнейшей сборки DDS -- - Дедупликация через ReplacingMergeTree (последняя версия по src_ingest_ts) -- - Контроль качества: массив parse_errors для "грязных" данных -- - Разделение: валидные строки → основная таблица, ошибки → *_errors -- --- Поток данных: --- STG (*_raw) → MV → ODS (основная таблица + error_tables) +-- Важно: +-- - Наполнение ODS выполняется batch-процессом из sql/ods/20_stg_to_ods.sql +-- - Materialized View STG → ODS в target-state не используются -- ============================================================================ +-- ---------------------------------------------------------------------------- +-- Удаляем legacy MV STG → ODS (если ранее были созданы) +-- ---------------------------------------------------------------------------- +DROP TABLE IF EXISTS stg.mv_browser_raw_to_ods; +DROP TABLE IF EXISTS stg.mv_browser_raw_to_ods_errors; +DROP TABLE IF EXISTS stg.mv_location_raw_to_ods; +DROP TABLE IF EXISTS stg.mv_location_raw_to_ods_errors; +DROP TABLE IF EXISTS stg.mv_device_raw_to_ods; +DROP TABLE IF EXISTS stg.mv_device_raw_to_ods_errors; +DROP TABLE IF EXISTS stg.mv_geo_raw_to_ods; +DROP TABLE IF EXISTS stg.mv_geo_raw_to_ods_errors; + -- ============================================================================ -- BROWSER EVENTS -- ============================================================================ @@ -18,9 +31,6 @@ -- ---------------------------------------------------------------------------- -- Основная таблица: валидные строки (event_id IS NOT NULL) -- ---------------------------------------------------------------------------- --- ReplacingMergeTree: при мердже оставляет строку с максимальным src_ingest_ts --- allow_nullable_key = 1: разрешаем NULL в ключе (ClickHouse по умолчанию запрещает) --- ---------------------------------------------------------------------------- CREATE TABLE IF NOT EXISTS ods.browser_event ( event_id Nullable(UUID), -- UUID события (ключ) @@ -35,101 +45,28 @@ CREATE TABLE IF NOT EXISTS ods.browser_event src_raw String, -- Исходный JSON для аудита parse_errors Array(LowCardinality(String)) -- Ошибки парсинга (если есть) ) -ENGINE = ReplacingMergeTree(src_ingest_ts) -- Движок дедупликации по версии -PARTITION BY toYYYYMM(event_date) -- Партиции по месяцу для быстрой очистки -ORDER BY (event_id) -- Ключ сортировки (и дедупликации) -SETTINGS allow_nullable_key = 1; -- Разрешаем NULL в ключе (для "битых" данных) +ENGINE = ReplacingMergeTree(src_ingest_ts) +PARTITION BY toYYYYMM(event_date) +ORDER BY (event_id) +SETTINGS allow_nullable_key = 1; -- ---------------------------------------------------------------------------- --- Таблица ошибок: строки с невалидными ключами (event_id IS NULL) --- ---------------------------------------------------------------------------- --- Сохраняем полную информацию для анализа проблем с данными --- MergeTree без Replacing: сохраняем все ошибки (не дедуплицируем) +-- Таблица ошибок: строки с невалидными ключами/критичными ошибками browser -- ---------------------------------------------------------------------------- CREATE TABLE IF NOT EXISTS ods.browser_event_errors ( ingest_ts DateTime64(3), -- Время вставки в ClickHouse kafka_topic LowCardinality(String), -- Топик Kafka (для отслеживания источника) kafka_partition Int32, -- Партиция Kafka - kafka_offset Int64, -- Смещение Kafka (уникальный идентификатор сообщения) + kafka_offset Int64, -- Смещение Kafka (идентификатор сообщения) kafka_ts DateTime64(3), -- Время из Kafka raw String, -- Исходный JSON - error_reason LowCardinality(String) -- Описание ошибки (что именно не распарсилось) + error_reason LowCardinality(String) -- Описание ошибки ) ENGINE = MergeTree PARTITION BY toYYYYMM(ingest_ts) ORDER BY (ingest_ts, kafka_topic, kafka_partition, kafka_offset); --- ---------------------------------------------------------------------------- --- MV: STG → ODS (основная таблица) --- ---------------------------------------------------------------------------- --- Фильтруем только валидные строки: WHERE event_id IS NOT NULL --- Парсим JSON, типизируем поля, собираем ошибки в массив parse_errors --- ---------------------------------------------------------------------------- -CREATE MATERIALIZED VIEW IF NOT EXISTS stg.mv_browser_raw_to_ods -TO ods.browser_event -AS -WITH - toUUIDOrNull(JSONExtractString(raw, 'event_id')) AS event_id, - parseDateTime64BestEffortOrNull(JSONExtractString(raw, 'event_timestamp'), 6) AS event_ts, - JSONExtractString(raw, 'event_type') AS event_type, - toUUIDOrNull(JSONExtractString(raw, 'click_id')) AS click_id, - JSONExtractString(raw, 'browser_name') AS browser_name, - JSONExtractString(raw, 'browser_user_agent') AS browser_user_agent, - JSONExtractString(raw, 'browser_language') AS browser_language -SELECT - event_id, - event_ts, - event_type, - click_id, - browser_name, - browser_user_agent, - browser_language, - ingest_ts AS src_ingest_ts, - raw AS src_raw, - -- Собираем ошибки парсинга в массив (пустые строки фильтруем) - arrayFilter(x -> x != '', [ - if(event_id IS NULL, 'bad_event_id', ''), - if(event_ts IS NULL, 'bad_event_timestamp', ''), - if(click_id IS NULL, 'bad_click_id', '') - ]) AS parse_errors -FROM stg.browser_raw -WHERE event_id IS NOT NULL; -- Только валидные строки (NULL → в error_tables) - --- ---------------------------------------------------------------------------- --- MV: STG → ODS (таблица ошибок) --- ---------------------------------------------------------------------------- --- Перенаправляем строки с ошибками парсинга в отдельную таблицу --- Это позволяет не терять данные и анализировать проблемы --- ---------------------------------------------------------------------------- -CREATE MATERIALIZED VIEW IF NOT EXISTS stg.mv_browser_raw_to_ods_errors -TO ods.browser_event_errors -AS -WITH - toUUIDOrNull(JSONExtractString(raw, 'event_id')) AS event_id, - parseDateTime64BestEffortOrNull(JSONExtractString(raw, 'event_timestamp'), 6) AS event_ts, - toUUIDOrNull(JSONExtractString(raw, 'click_id')) AS click_id, - arrayFilter(x -> x != '', [ - if(event_id IS NULL, 'bad_event_id', ''), - if(event_ts IS NULL, 'bad_event_timestamp', ''), - if(click_id IS NULL, 'bad_click_id', '') - ]) AS parse_errors -SELECT - ingest_ts, - kafka_topic, - kafka_partition, - kafka_offset, - kafka_ts, - raw, - arrayStringConcat(parse_errors, ',') AS error_reason -FROM stg.browser_raw -WHERE length(parse_errors) > 0 -- Есть хотя бы одна ошибка - AND ( - event_id IS NULL - OR event_ts IS NULL - OR click_id IS NULL - ); - -- ============================================================================ -- LOCATION EVENTS -- ============================================================================ @@ -137,9 +74,9 @@ WHERE length(parse_errors) > 0 -- Есть хотя бы одна ош CREATE TABLE IF NOT EXISTS ods.location_event ( event_id Nullable(UUID), - page_url Nullable(String), -- Полный URL страницы + page_url Nullable(String), -- Полный URL страницы page_url_path LowCardinality(Nullable(String)), -- Путь (/home, /product и т.д.) - referer_url Nullable(String), -- Откуда пришёл пользователь + referer_url Nullable(String), -- Откуда пришёл пользователь referer_medium LowCardinality(Nullable(String)), -- Тип referer (internal, search и т.д.) utm_medium LowCardinality(Nullable(String)), -- UTM medium (cpc, organic и т.д.) utm_source LowCardinality(Nullable(String)), -- UTM source (google, mailchimp и т.д.) @@ -150,7 +87,7 @@ CREATE TABLE IF NOT EXISTS ods.location_event parse_errors Array(LowCardinality(String)) ) ENGINE = ReplacingMergeTree(src_ingest_ts) -PARTITION BY toYYYYMM(toDate(src_ingest_ts)) -- Партиция по времени загрузки (нет event_date) +PARTITION BY toYYYYMM(toDate(src_ingest_ts)) ORDER BY (event_id) SETTINGS allow_nullable_key = 1; @@ -168,49 +105,6 @@ ENGINE = MergeTree PARTITION BY toYYYYMM(ingest_ts) ORDER BY (ingest_ts, kafka_topic, kafka_partition, kafka_offset); -CREATE MATERIALIZED VIEW IF NOT EXISTS stg.mv_location_raw_to_ods -TO ods.location_event -AS -WITH - toUUIDOrNull(JSONExtractString(raw, 'event_id')) AS event_id -SELECT - event_id, - JSONExtractString(raw, 'page_url') AS page_url, - JSONExtractString(raw, 'page_url_path') AS page_url_path, - JSONExtractString(raw, 'referer_url') AS referer_url, - JSONExtractString(raw, 'referer_medium') AS referer_medium, - JSONExtractString(raw, 'utm_medium') AS utm_medium, - JSONExtractString(raw, 'utm_source') AS utm_source, - JSONExtractString(raw, 'utm_content') AS utm_content, - JSONExtractString(raw, 'utm_campaign') AS utm_campaign, - ingest_ts AS src_ingest_ts, - raw AS src_raw, - arrayFilter(x -> x != '', [ - if(event_id IS NULL, 'bad_event_id', '') - ]) AS parse_errors -FROM stg.location_raw -WHERE event_id IS NOT NULL; - -CREATE MATERIALIZED VIEW IF NOT EXISTS stg.mv_location_raw_to_ods_errors -TO ods.location_event_errors -AS -WITH - toUUIDOrNull(JSONExtractString(raw, 'event_id')) AS event_id, - arrayFilter(x -> x != '', [ - if(event_id IS NULL, 'bad_event_id', '') - ]) AS parse_errors -SELECT - ingest_ts, - kafka_topic, - kafka_partition, - kafka_offset, - kafka_ts, - raw, - arrayStringConcat(parse_errors, ',') AS error_reason -FROM stg.location_raw -WHERE length(parse_errors) > 0 - AND event_id IS NULL; - -- ============================================================================ -- DEVICE EVENTS -- ============================================================================ @@ -218,13 +112,13 @@ WHERE length(parse_errors) > 0 CREATE TABLE IF NOT EXISTS ods.device_by_click ( click_id Nullable(UUID), - os Nullable(String), -- Полное название ОС + os Nullable(String), -- Полное название ОС os_name LowCardinality(Nullable(String)), -- Короткое название (Windows, iOS и т.д.) os_timezone LowCardinality(Nullable(String)), -- Таймзона пользователя device_type LowCardinality(Nullable(String)), -- Mobile, Computer, Tablet - device_is_mobile Nullable(UInt8), -- 1 = мобильное, 0 = десктоп - user_custom_id Nullable(String), -- Email или username - user_domain_id Nullable(UUID), -- UUID пользователя в системе + device_is_mobile Nullable(UInt8), -- 1 = мобильное, 0 = десктоп + user_custom_id Nullable(String), -- Email или username + user_domain_id Nullable(UUID), -- UUID пользователя в системе src_ingest_ts DateTime64(3), src_raw String, parse_errors Array(LowCardinality(String)) @@ -248,56 +142,6 @@ ENGINE = MergeTree PARTITION BY toYYYYMM(ingest_ts) ORDER BY (ingest_ts, kafka_topic, kafka_partition, kafka_offset); -CREATE MATERIALIZED VIEW IF NOT EXISTS stg.mv_device_raw_to_ods -TO ods.device_by_click -AS -WITH - toUUIDOrNull(JSONExtractString(raw, 'click_id')) AS click_id, - JSONExtract(raw, 'device_is_mobile', 'Nullable(UInt8)') AS device_is_mobile, - toUUIDOrNull(JSONExtractString(raw, 'user_domain_id')) AS user_domain_id -SELECT - click_id, - JSONExtractString(raw, 'os') AS os, - JSONExtractString(raw, 'os_name') AS os_name, - JSONExtractString(raw, 'os_timezone') AS os_timezone, - JSONExtractString(raw, 'device_type') AS device_type, - device_is_mobile, - JSONExtractString(raw, 'user_custom_id') AS user_custom_id, - user_domain_id, - ingest_ts AS src_ingest_ts, - raw AS src_raw, - arrayFilter(x -> x != '', [ - if(click_id IS NULL, 'bad_click_id', ''), - if(user_domain_id IS NULL, 'bad_user_domain_id', '') - ]) AS parse_errors -FROM stg.device_raw -WHERE click_id IS NOT NULL; - -CREATE MATERIALIZED VIEW IF NOT EXISTS stg.mv_device_raw_to_ods_errors -TO ods.device_by_click_errors -AS -WITH - toUUIDOrNull(JSONExtractString(raw, 'click_id')) AS click_id, - toUUIDOrNull(JSONExtractString(raw, 'user_domain_id')) AS user_domain_id, - arrayFilter(x -> x != '', [ - if(click_id IS NULL, 'bad_click_id', ''), - if(user_domain_id IS NULL, 'bad_user_domain_id', '') - ]) AS parse_errors -SELECT - ingest_ts, - kafka_topic, - kafka_partition, - kafka_offset, - kafka_ts, - raw, - arrayStringConcat(parse_errors, ',') AS error_reason -FROM stg.device_raw -WHERE length(parse_errors) > 0 - AND ( - click_id IS NULL - OR user_domain_id IS NULL - ); - -- ============================================================================ -- GEO EVENTS -- ============================================================================ @@ -305,12 +149,12 @@ WHERE length(parse_errors) > 0 CREATE TABLE IF NOT EXISTS ods.geo_by_click ( click_id Nullable(UUID), - geo_latitude Nullable(Float64), -- Широта - geo_longitude Nullable(Float64), -- Долгота - geo_country LowCardinality(Nullable(String)), -- Код страны (RU, US и т.д.) - geo_timezone LowCardinality(Nullable(String)), -- Таймзона - geo_region_name Nullable(String), -- Название региона/города - ip_address Nullable(String), -- IP адрес + geo_latitude Nullable(Float64), -- Широта + geo_longitude Nullable(Float64), -- Долгота + geo_country LowCardinality(Nullable(String)), -- Код страны (RU, US и т.д.) + geo_timezone LowCardinality(Nullable(String)), -- Таймзона + geo_region_name Nullable(String), -- Название региона/города + ip_address Nullable(String), -- IP адрес src_ingest_ts DateTime64(3), src_raw String, parse_errors Array(LowCardinality(String)) @@ -333,56 +177,3 @@ CREATE TABLE IF NOT EXISTS ods.geo_by_click_errors ENGINE = MergeTree PARTITION BY toYYYYMM(ingest_ts) ORDER BY (ingest_ts, kafka_topic, kafka_partition, kafka_offset); - -CREATE MATERIALIZED VIEW IF NOT EXISTS stg.mv_geo_raw_to_ods -TO ods.geo_by_click -AS -WITH - toUUIDOrNull(JSONExtractString(raw, 'click_id')) AS click_id, - toFloat64OrNull(JSONExtractString(raw, 'geo_latitude')) AS geo_latitude, - toFloat64OrNull(JSONExtractString(raw, 'geo_longitude')) AS geo_longitude -SELECT - click_id, - geo_latitude, - geo_longitude, - JSONExtractString(raw, 'geo_country') AS geo_country, - JSONExtractString(raw, 'geo_timezone') AS geo_timezone, - JSONExtractString(raw, 'geo_region_name') AS geo_region_name, - JSONExtractString(raw, 'ip_address') AS ip_address, - ingest_ts AS src_ingest_ts, - raw AS src_raw, - arrayFilter(x -> x != '', [ - if(click_id IS NULL, 'bad_click_id', ''), - if(geo_latitude IS NULL, 'bad_geo_latitude', ''), - if(geo_longitude IS NULL, 'bad_geo_longitude', '') - ]) AS parse_errors -FROM stg.geo_raw -WHERE click_id IS NOT NULL; - -CREATE MATERIALIZED VIEW IF NOT EXISTS stg.mv_geo_raw_to_ods_errors -TO ods.geo_by_click_errors -AS -WITH - toUUIDOrNull(JSONExtractString(raw, 'click_id')) AS click_id, - toFloat64OrNull(JSONExtractString(raw, 'geo_latitude')) AS geo_latitude, - toFloat64OrNull(JSONExtractString(raw, 'geo_longitude')) AS geo_longitude, - arrayFilter(x -> x != '', [ - if(click_id IS NULL, 'bad_click_id', ''), - if(geo_latitude IS NULL, 'bad_geo_latitude', ''), - if(geo_longitude IS NULL, 'bad_geo_longitude', '') - ]) AS parse_errors -SELECT - ingest_ts, - kafka_topic, - kafka_partition, - kafka_offset, - kafka_ts, - raw, - arrayStringConcat(parse_errors, ',') AS error_reason -FROM stg.geo_raw -WHERE length(parse_errors) > 0 - AND ( - click_id IS NULL - OR geo_latitude IS NULL - OR geo_longitude IS NULL - ); diff --git a/sql/dm/40_dds_to_dm.sql b/sql/dm/40_dds_to_dm.sql index 0edece8..58d20d7 100644 --- a/sql/dm/40_dds_to_dm.sql +++ b/sql/dm/40_dds_to_dm.sql @@ -85,10 +85,33 @@ 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 diff --git a/sql/ods/20_stg_to_ods.sql b/sql/ods/20_stg_to_ods.sql new file mode 100644 index 0000000..50b4694 --- /dev/null +++ b/sql/ods/20_stg_to_ods.sql @@ -0,0 +1,251 @@ +-- ============================================================================ +-- Batch-трансформация: STG → ODS +-- ============================================================================ +-- Назначение: +-- - Перенос типизации STG → ODS из Materialized View в управляемый batch +-- - Полная пересборка ODS для прозрачного мониторинга в Airflow +-- - Сохранение DQ-логики: parse_errors + отдельные *_errors таблицы +-- +-- Когда запускать: +-- - В DAG etl_pipeline перед ODS → DDS +-- - После загрузки очередного среза данных в Kafka/STG +-- ============================================================================ + +-- ---------------------------------------------------------------------------- +-- Подготовка: очищаем ODS перед полной пересборкой +-- ---------------------------------------------------------------------------- +TRUNCATE TABLE ods.browser_event; +TRUNCATE TABLE ods.browser_event_errors; +TRUNCATE TABLE ods.location_event; +TRUNCATE TABLE ods.location_event_errors; +TRUNCATE TABLE ods.device_by_click; +TRUNCATE TABLE ods.device_by_click_errors; +TRUNCATE TABLE ods.geo_by_click; +TRUNCATE TABLE ods.geo_by_click_errors; + +-- ============================================================================ +-- BROWSER EVENTS +-- ============================================================================ + +-- ---------------------------------------------------------------------------- +-- Основная таблица: валидные ключи event_id +-- ---------------------------------------------------------------------------- +INSERT INTO ods.browser_event +WITH + toUUIDOrNull(JSONExtractString(raw, 'event_id')) AS event_id, + parseDateTime64BestEffortOrNull(JSONExtractString(raw, 'event_timestamp'), 6) AS event_ts, + JSONExtractString(raw, 'event_type') AS event_type, + toUUIDOrNull(JSONExtractString(raw, 'click_id')) AS click_id, + JSONExtractString(raw, 'browser_name') AS browser_name, + JSONExtractString(raw, 'browser_user_agent') AS browser_user_agent, + JSONExtractString(raw, 'browser_language') AS browser_language +SELECT + event_id, + event_ts, + event_type, + click_id, + browser_name, + browser_user_agent, + browser_language, + ingest_ts AS src_ingest_ts, + raw AS src_raw, + arrayFilter(x -> x != '', [ + if(event_id IS NULL, 'bad_event_id', ''), + if(event_ts IS NULL, 'bad_event_timestamp', ''), + if(click_id IS NULL, 'bad_click_id', '') + ]) AS parse_errors +FROM stg.browser_raw +WHERE event_id IS NOT NULL; + +-- ---------------------------------------------------------------------------- +-- Таблица ошибок: строки с ошибками парсинга browser +-- ---------------------------------------------------------------------------- +INSERT INTO ods.browser_event_errors +WITH + toUUIDOrNull(JSONExtractString(raw, 'event_id')) AS event_id, + parseDateTime64BestEffortOrNull(JSONExtractString(raw, 'event_timestamp'), 6) AS event_ts, + toUUIDOrNull(JSONExtractString(raw, 'click_id')) AS click_id, + arrayFilter(x -> x != '', [ + if(event_id IS NULL, 'bad_event_id', ''), + if(event_ts IS NULL, 'bad_event_timestamp', ''), + if(click_id IS NULL, 'bad_click_id', '') + ]) AS parse_errors +SELECT + ingest_ts, + kafka_topic, + kafka_partition, + kafka_offset, + kafka_ts, + raw, + arrayStringConcat(parse_errors, ',') AS error_reason +FROM stg.browser_raw +WHERE length(parse_errors) > 0 + AND ( + event_id IS NULL + OR event_ts IS NULL + OR click_id IS NULL + ); + +-- ============================================================================ +-- LOCATION EVENTS +-- ============================================================================ + +-- ---------------------------------------------------------------------------- +-- Основная таблица: валидные ключи event_id +-- ---------------------------------------------------------------------------- +INSERT INTO ods.location_event +WITH + toUUIDOrNull(JSONExtractString(raw, 'event_id')) AS event_id +SELECT + event_id, + JSONExtractString(raw, 'page_url') AS page_url, + JSONExtractString(raw, 'page_url_path') AS page_url_path, + JSONExtractString(raw, 'referer_url') AS referer_url, + JSONExtractString(raw, 'referer_medium') AS referer_medium, + JSONExtractString(raw, 'utm_medium') AS utm_medium, + JSONExtractString(raw, 'utm_source') AS utm_source, + JSONExtractString(raw, 'utm_content') AS utm_content, + JSONExtractString(raw, 'utm_campaign') AS utm_campaign, + ingest_ts AS src_ingest_ts, + raw AS src_raw, + arrayFilter(x -> x != '', [ + if(event_id IS NULL, 'bad_event_id', '') + ]) AS parse_errors +FROM stg.location_raw +WHERE event_id IS NOT NULL; + +-- ---------------------------------------------------------------------------- +-- Таблица ошибок: строки с ошибками парсинга location +-- ---------------------------------------------------------------------------- +INSERT INTO ods.location_event_errors +WITH + toUUIDOrNull(JSONExtractString(raw, 'event_id')) AS event_id, + arrayFilter(x -> x != '', [ + if(event_id IS NULL, 'bad_event_id', '') + ]) AS parse_errors +SELECT + ingest_ts, + kafka_topic, + kafka_partition, + kafka_offset, + kafka_ts, + raw, + arrayStringConcat(parse_errors, ',') AS error_reason +FROM stg.location_raw +WHERE length(parse_errors) > 0 + AND event_id IS NULL; + +-- ============================================================================ +-- DEVICE EVENTS +-- ============================================================================ + +-- ---------------------------------------------------------------------------- +-- Основная таблица: валидные ключи click_id +-- ---------------------------------------------------------------------------- +INSERT INTO ods.device_by_click +WITH + toUUIDOrNull(JSONExtractString(raw, 'click_id')) AS click_id, + JSONExtract(raw, 'device_is_mobile', 'Nullable(UInt8)') AS device_is_mobile, + toUUIDOrNull(JSONExtractString(raw, 'user_domain_id')) AS user_domain_id +SELECT + click_id, + JSONExtractString(raw, 'os') AS os, + JSONExtractString(raw, 'os_name') AS os_name, + JSONExtractString(raw, 'os_timezone') AS os_timezone, + JSONExtractString(raw, 'device_type') AS device_type, + device_is_mobile, + JSONExtractString(raw, 'user_custom_id') AS user_custom_id, + user_domain_id, + ingest_ts AS src_ingest_ts, + raw AS src_raw, + arrayFilter(x -> x != '', [ + if(click_id IS NULL, 'bad_click_id', ''), + if(user_domain_id IS NULL, 'bad_user_domain_id', '') + ]) AS parse_errors +FROM stg.device_raw +WHERE click_id IS NOT NULL; + +-- ---------------------------------------------------------------------------- +-- Таблица ошибок: строки с ошибками парсинга device +-- ---------------------------------------------------------------------------- +INSERT INTO ods.device_by_click_errors +WITH + toUUIDOrNull(JSONExtractString(raw, 'click_id')) AS click_id, + toUUIDOrNull(JSONExtractString(raw, 'user_domain_id')) AS user_domain_id, + arrayFilter(x -> x != '', [ + if(click_id IS NULL, 'bad_click_id', ''), + if(user_domain_id IS NULL, 'bad_user_domain_id', '') + ]) AS parse_errors +SELECT + ingest_ts, + kafka_topic, + kafka_partition, + kafka_offset, + kafka_ts, + raw, + arrayStringConcat(parse_errors, ',') AS error_reason +FROM stg.device_raw +WHERE length(parse_errors) > 0 + AND ( + click_id IS NULL + OR user_domain_id IS NULL + ); + +-- ============================================================================ +-- GEO EVENTS +-- ============================================================================ + +-- ---------------------------------------------------------------------------- +-- Основная таблица: валидные ключи click_id +-- ---------------------------------------------------------------------------- +INSERT INTO ods.geo_by_click +WITH + toUUIDOrNull(JSONExtractString(raw, 'click_id')) AS click_id, + toFloat64OrNull(JSONExtractString(raw, 'geo_latitude')) AS geo_latitude, + toFloat64OrNull(JSONExtractString(raw, 'geo_longitude')) AS geo_longitude +SELECT + click_id, + geo_latitude, + geo_longitude, + JSONExtractString(raw, 'geo_country') AS geo_country, + JSONExtractString(raw, 'geo_timezone') AS geo_timezone, + JSONExtractString(raw, 'geo_region_name') AS geo_region_name, + JSONExtractString(raw, 'ip_address') AS ip_address, + ingest_ts AS src_ingest_ts, + raw AS src_raw, + arrayFilter(x -> x != '', [ + if(click_id IS NULL, 'bad_click_id', ''), + if(geo_latitude IS NULL, 'bad_geo_latitude', ''), + if(geo_longitude IS NULL, 'bad_geo_longitude', '') + ]) AS parse_errors +FROM stg.geo_raw +WHERE click_id IS NOT NULL; + +-- ---------------------------------------------------------------------------- +-- Таблица ошибок: строки с ошибками парсинга geo +-- ---------------------------------------------------------------------------- +INSERT INTO ods.geo_by_click_errors +WITH + toUUIDOrNull(JSONExtractString(raw, 'click_id')) AS click_id, + toFloat64OrNull(JSONExtractString(raw, 'geo_latitude')) AS geo_latitude, + toFloat64OrNull(JSONExtractString(raw, 'geo_longitude')) AS geo_longitude, + arrayFilter(x -> x != '', [ + if(click_id IS NULL, 'bad_click_id', ''), + if(geo_latitude IS NULL, 'bad_geo_latitude', ''), + if(geo_longitude IS NULL, 'bad_geo_longitude', '') + ]) AS parse_errors +SELECT + ingest_ts, + kafka_topic, + kafka_partition, + kafka_offset, + kafka_ts, + raw, + arrayStringConcat(parse_errors, ',') AS error_reason +FROM stg.geo_raw +WHERE length(parse_errors) > 0 + AND ( + click_id IS NULL + OR geo_latitude IS NULL + OR geo_longitude IS NULL + );