feat(airflow): реализован DAG kafka_load для загрузки в Kafka (фаза 2)

- Добавлен kafka-python==2.0.6 в airflow/requirements.txt
- Создан dags/utils/kafka_helpers.py с функциями:
  - check_kafka_ready() — проверка доступности брокера
  - prepare_topics() — создание/сброс топиков через KafkaAdminClient
  - load_jsonl() — загрузка данных через KafkaProducer (limit=0 = все)
  - validate_load_params(), check_input_files() — валидация
- Создан dags/kafka_load_dag.py с TaskGroup:
  - precheck: check_kafka, check_input_files, validate_load_params
  - ingest: prepare_topics, параллельная загрузка 4 потоков, verify_publish_counts
- Параметры DAG: limit (0 = все), reset_topics, load_* (выбор потоков)
- Обновлена документация: AGENTS.md, README.md, plans/runbook.md,
  plans/airflow_dags_plan.md, docs/ARCHITECTURE.md

Тестирование:
- Подключение к Kafka:  (kafka:29092 доступен, брокер 2.6.0)
- Загрузка данных:  (1000 сообщений — полный файл browser_events)
- Python синтаксис:  (py_compile проходит)
- Структура DAG:  (все 9 задач корректно определены)
This commit is contained in:
2026-02-08 18:13:22 +03:00
parent 12f35679f0
commit 10f5bc3510
10 changed files with 778 additions and 16 deletions
+322
View File
@@ -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
+1
View File
@@ -0,0 +1 @@
# Utils package для Airflow DAG'ов
+342
View File
@@ -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