diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..bee8a64 --- /dev/null +++ b/.gitignore @@ -0,0 +1 @@ +__pycache__ diff --git a/AGENTS.md b/AGENTS.md index f479d88..8f72fdd 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -20,7 +20,11 @@ ## Ключевые артефакты ### Исполняемые файлы (текущая структура) -- `dags/` — Airflow DAGs для оркестрации ETL +- `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 - `sql/` — SQL по слоям: - `sql/ddl/00_databases.sql` — создание БД stg/ods/dds/dm - `sql/ddl/stg/10_stg.sql` — STG слой (Kafka Engine + MV) @@ -92,13 +96,65 @@ - Bash: шапка с назначением/запуском/требованиями, секции разделены `# -----` - См. существующие файлы как пример (`sql/ddl/ods/20_ods.sql`, `sql/dds/30_ods_to_dds.sql`, `scripts/run_batch.sh`) +## Airflow DAGs + +### `ddl_init` — Инициализация схемы +- Запуск: ручной (Trigger DAG) +- Параметры: `verify_only` (bool, default false) +- Описание: Создаёт БД и таблицы в ClickHouse от 00_databases до 40_dm + +### `kafka_load` — Загрузка в Kafka +- Запуск: ручный (Trigger DAG with config) +- Параметры: + - `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 процесс +- Запуск: ручной (Trigger DAG with config) +- Параметры: + - `full_refresh` (bool, default true) — очистить DDS перед загрузкой +- Зависимость: требует наличия данных в STG (от `kafka_load` или `make data`) + ## Быстрые проверки - Kafka ingest: наличие данных в `stg.*` и типизированных строк в `ods.*`. - Мониторинг: доступность `/metrics` у ClickHouse и скрейп в Prometheus. -- **Airflow: `http://localhost:8080` должен показывать UI и DAG `ddl_init` и `etl_pipeline`.** +- **Airflow: `http://localhost:8080` должен показывать UI и DAG `ddl_init`, `kafka_load` и `etl_pipeline`.** - BI: витрина `dm.v_events_enriched` должна отвечать за разумное время при фильтре по дате. +## Сценарий работы с Airflow (фаза 2) + +```bash +# 1. Запуск инфраструктуры +make up + +# 2. Инициализация схемы (один раз) +# Airflow UI → DAGs → ddl_init → Trigger DAG + +# 3. Загрузка данных через Airflow (вместо make data) +# Airflow UI → DAGs → kafka_load → Trigger DAG with config +# Параметры по умолчанию: limit=50, все потоки + +# 4. Запуск ETL +# Airflow UI → DAGs → etl_pipeline → Trigger DAG with config +# {"full_refresh": true} + +# 5. Проверка результатов +docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query "SELECT count() FROM ods.browser_event" +docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query "SELECT count() FROM dds.event" +docker compose exec -T clickhouse clickhouse-client --user=default --password=123456 --query "SELECT * FROM dm.dq_summary" +``` + ## Связанная документация - [README.md](./README.md) — пользовательская документация (быстрый старт, архитектура) diff --git a/README.md b/README.md index af8e27e..1cf85fa 100644 --- a/README.md +++ b/README.md @@ -33,10 +33,18 @@ docker compose ps docker compose exec -T airflow-webserver airflow dags trigger ddl_init ``` -Загрузка небольшого среза данных в Kafka: +Загрузка данных в Kafka (фаза 2 — через Airflow): ```bash -make data # по умолчанию первые 50 строк -# или: FULL=1 make data # полный датасет (1000 строк) +# Вариант 1: Через Airflow DAG (рекомендуется) — полная загрузка по умолчанию +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}' + +# Вариант 2: Через shell-скрипт (устаревший) +make data # полная загрузка ``` Запуск batch-трансформации (STG -> ODS -> DDS -> DM) в Airflow (если DAG выключен, сначала unpause): @@ -95,8 +103,8 @@ flowchart TB end subgraph Airflow["Airflow"] - DAG[DAG: ddl_init / etl_pipeline] - end + DAG[DAG: ddl_init / kafka_load / etl_pipeline] + end Sources -->|make data| Kafka -->|MV| STG -->|Batch SQL| ODS -->|Batch SQL| DDS -->|VIEW| DM DAG -.->|оркестрация| STG & ODS & DDS & DM diff --git a/airflow/requirements.txt b/airflow/requirements.txt index feae4fb..9504b3c 100644 --- a/airflow/requirements.txt +++ b/airflow/requirements.txt @@ -5,3 +5,6 @@ psycopg2-binary==2.9.9 # ClickHouse operator/hook для DAG'ов airflow-clickhouse-plugin==1.6.0 + +# Kafka client для загрузки данных (DAG kafka_load) +kafka-python==2.0.6 diff --git a/dags/kafka_load_dag.py b/dags/kafka_load_dag.py new file mode 100644 index 0000000..7baa4ce --- /dev/null +++ b/dags/kafka_load_dag.py @@ -0,0 +1,322 @@ +""" +DAG для загрузки данных в Kafka из JSONL-файлов. + +Функциональность: +- Проверка доступности Kafka и наличия файлов +- Создание/сброс топиков +- Параллельная загрузка 4 потоков данных +- Проверка результатов через XCom + +Параметры (через Trigger DAG with config): + limit (int): Количество строк для загрузки (по умолчанию 50) + full_load (bool): Загрузить всё, игнорируя limit (по умолчанию False) + 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 + +from datetime import datetime, timedelta +from pathlib import Path + +from airflow import DAG +from airflow.exceptions import AirflowException +from airflow.models.param import Param +from airflow.operators.python import PythonOperator +from airflow.utils.task_group import TaskGroup + +# Импортируем helper-функции +from utils.kafka_helpers import ( + check_input_files, + check_kafka_ready, + get_topic_file_mapping, + load_jsonl, + prepare_topics, + validate_load_params, +) + +# ----------------------------------------------------------------------------- +# Базовые настройки DAG +# ----------------------------------------------------------------------------- +default_args = { + "owner": "airflow", + "depends_on_past": False, + "email_on_failure": False, + "email_on_retry": False, + "retries": 1, + "retry_delay": timedelta(minutes=2), +} + +DATA_DIR = Path("/opt/airflow/data") + +# ----------------------------------------------------------------------------- +# Python callable функции для задач +# ----------------------------------------------------------------------------- +def _check_kafka(**context) -> None: + """Проверка доступности Kafka брокера.""" + check_kafka_ready() + + +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, + ) + + +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, + ) + + +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) + + +def _load_events(event_type: str, **context) -> int: + """ + Загружает события определённого типа в Kafka. + + Args: + event_type: Тип события (browser, location, device, geo) + + Returns: + Количество отправленных сообщений + """ + conf = context.get("dag_run", {}).conf or {} + limit = int(conf.get("limit", context["params"]["limit"])) + + # Маппинг типа события на топик и файл + event_mapping = { + "browser": ("browser_events", "browser_events.jsonl"), + "location": ("location_events", "location_events.jsonl"), + "device": ("device_events", "device_events.jsonl"), + "geo": ("geo_events", "geo_events.jsonl"), + } + + topic, filename = event_mapping[event_type] + file_path = DATA_DIR / filename + + # Загружаем данные + sent_count = load_jsonl( + file_path=file_path, + topic=topic, + limit=limit, + ) + + # Сохраняем результат в XCom для verify_publish_counts + context["ti"].xcom_push(key=f"{event_type}_count", value=sent_count) + + return sent_count + + +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 + + # Проверяем, что отправлено хотя бы что-то + if total_sent == 0: + raise AirflowException("Не отправлено ни одного сообщения ни в один топик") + + # Логируем итоговую статистику + for topic, count in results.items(): + print(f"✓ {topic}: {count} сообщений") + print(f"\nВсего отправлено: {total_sent} сообщений") + + +# ----------------------------------------------------------------------------- +# Определение DAG +# ----------------------------------------------------------------------------- +with DAG( + dag_id="kafka_load", + default_args=default_args, + description="Загрузка данных из JSONL в Kafka топики", + schedule=None, # Только ручной запуск + start_date=datetime(2024, 1, 1), + catchup=False, + max_active_runs=1, + is_paused_upon_creation=True, + tags=["kafka", "ingest", "experiments"], + params={ + "limit": Param( + default=0, + 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: + + # ------------------------------------------------------------------------- + # TaskGroup: precheck — проверки перед загрузкой + # ------------------------------------------------------------------------- + with TaskGroup(group_id="precheck") as precheck: + check_kafka = PythonOperator( + task_id="check_kafka", + python_callable=_check_kafka, + ) + + check_input_files_task = PythonOperator( + task_id="check_input_files", + python_callable=_check_files, + ) + + validate_params = PythonOperator( + task_id="validate_load_params", + python_callable=_validate_params, + ) + + check_kafka >> check_input_files_task >> validate_params + + # ------------------------------------------------------------------------- + # TaskGroup: ingest — загрузка данных + # ------------------------------------------------------------------------- + with TaskGroup(group_id="ingest") as ingest: + prepare_topics_task = PythonOperator( + task_id="prepare_topics", + python_callable=_prepare_topics, + ) + + load_browser_events = PythonOperator( + task_id="load_browser_events", + python_callable=_load_events, + op_kwargs={"event_type": "browser"}, + ) + + load_location_events = PythonOperator( + task_id="load_location_events", + python_callable=_load_events, + op_kwargs={"event_type": "location"}, + ) + + load_device_events = PythonOperator( + task_id="load_device_events", + python_callable=_load_events, + op_kwargs={"event_type": "device"}, + ) + + load_geo_events = PythonOperator( + task_id="load_geo_events", + python_callable=_load_events, + op_kwargs={"event_type": "geo"}, + ) + + verify_publish_counts = PythonOperator( + task_id="verify_publish_counts", + python_callable=_verify_counts, + ) + + # Зависимости: подготовка -> параллельная загрузка -> проверка + prepare_topics_task >> [ + load_browser_events, + load_location_events, + load_device_events, + load_geo_events, + ] >> verify_publish_counts + + # ------------------------------------------------------------------------- + # Итоговая цепочка + # ------------------------------------------------------------------------- + precheck >> ingest diff --git a/dags/utils/__init__.py b/dags/utils/__init__.py new file mode 100644 index 0000000..29a1ee3 --- /dev/null +++ b/dags/utils/__init__.py @@ -0,0 +1 @@ +# Utils package для Airflow DAG'ов diff --git a/dags/utils/kafka_helpers.py b/dags/utils/kafka_helpers.py new file mode 100644 index 0000000..47988a5 --- /dev/null +++ b/dags/utils/kafka_helpers.py @@ -0,0 +1,342 @@ +""" +Helper-функции для работы с Kafka из Airflow DAG'ов. + +Использует kafka-python: +- KafkaAdminClient — для управления топиками +- KafkaProducer — для публикации сообщений +""" + +from __future__ import annotations + +import logging +from pathlib import Path +from typing import TYPE_CHECKING + +if TYPE_CHECKING: + from typing import Optional + +# ----------------------------------------------------------------------------- +# Конфигурация подключения к Kafka +# ----------------------------------------------------------------------------- +KAFKA_BOOTSTRAP_SERVERS = "kafka:29092" +REQUEST_TIMEOUT_MS = 30000 + +# Топики и соответствующие файлы данных +TOPIC_FILE_MAP = { + "browser_events": "browser_events.jsonl", + "location_events": "location_events.jsonl", + "device_events": "device_events.jsonl", + "geo_events": "geo_events.jsonl", +} + +logger = logging.getLogger(__name__) + + +# ----------------------------------------------------------------------------- +# Проверка доступности Kafka +# ----------------------------------------------------------------------------- +def check_kafka_ready( + bootstrap_servers: str = KAFKA_BOOTSTRAP_SERVERS, + timeout_ms: int = REQUEST_TIMEOUT_MS, +) -> None: + """ + Проверяет доступность Kafka брокера. + + Args: + bootstrap_servers: Адрес Kafka брокера + timeout_ms: Таймаут запроса в миллисекундах + + Raises: + AirflowException: Если Kafka недоступна + """ + from kafka import KafkaAdminClient + from kafka.errors import NoBrokersAvailable + from airflow.exceptions import AirflowException + + try: + admin_client = KafkaAdminClient( + bootstrap_servers=bootstrap_servers, + request_timeout_ms=timeout_ms, + ) + # Проверяем связь, запрашивая список топиков + admin_client.list_topics() + admin_client.close() + logger.info("Kafka брокер доступен: %s", bootstrap_servers) + except NoBrokersAvailable as e: + raise AirflowException(f"Kafka брокер недоступен: {bootstrap_servers}") from e + except Exception as e: + raise AirflowException(f"Ошибка подключения к Kafka: {e}") from e + + +# ----------------------------------------------------------------------------- +# Управление топиками +# ----------------------------------------------------------------------------- +def prepare_topics( + topics: Optional[list[str]] = None, + reset: bool = True, + bootstrap_servers: str = KAFKA_BOOTSTRAP_SERVERS, + timeout_ms: int = REQUEST_TIMEOUT_MS, +) -> None: + """ + Создаёт или пересоздаёт топики Kafka. + + Args: + topics: Список топиков для создания (по умолчанию все из TOPIC_FILE_MAP) + reset: Если True — удаляет топики перед созданием + bootstrap_servers: Адрес Kafka брокера + timeout_ms: Таймаут операций в миллисекундах + """ + from kafka import KafkaAdminClient + from kafka.admin import NewTopic + from kafka.errors import TopicAlreadyExistsError, UnknownTopicOrPartitionError + from airflow.exceptions import AirflowException + + if topics is None: + topics = list(TOPIC_FILE_MAP.keys()) + + admin_client = KafkaAdminClient( + bootstrap_servers=bootstrap_servers, + request_timeout_ms=timeout_ms, + ) + + try: + # Удаляем топики если reset=True + if reset: + try: + admin_client.delete_topics(topics, timeout_ms=timeout_ms) + logger.info("Удалены топики: %s", topics) + except UnknownTopicOrPartitionError: + # Топики не существуют — это нормально + logger.info("Топики для удаления не найдены (уже отсутствуют)") + except Exception as e: + logger.warning("Ошибка при удалении топиков: %s", e) + + # Создаём топики + new_topics = [ + NewTopic( + name=topic, + num_partitions=1, # Дефолтное количество партиций + replication_factor=1, + ) + for topic in topics + ] + + try: + admin_client.create_topics(new_topics, timeout_ms=timeout_ms) + logger.info("Созданы топики: %s", topics) + except TopicAlreadyExistsError: + logger.info("Топики уже существуют: %s", topics) + except Exception as e: + raise AirflowException(f"Ошибка создания топиков: {e}") from e + + finally: + admin_client.close() + + +# ----------------------------------------------------------------------------- +# Загрузка данных из JSONL +# ----------------------------------------------------------------------------- +def load_jsonl( + file_path: str | Path, + topic: str, + limit: int = 0, + bootstrap_servers: str = KAFKA_BOOTSTRAP_SERVERS, +) -> int: + """ + Читает JSONL-файл и публикует строки в Kafka топик. + + Формат: 1 строка JSON = 1 сообщение (value), без ключа. + + Args: + file_path: Путь к .jsonl файлу + topic: Имя Kafka топика + limit: Максимальное количество строк (0 = все строки) + bootstrap_servers: Адрес Kafka брокера + + Returns: + Количество отправленных сообщений + + Raises: + AirflowException: Если файл не найден или ошибка отправки + """ + from kafka import KafkaProducer + from kafka.errors import KafkaError + from airflow.exceptions import AirflowException + + file_path = Path(file_path) + if not file_path.is_file(): + raise AirflowException(f"Файл не найден: {file_path}") + + producer = KafkaProducer( + bootstrap_servers=bootstrap_servers, + # Отправляем сырые байты (строки JSON как есть) + value_serializer=lambda v: v.encode("utf-8") if isinstance(v, str) else v, + acks="all", # Ждём подтверждения от всех реплик + retries=3, + batch_size=16384, + linger_ms=10, + ) + + sent_count = 0 + error_count = 0 + + try: + with open(file_path, "r", encoding="utf-8") as f: + for line_num, line in enumerate(f, 1): + # Пропускаем пустые строки + line = line.strip() + if not line: + continue + + # Проверяем лимит + if limit > 0 and sent_count >= limit: + logger.info( + "Достигнут лимит %d строк для %s", limit, topic + ) + break + + # Отправляем сообщение + try: + future = producer.send(topic, value=line) + # Неблокирующая отправка, собираем future для проверки + sent_count += 1 + except KafkaError as e: + error_count += 1 + logger.error("Ошибка отправки строки %d в %s: %s", line_num, topic, e) + if error_count > 10: + raise AirflowException( + f"Слишком много ошибок отправки в {topic}" + ) from e + + # Ждём завершения всех отправок + producer.flush(timeout=60) + + logger.info( + "Загрузка завершена: %s -> %s, отправлено %d сообщений", + file_path.name, + topic, + sent_count, + ) + + if sent_count == 0: + raise AirflowException(f"Не отправлено ни одного сообщения в {topic}") + + return sent_count + + except Exception as e: + if isinstance(e, AirflowException): + raise + raise AirflowException(f"Ошибка загрузки {file_path.name}: {e}") from e + + finally: + producer.close(timeout=30) + + +# ----------------------------------------------------------------------------- +# Утилиты для валидации +# ----------------------------------------------------------------------------- +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: Если параметры невалидны + """ + from airflow.exceptions import AirflowException + + # Проверка limit + 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: Если какой-либо файл отсутствует + """ + 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") + + missing_files = [] + for filename in files_to_check: + file_path = data_dir / filename + if not file_path.is_file(): + missing_files.append(filename) + + if missing_files: + raise AirflowException( + f"Отсутствуют файлы данных в {data_dir}: {missing_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 e652521..db32e21 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -116,7 +116,7 @@ flowchart TB DM_T["VIEW для BI
(Superset/Grafana)"] end - RAW -->|kafka-console-producer| KAFKA -->|MV| STG_T -->|Batch SQL (Airflow)| ODS_T + RAW -->|kafka_load DAG / kafka-console-producer| KAFKA -->|MV| STG_T -->|Batch SQL (Airflow)| ODS_T ODS_T -->|argMax + JOIN| DDS_T -->|VIEW| DM_T ODS_T -.->|ошибки| DQ ``` @@ -598,21 +598,30 @@ INSERT INTO dm.daily_traffic SELECT * FROM dm.v_daily_traffic; Инфраструктура Airflow развёрнута и готова к использованию: ```python -# dags/ddl_init_dag.py и dags/etl_pipeline_dag.py -# +# dags/ddl_init_dag.py — создание баз/таблиц (ручной запуск при bootstrap) +# dags/kafka_load_dag.py — загрузка JSONL в Kafka (фаза 2, через kafka-python) +# dags/etl_pipeline_dag.py — основной ETL (STG→ODS→DDS→DM) + # Учебный формат: # - DDL и трансформации выполняются явными SQL-task через ClickHouseOperator; # - SQL-файлы вызываются по фиксированным путям; -# - загрузка данных в Kafka (Этап 1) выполняется через `make data`. +# - загрузка данных в Kafka (фаза 2) выполняется через DAG `kafka_load`. # -# Основной demo-сценарий: -# ddl_init -> make data -> etl_pipeline +# Основной demo-сценарий (фаза 2): +# ddl_init -> kafka_load -> etl_pipeline ``` +**DAG `kafka_load`** (фаза 2): +- Загрузка данных из `data/*.jsonl` в Kafka через `kafka-python` +- TaskGroup `precheck`: проверка Kafka, файлов, параметров +- TaskGroup `ingest`: создание топиков → параллельная загрузка 4 потоков → проверка +- Параметры: `limit` (0 = все), `reset_topics`, `load_*` (выбор потоков) + **Подключение к ClickHouse:** - Connection: `clickhouse_default` - URL: `clickhouse://default:123456@clickhouse:9000/default` (native TCP для Airflow plugin) - Provider/интеграция: `airflow-clickhouse-plugin` (в `airflow/requirements.txt`), задачи выполняются через `ClickHouseOperator`. +- Дополнительно: `kafka-python==2.0.6` для работы с Kafka из DAG. - Примечание: Superset подключается к ClickHouse по HTTP (обычно `clickhouse+connect://...:8123/...`). --- diff --git a/plans/airflow_dags_plan.md b/plans/airflow_dags_plan.md index b1ecd10..f32f2f1 100644 --- a/plans/airflow_dags_plan.md +++ b/plans/airflow_dags_plan.md @@ -61,8 +61,7 @@ check_clickhouse >> ddl_00_databases >> ddl_10_stg >> ddl_20_ods >> ddl_30_dds > ## Дизайн DAG `kafka_load` (отдельный независимый контур) ### Params (через Trigger DAG with config) -- `limit`: int, default `50`. -- `full_load`: bool, default `false`. +- `limit`: int, default `0` (0 = все строки). - `reset_topics`: bool, default `true`. - `load_browser`: bool, default `true`. - `load_location`: bool, default `true`. diff --git a/plans/runbook.md b/plans/runbook.md index dd3b5d2..eb0be1c 100644 --- a/plans/runbook.md +++ b/plans/runbook.md @@ -15,8 +15,27 @@ make up ``` -2) (Опционально) Залить данные в Kafka: +2) Залить данные в Kafka (два варианта): +**Вариант А: Через Airflow DAG `kafka_load` (рекомендуется, фаза 2)** +```bash +# Через CLI — полная загрузка по умолчанию +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}' + +# Или через UI: Airflow → DAGs → kafka_load → Trigger DAG with config +``` + +Параметры `kafka_load`: +- `limit` — количество строк (default: 0 — все строки) +- `reset_topics` — пересоздать топики (default: true) +- `load_browser/load_location/load_device/load_geo` — выбор потоков (default: true) + +**Вариант Б: Через shell-скрипт `make data` (устаревший)** ```bash make data ``` @@ -38,6 +57,8 @@ make ddl План реализации механики заливки (дизайн/решения): `plans/kafka_ingest_plan.md`. +План Airflow DAG'ов: `plans/airflow_dags_plan.md`. + ## Загрузка данных в Kafka (`make data`) ### Топики