diff --git a/airflow-docker/dags/csv_to_greenplum.py b/airflow-docker/dags/csv_to_greenplum.py new file mode 100644 index 0000000..f99eaa7 --- /dev/null +++ b/airflow-docker/dags/csv_to_greenplum.py @@ -0,0 +1,167 @@ +from __future__ import annotations + +import logging +import os +import random +from datetime import UTC, datetime, timedelta +from pathlib import Path +from typing import List + +import pandas as pd +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")) + + +def _create_table() -> None: + """Создаёт таблицу public.orders, если она ещё не существует.""" + ddl = """ + CREATE TABLE IF NOT EXISTS public.orders ( + order_id BIGINT, + order_ts TIMESTAMP NOT NULL, + customer_id BIGINT NOT NULL, + amount NUMERIC(12,2) NOT NULL + ) + WITH (appendonly=true, orientation=row, compresstype=zstd, compresslevel=1) + DISTRIBUTED BY (order_id); + """ + with get_gp_conn() as conn, conn.cursor() as cur: + cur.execute(ddl) + conn.commit() + + +def _generate_csv(rows: int, csv_dir: Path) -> str: + """Генерирует CSV c заказами с помощью pandas и сохраняет на диск.""" + csv_dir.mkdir(parents=True, exist_ok=True) + timestamp = datetime.now(UTC).strftime("%Y%m%d_%H%M%S") + csv_path = csv_dir / f"orders_{timestamp}.csv" + + # Генерируем данные в pandas-стиле + base_order_id = int(datetime.now(UTC).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.now(UTC), 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) + + +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) + ) + numeric_summary = df.describe(include="number") + logging.info("Числовая статистика:\n%s", numeric_summary.to_string()) + if df["order_ts"].notna().any(): + logging.info( + "Диапазон order_ts: %s → %s", + df["order_ts"].min().isoformat(), + df["order_ts"].max().isoformat(), + ) + + +def _load_csv(csv_path: str) -> None: + """Загружает CSV в Greenplum через временную таблицу и anti-join.""" + csv_file = Path(csv_path) + 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;" + ) + cur.copy_expert( + "COPY tmp_orders (order_id, order_ts, customer_id, amount) FROM STDIN WITH CSV HEADER", + f, + ) + + cur.execute("SELECT COUNT(*) FROM tmp_orders") + tmp_rows = cur.fetchone()[0] + + cur.execute( + """ + INSERT INTO public.orders(order_id, order_ts, customer_id, amount) + SELECT t.order_id, t.order_ts, t.customer_id, t.amount + FROM tmp_orders t + LEFT JOIN public.orders o ON o.order_id = t.order_id + WHERE o.order_id IS NULL + """ + ) + inserted = cur.rowcount if cur.rowcount != -1 else 0 + conn.commit() + + logging.info("Загружено строк: %s (прочитано из CSV: %s)", inserted, tmp_rows) + + +default_args = {"owner": "airflow", "retries": 1, "retry_delay": timedelta(seconds=30)} + +with DAG( + dag_id="csv_to_greenplum", + start_date=datetime(2017, 1, 1), + schedule=None, + catchup=False, + default_args=default_args, + tags=["demo", "greenplum", "csv"], +) as dag: + create_table = PythonOperator( + task_id="create_orders_table", + python_callable=_create_table, + ) + + generate_csv = PythonOperator( + task_id="generate_csv", + python_callable=_generate_csv, + op_kwargs={"rows": CSV_ROWS, "csv_dir": CSV_DIR}, + ) + + preview_csv = PythonOperator( + task_id="preview_csv", + python_callable=_preview_csv, + op_kwargs={ + "csv_path": "{{ ti.xcom_pull(task_ids='generate_csv') }}", + "sample_rows": 5, + }, + ) + + load_csv = PythonOperator( + task_id="load_csv_to_greenplum", + python_callable=_load_csv, + op_kwargs={ + "csv_path": "{{ ti.xcom_pull(task_ids='generate_csv') }}", + }, + ) + + create_table >> generate_csv >> preview_csv >> load_csv diff --git a/airflow-docker/dags/csv_to_greenplum_dq.py b/airflow-docker/dags/csv_to_greenplum_dq.py new file mode 100644 index 0000000..6893d5f --- /dev/null +++ b/airflow-docker/dags/csv_to_greenplum_dq.py @@ -0,0 +1,97 @@ +from __future__ import annotations + +import logging +from datetime import datetime, timedelta + +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 airflow import DAG + + +def _run_check(check_callable): + """ + Оборачивает проверку качества данных в контекст подключения к Greenplum. + + Этот DAG предназначен для автоматической проверки качества данных + после CSV-пайплайна в таблице public.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)} + +with DAG( + dag_id="csv_to_greenplum_dq", + start_date=datetime(2017, 1, 1), + schedule=None, + catchup=False, + default_args=default_args, + tags=["demo", "greenplum", "quality", "csv", "dq"], + description="Проверки качества данных после CSV → public.orders в 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 >> dq_summary diff --git a/airflow-docker/dags/ddl_greenplum_base.py b/airflow-docker/dags/ddl_greenplum_base.py new file mode 100644 index 0000000..cfd7e64 --- /dev/null +++ b/airflow-docker/dags/ddl_greenplum_base.py @@ -0,0 +1,32 @@ +from __future__ import annotations + +""" +Учебный DAG: применяет DDL для базовой таблицы orders в Greenplum. +Запускается вручную перед CSV‑пайплайном или после изменения схемы. +""" + +from datetime import datetime, timedelta + +from airflow.providers.postgres.operators.postgres import PostgresOperator + +from airflow import DAG + +GREENPLUM_CONN_ID = "greenplum_conn" + +default_args = {"owner": "airflow", "retries": 1, "retry_delay": timedelta(seconds=30)} + +with DAG( + dag_id="orders_base_ddl", + start_date=datetime(2017, 1, 1), + schedule=None, + catchup=False, + template_searchpath="/sql", + default_args=default_args, + tags=["demo", "greenplum", "ddl", "orders"], + description="Создаёт/обновляет базовую таблицу orders в схеме public", +) as dag: + apply_orders_ddl = PostgresOperator( + task_id="apply_orders_ddl", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="base/orders_ddl.sql", + ) diff --git a/airflow-docker/dags/helpers/greenplum.py b/airflow-docker/dags/helpers/greenplum.py new file mode 100644 index 0000000..0c78c6e --- /dev/null +++ b/airflow-docker/dags/helpers/greenplum.py @@ -0,0 +1,223 @@ +from __future__ import annotations + +""" +LEGACY: Вспомогательные функции для прямого подключения к Greenplum через psycopg2. +Внимание: этот модуль оставлен только для поддержки базового CSV-пайплайна. +В новых DAG (ODS/DDS/DM) используйте встроенный в Airflow PostgresOperator +и штатные механизмы XCom. +""" + +import logging +import os +from typing import List, Sequence, Tuple + +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", +) + +# Ожидаемая схема таблицы orders для проверки качества данных +EXPECTED_ORDERS_SCHEMA: List[Tuple[str, str]] = [ + ("order_id", "bigint"), + ("order_ts", "timestamp without time zone"), + ("customer_id", "bigint"), + ("amount", "numeric"), +] + + +def get_gp_conn(): + """ + Возвращает 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) + conn = hook.get_conn() + logging.info("✅ Подключение через Airflow Connection успешно") + return conn + except Exception as e: + logging.warning("⚠️ Не удалось подключиться через Airflow Connection: %s", e) + logging.info("🔄 Переключаемся на прямое подключение по ENV переменным") + # Фоллбек на прямое подключение по переменным окружения. + + # Прямое подключение по переменным окружения + conn_params = { + "dbname": os.getenv("GP_DB", "gp_dwh"), + "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/%s", + conn_params["host"], + conn_params["port"], + conn_params["dbname"], + ) + return psycopg2.connect(**conn_params) + + +def assert_orders_table_exists(conn) -> None: + """ + Проверяет наличие таблицы orders в схеме public. + + Args: + conn: Подключение к Greenplum + + Raises: + ValueError: Если таблица не найдена + """ + logging.info("🔍 Проверяем существование таблицы public.orders...") + with conn.cursor() as cur: + cur.execute( + """ + SELECT 1 + FROM pg_catalog.pg_tables + WHERE schemaname = 'public' AND tablename = 'orders' + """ + ) + if cur.fetchone() is None: + 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( + """ + SELECT column_name, data_type + FROM information_schema.columns + WHERE table_schema = 'public' AND table_name = 'orders' + ORDER BY ordinal_position + """ + ) + return cur.fetchall() + + +def assert_orders_schema(conn) -> None: + """ + Проверяет, что схема таблицы 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}." + ) + 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 не пустая. + + 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( + """ + SELECT COUNT(*) FROM ( + SELECT order_id + FROM public.orders + GROUP BY order_id + HAVING COUNT(*) > 1 + ) d + """ + ) + return cur.fetchone()[0] + + +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} шт.) — проверь загрузку данных." + ) + logging.info("✅ Дубликаты не обнаружены") diff --git a/airflow-docker/sql/base/orders_ddl.sql b/airflow-docker/sql/base/orders_ddl.sql new file mode 100644 index 0000000..f6b4caf --- /dev/null +++ b/airflow-docker/sql/base/orders_ddl.sql @@ -0,0 +1,11 @@ +-- DDL для базовой таблицы orders, которую использует CSV‑pipeline. +-- Выполняется идемпотентно: таблица создаётся, если ещё не существует. + +CREATE TABLE IF NOT EXISTS public.orders ( + order_id BIGINT, + order_ts TIMESTAMP NOT NULL, + customer_id BIGINT NOT NULL, + amount NUMERIC(12,2) NOT NULL +) +WITH (appendonly=true, orientation=row, compresstype=zstd, compresslevel=1) +DISTRIBUTED BY (order_id); diff --git a/airflow-docker/tests/test_greenplum_helpers.py b/airflow-docker/tests/test_greenplum_helpers.py new file mode 100644 index 0000000..a60876c --- /dev/null +++ b/airflow-docker/tests/test_greenplum_helpers.py @@ -0,0 +1,178 @@ +from __future__ import annotations + +from dataclasses import dataclass +from typing import Any, List, Sequence + +import pytest + +import airflow.dags.helpers.greenplum as greenplum +from tests.conftest import patch_postgres_hook + + +@dataclass +class FakeCursor: + fetchone_value: Any = None + fetchall_value: Sequence[Any] | None = None + rowcount: int | None = None + + def __post_init__(self) -> None: + self.queries: List[Any] = [] + + def execute(self, query: str, params: Any | None = None) -> None: + self.queries.append((query, params)) + + def fetchone(self) -> Any: + return self.fetchone_value + + def fetchall(self) -> Sequence[Any] | None: + return self.fetchall_value + + def __enter__(self) -> FakeCursor: + return self + + def __exit__(self, exc_type, exc, tb) -> None: + return None + + +class FakeConn: + def __init__(self, cursors: Sequence[FakeCursor]) -> None: + self._cursors = list(cursors) + self._index = 0 + self.commits = 0 + + def cursor(self) -> FakeCursor: + cursor = self._cursors[self._index] + self._index += 1 + return cursor + + def commit(self) -> None: + self.commits += 1 + + +def test_get_gp_conn_uses_airflow_hook(monkeypatch) -> None: + class FakeHook: + def __init__(self, postgres_conn_id: str) -> None: + self.postgres_conn_id = postgres_conn_id + + def get_conn(self) -> str: + return "hook_connection" + + patch_postgres_hook(monkeypatch, FakeHook) + monkeypatch.setattr(greenplum, "GP_CONN_ID", "demo_conn", raising=False) + monkeypatch.setattr(greenplum, "GP_USE_AIRFLOW_CONN", True, raising=False) + + conn = greenplum.get_gp_conn() + + assert conn == "hook_connection" + + +def test_get_gp_conn_fallback_to_psycopg(monkeypatch) -> None: + class BrokenHook: + def __init__(self, postgres_conn_id: str) -> None: + self.postgres_conn_id = postgres_conn_id + + def get_conn(self): + raise RuntimeError("boom") + + patch_postgres_hook(monkeypatch, BrokenHook) + monkeypatch.setattr(greenplum, "GP_USE_AIRFLOW_CONN", True, raising=False) + monkeypatch.setattr(greenplum, "GP_CONN_ID", "demo_conn", raising=False) + monkeypatch.setenv("GP_DB", "demo_db") + monkeypatch.setenv("GP_USER", "demo_user") + monkeypatch.setenv("GP_PASSWORD", "secret") + monkeypatch.setenv("GP_HOST", "greenplum-host") + monkeypatch.setenv("GP_PORT", "5434") + + captured_kwargs = {} + + def fake_connect(**kwargs): + captured_kwargs.update(kwargs) + return "psycopg_connection" + + monkeypatch.setattr(greenplum.psycopg2, "connect", fake_connect) + + conn = greenplum.get_gp_conn() + + assert conn == "psycopg_connection" + assert captured_kwargs == { + "dbname": "demo_db", + "user": "demo_user", + "password": "secret", + "host": "greenplum-host", + "port": 5434, + } + + +def test_get_gp_conn_without_airflow(monkeypatch) -> None: + monkeypatch.setattr(greenplum, "GP_USE_AIRFLOW_CONN", False, raising=False) + monkeypatch.setenv("GP_DB", "demo_db") + monkeypatch.setenv("GP_USER", "demo_user") + monkeypatch.setenv("GP_PASSWORD", "secret") + monkeypatch.setenv("GP_HOST", "greenplum-host") + monkeypatch.setenv("GP_PORT", "5435") + + captured_kwargs = {} + + def fake_connect(**kwargs): + captured_kwargs.update(kwargs) + return "direct_psycopg" + + monkeypatch.setattr(greenplum.psycopg2, "connect", fake_connect) + + conn = greenplum.get_gp_conn() + + assert conn == "direct_psycopg" + assert captured_kwargs["port"] == 5435 + + +def test_assert_orders_table_exists_ok() -> None: + conn = FakeConn([FakeCursor(fetchone_value=(1,))]) + + greenplum.assert_orders_table_exists(conn) + + +def test_assert_orders_table_exists_missing() -> None: + conn = FakeConn([FakeCursor(fetchone_value=None)]) + + with pytest.raises(ValueError): + greenplum.assert_orders_table_exists(conn) + + +def test_assert_orders_schema_ok() -> None: + expected = list(greenplum.EXPECTED_ORDERS_SCHEMA) + conn = FakeConn([FakeCursor(fetchall_value=expected)]) + + greenplum.assert_orders_schema(conn) + + +def test_assert_orders_schema_mismatch() -> None: + conn = FakeConn([FakeCursor(fetchall_value=[("order_id", "bigint")])]) + + with pytest.raises(ValueError): + greenplum.assert_orders_schema(conn) + + +def test_assert_orders_have_rows_ok() -> None: + conn = FakeConn([FakeCursor(fetchone_value=(5,))]) + + greenplum.assert_orders_have_rows(conn) + + +def test_assert_orders_have_rows_empty() -> None: + conn = FakeConn([FakeCursor(fetchone_value=(0,))]) + + with pytest.raises(ValueError): + greenplum.assert_orders_have_rows(conn) + + +def test_assert_orders_no_duplicates_ok() -> None: + conn = FakeConn([FakeCursor(fetchone_value=(0,))]) + + greenplum.assert_orders_no_duplicates(conn) + + +def test_assert_orders_no_duplicates_detected() -> None: + conn = FakeConn([FakeCursor(fetchone_value=(3,))]) + + with pytest.raises(ValueError): + greenplum.assert_orders_no_duplicates(conn)