diff --git a/README.md b/README.md index f991096..2bd1cb8 100644 --- a/README.md +++ b/README.md @@ -9,7 +9,7 @@ Фокус проекта: быстро показать работающий end-to-end сценарий и понятным языком объяснить, как устроены слои и почему пайплайн не падает на "грязных" данных. Коротко про поток: -`data/*.jsonl` -> Kafka (1 строка = 1 сообщение) -> ClickHouse `stg` (сырые JSON) -> Airflow batch `stg -> ods -> dds -> dm` -> Superset. +`data/*.jsonl` -> Airflow DAG `kafka_load` -> Kafka (1 строка = 1 сообщение) -> ClickHouse `stg` (сырые JSON) -> Airflow DAG `etl_pipeline` (`stg -> ods -> dds -> dm`) -> Superset. --- @@ -33,18 +33,15 @@ docker compose ps docker compose exec -T airflow-webserver airflow dags trigger ddl_init ``` -Загрузка данных в Kafka (фаза 2 — через Airflow): +Загрузка данных в Kafka через Airflow DAG: ```bash -# Вариант 1: Через Airflow DAG (рекомендуется) — полная загрузка по умолчанию +# Полная загрузка (по умолчанию limit=0) docker compose exec -T airflow-webserver airflow dags trigger kafka_load \ --conf '{"reset_topics": true}' # Ограниченная загрузка — первые 100 строк docker compose exec -T airflow-webserver airflow dags trigger kafka_load \ --conf '{"limit": 100, "reset_topics": true}' - -# Вариант 2: Через shell-скрипт (устаревший) -make data # полная загрузка ``` Запуск batch-трансформации (STG -> ODS -> DDS -> DM) в Airflow (если DAG выключен, сначала unpause): @@ -106,7 +103,7 @@ flowchart TB DAG[DAG: ddl_init / kafka_load / etl_pipeline] end - Sources -->|kafka_load / make data| Kafka -->|MV| STG -->|Batch SQL| ODS -->|Batch SQL| DDS -->|VIEW| DM + Sources -->|kafka_load| Kafka -->|MV| STG -->|Batch SQL| ODS -->|Batch SQL| DDS -->|VIEW| DM DAG -.->|оркестрация| STG & ODS & DDS & DM ``` @@ -131,14 +128,14 @@ flowchart TB │ ├── ods/ # Batch SQL: STG -> ODS │ ├── dds/ # Batch SQL: ODS -> DDS │ └── dm/ # Batch SQL: DDS -> DM -├── scripts/ # Автоматизация (apply ddl, load data, run batch) +├── scripts/ # Служебные shell-скрипты (legacy fallback, не основной путь) ├── airflow/ # Конфигурация Airflow │ └── requirements.txt ├── docs/ # Документация │ └── ARCHITECTURE.md # Подробное описание слоёв ├── data/ # Исходные JSONL файлы ├── docker-compose.yml -└── Makefile # Команды: up, ddl, data, transform +└── Makefile # Команды: up, ddl, transform ``` --- @@ -149,8 +146,6 @@ flowchart TB |---------|----------| | `make up` | Поднять инфраструктуру | | `make ddl` | Применить DDL в ClickHouse (вне Airflow) | -| `make data` | Загрузить данные в Kafka (50 строк) | -| `FULL=1 make data` | Загрузить полный датасет | | `make transform` | Запустить batch-процесс `STG -> ODS -> DDS -> DM` (вне Airflow) | Примечания про сохранность данных: @@ -214,7 +209,7 @@ flowchart LR ## Частые проблемы - `etl_pipeline` падает с сообщением про схему: сначала запустите `ddl_init`. -- После `docker compose down -v` схема и данные исчезнут: нужно заново `ddl_init` и `make data`. +- После `docker compose down -v` схема и данные исчезнут: нужно заново запустить `ddl_init`, затем `kafka_load`, затем `etl_pipeline`. - Подключения используют разные протоколы: - Airflow (ClickHouseOperator) ходит в ClickHouse по native TCP (порт `9000` внутри сети Docker). - Superset (clickhouse-connect) ходит по HTTP (порт `8123` внутри сети Docker). @@ -223,13 +218,13 @@ flowchart LR ## Статус проекта -Реализовано (Этап 1): +Реализовано: - DAG `ddl_init`: последовательное применение DDL + проверка схемы. +- DAG `kafka_load`: ingest из `.jsonl` в Kafka через `kafka-python` (параметры `limit`, `reset_topics`). - DAG `etl_pipeline`: precheck, ожидание данных в STG, batch-пересчёт ODS/DDS/DM, базовые проверки. - Устойчивость к "грязным" данным: ошибки парсинга сохраняются в ODS, а не валят ingest. В планах (не требуется для MVP задания): -- DAG `kafka_load` (чистый ingest из `.jsonl` в Kafka средствами Airflow). - Инкрементальный batch (watermark вместо `full_refresh`). - DQ мониторинг по расписанию. diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 18ce86b..a89d1d4 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -67,7 +67,7 @@ flowchart LR DM1[v_events_enriched] DM2[v_daily_traffic] DM3[v_utm_effectiveness] - DM4[v_top_pages] + DM4[v_top_pages_daily] end BE --> K1 --> S1 @@ -116,7 +116,7 @@ flowchart TB DM_T["VIEW для BI
(Superset/Grafana)"] end - RAW -->|kafka_load DAG / kafka-console-producer| KAFKA -->|MV| STG_T -->|Batch SQL (Airflow)| ODS_T + RAW -->|Airflow DAG kafka_load| KAFKA -->|MV| STG_T -->|Batch SQL (Airflow)| ODS_T ODS_T -->|argMax + JOIN| DDS_T -->|VIEW| DM_T ODS_T -.->|ошибки| DQ ``` @@ -167,7 +167,7 @@ CREATE TABLE stg.browser_raw ( | `location_event` | event_id | ReplacingMergeTree(src_ingest_ts) | Данные страниц | | `device_by_click` | click_id | ReplacingMergeTree(src_ingest_ts) | Устройства | | `geo_by_click` | click_id | ReplacingMergeTree(src_ingest_ts) | Гео-данные | -| `*_errors` | — | MergeTree | Строки с битыми ключами | +| `*_errors` | — | MergeTree | Строки с критичными ошибками парсинга | **Batch шаг наполнения ODS:** @@ -179,8 +179,8 @@ CREATE TABLE stg.browser_raw ( | `INSERT ... SELECT` в `*_errors` | Сохранение строк с критичными ошибками парсинга | **Логика разделения:** -- **Основная таблица**: строки с валидными ключами (`WHERE key IS NOT NULL`) -- **Таблица ошибок**: строки с невалидными ключами (`WHERE key IS NULL`) +- **Основная таблица**: строки с валидным business key (`WHERE key IS NOT NULL`). +- **Таблица ошибок**: строки с критичными ошибками (невалидный key, timestamp/координаты и т.п.) через отдельный `INSERT ... SELECT`. **Пример структуры:** ```sql @@ -365,7 +365,7 @@ FROM ... ```mermaid sequenceDiagram participant User as Пользователь - participant Make as Makefile + participant Compose as Docker Compose participant Airflow as Airflow participant K as Kafka participant CH as ClickHouse @@ -374,24 +374,24 @@ sequenceDiagram participant DDS as dds.* participant DM as dm.* - User->>Make: make up - Make->>K: docker compose up kafka - Make->>CH: docker compose up clickhouse - K-->>User: ✅ Инфраструктура готова + User->>Compose: make up + Compose->>K: docker compose up -d kafka + Compose->>CH: docker compose up -d clickhouse + Compose->>Airflow: docker compose up -d airflow-* + Compose-->>User: ✅ Инфраструктура готова - 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 (таблицы ODS + drop legacy MV) - Make->>CH: sql/ddl/dds/30_dds.sql - Make->>CH: sql/ddl/dm/40_dm.sql + User->>Airflow: Trigger ddl_init + Airflow->>CH: sql/ddl/00_databases.sql + Airflow->>CH: sql/ddl/stg/10_stg.sql (Kafka Engine + MV) + Airflow->>CH: sql/ddl/ods/20_ods.sql + Airflow->>CH: sql/ddl/dds/30_dds.sql + Airflow->>CH: sql/ddl/dm/40_dm.sql CH-->>User: ✅ Структура БД создана - User->>Make: make data - Make->>K: load_kafka_data.sh - K->>K: Создание топиков + User->>Airflow: Trigger kafka_load + Airflow->>K: precheck + prepare_topics loop 4 файла - Make->>K: kafka-console-producer + Airflow->>K: KafkaProducer.send(topic, json_line) end K->>CH: Потребление сообщений CH->>STG: INSERT через MV @@ -535,11 +535,11 @@ flowchart LR ### Обработка ошибок в ODS -**Проблема:** Грязные данные с невалидными ключами (NULL event_id/click_id). +**Проблема:** Грязные данные могут содержать не только невалидные ключи, но и невалидные timestamp/координаты/ID. **Решение:** -1. **Основная таблица**: только валидные строки (`WHERE key IS NOT NULL`) -2. **Таблица ошибок**: строки с невалидными ключами через отдельные `INSERT ... SELECT` +1. **Основная таблица**: строки с валидным business key (`WHERE key IS NOT NULL`) +2. **Таблица ошибок**: строки с критичными ошибками парсинга через отдельные `INSERT ... SELECT` 3. **DQ-метрики**: массив `parse_errors` для аудита ```sql @@ -549,7 +549,8 @@ SELECT ... FROM stg.browser_raw WHERE event_id IS NOT NULL; -- Таблица ошибок INSERT INTO ods.browser_event_errors -SELECT ... FROM stg.browser_raw WHERE event_id IS NULL; +SELECT ... FROM stg.browser_raw +WHERE event_id IS NULL OR event_ts IS NULL OR click_id IS NULL; ``` ### Partial data в DDS @@ -599,19 +600,19 @@ INSERT INTO dm.daily_traffic SELECT * FROM dm.v_daily_traffic; ```python # dags/ddl_init_dag.py — создание баз/таблиц (ручной запуск при bootstrap) -# dags/kafka_load_dag.py — загрузка JSONL в Kafka (фаза 2, через kafka-python) +# dags/kafka_load_dag.py — загрузка JSONL в Kafka (через kafka-python) # dags/etl_pipeline_dag.py — основной ETL (STG→ODS→DDS→DM) # Учебный формат: # - DDL и трансформации выполняются явными SQL-task через ClickHouseOperator; # - SQL-файлы вызываются по фиксированным путям; -# - загрузка данных в Kafka (фаза 2) выполняется через DAG `kafka_load`. +# - загрузка данных в Kafka выполняется через DAG `kafka_load`. # -# Основной demo-сценарий (фаза 2): +# Основной demo-сценарий: # ddl_init -> kafka_load -> etl_pipeline ``` -**DAG `kafka_load`** (фаза 2): +**DAG `kafka_load`**: - Загрузка данных из `data/*.jsonl` в Kafka через `kafka-python` - TaskGroup `precheck`: проверка Kafka, файлов, параметров - TaskGroup `ingest`: создание топиков → параллельная загрузка 4 потоков → проверка diff --git a/plans/runbook.md b/plans/runbook.md index ede4616..d380dae 100644 --- a/plans/runbook.md +++ b/plans/runbook.md @@ -1,107 +1,87 @@ -# Runbook: запуск демо и загрузка данных +# Runbook: Airflow-first запуск демо -Этот документ фиксирует порядок действий и `make`‑таргеты. Он не описывает внутренности ClickHouse‑слоёв (это в `plans/clickhouse_ddl.md`). +Документ фиксирует канонический пользовательский сценарий через Airflow DAG'и. +Внутренности слоёв ClickHouse описаны в `plans/clickhouse_ddl.md` и `docs/ARCHITECTURE.md`. ## Предпосылки - Docker + Docker Compose. -- Доступ к Docker daemon (если `docker compose ...` пишет `permission denied ... /var/run/docker.sock`, добавьте пользователя в группу `docker` или запускайте команды с правами, принятыми в вашей среде). +- Доступ к Docker daemon. -## Быстрый сценарий +## Канонический сценарий (через Airflow) 1) Поднять инфраструктуру: ```bash make up +docker compose ps ``` -2) Залить данные в Kafka (два варианта): +2) Инициализировать схему ClickHouse: -**Вариант А: Через Airflow DAG `kafka_load` (рекомендуется, фаза 2)** ```bash -# Через CLI — полная загрузка по умолчанию +docker compose exec -T airflow-webserver airflow dags trigger ddl_init +``` + +3) Загрузить данные в Kafka через DAG `kafka_load`: + +```bash +# Рекомендуется для демо: небольшой срез +docker compose exec -T airflow-webserver airflow dags trigger kafka_load \ + --conf '{"limit": 50, "reset_topics": true}' + +# Полная загрузка (limit=0 по умолчанию) docker compose exec -T airflow-webserver airflow dags trigger kafka_load \ --conf '{"reset_topics": true}' - -# Ограниченная загрузка — первые 100 строк -docker compose exec -T airflow-webserver airflow dags trigger kafka_load \ - --conf '{"limit": 100, "reset_topics": true}' - -# Или через UI: Airflow → DAGs → kafka_load → Trigger DAG with config ``` -Параметры `kafka_load`: -- `limit` — количество строк (default: 0 — все строки) -- `reset_topics` — пересоздать топики (default: true) - -**Вариант Б: Через shell-скрипт `make data` (устаревший)** -```bash -make data -``` - -`make data` не зависит от ClickHouse/DDL — достаточно, чтобы Kafka была поднята. - -3) Применить DDL в ClickHouse: +4) Запустить ETL: ```bash -make ddl +docker compose exec -T airflow-webserver airflow dags trigger etl_pipeline \ + --conf '{"full_refresh": true}' ``` -## Make таргеты - -- `make up` — `docker compose up -d` (поднимает весь стек из `docker-compose.yml`). -- `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`. - -План Airflow DAG'ов: `plans/airflow_dags_plan.md`. - -## Загрузка данных в Kafka (`make data`) - -### Топики - -Скрипт использует фиксированный маппинг: - -- `data/browser_events.jsonl` → `browser_events` -- `data/location_events.jsonl` → `location_events` -- `data/device_events.jsonl` → `device_events` -- `data/geo_events.jsonl` → `geo_events` - -### Режимы загрузки - -- По умолчанию — “debug срез”: первые 50 строк каждого файла. -- Полная загрузка — весь файл. - -Параметры (env): - -- `LIMIT` — сколько строк брать из каждого `.jsonl`. По умолчанию загружаются все записи (весь файл). - Для ограничения используйте `LIMIT=50` или `LIMIT=100`. -- `RESET_TOPICS` — если `RESET_TOPICS=1` (по умолчанию), топики удаляются и создаются заново с теми же именами. -- `BOOTSTRAP_SERVER` — bootstrap для Kafka *изнутри kafka‑контейнера* (по умолчанию `kafka:29092`). - -Примеры: +5) Проверить результат: ```bash -# Загрузить все данные (по умолчанию) -make data - -# Быстрый тест — 50 строк на поток -LIMIT=50 make data - -# Ограниченная загрузка — 100 строк на поток -LIMIT=100 make data - -# Дозалить данные без пересоздания топиков -RESET_TOPICS=0 make data +docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 \ + --query "SELECT count() FROM stg.browser_raw" +docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 \ + --query "SELECT count() FROM ods.browser_event" +docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 \ + --query "SELECT count() FROM dds.event" +docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 \ + --query "SELECT * FROM dm.dq_summary ORDER BY layer, table_name, check_name" ``` -## Применение DDL в ClickHouse (`make ddl`) +## Параметры DAG'ов -Скрипт исполняет SQL-файлы из `sql/ddl/*` по фиксированному порядку. Для `ENGINE = Kafka` важно, чтобы `kafka_broker_list` был доступен из контейнера ClickHouse. +- `ddl_init`: + - `verify_only` (bool, default: `false`) — только проверка схемы, без применения DDL. +- `kafka_load`: + - `limit` (int, default: `0`) — количество строк на поток (`0` = весь файл). + - `reset_topics` (bool, default: `true`) — удалить и создать топики заново. +- `etl_pipeline`: + - `full_refresh` (bool, default: `true`) — очищать DDS перед загрузкой. + - `wait_stg_timeout_sec` (int, default: `600`) — таймаут ожидания данных в STG. -В текущем compose: +## Роль Make-таргетов -- для соединений “контейнер → Kafka” используйте `kafka:29092`; -- `localhost:9092` подходит только для клиентов на хосте. +- `make up` — основной способ поднять стек. +- `make ddl`, `make transform` — технический fallback для низкоуровневой диагностики вне Airflow. +- Загрузка в Kafka в runbook выполняется только через DAG `kafka_load`. + +## Важные примечания + +- Для связей контейнеров используйте `kafka:29092` (не `localhost:9092`). +- После `docker compose down -v` нужно заново выполнить: + 1. `ddl_init` + 2. `kafka_load` + 3. `etl_pipeline` + +## Связанные документы + +- План ingest: `plans/kafka_ingest_plan.md` +- План DAG'ов: `plans/airflow_dags_plan.md` +- Архитектура: `docs/ARCHITECTURE.md`