Files
clickstream-ch-kafka-supers…/README.md
T
ddadmin a185a56762 docs(readme): add DBeaver guide and move DE-task to docs
- Why:
  - simplify first data checks for interview/demo audience
  - keep task reference in a stable docs location
- What:
  - add short DBeaver connection section with ready-to-use params and quick SQL checks
  - update DE-task links in README to docs path
  - move DE-task from data/ to docs/
- Check:
  - README links resolve to docs/DE-task.md
  - git shows file move data/DE-task.md -> docs/DE-task.md
2026-02-08 19:57:11 +03:00

259 lines
10 KiB
Markdown

# ClickHouse Mini DWH для кликстрима
[![Stack](https://img.shields.io/badge/stack-Kafka%20%7C%20ClickHouse%20%7C%20Airflow%20%7C%20Superset-blue)](./docker-compose.yml)
[![Layers](https://img.shields.io/badge/layers-STG%20→%20ODS%20→%20DDS%20→%20DM-green)](./docs/ARCHITECTURE.md)
[![License](https://img.shields.io/badge/license-Educational-orange)]()
Мини-демо для решения задания [DE-task.md](./docs/DE-task.md): развернуть инфраструктуру на своей машине, прогнать кликстрим через Kafka в ClickHouse, сделать регулярный расчёт в Airflow и подготовить витрины под дашборд.
Фокус проекта: быстро показать работающий end-to-end сценарий и понятным языком объяснить, как устроены слои и почему пайплайн не падает на "грязных" данных.
Коротко про поток:
`data/*.jsonl` -> Airflow DAG `kafka_load` -> Kafka (1 строка = 1 сообщение) -> ClickHouse `stg` (сырые JSON) -> Airflow DAG `etl_pipeline` (`stg -> ods -> dds -> dm`) -> Superset.
---
## Быстрый старт (демо-сценарий)
```bash
# 1) Поднять инфраструктуру
make up
# Проверить статусы контейнеров
docker compose ps
```
Дальше основной путь идёт через Airflow (как в задании).
1. Открыть Airflow UI: `http://localhost:8080` (admin/admin)
2. Включить (unpause) и запустить `ddl_init` (создаёт базы/таблицы/VIEW в ClickHouse)
Опционально можно триггернуть DAG из CLI (удобно для CI/скрипта):
```bash
docker compose exec -T airflow-webserver airflow dags trigger ddl_init
```
Загрузка данных в Kafka через Airflow DAG:
```bash
# Полная загрузка (по умолчанию 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}'
```
Запуск batch-трансформации (STG -> ODS -> DDS -> DM) в Airflow (если DAG выключен, сначала unpause):
```bash
docker compose exec -T airflow-webserver airflow dags trigger etl_pipeline \
--conf '{"full_refresh": true}'
```
Smoke-check результата в ClickHouse:
```bash
docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query \
"SELECT 'ods.browser_event' AS t, count() AS rows FROM ods.browser_event"
docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query \
"SELECT 'dds.click' AS t, count() AS rows FROM dds.click"
docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query \
"SELECT 'dds.event' AS t, count() AS rows FROM dds.event"
docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query \
"SELECT 'dm.dq_summary' AS t, count() AS rows FROM dm.dq_summary"
```
---
## Доступные сервисы
| Сервис | URL | Назначение |
|--------|-----|------------|
| ClickHouse HTTP | http://localhost:9123/play | SQL-запросы |
| Kafka UI | http://localhost:8082 | Просмотр топиков |
| Airflow | http://localhost:8080 | Оркестрация ETL (admin/admin) |
| Superset | http://localhost:8088 | BI-дашборды |
| Prometheus | http://localhost:9090 | Метрики |
| Grafana | http://localhost:3000 | Визуализация метрик |
---
## Подключение DBeaver (кратко)
После прогона `ddl_init -> kafka_load -> etl_pipeline` можно быстро проверить витрины в DBeaver (удобно для демо бизнесу).
1. `Database -> New Database Connection -> ClickHouse`
2. Параметры:
- `Host`: `localhost`
- `Port`: `9123` (HTTP)
- `Database`: `default`
- `Username`: `default`
- `Password`: `123456`
3. Нажать `Test Connection` -> `Finish`
Если ваш драйвер просит native-протокол, используйте порт `8002`.
Полезные быстрые запросы для первичного анализа:
```sql
SELECT count() AS rows FROM dm.v_events_enriched;
SELECT * FROM dm.v_daily_traffic ORDER BY event_date DESC LIMIT 20;
SELECT * FROM dm.v_utm_effectiveness ORDER BY clicks DESC LIMIT 20;
```
---
## Архитектура (в двух словах)
```mermaid
flowchart LR
subgraph AF["Airflow"]
D1["ddl_init"]
D2["kafka_load"]
D3["etl_pipeline"]
end
subgraph Kafka["Kafka"]
Topics[4 топика]
end
subgraph CH["ClickHouse"]
STG["STG: сырые данные"]
ODS["ODS: типизация + DQ"]
DDS["DDS: сущности"]
DM["DM: витрины VIEW"]
end
D2 -->|загрузка JSONL| Kafka -->|Kafka MV| STG
STG -->|batch| ODS -->|batch| DDS -->|VIEW| DM
D1 -.->|DDL| CH
D3 -.->|batch| ODS & DDS
```
Особенность задания про "грязные данные": парсинг не валит pipeline, ошибки фиксируются в `ods.*_errors` и в поле `parse_errors`.
[Подробное описание архитектуры →](./docs/ARCHITECTURE.md)
---
## Структура проекта
```
.
├── sql/
│ ├── ddl/ # DDL по слоям
│ │ ├── 00_databases.sql
│ │ ├── stg/10_stg.sql
│ │ ├── 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/ # Служебные shell-скрипты (legacy fallback, не основной путь)
├── airflow/ # Конфигурация Airflow
│ ├── dags/ # Airflow DAGs для оркестрации
│ └── requirements.txt
├── docs/ # Документация
│ └── ARCHITECTURE.md # Подробное описание слоёв
├── data/ # Исходные JSONL файлы
├── docker-compose.yml
└── Makefile # Команды: up, ddl, transform
```
---
## Команды Makefile
| Команда | Описание |
|---------|----------|
| `make up` | Поднять инфраструктуру |
| `make ddl` | Применить DDL в ClickHouse (вне Airflow) |
| `make transform` | Запустить batch-процесс `STG -> ODS -> DDS -> DM` (вне Airflow) |
Примечания про сохранность данных:
- Данные ClickHouse сохраняются в Docker volume `clickhouse-data`.
- Данные Kafka сохраняются в Docker volume `kafka-data`.
- `docker compose down` сохраняет named volumes, `docker compose down -v` удаляет их (и данные пропадут).
---
## Ключи данных (как джойним)
```mermaid
flowchart LR
subgraph Sources["Источники"]
BE["browser_events (event_id, click_id)"]
LE["location_events (event_id)"]
DE["device_events (click_id)"]
GE["geo_events (click_id)"]
end
subgraph DDS["DDS"]
EV["event (event_id PK)"]
CL["click (click_id PK)"]
end
subgraph DM["DM"]
V1[v_events_enriched]
V2[v_daily_traffic]
V3[v_utm_effectiveness]
end
BE -->|event_id| EV
LE -->|event_id| EV
BE -->|click_id| CL
DE -->|click_id| CL
GE -->|click_id| CL
EV -->|LEFT JOIN click_id| V1
CL --> V1
EV --> V2 & V3
CL --> V2 & V3
```
---
## Дашборд в Superset (опционально, но полезно)
1. Открыть `http://localhost:8088`
2. Database -> Add:
- URI: `clickhouse+connect://default:123456@clickhouse:8123/default`
3. Создать datasets из `dm.v_*` (VIEW) и собрать несколько графиков
Идеи графиков под задание:
- Трафик по дням: `dm.v_daily_traffic` (events, uniq_users)
- Эффективность UTM: `dm.v_utm_effectiveness` (clicks, purchases)
- Популярные страницы: `dm.v_top_pages_daily` (pageviews)
- Качество данных: `dm.v_dq_errors_daily` (rows_cnt по error_code)
---
## Частые проблемы
- `etl_pipeline` падает с сообщением про схему: сначала запустите `ddl_init`.
- После `docker compose down -v` схема и данные исчезнут: нужно заново запустить `ddl_init`, затем `kafka_load`, затем `etl_pipeline`.
- Подключения используют разные протоколы:
- Airflow (ClickHouseOperator) ходит в ClickHouse по native TCP (порт `9000` внутри сети Docker).
- Superset (clickhouse-connect) ходит по HTTP (порт `8123` внутри сети Docker).
---
## Статус проекта
Реализовано:
- 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 задания):
- Инкрементальный batch (watermark вместо `full_refresh`).
- DQ мониторинг по расписанию.
---
## Документация
- [Архитектура и слои](./docs/ARCHITECTURE.md) — подробное описание STG/ODS/DDS/DM, ER-диаграммы, обоснование решений
- [DE-task.md](./docs/DE-task.md) — исходное задание