docs(docs): align docs with airflow-first ingest workflow

- Why:\n  - User-facing docs mixed Airflow and legacy CLI ingest paths and caused confusion\n- What:\n  - Rework README quick start and status to use DAG chain ddl_init -> kafka_load -> etl_pipeline\n  - Rewrite runbook as canonical Airflow-first execution flow\n  - Sync architecture diagrams/sequence and DQ wording with current SQL and DAG behavior\n- Check:\n  - Verified updated sections and removed stale markers with rg in README.md, docs/ARCHITECTURE.md, plans/runbook.md
This commit is contained in:
2026-02-08 18:52:44 +03:00
parent d7588a8caa
commit 0b75da9c08
3 changed files with 95 additions and 119 deletions
+9 -14
View File
@@ -9,7 +9,7 @@
Фокус проекта: быстро показать работающий end-to-end сценарий и понятным языком объяснить, как устроены слои и почему пайплайн не падает на "грязных" данных. Фокус проекта: быстро показать работающий 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 docker compose exec -T airflow-webserver airflow dags trigger ddl_init
``` ```
Загрузка данных в Kafka (фаза 2 — через Airflow): Загрузка данных в Kafka через Airflow DAG:
```bash ```bash
# Вариант 1: Через Airflow DAG (рекомендуется) — полная загрузка по умолчанию # Полная загрузка (по умолчанию limit=0)
docker compose exec -T airflow-webserver airflow dags trigger kafka_load \ docker compose exec -T airflow-webserver airflow dags trigger kafka_load \
--conf '{"reset_topics": true}' --conf '{"reset_topics": true}'
# Ограниченная загрузка — первые 100 строк # Ограниченная загрузка — первые 100 строк
docker compose exec -T airflow-webserver airflow dags trigger kafka_load \ docker compose exec -T airflow-webserver airflow dags trigger kafka_load \
--conf '{"limit": 100, "reset_topics": true}' --conf '{"limit": 100, "reset_topics": true}'
# Вариант 2: Через shell-скрипт (устаревший)
make data # полная загрузка
``` ```
Запуск batch-трансформации (STG -> ODS -> DDS -> DM) в Airflow (если DAG выключен, сначала unpause): Запуск batch-трансформации (STG -> ODS -> DDS -> DM) в Airflow (если DAG выключен, сначала unpause):
@@ -106,7 +103,7 @@ flowchart TB
DAG[DAG: ddl_init / kafka_load / etl_pipeline] DAG[DAG: ddl_init / kafka_load / etl_pipeline]
end 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 DAG -.->|оркестрация| STG & ODS & DDS & DM
``` ```
@@ -131,14 +128,14 @@ flowchart TB
│ ├── ods/ # Batch SQL: STG -> ODS │ ├── ods/ # Batch SQL: STG -> ODS
│ ├── dds/ # Batch SQL: ODS -> DDS │ ├── dds/ # Batch SQL: ODS -> DDS
│ └── dm/ # Batch SQL: DDS -> DM │ └── dm/ # Batch SQL: DDS -> DM
├── scripts/ # Автоматизация (apply ddl, load data, run batch) ├── scripts/ # Служебные shell-скрипты (legacy fallback, не основной путь)
├── airflow/ # Конфигурация Airflow ├── airflow/ # Конфигурация Airflow
│ └── requirements.txt │ └── requirements.txt
├── docs/ # Документация ├── docs/ # Документация
│ └── ARCHITECTURE.md # Подробное описание слоёв │ └── ARCHITECTURE.md # Подробное описание слоёв
├── data/ # Исходные JSONL файлы ├── data/ # Исходные JSONL файлы
├── docker-compose.yml ├── docker-compose.yml
└── Makefile # Команды: up, ddl, data, transform └── Makefile # Команды: up, ddl, transform
``` ```
--- ---
@@ -149,8 +146,6 @@ flowchart TB
|---------|----------| |---------|----------|
| `make up` | Поднять инфраструктуру | | `make up` | Поднять инфраструктуру |
| `make ddl` | Применить DDL в ClickHouse (вне Airflow) | | `make ddl` | Применить DDL в ClickHouse (вне Airflow) |
| `make data` | Загрузить данные в Kafka (50 строк) |
| `FULL=1 make data` | Загрузить полный датасет |
| `make transform` | Запустить batch-процесс `STG -> ODS -> DDS -> DM` (вне Airflow) | | `make transform` | Запустить batch-процесс `STG -> ODS -> DDS -> DM` (вне Airflow) |
Примечания про сохранность данных: Примечания про сохранность данных:
@@ -214,7 +209,7 @@ flowchart LR
## Частые проблемы ## Частые проблемы
- `etl_pipeline` падает с сообщением про схему: сначала запустите `ddl_init`. - `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). - Airflow (ClickHouseOperator) ходит в ClickHouse по native TCP (порт `9000` внутри сети Docker).
- Superset (clickhouse-connect) ходит по HTTP (порт `8123` внутри сети Docker). - Superset (clickhouse-connect) ходит по HTTP (порт `8123` внутри сети Docker).
@@ -223,13 +218,13 @@ flowchart LR
## Статус проекта ## Статус проекта
Реализовано (Этап 1): Реализовано:
- DAG `ddl_init`: последовательное применение DDL + проверка схемы. - DAG `ddl_init`: последовательное применение DDL + проверка схемы.
- DAG `kafka_load`: ingest из `.jsonl` в Kafka через `kafka-python` (параметры `limit`, `reset_topics`).
- DAG `etl_pipeline`: precheck, ожидание данных в STG, batch-пересчёт ODS/DDS/DM, базовые проверки. - DAG `etl_pipeline`: precheck, ожидание данных в STG, batch-пересчёт ODS/DDS/DM, базовые проверки.
- Устойчивость к "грязным" данным: ошибки парсинга сохраняются в ODS, а не валят ingest. - Устойчивость к "грязным" данным: ошибки парсинга сохраняются в ODS, а не валят ingest.
В планах (не требуется для MVP задания): В планах (не требуется для MVP задания):
- DAG `kafka_load` (чистый ingest из `.jsonl` в Kafka средствами Airflow).
- Инкрементальный batch (watermark вместо `full_refresh`). - Инкрементальный batch (watermark вместо `full_refresh`).
- DQ мониторинг по расписанию. - DQ мониторинг по расписанию.
+29 -28
View File
@@ -67,7 +67,7 @@ flowchart LR
DM1[v_events_enriched] DM1[v_events_enriched]
DM2[v_daily_traffic] DM2[v_daily_traffic]
DM3[v_utm_effectiveness] DM3[v_utm_effectiveness]
DM4[v_top_pages] DM4[v_top_pages_daily]
end end
BE --> K1 --> S1 BE --> K1 --> S1
@@ -116,7 +116,7 @@ flowchart TB
DM_T["VIEW для BI<br/>(Superset/Grafana)"] DM_T["VIEW для BI<br/>(Superset/Grafana)"]
end 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 -->|argMax + JOIN| DDS_T -->|VIEW| DM_T
ODS_T -.->|ошибки| DQ ODS_T -.->|ошибки| DQ
``` ```
@@ -167,7 +167,7 @@ CREATE TABLE stg.browser_raw (
| `location_event` | event_id | ReplacingMergeTree(src_ingest_ts) | Данные страниц | | `location_event` | event_id | ReplacingMergeTree(src_ingest_ts) | Данные страниц |
| `device_by_click` | click_id | ReplacingMergeTree(src_ingest_ts) | Устройства | | `device_by_click` | click_id | ReplacingMergeTree(src_ingest_ts) | Устройства |
| `geo_by_click` | click_id | ReplacingMergeTree(src_ingest_ts) | Гео-данные | | `geo_by_click` | click_id | ReplacingMergeTree(src_ingest_ts) | Гео-данные |
| `*_errors` | — | MergeTree | Строки с битыми ключами | | `*_errors` | — | MergeTree | Строки с критичными ошибками парсинга |
**Batch шаг наполнения ODS:** **Batch шаг наполнения ODS:**
@@ -179,8 +179,8 @@ CREATE TABLE stg.browser_raw (
| `INSERT ... SELECT` в `*_errors` | Сохранение строк с критичными ошибками парсинга | | `INSERT ... SELECT` в `*_errors` | Сохранение строк с критичными ошибками парсинга |
**Логика разделения:** **Логика разделения:**
- **Основная таблица**: строки с валидными ключами (`WHERE key IS NOT NULL`) - **Основная таблица**: строки с валидным business key (`WHERE key IS NOT NULL`).
- **Таблица ошибок**: строки с невалидными ключами (`WHERE key IS NULL`) - **Таблица ошибок**: строки с критичными ошибками (невалидный key, timestamp/координаты и т.п.) через отдельный `INSERT ... SELECT`.
**Пример структуры:** **Пример структуры:**
```sql ```sql
@@ -365,7 +365,7 @@ FROM ...
```mermaid ```mermaid
sequenceDiagram sequenceDiagram
participant User as Пользователь participant User as Пользователь
participant Make as Makefile participant Compose as Docker Compose
participant Airflow as Airflow participant Airflow as Airflow
participant K as Kafka participant K as Kafka
participant CH as ClickHouse participant CH as ClickHouse
@@ -374,24 +374,24 @@ sequenceDiagram
participant DDS as dds.* participant DDS as dds.*
participant DM as dm.* participant DM as dm.*
User->>Make: make up User->>Compose: make up
Make->>K: docker compose up kafka Compose->>K: docker compose up -d kafka
Make->>CH: docker compose up clickhouse Compose->>CH: docker compose up -d clickhouse
K-->>User: ✅ Инфраструктура готова Compose->>Airflow: docker compose up -d airflow-*
Compose-->>User: ✅ Инфраструктура готова
User->>Make: make ddl User->>Airflow: Trigger ddl_init
Make->>CH: sql/ddl/00_databases.sql Airflow->>CH: sql/ddl/00_databases.sql
Make->>CH: sql/ddl/stg/10_stg.sql (Kafka Engine) Airflow->>CH: sql/ddl/stg/10_stg.sql (Kafka Engine + MV)
Make->>CH: sql/ddl/ods/20_ods.sql (таблицы ODS + drop legacy MV) Airflow->>CH: sql/ddl/ods/20_ods.sql
Make->>CH: sql/ddl/dds/30_dds.sql Airflow->>CH: sql/ddl/dds/30_dds.sql
Make->>CH: sql/ddl/dm/40_dm.sql Airflow->>CH: sql/ddl/dm/40_dm.sql
CH-->>User: ✅ Структура БД создана CH-->>User: ✅ Структура БД создана
User->>Make: make data User->>Airflow: Trigger kafka_load
Make->>K: load_kafka_data.sh Airflow->>K: precheck + prepare_topics
K->>K: Создание топиков
loop 4 файла loop 4 файла
Make->>K: kafka-console-producer Airflow->>K: KafkaProducer.send(topic, json_line)
end end
K->>CH: Потребление сообщений K->>CH: Потребление сообщений
CH->>STG: INSERT через MV CH->>STG: INSERT через MV
@@ -535,11 +535,11 @@ flowchart LR
### Обработка ошибок в ODS ### Обработка ошибок в ODS
**Проблема:** Грязные данные с невалидными ключами (NULL event_id/click_id). **Проблема:** Грязные данные могут содержать не только невалидные ключи, но и невалидные timestamp/координаты/ID.
**Решение:** **Решение:**
1. **Основная таблица**: только валидные строки (`WHERE key IS NOT NULL`) 1. **Основная таблица**: строки с валидным business key (`WHERE key IS NOT NULL`)
2. **Таблица ошибок**: строки с невалидными ключами через отдельные `INSERT ... SELECT` 2. **Таблица ошибок**: строки с критичными ошибками парсинга через отдельные `INSERT ... SELECT`
3. **DQ-метрики**: массив `parse_errors` для аудита 3. **DQ-метрики**: массив `parse_errors` для аудита
```sql ```sql
@@ -549,7 +549,8 @@ SELECT ... FROM stg.browser_raw WHERE event_id IS NOT NULL;
-- Таблица ошибок -- Таблица ошибок
INSERT INTO ods.browser_event_errors 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 ### Partial data в DDS
@@ -599,19 +600,19 @@ INSERT INTO dm.daily_traffic SELECT * FROM dm.v_daily_traffic;
```python ```python
# dags/ddl_init_dag.py — создание баз/таблиц (ручной запуск при bootstrap) # 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) # dags/etl_pipeline_dag.py — основной ETL (STG→ODS→DDS→DM)
# Учебный формат: # Учебный формат:
# - DDL и трансформации выполняются явными SQL-task через ClickHouseOperator; # - DDL и трансформации выполняются явными SQL-task через ClickHouseOperator;
# - SQL-файлы вызываются по фиксированным путям; # - SQL-файлы вызываются по фиксированным путям;
# - загрузка данных в Kafka (фаза 2) выполняется через DAG `kafka_load`. # - загрузка данных в Kafka выполняется через DAG `kafka_load`.
# #
# Основной demo-сценарий (фаза 2): # Основной demo-сценарий:
# ddl_init -> kafka_load -> etl_pipeline # ddl_init -> kafka_load -> etl_pipeline
``` ```
**DAG `kafka_load`** (фаза 2): **DAG `kafka_load`**:
- Загрузка данных из `data/*.jsonl` в Kafka через `kafka-python` - Загрузка данных из `data/*.jsonl` в Kafka через `kafka-python`
- TaskGroup `precheck`: проверка Kafka, файлов, параметров - TaskGroup `precheck`: проверка Kafka, файлов, параметров
- TaskGroup `ingest`: создание топиков → параллельная загрузка 4 потоков → проверка - TaskGroup `ingest`: создание топиков → параллельная загрузка 4 потоков → проверка
+57 -77
View File
@@ -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 + Docker Compose.
- Доступ к Docker daemon (если `docker compose ...` пишет `permission denied ... /var/run/docker.sock`, добавьте пользователя в группу `docker` или запускайте команды с правами, принятыми в вашей среде). - Доступ к Docker daemon.
## Быстрый сценарий ## Канонический сценарий (через Airflow)
1) Поднять инфраструктуру: 1) Поднять инфраструктуру:
```bash ```bash
make up make up
docker compose ps
``` ```
2) Залить данные в Kafka (два варианта): 2) Инициализировать схему ClickHouse:
**Вариант А: Через Airflow DAG `kafka_load` (рекомендуется, фаза 2)**
```bash ```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 \ docker compose exec -T airflow-webserver airflow dags trigger kafka_load \
--conf '{"reset_topics": true}' --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`: 4) Запустить ETL:
- `limit` — количество строк (default: 0 — все строки)
- `reset_topics` — пересоздать топики (default: true)
**Вариант Б: Через shell-скрипт `make data` (устаревший)**
```bash
make data
```
`make data` не зависит от ClickHouse/DDL — достаточно, чтобы Kafka была поднята.
3) Применить DDL в ClickHouse:
```bash ```bash
make ddl docker compose exec -T airflow-webserver airflow dags trigger etl_pipeline \
--conf '{"full_refresh": true}'
``` ```
## Make таргеты 5) Проверить результат:
- `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`).
Примеры:
```bash ```bash
# Загрузить все данные (по умолчанию) docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 \
make data --query "SELECT count() FROM stg.browser_raw"
docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 \
# Быстрый тест — 50 строк на поток --query "SELECT count() FROM ods.browser_event"
LIMIT=50 make data docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 \
--query "SELECT count() FROM dds.event"
# Ограниченная загрузка — 100 строк на поток docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 \
LIMIT=100 make data --query "SELECT * FROM dm.dq_summary ORDER BY layer, table_name, check_name"
# Дозалить данные без пересоздания топиков
RESET_TOPICS=0 make data
``` ```
## Применение 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`; - `make up` — основной способ поднять стек.
- `localhost:9092` подходит только для клиентов на хосте. - `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`