Merge branch 'feature/airflow-orchestration'
This commit is contained in:
@@ -20,11 +20,11 @@
|
||||
## Ключевые артефакты
|
||||
|
||||
### Исполняемые файлы (текущая структура)
|
||||
- `dags/` — Airflow DAGs для оркестрации ETL:
|
||||
- `dags/ddl_init_dag.py` — инициализация схемы ClickHouse
|
||||
- `dags/etl_pipeline_dag.py` — ETL процесс STG → ODS → DDS → DM
|
||||
- `dags/kafka_load_dag.py` — загрузка данных в Kafka из JSONL
|
||||
- `dags/utils/kafka_helpers.py` — helper-функции для работы с Kafka
|
||||
- `airflow/dags/` — Airflow DAGs для оркестрации ETL:
|
||||
- `airflow/dags/ddl_init_dag.py` — инициализация схемы ClickHouse
|
||||
- `airflow/dags/etl_pipeline_dag.py` — ETL процесс STG → ODS → DDS → DM
|
||||
- `airflow/dags/kafka_load_dag.py` — загрузка данных в Kafka из JSONL
|
||||
- `airflow/dags/utils/kafka_helpers.py` — helper-функции для работы с Kafka
|
||||
- `sql/` — SQL по слоям:
|
||||
- `sql/ddl/00_databases.sql` — создание БД stg/ods/dds/dm
|
||||
- `sql/ddl/stg/10_stg.sql` — STG слой (Kafka Engine + MV)
|
||||
|
||||
@@ -115,7 +115,6 @@ flowchart LR
|
||||
|
||||
```
|
||||
.
|
||||
├── dags/ # Airflow DAGs для оркестрации
|
||||
├── sql/
|
||||
│ ├── ddl/ # DDL по слоям
|
||||
│ │ ├── 00_databases.sql
|
||||
@@ -128,6 +127,7 @@ flowchart LR
|
||||
│ └── dm/ # Batch SQL: DDS -> DM
|
||||
├── scripts/ # Служебные shell-скрипты (legacy fallback, не основной путь)
|
||||
├── airflow/ # Конфигурация Airflow
|
||||
│ ├── dags/ # Airflow DAGs для оркестрации
|
||||
│ └── requirements.txt
|
||||
├── docs/ # Документация
|
||||
│ └── ARCHITECTURE.md # Подробное описание слоёв
|
||||
|
||||
@@ -38,7 +38,19 @@ default_args = {
|
||||
# -----------------------------------------------------------------------------
|
||||
# SQL-файлы проекта
|
||||
# -----------------------------------------------------------------------------
|
||||
SQL_ROOT = Path(__file__).resolve().parents[1] / "sql"
|
||||
def resolve_sql_root() -> Path:
|
||||
"""Определяет корень SQL для контейнера и локального запуска."""
|
||||
candidates = (
|
||||
Path(__file__).resolve().parents[1] / "sql", # /opt/airflow/sql в контейнере
|
||||
Path(__file__).resolve().parents[2] / "sql", # <repo>/sql при локальном запуске
|
||||
)
|
||||
for candidate in candidates:
|
||||
if candidate.is_dir():
|
||||
return candidate
|
||||
return candidates[0]
|
||||
|
||||
|
||||
SQL_ROOT = resolve_sql_root()
|
||||
|
||||
|
||||
def load_sql_statements(relative_path: str) -> tuple[str, ...]:
|
||||
@@ -41,7 +41,19 @@ default_args = {
|
||||
# -----------------------------------------------------------------------------
|
||||
# SQL-файлы проекта
|
||||
# -----------------------------------------------------------------------------
|
||||
SQL_ROOT = Path(__file__).resolve().parents[1] / "sql"
|
||||
def resolve_sql_root() -> Path:
|
||||
"""Определяет корень SQL для контейнера и локального запуска."""
|
||||
candidates = (
|
||||
Path(__file__).resolve().parents[1] / "sql", # /opt/airflow/sql в контейнере
|
||||
Path(__file__).resolve().parents[2] / "sql", # <repo>/sql при локальном запуске
|
||||
)
|
||||
for candidate in candidates:
|
||||
if candidate.is_dir():
|
||||
return candidate
|
||||
return candidates[0]
|
||||
|
||||
|
||||
SQL_ROOT = resolve_sql_root()
|
||||
|
||||
|
||||
def load_sql_statements(relative_path: str) -> tuple[str, ...]:
|
||||
Binary file not shown.
Binary file not shown.
+3
-3
@@ -108,7 +108,7 @@ services:
|
||||
ports:
|
||||
- "8080:8080"
|
||||
volumes:
|
||||
- ./dags:/opt/airflow/dags
|
||||
- ./airflow/dags:/opt/airflow/dags
|
||||
- ./sql:/opt/airflow/sql:ro
|
||||
- ./data:/opt/airflow/data
|
||||
networks:
|
||||
@@ -133,7 +133,7 @@ services:
|
||||
airflow scheduler
|
||||
"
|
||||
volumes:
|
||||
- ./dags:/opt/airflow/dags
|
||||
- ./airflow/dags:/opt/airflow/dags
|
||||
- ./sql:/opt/airflow/sql:ro
|
||||
- ./data:/opt/airflow/data
|
||||
networks:
|
||||
@@ -155,7 +155,7 @@ services:
|
||||
environment:
|
||||
<<: *airflow-default-env
|
||||
volumes:
|
||||
- ./dags:/opt/airflow/dags
|
||||
- ./airflow/dags:/opt/airflow/dags
|
||||
- ./sql:/opt/airflow/sql:ro
|
||||
- ./data:/opt/airflow/data
|
||||
networks:
|
||||
|
||||
@@ -578,9 +578,9 @@ INSERT INTO dm.daily_traffic SELECT * FROM dm.v_daily_traffic;
|
||||
Инфраструктура Airflow развёрнута и готова к использованию:
|
||||
|
||||
```python
|
||||
# dags/ddl_init_dag.py — создание баз/таблиц (ручной запуск при bootstrap)
|
||||
# dags/kafka_load_dag.py — загрузка JSONL в Kafka (через kafka-python)
|
||||
# dags/etl_pipeline_dag.py — основной ETL (STG→ODS→DDS→DM)
|
||||
# airflow/dags/ddl_init_dag.py — создание баз/таблиц (ручной запуск при bootstrap)
|
||||
# airflow/dags/kafka_load_dag.py — загрузка JSONL в Kafka (через kafka-python)
|
||||
# airflow/dags/etl_pipeline_dag.py — основной ETL (STG→ODS→DDS→DM)
|
||||
|
||||
# Учебный формат:
|
||||
# - DDL и трансформации выполняются явными SQL-task через ClickHouseOperator;
|
||||
|
||||
@@ -146,18 +146,18 @@ ddl_init -> kafka_load -> etl_pipeline
|
||||
- публикации строк из `.jsonl` (`1 строка = 1 message value`).
|
||||
|
||||
### Общие helper-функции
|
||||
- `dags/utils/clickhouse_helpers.py`:
|
||||
- `airflow/dags/utils/clickhouse_helpers.py`:
|
||||
- `execute_sql(sql: str) -> None`
|
||||
- `execute_sql_file(path: str) -> None`
|
||||
- `fetch_one(sql: str) -> tuple`
|
||||
- `dags/utils/kafka_helpers.py`:
|
||||
- `airflow/dags/utils/kafka_helpers.py`:
|
||||
- `prepare_topics(reset: bool) -> None`
|
||||
- `load_jsonl(file_path: str, topic: str, limit: int) -> int`
|
||||
- `check_kafka_ready() -> None`
|
||||
|
||||
## Структура файлов
|
||||
```text
|
||||
dags/
|
||||
airflow/dags/
|
||||
├── __init__.py
|
||||
├── ddl_init_dag.py # отдельный DAG для DDL (обязателен)
|
||||
├── kafka_load_dag.py # отдельный DAG для ingest в Kafka (обязателен)
|
||||
|
||||
Reference in New Issue
Block a user