feat(airflow): configure ClickHouse connection and update infrastructure
Update Airflow configuration to integrate with ClickHouse DWH instead of PostgreSQL training database. Changes include: - Switch Airflow dependencies from PostgreSQL to ClickHouse connector - Update docker-compose to use ClickHouse connection and correct Dockerfile - Refactor airflow/requirements.txt to include only essential packages - Add DAGs directory for ETL pipeline orchestration - Update documentation to reflect Airflow integration and access credentials - Adjust service dependencies to wait for ClickHouse startup
This commit is contained in:
@@ -12,6 +12,7 @@
|
|||||||
|
|
||||||
- Kafka (источник событий, 1 JSON message = 1 event)
|
- Kafka (источник событий, 1 JSON message = 1 event)
|
||||||
- ClickHouse (STG → ODS → DDS → DM)
|
- ClickHouse (STG → ODS → DDS → DM)
|
||||||
|
- **Airflow** (оркестрация ETL-пайплайна)
|
||||||
- Superset (BI поверх витрин)
|
- Superset (BI поверх витрин)
|
||||||
- Prometheus + Grafana (мониторинг)
|
- Prometheus + Grafana (мониторинг)
|
||||||
- простой инструмент/скрипт, который читает `.jsonl` и пишет события в Kafka
|
- простой инструмент/скрипт, который читает `.jsonl` и пишет события в Kafka
|
||||||
@@ -19,6 +20,7 @@
|
|||||||
## Ключевые артефакты
|
## Ключевые артефакты
|
||||||
|
|
||||||
### Исполняемые файлы (текущая структура)
|
### Исполняемые файлы (текущая структура)
|
||||||
|
- `dags/` — Airflow DAGs для оркестрации ETL
|
||||||
- `ddl/` — SQL для создания объектов БД:
|
- `ddl/` — SQL для создания объектов БД:
|
||||||
- `00_databases.sql` — создание БД stg/ods/dds/dm
|
- `00_databases.sql` — создание БД stg/ods/dds/dm
|
||||||
- `10_stg.sql` — STG слой (Kafka Engine + MV)
|
- `10_stg.sql` — STG слой (Kafka Engine + MV)
|
||||||
@@ -32,6 +34,8 @@
|
|||||||
- `apply_clickhouse_ddl.sh` — применение DDL
|
- `apply_clickhouse_ddl.sh` — применение DDL
|
||||||
- `load_kafka_data.sh` — загрузка в Kafka
|
- `load_kafka_data.sh` — загрузка в Kafka
|
||||||
- `run_batch.sh` — запуск batch-процесса
|
- `run_batch.sh` — запуск batch-процесса
|
||||||
|
- `airflow/` — конфигурация Airflow:
|
||||||
|
- `requirements.txt` — зависимости (clickhouse-connect и др.)
|
||||||
|
|
||||||
### Планы и документация (legacy)
|
### Планы и документация (legacy)
|
||||||
- `plans/clickhouse_ddl.md` — исходный план (inline DDL, legacy)
|
- `plans/clickhouse_ddl.md` — исходный план (inline DDL, legacy)
|
||||||
@@ -67,6 +71,7 @@
|
|||||||
- ClickHouse HTTP: `localhost:9123`
|
- ClickHouse HTTP: `localhost:9123`
|
||||||
- Kafka: `localhost:9092`
|
- Kafka: `localhost:9092`
|
||||||
- Kafka UI: `http://localhost:8082`
|
- Kafka UI: `http://localhost:8082`
|
||||||
|
- **Airflow: `http://localhost:8080` (admin/admin)**
|
||||||
- Prometheus: `http://localhost:9090`
|
- Prometheus: `http://localhost:9090`
|
||||||
- Grafana: `http://localhost:3000`
|
- Grafana: `http://localhost:3000`
|
||||||
|
|
||||||
@@ -89,6 +94,7 @@
|
|||||||
|
|
||||||
- Kafka ingest: наличие данных в `stg.*` и типизированных строк в `ods.*`.
|
- Kafka ingest: наличие данных в `stg.*` и типизированных строк в `ods.*`.
|
||||||
- Мониторинг: доступность `/metrics` у ClickHouse и скрейп в Prometheus.
|
- Мониторинг: доступность `/metrics` у ClickHouse и скрейп в Prometheus.
|
||||||
|
- **Airflow: `http://localhost:8080` должен показывать UI и DAG `etl_pipeline`.**
|
||||||
- BI: витрина `dm.v_events_enriched` должна отвечать за разумное время при фильтре по дате.
|
- BI: витрина `dm.v_events_enriched` должна отвечать за разумное время при фильтре по дате.
|
||||||
|
|
||||||
## Связанная документация
|
## Связанная документация
|
||||||
|
|||||||
+1
-1
@@ -5,7 +5,7 @@ FROM apache/airflow:2.9.3
|
|||||||
ENV AIRFLOW_HOME=/opt/airflow
|
ENV AIRFLOW_HOME=/opt/airflow
|
||||||
|
|
||||||
# Copy the project requirements file
|
# Copy the project requirements file
|
||||||
COPY --chown=airflow:0 requirements.txt ${AIRFLOW_HOME}/requirements.txt
|
COPY --chown=airflow:0 airflow/requirements.txt ${AIRFLOW_HOME}/requirements.txt
|
||||||
|
|
||||||
# Install the required Python packages
|
# Install the required Python packages
|
||||||
RUN pip install --no-cache-dir -r ${AIRFLOW_HOME}/requirements.txt
|
RUN pip install --no-cache-dir -r ${AIRFLOW_HOME}/requirements.txt
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
# ClickHouse Mini DWH для кликстрима
|
# ClickHouse Mini DWH для кликстрима
|
||||||
|
|
||||||
[](./docker-compose.yml)
|
[](./docker-compose.yml)
|
||||||
[](./docs/ARCHITECTURE.md)
|
[](./docs/ARCHITECTURE.md)
|
||||||
[]()
|
[]()
|
||||||
|
|
||||||
@@ -57,6 +57,7 @@ docker compose exec clickhouse clickhouse-client \
|
|||||||
|--------|-----|------------|
|
|--------|-----|------------|
|
||||||
| ClickHouse HTTP | http://localhost:9123/play | SQL-запросы |
|
| ClickHouse HTTP | http://localhost:9123/play | SQL-запросы |
|
||||||
| Kafka UI | http://localhost:8082 | Просмотр топиков |
|
| Kafka UI | http://localhost:8082 | Просмотр топиков |
|
||||||
|
| Airflow | http://localhost:8080 | Оркестрация ETL (admin/admin) |
|
||||||
| Superset | http://localhost:8088 | BI-дашборды |
|
| Superset | http://localhost:8088 | BI-дашборды |
|
||||||
| Prometheus | http://localhost:9090 | Метрики |
|
| Prometheus | http://localhost:9090 | Метрики |
|
||||||
| Grafana | http://localhost:3000 | Визуализация метрик |
|
| Grafana | http://localhost:3000 | Визуализация метрик |
|
||||||
@@ -85,7 +86,12 @@ flowchart TB
|
|||||||
DM["DM — витрины"]
|
DM["DM — витрины"]
|
||||||
end
|
end
|
||||||
|
|
||||||
|
subgraph Airflow["⚙️ Airflow"]
|
||||||
|
DAG[ETL DAGs]
|
||||||
|
end
|
||||||
|
|
||||||
Sources -->|make data| Kafka -->|MV| STG -->|MV| ODS -->|Batch SQL| DDS -->|VIEW| DM
|
Sources -->|make data| Kafka -->|MV| STG -->|MV| ODS -->|Batch SQL| DDS -->|VIEW| DM
|
||||||
|
DAG -.->|оркестрация| ODS & DDS & DM
|
||||||
```
|
```
|
||||||
|
|
||||||
**Поток данных:**
|
**Поток данных:**
|
||||||
@@ -102,9 +108,12 @@ flowchart TB
|
|||||||
|
|
||||||
```
|
```
|
||||||
.
|
.
|
||||||
|
├── dags/ # Airflow DAGs для оркестрации
|
||||||
├── ddl/ # SQL для создания объектов (00_databases → 40_dm)
|
├── ddl/ # SQL для создания объектов (00_databases → 40_dm)
|
||||||
├── jobs/ # Batch-трансформации (ODS→DDS, DDS→DM)
|
├── jobs/ # Batch-трансформации (ODS→DDS, DDS→DM)
|
||||||
├── scripts/ # Автоматизация (apply ddl, load data, run batch)
|
├── scripts/ # Автоматизация (apply ddl, load data, run batch)
|
||||||
|
├── airflow/ # Конфигурация Airflow
|
||||||
|
│ └── requirements.txt
|
||||||
├── docs/ # Документация
|
├── docs/ # Документация
|
||||||
│ └── ARCHITECTURE.md # Подробное описание слоёв
|
│ └── ARCHITECTURE.md # Подробное описание слоёв
|
||||||
├── data/ # Исходные JSONL файлы
|
├── data/ # Исходные JSONL файлы
|
||||||
@@ -187,7 +196,10 @@ flowchart LR
|
|||||||
|
|
||||||
## 🔮 Развитие проекта
|
## 🔮 Развитие проекта
|
||||||
|
|
||||||
- [ ] **Airflow** — оркестрация batch-процесса
|
### ✅ Реализовано
|
||||||
|
- [x] **Airflow** — оркестрация batch-процесса (инфраструктура готова, DAGs в разработке)
|
||||||
|
|
||||||
|
### 📋 В планах
|
||||||
- [ ] **Инкрементальный batch** — watermark-based загрузка
|
- [ ] **Инкрементальный batch** — watermark-based загрузка
|
||||||
- [ ] **Материализация витрин** — для тяжёлых агрегаций
|
- [ ] **Материализация витрин** — для тяжёлых агрегаций
|
||||||
- [ ] **DQ мониторинг** — алерты на ошибки парсинга
|
- [ ] **DQ мониторинг** — алерты на ошибки парсинга
|
||||||
|
|||||||
@@ -1,14 +1,11 @@
|
|||||||
# Optimized requirements for Airflow Docker setup
|
# Airflow requirements для ClickHouse DWH проекта
|
||||||
# Only includes packages actually used in DAGs
|
|
||||||
|
|
||||||
# Core database connector
|
# Core database connector для metadata
|
||||||
psycopg2-binary==2.9.9
|
psycopg2-binary==2.9.9
|
||||||
|
|
||||||
# Data processing (used in data_processing_dag.py and file_operations_dag.py)
|
# ClickHouse provider для ETL
|
||||||
|
# Примечание: официальный провайдер deprecated, используем clickhouse-connect
|
||||||
|
clickhouse-connect==0.8.0
|
||||||
|
|
||||||
|
# Для работы с данными
|
||||||
pandas==2.1.4
|
pandas==2.1.4
|
||||||
|
|
||||||
# Data generation library (used in file_operations_dag.py)
|
|
||||||
mimesis==15.1.0
|
|
||||||
|
|
||||||
# Airflow PostgreSQL provider (used in sql_basic_dag.py and data_processing_dag.py)
|
|
||||||
apache-airflow-providers-postgres==5.11.1
|
|
||||||
|
|||||||
@@ -0,0 +1,46 @@
|
|||||||
|
"""
|
||||||
|
ETL Pipeline DAG для ClickHouse DWH
|
||||||
|
|
||||||
|
Шаблон DAG для оркестрации пайплайна данных.
|
||||||
|
Полная реализация будет добавлена позже.
|
||||||
|
|
||||||
|
Пайплайн:
|
||||||
|
1. DDL - создание структуры БД
|
||||||
|
2. Load - загрузка данных в Kafka
|
||||||
|
3. Transform - batch трансформация ODS → DDS → DM
|
||||||
|
"""
|
||||||
|
|
||||||
|
from datetime import datetime, timedelta
|
||||||
|
from airflow import DAG
|
||||||
|
from airflow.operators.bash import BashOperator
|
||||||
|
from airflow.operators.empty import EmptyOperator
|
||||||
|
|
||||||
|
# Базовые настройки DAG
|
||||||
|
default_args = {
|
||||||
|
"owner": "airflow",
|
||||||
|
"depends_on_past": False,
|
||||||
|
"email_on_failure": False,
|
||||||
|
"email_on_retry": False,
|
||||||
|
"retries": 1,
|
||||||
|
"retry_delay": timedelta(minutes=5),
|
||||||
|
}
|
||||||
|
|
||||||
|
with DAG(
|
||||||
|
dag_id="etl_pipeline",
|
||||||
|
default_args=default_args,
|
||||||
|
description="ETL pipeline для ClickHouse DWH",
|
||||||
|
schedule=None, # Запуск только вручную (пока)
|
||||||
|
start_date=datetime(2024, 1, 1),
|
||||||
|
catchup=False,
|
||||||
|
tags=["etl", "clickhouse", "dwh"],
|
||||||
|
) as dag:
|
||||||
|
|
||||||
|
# TODO: добавить задачи пайплайна
|
||||||
|
# - ddl: создание структуры БД
|
||||||
|
# - load: загрузка данных в Kafka
|
||||||
|
# - transform: batch трансформация
|
||||||
|
|
||||||
|
start = EmptyOperator(task_id="start")
|
||||||
|
end = EmptyOperator(task_id="end")
|
||||||
|
|
||||||
|
start >> end
|
||||||
+19
-19
@@ -2,13 +2,9 @@ x-airflow-env: &airflow-default-env
|
|||||||
AIRFLOW__CORE__LOAD_EXAMPLES: "False"
|
AIRFLOW__CORE__LOAD_EXAMPLES: "False"
|
||||||
AIRFLOW__DATABASE__SQL_ALCHEMY_CONN: postgresql+psycopg2://airflow:airflow@postgres-metadata:5432/airflow
|
AIRFLOW__DATABASE__SQL_ALCHEMY_CONN: postgresql+psycopg2://airflow:airflow@postgres-metadata:5432/airflow
|
||||||
AIRFLOW__CORE__EXECUTOR: LocalExecutor
|
AIRFLOW__CORE__EXECUTOR: LocalExecutor
|
||||||
AIRFLOW__WEBSERVER__SECRET_KEY: replace-me-with-random-string
|
AIRFLOW__WEBSERVER__SECRET_KEY: ${AIRFLOW_SECRET_KEY:-replace-me-with-random-string}
|
||||||
POSTGRES_TRAINING_HOST: postgres-training
|
# ClickHouse connection для ETL
|
||||||
POSTGRES_TRAINING_PORT: "5432"
|
AIRFLOW_CONN_CLICKHOUSE_DEFAULT: clickhouse://default:123456@clickhouse:8123/default
|
||||||
POSTGRES_TRAINING_DB: training
|
|
||||||
POSTGRES_TRAINING_USER: student
|
|
||||||
POSTGRES_TRAINING_PASSWORD: student
|
|
||||||
AIRFLOW_CONN_POSTGRES_TRAINING: postgresql://student:student@postgres-training:5432/training
|
|
||||||
|
|
||||||
|
|
||||||
services:
|
services:
|
||||||
@@ -98,8 +94,10 @@ services:
|
|||||||
|
|
||||||
# Airflow webserver
|
# Airflow webserver
|
||||||
airflow-webserver:
|
airflow-webserver:
|
||||||
build: .
|
build:
|
||||||
image: airflow-optimized:2.9.2
|
context: .
|
||||||
|
dockerfile: Dockerfile.airflow
|
||||||
|
image: airflow-optimized:2.9.3
|
||||||
environment:
|
environment:
|
||||||
<<: *airflow-default-env
|
<<: *airflow-default-env
|
||||||
command: >
|
command: >
|
||||||
@@ -118,13 +116,14 @@ services:
|
|||||||
condition: service_completed_successfully
|
condition: service_completed_successfully
|
||||||
postgres-metadata:
|
postgres-metadata:
|
||||||
condition: service_healthy
|
condition: service_healthy
|
||||||
postgres-training:
|
clickhouse:
|
||||||
condition: service_healthy
|
condition: service_started
|
||||||
|
|
||||||
# Airflow scheduler
|
# Airflow scheduler
|
||||||
airflow-scheduler:
|
airflow-scheduler:
|
||||||
build: .
|
build:
|
||||||
# image: airflow-optimized:2.9.2
|
context: .
|
||||||
|
dockerfile: Dockerfile.airflow
|
||||||
environment:
|
environment:
|
||||||
<<: *airflow-default-env
|
<<: *airflow-default-env
|
||||||
command: >
|
command: >
|
||||||
@@ -141,14 +140,15 @@ services:
|
|||||||
condition: service_completed_successfully
|
condition: service_completed_successfully
|
||||||
postgres-metadata:
|
postgres-metadata:
|
||||||
condition: service_healthy
|
condition: service_healthy
|
||||||
postgres-training:
|
clickhouse:
|
||||||
condition: service_healthy
|
condition: service_started
|
||||||
|
|
||||||
# Airflow init to create admin user
|
# Airflow init to create admin user
|
||||||
airflow-init:
|
airflow-init:
|
||||||
build: .
|
build:
|
||||||
|
context: .
|
||||||
|
dockerfile: Dockerfile.airflow
|
||||||
user: "0:0"
|
user: "0:0"
|
||||||
# image: airflow-optimized:2.9.3
|
|
||||||
environment:
|
environment:
|
||||||
<<: *airflow-default-env
|
<<: *airflow-default-env
|
||||||
volumes:
|
volumes:
|
||||||
@@ -164,8 +164,8 @@ services:
|
|||||||
umask 000 &&
|
umask 000 &&
|
||||||
su -s /bin/bash airflow -c 'airflow db migrate' &&
|
su -s /bin/bash airflow -c 'airflow db migrate' &&
|
||||||
su -s /bin/bash airflow -c 'airflow users create --username admin --password admin --firstname Admin --lastname User --role Admin --email admin@example.org' &&
|
su -s /bin/bash airflow -c 'airflow users create --username admin --password admin --firstname Admin --lastname User --role Admin --email admin@example.org' &&
|
||||||
su -s /bin/bash airflow -c 'airflow connections delete postgres_training || true' &&
|
su -s /bin/bash airflow -c 'airflow connections delete clickhouse_default || true' &&
|
||||||
su -s /bin/bash airflow -c \"airflow connections add 'postgres_training' --conn-uri \\\"$${AIRFLOW_CONN_POSTGRES_TRAINING}\\\" --conn-description 'Training exercises database'\"
|
su -s /bin/bash airflow -c \"airflow connections add 'clickhouse_default' --conn-uri \\\"$${AIRFLOW_CONN_CLICKHOUSE_DEFAULT}\\\" --conn-description 'ClickHouse DWH'\"
|
||||||
"
|
"
|
||||||
depends_on:
|
depends_on:
|
||||||
postgres-metadata:
|
postgres-metadata:
|
||||||
|
|||||||
@@ -586,9 +586,11 @@ INSERT INTO dm.daily_traffic SELECT * FROM dm.v_daily_traffic;
|
|||||||
|
|
||||||
### Airflow-оркестрация
|
### Airflow-оркестрация
|
||||||
|
|
||||||
|
Инфраструктура Airflow развёрнута и готова к использованию:
|
||||||
|
|
||||||
```python
|
```python
|
||||||
# dag.py
|
# dags/etl_pipeline_dag.py
|
||||||
with DAG('clickhouse_etl'):
|
with DAG('etl_pipeline'):
|
||||||
ddl = BashOperator(task_id='ddl', bash_command='make ddl')
|
ddl = BashOperator(task_id='ddl', bash_command='make ddl')
|
||||||
load = BashOperator(task_id='load', bash_command='make data')
|
load = BashOperator(task_id='load', bash_command='make data')
|
||||||
transform = BashOperator(task_id='transform', bash_command='make transform')
|
transform = BashOperator(task_id='transform', bash_command='make transform')
|
||||||
@@ -596,6 +598,11 @@ with DAG('clickhouse_etl'):
|
|||||||
ddl >> load >> transform
|
ddl >> load >> transform
|
||||||
```
|
```
|
||||||
|
|
||||||
|
**Подключение к ClickHouse:**
|
||||||
|
- Connection: `clickhouse_default`
|
||||||
|
- URL: `clickhouse://default:123456@clickhouse:8123/default`
|
||||||
|
- Provider: `clickhouse-connect` (в `airflow/requirements.txt`)
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
## Полезные запросы
|
## Полезные запросы
|
||||||
|
|||||||
Reference in New Issue
Block a user