refactor(airflow): move DAGs to airflow/dags and update paths

- Why:
  - keep Airflow artifacts under a single airflow/ directory
  - align repository layout with intended project structure
- What:
  - move dags/ to airflow/dags/ and update compose mounts
  - make SQL root resolution work in container and local runs
  - update DAG path references in README, AGENTS, ARCHITECTURE, and plans
  - remove tracked Python cache artifacts from old DAG location
- Check:
  - airflow dags list
  - airflow dags list-import-errors
  - e2e success: ddl_init, kafka_load(limit=50), etl_pipeline
This commit is contained in:
2026-02-08 19:34:10 +03:00
parent 44a75691e8
commit 869c189fe8
14 changed files with 41 additions and 17 deletions
View File
View File
+200
View File
@@ -0,0 +1,200 @@
"""
DAG инициализации DDL в ClickHouse.
Учебный формат:
- каждая операция DDL выполняется отдельной SQL-task;
- SQL-файлы вызываются явно по фиксированным путям;
- режим verify_only позволяет прогонять только проверки схемы.
"""
from __future__ import annotations
import re
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.empty import EmptyOperator
from airflow.operators.python import BranchPythonOperator, PythonOperator
from airflow.utils.trigger_rule import TriggerRule
from airflow_clickhouse_plugin.operators.clickhouse import ClickHouseOperator
# -----------------------------------------------------------------------------
# Базовые настройки DAG
# -----------------------------------------------------------------------------
default_args = {
"owner": "airflow",
"depends_on_past": False,
"email_on_failure": False,
"email_on_retry": False,
"retries": 1,
"retry_delay": timedelta(minutes=2),
}
# -----------------------------------------------------------------------------
# SQL-файлы проекта
# -----------------------------------------------------------------------------
def resolve_sql_root() -> Path:
"""Определяет корень SQL для контейнера и локального запуска."""
candidates = (
Path(__file__).resolve().parents[1] / "sql", # /opt/airflow/sql в контейнере
Path(__file__).resolve().parents[2] / "sql", # <repo>/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, ...]:
"""Читает SQL-файл и делит его на отдельные команды по ';'."""
file_path = SQL_ROOT / relative_path
if not file_path.is_file():
raise AirflowException(f"SQL-файл не найден: {file_path}")
sql_text = file_path.read_text(encoding="utf-8")
statements: list[str] = []
for segment in sql_text.split(";"):
# Убираем блочные и строковые комментарии, чтобы не отправлять "пустые" запросы.
no_block_comments = re.sub(r"/\*.*?\*/", "", segment, flags=re.S)
lines = [line for line in no_block_comments.splitlines() if not line.strip().startswith("--")]
cleaned = "\n".join(lines).strip()
if cleaned:
statements.append(cleaned)
if not statements:
raise AirflowException(f"SQL-файл пустой: {file_path}")
return tuple(statements)
# -----------------------------------------------------------------------------
# SQL-проверки
# -----------------------------------------------------------------------------
SQL_CHECK_CLICKHOUSE = "SELECT 1 AS ok"
SQL_VERIFY_SCHEMA = """
SELECT
(SELECT count() FROM system.tables WHERE database = 'stg' AND name = 'browser_raw') AS stg_browser_raw,
(SELECT count() FROM system.tables WHERE database = 'ods' AND name = 'browser_event') AS ods_browser_event,
(SELECT count() FROM system.tables WHERE database = 'dds' AND name = 'click') AS dds_click,
(SELECT count() FROM system.tables WHERE database = 'dds' AND name = 'event') AS dds_event,
(SELECT count() FROM system.tables WHERE database = 'dm' AND name = 'v_events_enriched') AS dm_v_events_enriched
"""
# -----------------------------------------------------------------------------
# Управляющие функции
# -----------------------------------------------------------------------------
def choose_ddl_mode(**context) -> str:
"""Выбирает ветку выполнения: full DDL или только verify."""
dag_run = context.get("dag_run")
conf = dag_run.conf if dag_run else {}
verify_only = bool(conf.get("verify_only", context["params"]["verify_only"]))
return "skip_ddl" if verify_only else "ddl_00_databases"
def assert_schema_ready(**context) -> None:
"""Проверяет результат финальной SQL-проверки схемы."""
ti = context["ti"]
result = ti.xcom_pull(task_ids="verify_schema_sql")
if not result or not result[0] or len(result[0]) != 5:
raise AirflowException(f"Некорректный результат проверки схемы: {result}")
if any(value == 0 for value in result[0]):
raise AirflowException(
"Схема применена не полностью. Проверьте таблицы/VIEW stg, ods, dds, dm."
)
with DAG(
dag_id="ddl_init",
description="Инициализация схемы ClickHouse (stg/ods/dds/dm)",
default_args=default_args,
schedule=None,
start_date=datetime(2024, 1, 1),
catchup=False,
max_active_runs=1,
is_paused_upon_creation=True,
tags=["ddl", "bootstrap", "clickhouse"],
params={
"verify_only": Param(False, type="boolean"),
},
) as dag:
check_clickhouse = ClickHouseOperator(
task_id="check_clickhouse",
sql=SQL_CHECK_CLICKHOUSE,
clickhouse_conn_id="clickhouse_default",
database="default",
)
choose_mode = BranchPythonOperator(
task_id="choose_mode",
python_callable=choose_ddl_mode,
)
ddl_00_databases = ClickHouseOperator(
task_id="ddl_00_databases",
sql=load_sql_statements("ddl/00_databases.sql"),
clickhouse_conn_id="clickhouse_default",
database="default",
)
ddl_10_stg = ClickHouseOperator(
task_id="ddl_10_stg",
sql=load_sql_statements("ddl/stg/10_stg.sql"),
clickhouse_conn_id="clickhouse_default",
database="default",
)
ddl_20_ods = ClickHouseOperator(
task_id="ddl_20_ods",
sql=load_sql_statements("ddl/ods/20_ods.sql"),
clickhouse_conn_id="clickhouse_default",
database="default",
)
ddl_30_dds = ClickHouseOperator(
task_id="ddl_30_dds",
sql=load_sql_statements("ddl/dds/30_dds.sql"),
clickhouse_conn_id="clickhouse_default",
database="default",
)
ddl_40_dm = ClickHouseOperator(
task_id="ddl_40_dm",
sql=load_sql_statements("ddl/dm/40_dm.sql"),
clickhouse_conn_id="clickhouse_default",
database="default",
)
skip_ddl = EmptyOperator(task_id="skip_ddl")
ddl_complete = EmptyOperator(
task_id="ddl_complete",
trigger_rule=TriggerRule.NONE_FAILED_MIN_ONE_SUCCESS,
)
verify_schema_sql = ClickHouseOperator(
task_id="verify_schema_sql",
sql=SQL_VERIFY_SCHEMA,
clickhouse_conn_id="clickhouse_default",
database="default",
)
verify_schema = PythonOperator(
task_id="verify_schema",
python_callable=assert_schema_ready,
)
check_clickhouse >> choose_mode
choose_mode >> skip_ddl >> ddl_complete
choose_mode >> ddl_00_databases >> ddl_10_stg >> ddl_20_ods >> ddl_30_dds >> ddl_40_dm >> ddl_complete
ddl_complete >> verify_schema_sql >> verify_schema
+369
View File
@@ -0,0 +1,369 @@
"""
DAG ETL-процесса STG -> ODS -> DDS -> DM для учебного проекта.
Принципы реализации:
- SQL выполняется явными task на ClickHouseOperator;
- SQL-файлы вызываются по фиксированным путям;
- Python используется только для управляющей логики (branch/wait/assert).
"""
from __future__ import annotations
import re
import time
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.empty import EmptyOperator
from airflow.operators.python import BranchPythonOperator, PythonOperator
from airflow.utils.task_group import TaskGroup
from airflow.utils.trigger_rule import TriggerRule
from airflow_clickhouse_plugin.hooks.clickhouse import ClickHouseHook
from airflow_clickhouse_plugin.operators.clickhouse import ClickHouseOperator
# -----------------------------------------------------------------------------
# Базовые настройки DAG
# -----------------------------------------------------------------------------
default_args = {
"owner": "airflow",
"depends_on_past": False,
"email_on_failure": False,
"email_on_retry": False,
"retries": 1,
"retry_delay": timedelta(minutes=2),
}
# -----------------------------------------------------------------------------
# SQL-файлы проекта
# -----------------------------------------------------------------------------
def resolve_sql_root() -> Path:
"""Определяет корень SQL для контейнера и локального запуска."""
candidates = (
Path(__file__).resolve().parents[1] / "sql", # /opt/airflow/sql в контейнере
Path(__file__).resolve().parents[2] / "sql", # <repo>/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, ...]:
"""Читает SQL-файл и делит его на отдельные команды по ';'."""
file_path = SQL_ROOT / relative_path
if not file_path.is_file():
raise AirflowException(f"SQL-файл не найден: {file_path}")
sql_text = file_path.read_text(encoding="utf-8")
statements: list[str] = []
for segment in sql_text.split(";"):
# Убираем блочные и строковые комментарии, чтобы не отправлять "пустые" запросы.
no_block_comments = re.sub(r"/\*.*?\*/", "", segment, flags=re.S)
lines = [line for line in no_block_comments.splitlines() if not line.strip().startswith("--")]
cleaned = "\n".join(lines).strip()
if cleaned:
statements.append(cleaned)
if not statements:
raise AirflowException(f"SQL-файл пустой: {file_path}")
return tuple(statements)
# -----------------------------------------------------------------------------
# SQL для проверок и технических шагов
# -----------------------------------------------------------------------------
SQL_CHECK_CLICKHOUSE = "SELECT 1 AS ok"
SQL_CHECK_SCHEMA_READY = """
SELECT
(SELECT count() FROM system.tables WHERE database = 'stg' AND name = 'browser_raw') AS stg_browser_raw,
(SELECT count() FROM system.tables WHERE database = 'ods' AND name = 'browser_event') AS ods_browser_event,
(SELECT count() FROM system.tables WHERE database = 'dds' AND name = 'event') AS dds_event,
(SELECT count() FROM system.tables WHERE database = 'dm' AND name = 'v_events_enriched') AS dm_v_events_enriched
"""
SQL_CHECK_ODS_QUALITY = """
SELECT
table_name,
total_rows,
rows_with_errors,
round(if(total_rows = 0, 0, rows_with_errors / total_rows * 100), 2) AS error_pct
FROM
(
SELECT
'browser_event' AS table_name,
toFloat64(count()) AS total_rows,
toFloat64(countIf(length(parse_errors) > 0)) AS rows_with_errors
FROM ods.browser_event
UNION ALL
SELECT
'location_event',
toFloat64(count()),
toFloat64(countIf(length(parse_errors) > 0))
FROM ods.location_event
UNION ALL
SELECT
'device_by_click',
toFloat64(count()),
toFloat64(countIf(length(parse_errors) > 0))
FROM ods.device_by_click
UNION ALL
SELECT
'geo_by_click',
toFloat64(count()),
toFloat64(countIf(length(parse_errors) > 0))
FROM ods.geo_by_click
UNION ALL
SELECT
'browser_event_errors',
toFloat64(count()),
toFloat64(count())
FROM ods.browser_event_errors
UNION ALL
SELECT
'location_event_errors',
toFloat64(count()),
toFloat64(count())
FROM ods.location_event_errors
UNION ALL
SELECT
'device_by_click_errors',
toFloat64(count()),
toFloat64(count())
FROM ods.device_by_click_errors
UNION ALL
SELECT
'geo_by_click_errors',
toFloat64(count()),
toFloat64(count())
FROM ods.geo_by_click_errors
)
ORDER BY table_name
"""
SQL_TRUNCATE_DDS_CLICK = "TRUNCATE TABLE dds.click"
SQL_TRUNCATE_DDS_EVENT = "TRUNCATE TABLE dds.event"
SQL_CHECK_DDS_INTEGRITY = """
SELECT
countIf(click_id IS NOT NULL AND click_id NOT IN (SELECT click_id FROM dds.click)) AS orphan_events
FROM dds.event
"""
SQL_VALIDATE_DM_SUMMARY = "SELECT count() AS dq_rows FROM dm.dq_summary"
# -----------------------------------------------------------------------------
# Управляющие функции
# -----------------------------------------------------------------------------
def assert_schema_ready(**context) -> None:
"""Падает, если DDL не применён полностью."""
ti = context["ti"]
result = ti.xcom_pull(task_ids="precheck.check_schema_ready_sql")
if not result or not result[0] or len(result[0]) != 4:
raise AirflowException(f"Некорректный результат check_schema_ready_sql: {result}")
if any(value == 0 for value in result[0]):
raise AirflowException(
"Схема не готова: сначала запустите DAG ddl_init, затем повторите etl_pipeline."
)
def wait_for_stg_data(**context) -> None:
"""
Ожидает появления строк в STG до заданного таймаута.
Таймаут берётся из dag_run.conf.wait_stg_timeout_sec (или legacy wait_ods_timeout_sec)
либо из params.
"""
dag_run = context.get("dag_run")
conf = dag_run.conf if dag_run else {}
timeout_sec = int(
conf.get(
"wait_stg_timeout_sec",
conf.get(
"wait_ods_timeout_sec",
context["params"]["wait_stg_timeout_sec"],
),
)
)
poll_interval_sec = 10
hook = ClickHouseHook(clickhouse_conn_id="clickhouse_default", database="default")
started = time.monotonic()
while True:
rows = hook.execute(
"""
SELECT
(SELECT count() FROM stg.browser_raw)
+ (SELECT count() FROM stg.location_raw)
+ (SELECT count() FROM stg.device_raw)
+ (SELECT count() FROM stg.geo_raw) AS stg_rows_total
"""
)
count_rows = int(rows[0][0]) if rows else 0
if count_rows > 0:
return
elapsed = int(time.monotonic() - started)
if elapsed >= timeout_sec:
raise AirflowException(
f"Таймаут ожидания STG истёк ({timeout_sec} сек). "
"Таблицы stg.*_raw всё ещё пусты."
)
time.sleep(poll_interval_sec)
def choose_full_refresh(**context) -> str:
"""Ветвление: делать TRUNCATE DDS или пропустить."""
dag_run = context.get("dag_run")
conf = dag_run.conf if dag_run else {}
full_refresh = bool(conf.get("full_refresh", context["params"]["full_refresh"]))
return "transform.truncate_dds_click" if full_refresh else "transform.skip_truncate"
def assert_dm_summary_not_empty(**context) -> None:
"""Проверяет, что dm.dq_summary заполнена после загрузки."""
ti = context["ti"]
result = ti.xcom_pull(task_ids="transform.validate_dm_summary_sql")
if not result or not result[0] or len(result[0]) != 1:
raise AirflowException(f"Некорректный результат validate_dm_summary_sql: {result}")
dq_rows = int(result[0][0])
if dq_rows <= 0:
raise AirflowException("dm.dq_summary пуста после load_dm_summary.")
with DAG(
dag_id="etl_pipeline",
description="ETL STG -> ODS -> DDS -> DM для demo-проекта",
default_args=default_args,
schedule=None,
start_date=datetime(2024, 1, 1),
catchup=False,
max_active_runs=1,
is_paused_upon_creation=True,
tags=["etl", "clickhouse", "demo"],
params={
"full_refresh": Param(True, type="boolean"),
"wait_stg_timeout_sec": Param(600, type="integer", minimum=30),
},
) as dag:
with TaskGroup(group_id="precheck") as precheck:
check_clickhouse = ClickHouseOperator(
task_id="check_clickhouse",
sql=SQL_CHECK_CLICKHOUSE,
clickhouse_conn_id="clickhouse_default",
database="default",
)
check_schema_ready_sql = ClickHouseOperator(
task_id="check_schema_ready_sql",
sql=SQL_CHECK_SCHEMA_READY,
clickhouse_conn_id="clickhouse_default",
database="default",
)
check_schema_ready = PythonOperator(
task_id="check_schema_ready",
python_callable=assert_schema_ready,
)
check_clickhouse >> check_schema_ready_sql >> check_schema_ready
with TaskGroup(group_id="transform") as transform:
wait_for_stg_data_task = PythonOperator(
task_id="wait_for_stg_data",
python_callable=wait_for_stg_data,
)
load_ods = ClickHouseOperator(
task_id="load_ods",
sql=load_sql_statements("ods/20_stg_to_ods.sql"),
clickhouse_conn_id="clickhouse_default",
database="default",
)
check_ods_quality = ClickHouseOperator(
task_id="check_ods_quality",
sql=SQL_CHECK_ODS_QUALITY,
clickhouse_conn_id="clickhouse_default",
database="default",
)
choose_refresh_mode = BranchPythonOperator(
task_id="choose_refresh_mode",
python_callable=choose_full_refresh,
)
truncate_dds_click = ClickHouseOperator(
task_id="truncate_dds_click",
sql=SQL_TRUNCATE_DDS_CLICK,
clickhouse_conn_id="clickhouse_default",
database="default",
)
truncate_dds_event = ClickHouseOperator(
task_id="truncate_dds_event",
sql=SQL_TRUNCATE_DDS_EVENT,
clickhouse_conn_id="clickhouse_default",
database="default",
)
skip_truncate = EmptyOperator(task_id="skip_truncate")
truncate_complete = EmptyOperator(
task_id="truncate_complete",
trigger_rule=TriggerRule.NONE_FAILED_MIN_ONE_SUCCESS,
)
load_dds = ClickHouseOperator(
task_id="load_dds",
sql=load_sql_statements("dds/30_ods_to_dds.sql"),
clickhouse_conn_id="clickhouse_default",
database="default",
)
check_dds_integrity = ClickHouseOperator(
task_id="check_dds_integrity",
sql=SQL_CHECK_DDS_INTEGRITY,
clickhouse_conn_id="clickhouse_default",
database="default",
)
load_dm_summary = ClickHouseOperator(
task_id="load_dm_summary",
sql=load_sql_statements("dm/40_dds_to_dm.sql"),
clickhouse_conn_id="clickhouse_default",
database="default",
)
validate_dm_summary_sql = ClickHouseOperator(
task_id="validate_dm_summary_sql",
sql=SQL_VALIDATE_DM_SUMMARY,
clickhouse_conn_id="clickhouse_default",
database="default",
)
validate_dm_summary = PythonOperator(
task_id="validate_dm_summary",
python_callable=assert_dm_summary_not_empty,
)
wait_for_stg_data_task >> load_ods >> check_ods_quality >> choose_refresh_mode
choose_refresh_mode >> truncate_dds_click >> truncate_dds_event >> truncate_complete
choose_refresh_mode >> skip_truncate >> truncate_complete
truncate_complete >> load_dds >> check_dds_integrity >> load_dm_summary >> validate_dm_summary_sql >> validate_dm_summary
precheck >> transform
+234
View File
@@ -0,0 +1,234 @@
"""
DAG для загрузки данных в Kafka из JSONL-файлов.
Функциональность:
- Проверка доступности Kafka и наличия файлов
- Создание/сброс топиков
- Параллельная загрузка 4 потоков данных
- Проверка результатов через XCom
Параметры (через Trigger DAG with config):
limit (int): Количество строк для загрузки (по умолчанию 0 = все)
reset_topics (bool): Пересоздать топики (по умолчанию 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,
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:
"""Проверка наличия входных файлов."""
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"]))
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"]))
prepare_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:
"""Проверяет, что все загрузки отправили сообщения."""
ti = context["ti"]
# Собираем результаты из XCom
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:
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 = все строки, по умолчанию)",
),
"reset_topics": Param(
default=True,
type="boolean",
description="Пересоздать топики перед загрузкой",
),
},
) 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'ов
+281
View File
@@ -0,0 +1,281 @@
"""
Helper-функции для работы с Kafka из Airflow DAG'ов.
Использует kafka-python:
- KafkaAdminClient — для управления топиками
- KafkaProducer — для публикации сообщений
"""
from __future__ import annotations
import logging
from pathlib import Path
# -----------------------------------------------------------------------------
# Конфигурация подключения к 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: list[str] | None = 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,
) -> None:
"""
Валидирует параметры загрузки данных.
Args:
limit: Количество строк для загрузки
Raises:
AirflowException: Если параметры невалидны
"""
from airflow.exceptions import AirflowException
# Проверка limit
if not isinstance(limit, int) or limit < 0:
raise AirflowException(f"limit должен быть неотрицательным int, получено: {limit}")
def check_input_files(
data_dir: str | Path = "/opt/airflow/data",
) -> None:
"""
Проверяет наличие необходимых JSONL-файлов.
Args:
data_dir: Директория с данными
Raises:
AirflowException: Если какой-либо файл отсутствует
"""
from airflow.exceptions import AirflowException
data_dir = Path(data_dir)
files_to_check = list(TOPIC_FILE_MAP.values())
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)