diff --git a/AGENTS.md b/AGENTS.md index 8f72fdd..d03e985 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -108,15 +108,12 @@ - Параметры: - `limit` (int, default 0) — количество строк (0 = все) - `reset_topics` (bool, default true) — пересоздать топики - - `load_browser/load_location/load_device/load_geo` — выбор потоков - Примеры запуска: ```json // Полная загрузка (по умолчанию) {} // Ограниченная загрузка — 100 строк {"limit": 100} - // Только browser_events - {"limit": 100, "load_location": false, "load_device": false, "load_geo": false} ``` ### `etl_pipeline` — ETL процесс @@ -143,7 +140,7 @@ make up # 3. Загрузка данных через Airflow (вместо make data) # Airflow UI → DAGs → kafka_load → Trigger DAG with config -# Параметры по умолчанию: limit=50, все потоки +# Параметры по умолчанию: limit=0 (полная загрузка), reset_topics=true # 4. Запуск ETL # Airflow UI → DAGs → etl_pipeline → Trigger DAG with config diff --git a/README.md b/README.md index 1cf85fa..f991096 100644 --- a/README.md +++ b/README.md @@ -106,7 +106,7 @@ flowchart TB DAG[DAG: ddl_init / kafka_load / etl_pipeline] end - Sources -->|make data| Kafka -->|MV| STG -->|Batch SQL| ODS -->|Batch SQL| DDS -->|VIEW| DM + Sources -->|kafka_load / make data| Kafka -->|MV| STG -->|Batch SQL| ODS -->|Batch SQL| DDS -->|VIEW| DM DAG -.->|оркестрация| STG & ODS & DDS & DM ``` diff --git a/dags/kafka_load_dag.py b/dags/kafka_load_dag.py index 7baa4ce..484f8c8 100644 --- a/dags/kafka_load_dag.py +++ b/dags/kafka_load_dag.py @@ -8,13 +8,8 @@ DAG для загрузки данных в Kafka из JSONL-файлов. - Проверка результатов через XCom Параметры (через Trigger DAG with config): - limit (int): Количество строк для загрузки (по умолчанию 50) - full_load (bool): Загрузить всё, игнорируя limit (по умолчанию False) + limit (int): Количество строк для загрузки (по умолчанию 0 = все) reset_topics (bool): Пересоздать топики (по умолчанию True) - load_browser (bool): Загружать browser_events (по умолчанию True) - load_location (bool): Загружать location_events (по умолчанию True) - load_device (bool): Загружать device_events (по умолчанию True) - load_geo (bool): Загружать geo_events (по умолчанию True) """ from __future__ import annotations @@ -32,7 +27,6 @@ from airflow.utils.task_group import TaskGroup from utils.kafka_helpers import ( check_input_files, check_kafka_ready, - get_topic_file_mapping, load_jsonl, prepare_topics, validate_load_params, @@ -62,58 +56,22 @@ def _check_kafka(**context) -> None: def _check_files(**context) -> None: """Проверка наличия входных файлов.""" - conf = context.get("dag_run", {}).conf or {} - load_browser = bool(conf.get("load_browser", context["params"]["load_browser"])) - load_location = bool(conf.get("load_location", context["params"]["load_location"])) - load_device = bool(conf.get("load_device", context["params"]["load_device"])) - load_geo = bool(conf.get("load_geo", context["params"]["load_geo"])) - - check_input_files( - data_dir=DATA_DIR, - load_browser=load_browser, - load_location=load_location, - load_device=load_device, - load_geo=load_geo, - ) + check_input_files(data_dir=DATA_DIR) def _validate_params(**context) -> None: """Валидация параметров загрузки.""" conf = context.get("dag_run", {}).conf or {} limit = int(conf.get("limit", context["params"]["limit"])) - load_browser = bool(conf.get("load_browser", context["params"]["load_browser"])) - load_location = bool(conf.get("load_location", context["params"]["load_location"])) - load_device = bool(conf.get("load_device", context["params"]["load_device"])) - load_geo = bool(conf.get("load_geo", context["params"]["load_geo"])) - - validate_load_params( - limit=limit, - load_browser=load_browser, - load_location=load_location, - load_device=load_device, - load_geo=load_geo, - ) + validate_load_params(limit=limit) def _prepare_topics(**context) -> None: """Подготовка топиков Kafka (создание/сброс).""" conf = context.get("dag_run", {}).conf or {} reset_topics = bool(conf.get("reset_topics", context["params"]["reset_topics"])) - load_browser = bool(conf.get("load_browser", context["params"]["load_browser"])) - load_location = bool(conf.get("load_location", context["params"]["load_location"])) - load_device = bool(conf.get("load_device", context["params"]["load_device"])) - load_geo = bool(conf.get("load_geo", context["params"]["load_geo"])) - # Получаем список топиков для загрузки - mapping = get_topic_file_mapping( - load_browser=load_browser, - load_location=load_location, - load_device=load_device, - load_geo=load_geo, - ) - topics = list(mapping.keys()) - - prepare_topics(topics=topics, reset=reset_topics) + prepare_topics(reset=reset_topics) def _load_events(event_type: str, **context) -> int: @@ -155,37 +113,16 @@ def _load_events(event_type: str, **context) -> int: def _verify_counts(**context) -> None: """Проверяет, что все загрузки отправили сообщения.""" - conf = context.get("dag_run", {}).conf or {} - load_browser = bool(conf.get("load_browser", context["params"]["load_browser"])) - load_location = bool(conf.get("load_location", context["params"]["load_location"])) - load_device = bool(conf.get("load_device", context["params"]["load_device"])) - load_geo = bool(conf.get("load_geo", context["params"]["load_geo"])) - ti = context["ti"] # Собираем результаты из XCom - results = {} - total_sent = 0 - - if load_browser: - count = ti.xcom_pull(task_ids="ingest.load_browser_events", key="browser_count") - results["browser_events"] = count or 0 - total_sent += count or 0 - - if load_location: - count = ti.xcom_pull(task_ids="ingest.load_location_events", key="location_count") - results["location_events"] = count or 0 - total_sent += count or 0 - - if load_device: - count = ti.xcom_pull(task_ids="ingest.load_device_events", key="device_count") - results["device_events"] = count or 0 - total_sent += count or 0 - - if load_geo: - count = ti.xcom_pull(task_ids="ingest.load_geo_events", key="geo_count") - results["geo_events"] = count or 0 - total_sent += count or 0 + results = { + "browser_events": ti.xcom_pull(task_ids="ingest.load_browser_events", key="browser_count") or 0, + "location_events": ti.xcom_pull(task_ids="ingest.load_location_events", key="location_count") or 0, + "device_events": ti.xcom_pull(task_ids="ingest.load_device_events", key="device_count") or 0, + "geo_events": ti.xcom_pull(task_ids="ingest.load_geo_events", key="geo_count") or 0, + } + total_sent = sum(results.values()) # Проверяем, что отправлено хотя бы что-то if total_sent == 0: @@ -216,36 +153,11 @@ with DAG( type="integer", description="Количество строк для загрузки (0 = все строки, по умолчанию)", ), - "full_load": Param( - default=False, - type="boolean", - description="Загрузить все данные, игнорируя limit", - ), "reset_topics": Param( default=True, type="boolean", description="Пересоздать топики перед загрузкой", ), - "load_browser": Param( - default=True, - type="boolean", - description="Загружать browser_events", - ), - "load_location": Param( - default=True, - type="boolean", - description="Загружать location_events", - ), - "load_device": Param( - default=True, - type="boolean", - description="Загружать device_events", - ), - "load_geo": Param( - default=True, - type="boolean", - description="Загружать geo_events", - ), }, ) as dag: diff --git a/dags/utils/kafka_helpers.py b/dags/utils/kafka_helpers.py index 47988a5..ac142ca 100644 --- a/dags/utils/kafka_helpers.py +++ b/dags/utils/kafka_helpers.py @@ -10,10 +10,6 @@ from __future__ import annotations import logging from pathlib import Path -from typing import TYPE_CHECKING - -if TYPE_CHECKING: - from typing import Optional # ----------------------------------------------------------------------------- # Конфигурация подключения к Kafka @@ -72,7 +68,7 @@ def check_kafka_ready( # Управление топиками # ----------------------------------------------------------------------------- def prepare_topics( - topics: Optional[list[str]] = None, + topics: list[str] | None = None, reset: bool = True, bootstrap_servers: str = KAFKA_BOOTSTRAP_SERVERS, timeout_ms: int = REQUEST_TIMEOUT_MS, @@ -237,21 +233,12 @@ def load_jsonl( # ----------------------------------------------------------------------------- def validate_load_params( limit: int, - load_browser: bool, - load_location: bool, - load_device: bool, - load_geo: bool, ) -> None: """ Валидирует параметры загрузки данных. Args: limit: Количество строк для загрузки - full_load: Флаг полной загрузки - load_browser: Загружать browser_events - load_location: Загружать location_events - load_device: Загружать device_events - load_geo: Загружать geo_events Raises: AirflowException: Если параметры невалидны @@ -262,27 +249,15 @@ def validate_load_params( if not isinstance(limit, int) or limit < 0: raise AirflowException(f"limit должен быть неотрицательным int, получено: {limit}") - # Проверка что хотя бы один поток выбран - if not any([load_browser, load_location, load_device, load_geo]): - raise AirflowException("Должен быть выбран хотя бы один поток для загрузки") - def check_input_files( data_dir: str | Path = "/opt/airflow/data", - load_browser: bool = True, - load_location: bool = True, - load_device: bool = True, - load_geo: bool = True, ) -> None: """ Проверяет наличие необходимых JSONL-файлов. Args: data_dir: Директория с данными - load_browser: Проверять browser_events.jsonl - load_location: Проверять location_events.jsonl - load_device: Проверять device_events.jsonl - load_geo: Проверять geo_events.jsonl Raises: AirflowException: Если какой-либо файл отсутствует @@ -290,16 +265,7 @@ def check_input_files( from airflow.exceptions import AirflowException data_dir = Path(data_dir) - files_to_check = [] - - if load_browser: - files_to_check.append("browser_events.jsonl") - if load_location: - files_to_check.append("location_events.jsonl") - if load_device: - files_to_check.append("device_events.jsonl") - if load_geo: - files_to_check.append("geo_events.jsonl") + files_to_check = list(TOPIC_FILE_MAP.values()) missing_files = [] for filename in files_to_check: @@ -313,30 +279,3 @@ def check_input_files( ) logger.info("Все необходимые файлы найдены: %s", files_to_check) - - -def get_topic_file_mapping( - load_browser: bool = True, - load_location: bool = True, - load_device: bool = True, - load_geo: bool = True, -) -> dict[str, str]: - """ - Возвращает маппинг топиков на файлы для выбранных потоков. - - Returns: - Словарь {topic_name: filename} - """ - result = {} - flags = { - "browser_events": load_browser, - "location_events": load_location, - "device_events": load_device, - "geo_events": load_geo, - } - - for topic, filename in TOPIC_FILE_MAP.items(): - if flags.get(topic, True): - result[topic] = filename - - return result diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index db32e21..18ce86b 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -615,7 +615,7 @@ INSERT INTO dm.daily_traffic SELECT * FROM dm.v_daily_traffic; - Загрузка данных из `data/*.jsonl` в Kafka через `kafka-python` - TaskGroup `precheck`: проверка Kafka, файлов, параметров - TaskGroup `ingest`: создание топиков → параллельная загрузка 4 потоков → проверка -- Параметры: `limit` (0 = все), `reset_topics`, `load_*` (выбор потоков) +- Параметры: `limit` (0 = все), `reset_topics` **Подключение к ClickHouse:** - Connection: `clickhouse_default` diff --git a/plans/airflow_dags_plan.md b/plans/airflow_dags_plan.md index f32f2f1..2b69ed5 100644 --- a/plans/airflow_dags_plan.md +++ b/plans/airflow_dags_plan.md @@ -63,17 +63,13 @@ check_clickhouse >> ddl_00_databases >> ddl_10_stg >> ddl_20_ods >> ddl_30_dds > ### Params (через Trigger DAG with config) - `limit`: int, default `0` (0 = все строки). - `reset_topics`: bool, default `true`. -- `load_browser`: bool, default `true`. -- `load_location`: bool, default `true`. -- `load_device`: bool, default `true`. -- `load_geo`: bool, default `true`. ### TaskGroup `precheck` | Task ID | Что делает | Реализация | |---------|------------|------------| | `check_kafka` | Проверка доступности Kafka broker | `PythonOperator` + `kafka-python` | | `check_input_files` | Проверка наличия `data/*_events.jsonl` | `PythonOperator` | -| `validate_load_params` | Валидация параметров загрузки (`limit`, `full_load`, флаги потоков) | `PythonOperator` | +| `validate_load_params` | Валидация параметров загрузки (`limit`) | `PythonOperator` | ### TaskGroup `ingest` | Task ID | Что делает | Реализация | @@ -83,7 +79,7 @@ check_clickhouse >> ddl_00_databases >> ddl_10_stg >> ddl_20_ods >> ddl_30_dds > | `load_location_events` | Публикация `location_events.jsonl` | `PythonOperator` + KafkaProducer | | `load_device_events` | Публикация `device_events.jsonl` | `PythonOperator` + KafkaProducer | | `load_geo_events` | Публикация `geo_events.jsonl` | `PythonOperator` + KafkaProducer | -| `verify_publish_counts` | Проверка, что отправлено > 0 сообщений в выбранные потоки | `PythonOperator` (по XCom) | +| `verify_publish_counts` | Проверка, что отправлено > 0 сообщений в сумме | `PythonOperator` (по XCom) | Зависимости: ```text @@ -135,7 +131,7 @@ ddl_init -> kafka_load -> etl_pipeline ``` Экспериментальные сценарии: -- `kafka_load` отдельно: проверить разные наборы/параметры загрузки без запуска transform. +- `kafka_load` отдельно: проверить разные параметры загрузки без запуска transform. - `etl_pipeline` отдельно: повторно пересчитать DDS/DM по уже загруженным данным. ## Техническая реализация (приземленно) @@ -156,7 +152,7 @@ ddl_init -> kafka_load -> etl_pipeline - `fetch_one(sql: str) -> tuple` - `dags/utils/kafka_helpers.py`: - `prepare_topics(reset: bool) -> None` - - `load_jsonl(file_path: str, topic: str, limit: int, full_load: bool) -> int` + - `load_jsonl(file_path: str, topic: str, limit: int) -> int` - `check_kafka_ready() -> None` ## Структура файлов @@ -181,7 +177,7 @@ dags/ 2. Этап 2: - Реализовать отдельный DAG `kafka_load` на `kafka-python`. - Перенести загрузку из `make data` в `kafka_load` (функциональный паритет). - - Добавить в `kafka_load` расширенные параметры экспериментов (выбор потоков, сценарии reset/no-reset). + - Добавить в `kafka_load` параметры экспериментов по объёму и reset/no-reset. - Добавить опциональный параметр автотриггера `etl_pipeline` (по умолчанию `false`). 3. Этап 3: - Добавить `dq_monitor` и alert callback (email/Slack/webhook). @@ -225,7 +221,7 @@ docker compose exec -T clickhouse clickhouse-client --user=default --password=12 docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query "SELECT count() FROM dm.dq_summary" # 7) Этап 2: после реализации kafka_load -# Trigger DAG kafka_load {"limit": 50, "full_load": false, "reset_topics": true} +# Trigger DAG kafka_load {"limit": 50, "reset_topics": true} # Trigger DAG etl_pipeline {"full_refresh": true} ``` diff --git a/plans/runbook.md b/plans/runbook.md index eb0be1c..ede4616 100644 --- a/plans/runbook.md +++ b/plans/runbook.md @@ -33,7 +33,6 @@ docker compose exec -T airflow-webserver airflow dags trigger kafka_load \ Параметры `kafka_load`: - `limit` — количество строк (default: 0 — все строки) - `reset_topics` — пересоздать топики (default: true) -- `load_browser/load_location/load_device/load_geo` — выбор потоков (default: true) **Вариант Б: Через shell-скрипт `make data` (устаревший)** ```bash