diff --git a/airflow/dags/csv_to_greenplum.py b/airflow/dags/csv_to_greenplum.py index cb9986d..9289403 100644 --- a/airflow/dags/csv_to_greenplum.py +++ b/airflow/dags/csv_to_greenplum.py @@ -8,10 +8,11 @@ from pathlib import Path from typing import List import pandas as pd -from airflow import DAG from airflow.operators.python import PythonOperator from helpers.greenplum import get_gp_conn +from airflow import DAG + CSV_DIR = Path(os.getenv("CSV_DIR", "/opt/airflow/data")) CSV_ROWS = int(os.getenv("CSV_ROWS", "1000")) @@ -41,32 +42,30 @@ def _generate_csv(rows: int, csv_dir: Path) -> str: # Генерируем данные в pandas-стиле base_order_id = int(datetime.utcnow().timestamp() * 1_000) - + # Создаём DataFrame с использованием pandas методов - df = pd.DataFrame({ - # Уникальные order_id начиная с базового значения - "order_id": pd.Series(range(base_order_id, base_order_id + rows), dtype="int64"), - - # Временные метки с интервалом в 1 секунду в обратном порядке - "order_ts": pd.date_range( - end=datetime.utcnow(), - periods=rows, - freq="1S" - ).sort_values(ascending=False), - - # Случайные customer_id от 1 до 1000 - "customer_id": pd.Series( - random.choices(range(1, 1001), k=rows), - dtype="int64" - ), - - # Случайные суммы от 10 до 500 с округлением до 2 знаков - "amount": pd.Series( - [round(random.uniform(10, 500), 2) for _ in range(rows)], - dtype="float64" - ) - }) - + df = pd.DataFrame( + { + # Уникальные order_id начиная с базового значения + "order_id": pd.Series( + range(base_order_id, base_order_id + rows), dtype="int64" + ), + # Временные метки с интервалом в 1 секунду в обратном порядке + "order_ts": pd.date_range( + end=datetime.utcnow(), periods=rows, freq="1S" + ).sort_values(ascending=False), + # Случайные customer_id от 1 до 1000 + "customer_id": pd.Series( + random.choices(range(1, 1001), k=rows), dtype="int64" + ), + # Случайные суммы от 10 до 500 с округлением до 2 знаков + "amount": pd.Series( + [round(random.uniform(10, 500), 2) for _ in range(rows)], + dtype="float64", + ), + } + ) + # Сохраняем CSV без индекса df.to_csv(csv_path, index=False) logging.info("CSV сохранён: %s (строк: %s)", csv_path, len(df)) @@ -77,7 +76,9 @@ def _preview_csv(csv_path: str, sample_rows: int = 5) -> None: """Отображает предпросмотр CSV через pandas (head и describe).""" df = pd.read_csv(csv_path) df["order_ts"] = pd.to_datetime(df["order_ts"], errors="coerce") - logging.info("Первые %s строк:\n%s", sample_rows, df.head(sample_rows).to_string(index=False)) + logging.info( + "Первые %s строк:\n%s", sample_rows, df.head(sample_rows).to_string(index=False) + ) numeric_summary = df.describe(include="number") logging.info("Числовая статистика:\n%s", numeric_summary.to_string()) if df["order_ts"].notna().any(): @@ -94,8 +95,14 @@ def _load_csv(csv_path: str) -> None: if not csv_file.exists(): raise FileNotFoundError(f"CSV не найден: {csv_file}") - with get_gp_conn() as conn, conn.cursor() as cur, csv_file.open("r", encoding="utf-8") as f: - cur.execute("CREATE TEMP TABLE tmp_orders (LIKE public.orders INCLUDING DEFAULTS) ON COMMIT DROP;") + with ( + get_gp_conn() as conn, + conn.cursor() as cur, + csv_file.open("r", encoding="utf-8") as f, + ): + cur.execute( + "CREATE TEMP TABLE tmp_orders (LIKE public.orders INCLUDING DEFAULTS) ON COMMIT DROP;" + ) cur.copy_expert( "COPY tmp_orders (order_id, order_ts, customer_id, amount) FROM STDIN WITH CSV HEADER", f, diff --git a/airflow/dags/data_quality_greenplum.py b/airflow/dags/data_quality_greenplum.py index 697c0fd..8fab0e7 100644 --- a/airflow/dags/data_quality_greenplum.py +++ b/airflow/dags/data_quality_greenplum.py @@ -3,38 +3,35 @@ from __future__ import annotations import logging from datetime import datetime, timedelta -from airflow import DAG from airflow.operators.python import PythonOperator +from helpers.greenplum import (assert_orders_have_rows, + assert_orders_no_duplicates, + assert_orders_schema, + assert_orders_table_exists, get_gp_conn) -from helpers.greenplum import ( - assert_orders_have_rows, - assert_orders_no_duplicates, - assert_orders_schema, - assert_orders_table_exists, - get_gp_conn, -) +from airflow import DAG def _run_check(check_callable): """ Оборачивает проверку качества данных в контекст подключения к Greenplum. - + Этот DAG предназначен для автоматической проверки качества данных в таблице orders: 1. Проверяет существование таблицы 2. Проверяет соответствие схемы 3. Проверяет наличие данных 4. Проверяет отсутствие дубликатов - + Args: check_callable: Функция проверки, принимающая подключение к БД """ # Получаем имя функции для логов check_name = check_callable.__name__.replace("assert_", "") logging.info("🚀 Запуск проверки: %s", check_name) - + with get_gp_conn() as conn: check_callable(conn) - + logging.info("✅ Проверка пройдена: %s", check_name) @@ -64,28 +61,28 @@ with DAG( python_callable=_run_check, op_args=[assert_orders_table_exists], ) - + # Задача 2: Проверка соответствия схемы таблицы check_schema = PythonOperator( task_id="check_orders_schema", python_callable=_run_check, op_args=[assert_orders_schema], ) - + # Задача 3: Проверка наличия данных check_has_rows = PythonOperator( task_id="check_orders_has_rows", python_callable=_run_check, op_args=[assert_orders_have_rows], ) - + # Задача 4: Проверка отсутствия дубликатов check_no_duplicates = PythonOperator( task_id="check_order_duplicates", python_callable=_run_check, op_args=[assert_orders_no_duplicates], ) - + # Задача 5: Итоговая сводка dq_summary = PythonOperator( task_id="data_quality_summary", diff --git a/airflow/dags/helpers/greenplum.py b/airflow/dags/helpers/greenplum.py index 90e1185..841e4fa 100644 --- a/airflow/dags/helpers/greenplum.py +++ b/airflow/dags/helpers/greenplum.py @@ -9,7 +9,11 @@ import psycopg2 # Настройки для подключения к Greenplum. По умолчанию используем Airflow Connection, # но при проблемах можно переключиться на ENV-подключение, установив GP_USE_AIRFLOW_CONN=false. GP_CONN_ID = os.getenv("GP_CONN_ID", "greenplum_conn") -GP_USE_AIRFLOW_CONN = os.getenv("GP_USE_AIRFLOW_CONN", "true").lower() in ("1", "true", "yes") +GP_USE_AIRFLOW_CONN = os.getenv("GP_USE_AIRFLOW_CONN", "true").lower() in ( + "1", + "true", + "yes", +) # Ожидаемая схема таблицы orders для проверки качества данных EXPECTED_ORDERS_SCHEMA: List[Tuple[str, str]] = [ @@ -23,11 +27,11 @@ EXPECTED_ORDERS_SCHEMA: List[Tuple[str, str]] = [ def get_gp_conn(): """ Возвращает psycopg2 connection к Greenplum. - + Приоритет подключения: 1. Через Airflow Connection (если настроено и доступно) 2. Прямое подключение по переменным окружения (фоллбек) - + Returns: psycopg2 connection object """ @@ -52,17 +56,19 @@ def get_gp_conn(): "host": os.getenv("GP_HOST", "greenplum"), "port": int(os.getenv("GP_PORT", "5432")), } - logging.info("🔗 Подключение к Greenplum: %s:%s", conn_params["host"], conn_params["port"]) + logging.info( + "🔗 Подключение к Greenplum: %s:%s", conn_params["host"], conn_params["port"] + ) return psycopg2.connect(**conn_params) def assert_orders_table_exists(conn) -> None: """ Проверяет наличие таблицы orders в схеме public. - + Args: conn: Подключение к Greenplum - + Raises: ValueError: Если таблица не найдена """ @@ -76,17 +82,19 @@ def assert_orders_table_exists(conn) -> None: """ ) if cur.fetchone() is None: - raise ValueError("❌ Таблица public.orders не найдена; запусти DAG csv_to_greenplum.") + raise ValueError( + "❌ Таблица public.orders не найдена; запусти DAG csv_to_greenplum." + ) logging.info("✅ Таблица public.orders существует") def fetch_orders_schema(conn) -> Sequence[Tuple[str, str]]: """ Получает схему таблицы orders из information_schema. - + Args: conn: Подключение к Greenplum - + Returns: Список кортежей (имя_колонки, тип_данных) """ @@ -105,10 +113,10 @@ def fetch_orders_schema(conn) -> Sequence[Tuple[str, str]]: def assert_orders_schema(conn) -> None: """ Проверяет, что схема таблицы orders соответствует ожидаемой. - + Args: conn: Подключение к Greenplum - + Raises: ValueError: Если схема не соответствует ожидаемой """ @@ -116,19 +124,21 @@ def assert_orders_schema(conn) -> None: schema = fetch_orders_schema(conn) logging.info("📊 Фактическая схема: %s", list(schema)) logging.info("📊 Ожидаемая схема: %s", EXPECTED_ORDERS_SCHEMA) - + if list(schema) != EXPECTED_ORDERS_SCHEMA: - raise ValueError(f"❌ Неожиданная схема orders: {schema}. Ожидали {EXPECTED_ORDERS_SCHEMA}.") + raise ValueError( + f"❌ Неожиданная схема orders: {schema}. Ожидали {EXPECTED_ORDERS_SCHEMA}." + ) logging.info("✅ Схема таблицы orders соответствует ожиданиям") def fetch_orders_count(conn) -> int: """ Получает количество строк в таблице orders. - + Args: conn: Подключение к Greenplum - + Returns: Количество строк в таблице """ @@ -140,29 +150,31 @@ def fetch_orders_count(conn) -> int: def assert_orders_have_rows(conn) -> None: """ Проверяет, что таблица orders не пустая. - + Args: conn: Подключение к Greenplum - + Raises: ValueError: Если таблица пустая """ logging.info("📊 Проверяем наличие данных в таблице orders...") row_count = fetch_orders_count(conn) logging.info("📈 Количество строк в orders: %s", row_count) - + if row_count <= 0: - raise ValueError("❌ Таблица public.orders пустая — запусти DAG csv_to_greenplum перед проверкой.") + raise ValueError( + "❌ Таблица public.orders пустая — запусти DAG csv_to_greenplum перед проверкой." + ) logging.info("✅ Таблица orders содержит данные (%s строк)", row_count) def fetch_orders_duplicates(conn) -> int: """ Подсчитывает количество дубликатов по order_id. - + Args: conn: Подключение к Greenplum - + Returns: Количество дублирующихся order_id """ @@ -183,17 +195,19 @@ def fetch_orders_duplicates(conn) -> int: def assert_orders_no_duplicates(conn) -> None: """ Проверяет, что в таблице нет дублей по order_id. - + Args: conn: Подключение к Greenplum - + Raises: ValueError: Если обнаружены дубликаты """ logging.info("🔍 Проверяем отсутствие дубликатов по order_id...") duplicates = fetch_orders_duplicates(conn) logging.info("📊 Найдено дубликатов: %s", duplicates) - + if duplicates: - raise ValueError(f"❌ Обнаружены дубли по order_id ({duplicates} шт.) — проверь загрузку данных.") + raise ValueError( + f"❌ Обнаружены дубли по order_id ({duplicates} шт.) — проверь загрузку данных." + ) logging.info("✅ Дубликаты не обнаружены") diff --git a/tests/conftest.py b/tests/conftest.py index acbe54b..1d5d066 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -6,7 +6,6 @@ from pathlib import Path from types import ModuleType from typing import Type - PROJECT_ROOT = Path(__file__).resolve().parents[1] if str(PROJECT_ROOT) not in sys.path: sys.path.append(str(PROJECT_ROOT)) diff --git a/tests/test_dags_smoke.py b/tests/test_dags_smoke.py index a9955ff..e1ec388 100644 --- a/tests/test_dags_smoke.py +++ b/tests/test_dags_smoke.py @@ -14,7 +14,9 @@ def _airflow_available() -> bool: return hasattr(af, "DAG") -pytestmark = pytest.mark.skipif(not _airflow_available(), reason="Airflow is not installed for DAG smoke tests") +pytestmark = pytest.mark.skipif( + not _airflow_available(), reason="Airflow is not installed for DAG smoke tests" +) def _load_dag(module_name: str):