diff --git a/AGENTS.md b/AGENTS.md index d03e985..45d58e6 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -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) diff --git a/README.md b/README.md index 83a1ebf..cb84c58 100644 --- a/README.md +++ b/README.md @@ -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 # Подробное описание слоёв diff --git a/dags/.gitkeep b/airflow/dags/.gitkeep similarity index 100% rename from dags/.gitkeep rename to airflow/dags/.gitkeep diff --git a/dags/__init__.py b/airflow/dags/__init__.py similarity index 100% rename from dags/__init__.py rename to airflow/dags/__init__.py diff --git a/dags/ddl_init_dag.py b/airflow/dags/ddl_init_dag.py similarity index 92% rename from dags/ddl_init_dag.py rename to airflow/dags/ddl_init_dag.py index 9529177..ff34b0a 100644 --- a/dags/ddl_init_dag.py +++ b/airflow/dags/ddl_init_dag.py @@ -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", # /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, ...]: diff --git a/dags/etl_pipeline_dag.py b/airflow/dags/etl_pipeline_dag.py similarity index 95% rename from dags/etl_pipeline_dag.py rename to airflow/dags/etl_pipeline_dag.py index eaff6c2..ce43d2b 100644 --- a/dags/etl_pipeline_dag.py +++ b/airflow/dags/etl_pipeline_dag.py @@ -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", # /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, ...]: diff --git a/dags/kafka_load_dag.py b/airflow/dags/kafka_load_dag.py similarity index 100% rename from dags/kafka_load_dag.py rename to airflow/dags/kafka_load_dag.py diff --git a/dags/utils/__init__.py b/airflow/dags/utils/__init__.py similarity index 100% rename from dags/utils/__init__.py rename to airflow/dags/utils/__init__.py diff --git a/dags/utils/kafka_helpers.py b/airflow/dags/utils/kafka_helpers.py similarity index 100% rename from dags/utils/kafka_helpers.py rename to airflow/dags/utils/kafka_helpers.py diff --git a/dags/__pycache__/ddl_init_dag.cpython-312.pyc b/dags/__pycache__/ddl_init_dag.cpython-312.pyc deleted file mode 100644 index 9f0a801..0000000 Binary files a/dags/__pycache__/ddl_init_dag.cpython-312.pyc and /dev/null differ diff --git a/dags/__pycache__/etl_pipeline_dag.cpython-312.pyc b/dags/__pycache__/etl_pipeline_dag.cpython-312.pyc deleted file mode 100644 index aa2e353..0000000 Binary files a/dags/__pycache__/etl_pipeline_dag.cpython-312.pyc and /dev/null differ diff --git a/docker-compose.yml b/docker-compose.yml index 33bce41..d4804a2 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -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: diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 6503cfd..56b7362 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -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; diff --git a/plans/airflow_dags_plan.md b/plans/airflow_dags_plan.md index 2b69ed5..90d3868 100644 --- a/plans/airflow_dags_plan.md +++ b/plans/airflow_dags_plan.md @@ -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 (обязателен)