diff --git a/AGENTS.md b/AGENTS.md
index 20becea..f479d88 100644
--- a/AGENTS.md
+++ b/AGENTS.md
@@ -84,6 +84,7 @@
- Держать изменения минимальными и по теме задания (инфра, схема, ingest, витрины).
- Не коммитить секреты. Если требуется пароль/ключи — использовать `.env` и примеры `.env.example`.
+- Оформлять коммиты по правилам из [COMMIT_RULES.md](./docs/COMMIT_RULES.md).
- README/планы обновлять вместе с изменениями инфраструктуры/DDL.
- Для спорных или меняющихся API (особенно Airflow/operators/providers) проверять актуальную документацию через `context7` и фиксировать решение в коде/документации.
- **Комментарии в коде — на русском языке**:
@@ -103,6 +104,7 @@
- [README.md](./README.md) — пользовательская документация (быстрый старт, архитектура)
- [docs/ARCHITECTURE.md](./docs/ARCHITECTURE.md) — подробное описание слоёв и технических решений
- [data/DE-task.md](./data/DE-task.md) — исходное задание
+- [COMMIT_RULES.md](./docs/COMMIT_RULES.md) — правила оформления коммитов
## Примечания по текущему состоянию (если что-то “не встаёт”)
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/docs/COMMIT_RULES.md b/docs/COMMIT_RULES.md
new file mode 100644
index 0000000..269cb1d
--- /dev/null
+++ b/docs/COMMIT_RULES.md
@@ -0,0 +1,78 @@
+# Правила оформления коммитов
+
+Документ задаёт единый стиль коммитов для всех участников проекта.
+
+## Язык
+
+- Язык коммитов: русский.
+- Технические термины (`Airflow`, `ClickHouse`, `Kafka`, `MV`, `DDL`) допускаются на английском.
+
+## Формат заголовка
+
+- Формат: `(): <краткое действие>`
+- Максимальная длина заголовка: 72 символа.
+- Заголовок пишется в повелительном стиле, без точки в конце.
+
+### Разрешённые `type`
+
+- `feat` — новая функциональность
+- `fix` — исправление ошибки
+- `refactor` — изменение структуры без смены поведения
+- `docs` — документация
+- `test` — тесты/проверки
+- `chore` — сервисные изменения (конфиги, скрипты, хуки)
+- `ci` — CI/CD
+- `perf` — оптимизация производительности
+- `revert` — откат коммита
+
+### Рекомендуемые `scope` для этого репозитория
+
+- `airflow`
+- `stg`
+- `ods`
+- `dds`
+- `dm`
+- `kafka`
+- `superset`
+- `monitoring`
+- `scripts`
+- `docs`
+- `infra`
+
+## Структура тела коммита
+
+Если изменение не тривиальное, тело коммита обязательно. Для удобства чтения используйте буллеты.
+
+Рекомендуемый шаблон:
+
+```text
+(): <краткое действие>
+
+- Зачем:
+ - причина изменения
+- Что сделано:
+ - ключевое изменение 1
+ - ключевое изменение 2
+- Проверка:
+ - как проверено (команда/тест/смоук-чек)
+```
+
+## Размер и границы коммита
+
+- Один коммит = одна логическая задача.
+- Не смешивать в одном коммите функциональные изменения и крупный рефакторинг без необходимости.
+- Документацию обновлять в том же коммите, где меняется поведение пайплайна или инфраструктуры.
+
+## Ломающие изменения
+
+- Для ломающих изменений используйте `!` в заголовке:
+ - `feat(ods)!: изменить контракт таблицы browser_event`
+- Добавляйте footer:
+ - `BREAKING CHANGE: ...`
+
+## Примеры
+
+- `feat(ods): перенести STG->ODS в batch шаг Airflow`
+- `fix(airflow): ждать данные в STG перед load_ods`
+- `docs(architecture): обновить схему потока после миграции ODS`
+- `chore(scripts): синхронизировать make transform с новым ETL`
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
+ );