From 772804bda0c51520bc49fbccb68f7ea17d7d45be Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Tue, 14 Oct 2025 23:03:00 +0300 Subject: [PATCH] =?UTF-8?q?=D0=A3=D0=BB=D1=83=D1=87=D1=88=D0=B5=D0=BD?= =?UTF-8?q?=D0=B8=D1=8F=20=D1=87=D0=B8=D1=82=D0=B0=D0=B5=D0=BC=D0=BE=D1=81?= =?UTF-8?q?=D1=82=D0=B8=20=D0=BA=D0=BE=D0=B4=D0=B0=20dsg?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../airflow/dags/csv_to_greenplum.py | 44 ++++-- .../airflow/dags/data_quality_greenplum.py | 46 +++++- .../airflow/dags/helpers/greenplum.py | 132 +++++++++++++++--- 3 files changed, 185 insertions(+), 37 deletions(-) diff --git a/airflow-greenplum/airflow/dags/csv_to_greenplum.py b/airflow-greenplum/airflow/dags/csv_to_greenplum.py index fe51ac0..cb9986d 100644 --- a/airflow-greenplum/airflow/dags/csv_to_greenplum.py +++ b/airflow-greenplum/airflow/dags/csv_to_greenplum.py @@ -39,21 +39,35 @@ def _generate_csv(rows: int, csv_dir: Path) -> str: timestamp = datetime.utcnow().strftime("%Y%m%d_%H%M%S") csv_path = csv_dir / f"orders_{timestamp}.csv" - base_order_id = int(datetime.utcnow().timestamp()) - order_ids: List[int] = list(range(base_order_id * 1_000, base_order_id * 1_000 + rows)) - now = datetime.utcnow() - order_ts = [(now - timedelta(seconds=i)).isoformat() for i in range(rows)] - customer_ids = [random.randint(1, 1_000) for _ in range(rows)] - amounts = [round(random.uniform(10, 500), 2) for _ in range(rows)] - - df = pd.DataFrame( - { - "order_id": order_ids, - "order_ts": order_ts, - "customer_id": customer_ids, - "amount": amounts, - } - ) + # Генерируем данные в 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" + ) + }) + + # Сохраняем CSV без индекса df.to_csv(csv_path, index=False) logging.info("CSV сохранён: %s (строк: %s)", csv_path, len(df)) return str(csv_path) diff --git a/airflow-greenplum/airflow/dags/data_quality_greenplum.py b/airflow-greenplum/airflow/dags/data_quality_greenplum.py index 409f3fe..697c0fd 100644 --- a/airflow-greenplum/airflow/dags/data_quality_greenplum.py +++ b/airflow-greenplum/airflow/dags/data_quality_greenplum.py @@ -1,5 +1,6 @@ from __future__ import annotations +import logging from datetime import datetime, timedelta from airflow import DAG @@ -15,9 +16,35 @@ from helpers.greenplum import ( 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) + + +def _log_dq_summary(): + """ + Логирует итоговую сводку по качеству данных. + Эта задача выполняется после всех проверок и показывает общий результат. + """ + logging.info("🎉 Все проверки качества данных пройдены успешно!") + logging.info("📊 Качество данных в таблице orders соответствует требованиям.") default_args = {"owner": "airflow", "retries": 1, "retry_delay": timedelta(seconds=30)} @@ -29,26 +56,41 @@ with DAG( catchup=False, default_args=default_args, tags=["demo", "greenplum", "quality"], + description="Автоматизированные проверки качества данных в Greenplum", ) as dag: + # Задача 1: Проверка существования таблицы check_exists = PythonOperator( task_id="check_orders_table_exists", 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", + python_callable=_log_dq_summary, + ) - check_exists >> check_schema >> check_has_rows >> check_no_duplicates + # Определяем последовательность выполнения задач + check_exists >> check_schema >> check_has_rows >> check_no_duplicates >> dq_summary diff --git a/airflow-greenplum/airflow/dags/helpers/greenplum.py b/airflow-greenplum/airflow/dags/helpers/greenplum.py index 0709bd8..90e1185 100644 --- a/airflow-greenplum/airflow/dags/helpers/greenplum.py +++ b/airflow-greenplum/airflow/dags/helpers/greenplum.py @@ -1,5 +1,6 @@ from __future__ import annotations +import logging import os from typing import List, Sequence, Tuple @@ -10,6 +11,7 @@ import psycopg2 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") +# Ожидаемая схема таблицы orders для проверки качества данных EXPECTED_ORDERS_SCHEMA: List[Tuple[str, str]] = [ ("order_id", "bigint"), ("order_ts", "timestamp without time zone"), @@ -19,28 +21,52 @@ EXPECTED_ORDERS_SCHEMA: List[Tuple[str, str]] = [ def get_gp_conn(): - """Возвращает psycopg2 connection к Greenplum (через Airflow Connection или напрямую по ENV).""" + """ + Возвращает psycopg2 connection к Greenplum. + + Приоритет подключения: + 1. Через Airflow Connection (если настроено и доступно) + 2. Прямое подключение по переменным окружения (фоллбек) + + Returns: + psycopg2 connection object + """ if GP_USE_AIRFLOW_CONN: try: from airflow.providers.postgres.hooks.postgres import PostgresHook hook = PostgresHook(postgres_conn_id=GP_CONN_ID) - return hook.get_conn() - except Exception: + conn = hook.get_conn() + logging.info("✅ Подключение через Airflow Connection успешно") + return conn + except Exception as e: + logging.warning("⚠️ Не удалось подключиться через Airflow Connection: %s", e) + logging.info("🔄 Переключаемся на прямое подключение по ENV переменным") # Фоллбек на прямое подключение по переменным окружения. - pass - return psycopg2.connect( - dbname=os.getenv("GP_DB", "gpadmin"), - user=os.getenv("GP_USER", "gpadmin"), - password=os.getenv("GP_PASSWORD", ""), - host=os.getenv("GP_HOST", "greenplum"), - port=int(os.getenv("GP_PORT", "5432")), - ) + # Прямое подключение по переменным окружения + conn_params = { + "dbname": os.getenv("GP_DB", "gpadmin"), + "user": os.getenv("GP_USER", "gpadmin"), + "password": os.getenv("GP_PASSWORD", ""), + "host": os.getenv("GP_HOST", "greenplum"), + "port": int(os.getenv("GP_PORT", "5432")), + } + 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.""" + """ + Проверяет наличие таблицы orders в схеме public. + + Args: + conn: Подключение к Greenplum + + Raises: + ValueError: Если таблица не найдена + """ + logging.info("🔍 Проверяем существование таблицы public.orders...") with conn.cursor() as cur: cur.execute( """ @@ -50,10 +76,20 @@ 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: + Список кортежей (имя_колонки, тип_данных) + """ with conn.cursor() as cur: cur.execute( """ @@ -67,25 +103,69 @@ def fetch_orders_schema(conn) -> Sequence[Tuple[str, str]]: def assert_orders_schema(conn) -> None: - """Проверяет, что схема таблицы orders соответствует ожидаемой.""" + """ + Проверяет, что схема таблицы orders соответствует ожидаемой. + + Args: + conn: Подключение к Greenplum + + Raises: + ValueError: Если схема не соответствует ожидаемой + """ + logging.info("📋 Проверяем схему таблицы orders...") 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: + Количество строк в таблице + """ with conn.cursor() as cur: cur.execute("SELECT COUNT(*) FROM public.orders") return cur.fetchone()[0] def assert_orders_have_rows(conn) -> None: - """Проверяет, что таблица orders не пустая.""" - if fetch_orders_count(conn) <= 0: - raise ValueError("Таблица public.orders пустая — запусти DAG csv_to_greenplum перед проверкой.") + """ + Проверяет, что таблица 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 перед проверкой.") + logging.info("✅ Таблица orders содержит данные (%s строк)", row_count) def fetch_orders_duplicates(conn) -> int: + """ + Подсчитывает количество дубликатов по order_id. + + Args: + conn: Подключение к Greenplum + + Returns: + Количество дублирующихся order_id + """ with conn.cursor() as cur: cur.execute( """ @@ -101,7 +181,19 @@ def fetch_orders_duplicates(conn) -> int: def assert_orders_no_duplicates(conn) -> None: - """Проверяет, что в таблице нет дублей по order_id.""" + """ + Проверяет, что в таблице нет дублей по 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("✅ Дубликаты не обнаружены")