fix(airflow): simplify kafka_load params and align docs

- Why:

  - For DE task we only need full ingest or limit-based sample.

  - load_* and full_load params were redundant and unclear in current flow.

- What:

  - Remove full_load and load_* params from kafka_load DAG contract.

  - Simplify kafka helpers (validate/check files) to fixed 4-stream ingest.

  - Sync AGENTS, README, runbook, architecture and airflow plan docs.

- Check:

  - python3 -m py_compile dags/kafka_load_dag.py dags/utils/kafka_helpers.py

  - Airflow smoke/full runs: ddl_init -> kafka_load -> etl_pipeline (all success).

  - Legacy path: make data && make transform (success).
This commit is contained in:
2026-02-08 18:33:54 +03:00
parent 10f5bc3510
commit 1b9991f595
7 changed files with 22 additions and 179 deletions
+1 -4
View File
@@ -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
+1 -1
View File
@@ -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
```
+11 -99
View File
@@ -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:
+2 -63
View File
@@ -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
+1 -1
View File
@@ -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`
+6 -10
View File
@@ -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}
```
-1
View File
@@ -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