From a6a3559ab971d8387e4ec9bb921d8cd5058085df Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sun, 8 Mar 2026 19:12:13 +0300 Subject: [PATCH] =?UTF-8?q?refactor(csv):=20=D1=83=D0=B4=D0=B0=D0=BB=D0=B5?= =?UTF-8?q?=D0=BD=20=D0=BB=D0=B5=D0=B3=D0=B0=D1=81=D0=B8=20CSV-=D0=BF?= =?UTF-8?q?=D0=B0=D0=B9=D0=BF=D0=BB=D0=B0=D0=B9=D0=BD=20=D0=B8=20=D1=81?= =?UTF-8?q?=D0=B2=D1=8F=D0=B7=D0=B0=D0=BD=D0=BD=D1=8B=D0=B5=20=D1=81=20?= =?UTF-8?q?=D0=BD=D0=B8=D0=BC=20=D1=84=D0=B0=D0=B9=D0=BB=D1=8B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - Пример базовой загрузки CSV перенесен в отдельный репозиторий `airflow-manual` для разделения учебных треков. - Что: - удалены DAG-файлы `csv_to_greenplum` и вспомогательные скрипты `helpers/greenplum.py`, `orders_ddl.sql`. - из `docker-compose.yml` и `.env.example` удалены переменные и тома (`airflow_data`), необходимые для CSV. - очищена документация (`README.md`, `TESTING.md`, `educational-tasks.md`) и тесты (`test_dags_smoke.py`, `conftest.py`). - отмечен выполненным 'Этап 1' в `TODO.md`. - Проверка: - `make test` проходит успешно (smoke-тесты оставшихся DAG-ов не затронуты). --- .env.example | 4 - AGENTS.md | 6 +- README.md | 7 - TESTING.md | 13 -- TODO.md | 6 +- airflow/dags/csv_to_greenplum.py | 167 -------------------- airflow/dags/csv_to_greenplum_dq.py | 97 ------------ airflow/dags/ddl_greenplum_base.py | 32 ---- airflow/dags/helpers/greenplum.py | 223 --------------------------- docker-compose.yml | 7 +- docs/internal/architecture_review.md | 4 +- docs/stack.md | 2 +- educational-tasks.md | 65 ++------ plans/dockerfile-improvements.md | 2 +- sql/base/orders_ddl.sql | 11 -- sql/ddl_gp.sql | 5 +- tests/conftest.py | 11 -- tests/test_dags_smoke.py | 47 ------ tests/test_greenplum_helpers.py | 178 --------------------- 19 files changed, 21 insertions(+), 866 deletions(-) delete mode 100644 airflow/dags/csv_to_greenplum.py delete mode 100644 airflow/dags/csv_to_greenplum_dq.py delete mode 100644 airflow/dags/ddl_greenplum_base.py delete mode 100644 airflow/dags/helpers/greenplum.py delete mode 100644 sql/base/orders_ddl.sql delete mode 100644 tests/test_greenplum_helpers.py diff --git a/.env.example b/.env.example index a5de38d..5680198 100644 --- a/.env.example +++ b/.env.example @@ -43,7 +43,3 @@ GP_USE_AIRFLOW_CONN=true PXF_SEED_OVERWRITE=0 # PXF_SYNC_ON_START=1 — выполнять `pxf cluster sync` при старте контейнера (дольше, но гарантирует актуальные конфиги) PXF_SYNC_ON_START=0 - -# CSV Pipeline -CSV_DIR=/opt/airflow/data -CSV_ROWS=1000 diff --git a/AGENTS.md b/AGENTS.md index 8ce571b..b86b3de 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -9,7 +9,7 @@ - **Фокус на «Почему»:** При использовании специфичных паттернов DWH (например, `delete + insert` для инкремента в Greenplum вместо `merge`) — добавляйте краткий комментарий, объясняющий этот выбор студентам. ## 2. Карта проекта (Навигация для Агента) -- `airflow/dags/` — DAG-файлы (напр. `csv_to_greenplum.py`). +- `airflow/dags/` — DAG-файлы (напр. `bookings_to_gp_stage.py`). - `sql/` — DDL и SQL-скрипты. Разделены на слои: src/ (исходные системы), stg/ (стейджинг), ods/ (операционное хранилище), dds/ (детальное хранилище), dm/ (слой витрин). - *Правило ИИ:* DDL таблиц хранится строго рядом с объектом (напр. `sql/stg/bookings_ddl.sql`). - `docs/internal/naming_conventions.md` — Единый источник истины для нейминга служебных и SCD-полей. *Правило ИИ: Всегда сверяться с этим файлом при генерации новых DDL/SQL.* @@ -54,7 +54,7 @@ ## Тестирование - Тесты лежат в `tests/` (pytest). Запуск: `make test`. -- Есть юнит‑тесты для `helpers/greenplum.py` и smoke‑тесты DAG‑структуры (`tests/test_dags_smoke.py`). +- Есть smoke‑тесты DAG‑структуры (`tests/test_dags_smoke.py`). - Smoke‑тесты DAG автоматически пропускаются, если Airflow не установлен в venv. - Для ручного прогона стенда см. `TESTING.md` (пошаговый чек‑лист для студентов). - Для программной проверки DAG (без браузера) — см. `docs/agent-dag-testing.md`: CLI, REST API, проверка параллельности, запросы в Greenplum. @@ -66,7 +66,7 @@ - При изменении схемы/поведения — обновляйте `README.md` и `sql/ddl_gp.sql`. ## Безопасность и конфигурация -- Все настройки — через `.env`; креды в коде не хардкодим. Частые переменные: `GP_*`, `PG_*`, `AIRFLOW_*`, `CSV_*`. +- Все настройки — через `.env`; креды в коде не хардкодим. Частые переменные: `GP_*`, `PG_*`, `AIRFLOW_*`. - `make clean` удаляет тома — предупреждайте студентов, что данные пропадут. ## Для агента (особенности аудитории) diff --git a/README.md b/README.md index 6424488..0b9e128 100644 --- a/README.md +++ b/README.md @@ -22,7 +22,6 @@ - **Greenplum** (single‑node для обучения; внешний порт по умолчанию `5435`) - **bookings-db** (Postgres с демо‑БД `demo`; внешний порт по умолчанию `5434`) - **PXF** как “транспорт” между Postgres и Greenplum (уже настроен в образе) -- Побочный пример: загрузка данных через **pandas/CSV** (`csv_to_greenplum`) ## Требования @@ -104,12 +103,6 @@ SELECT COUNT(*) FROM dds.fact_flight_sales; - `bookings_dds_ddl` — создаёт/обновляет DDS-таблицы (`dim_*`, `fact_flight_sales`) по домену bookings. - `bookings_to_gp_dds` — загружает данные из ODS в DDS (SCD1/SCD2 + факт) и выполняет DQ‑проверки. -Вспомогательные (побочный трек с CSV): - -- `orders_base_ddl` — создаёт таблицу `public.orders` для CSV‑пайплайна; -- `csv_to_greenplum` — pandas → CSV → Greenplum (пример загрузки без источника‑БД); -- `csv_to_greenplum_dq` — проверки качества данных для `public.orders`. - ## Полезные команды ```bash diff --git a/TESTING.md b/TESTING.md index fb5f99d..f1370e1 100644 --- a/TESTING.md +++ b/TESTING.md @@ -28,14 +28,6 @@ 1. Открыть http://localhost:8080 (admin/admin). 2. (опционально) Зайти в Admin → Connections и убедиться, что DAG’и видят подключения: - `greenplum_conn` и `bookings_db` задаются через переменные `AIRFLOW_CONN_...` в docker-compose и могут не отображаться в списке, но `airflow connections get greenplum_conn` / `bookings_db` внутри контейнера должны отрабатывать без ошибок. -3. DAG `csv_to_greenplum`: - - Включить переключатель. - - Нажать «Trigger DAG». - - Контроль: все таски Success, в `data/` появился CSV, в логах `load_csv_to_greenplum` видно `INSERT`. - - В Greenplum (см. п.5) убедиться в наличии строк `(SELECT COUNT(*) ...)`. -4. DAG `csv_to_greenplum_dq`: - - Запустить вручную после первого DAG. - - Проверить, что все 5 задач Success и логи содержат `Проверка пройдена`. - DAG `bookings_to_gp_stage` (полная проверка цепочки bookings → Greenplum STG): - предварительно выполнить один раз: `make bookings-init` (установка демобазы `demo` в контейнере `bookings-db`) и `make ddl-gp` (создаёт STG/ODS/DDS слои в Greenplum, включая внешние `*_ext` через PXF); @@ -59,18 +51,13 @@ - `docker compose exec greenplum bash -lc "su - gpadmin -c '/usr/local/pxf/bin/pxf cluster status'"` - Команды внутри psql: - `\dt public.*` — таблицы схему public. - - `SELECT COUNT(*) FROM public.orders;` — оценка объёма. - - `SELECT * FROM public.orders LIMIT 5;` — визуальная проверка. - - `SELECT order_id FROM public.orders GROUP BY 1 HAVING COUNT(*) > 1;` — поиск дублей. - (после настройки PXF) `SELECT COUNT(*) FROM public.ext_bookings_bookings;` — проверка чтения из демо-БД bookings через PXF. - (после настройки PXF) `SELECT * FROM public.ext_bookings_bookings LIMIT 5;` — визуальное сравнение с таблицей `bookings.bookings` в исходной БД. - Завершить `\q`. ## 6. Негативные сценарии и fallback -- **Пустая таблица**: запустить `csv_to_greenplum_dq` до `csv_to_greenplum`. Ожидается ошибка на таске `check_orders_has_rows`. - **Проблемы с подключением**: временно изменить `GP_HOST` или `GP_PORT` на несуществующий, перезапустить `make up`, убедиться, что DAG падает с понятной ошибкой (`psycopg2.OperationalError`). - **Fallback без Airflow Connection**: установить `GP_USE_AIRFLOW_CONN=false`, перезапустить стек (`make down && make up`), удостовериться, что загрузка и DQ работают через ENV. -- **Дубликаты**: дважды вызвать `csv_to_greenplum` — ожидаем, что количество строк в `public.orders` не увеличится на размер CSV, а DAG `csv_to_greenplum_dq` не найдёт дублей. - **PXF и демобаза bookings** (после настройки PXF и выполнения `make ddl-gp`): временно остановить `bookings-db` (`docker compose stop bookings-db`) и попробовать выполнить `SELECT COUNT(*) FROM public.ext_bookings_bookings;` в `make gp-psql` — ожидается ошибка подключения. Затем запустить `bookings-db` (`docker compose start bookings-db`) и убедиться, что запрос снова работает. ## 7. Быстрый reset (если «что-то сломалось») diff --git a/TODO.md b/TODO.md index 1ae9581..5b20d21 100644 --- a/TODO.md +++ b/TODO.md @@ -16,11 +16,11 @@ **Инструмент:** Sonnet / Gemini / ChatGPT — механическая работа, перенос файлов. -- [ ] Перенести в [airflow-manual](https://github.com/dementev-dev/airflow-manual): +- [x] Перенести в [airflow-manual](https://github.com/dementev-dev/airflow-manual): `csv_to_greenplum.py`, `csv_to_greenplum_dq.py`, `ddl_greenplum_base.py`, `helpers/greenplum.py`, `sql/base/orders_ddl.sql`, связанные тесты -- [ ] Убрать CSV-зависимости из docker-compose / .env (`CSV_DIR`, `CSV_ROWS`) -- [ ] Обновить README (убрать упоминания CSV-пайплайна) +- [x] Убрать CSV-зависимости из docker-compose / .env (`CSV_DIR`, `CSV_ROWS`) +- [x] Обновить README (убрать упоминания CSV-пайплайна) ### Этап 1.5. Полировка эталона diff --git a/airflow/dags/csv_to_greenplum.py b/airflow/dags/csv_to_greenplum.py deleted file mode 100644 index f99eaa7..0000000 --- a/airflow/dags/csv_to_greenplum.py +++ /dev/null @@ -1,167 +0,0 @@ -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/dags/csv_to_greenplum_dq.py b/airflow/dags/csv_to_greenplum_dq.py deleted file mode 100644 index 6893d5f..0000000 --- a/airflow/dags/csv_to_greenplum_dq.py +++ /dev/null @@ -1,97 +0,0 @@ -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/dags/ddl_greenplum_base.py b/airflow/dags/ddl_greenplum_base.py deleted file mode 100644 index cfd7e64..0000000 --- a/airflow/dags/ddl_greenplum_base.py +++ /dev/null @@ -1,32 +0,0 @@ -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/dags/helpers/greenplum.py b/airflow/dags/helpers/greenplum.py deleted file mode 100644 index 0c78c6e..0000000 --- a/airflow/dags/helpers/greenplum.py +++ /dev/null @@ -1,223 +0,0 @@ -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/docker-compose.yml b/docker-compose.yml index f792e6e..d435488 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -11,7 +11,6 @@ x-airflow-common-env: &airflow-env x-airflow-common-volumes: &airflow-volumes - ./airflow/dags:/opt/airflow/dags - ./sql:/sql:ro - - airflow_data:/opt/airflow/data - airflow_logs:/opt/airflow/logs x-airflow-common-depends: &airflow-depends @@ -110,7 +109,6 @@ services: volumes: - ./airflow/dags:/opt/airflow/dags - ./sql:/sql:ro - - airflow_data:/opt/airflow/data - airflow_logs:/opt/airflow/logs - ./airflow/requirements.txt:/opt/airflow/requirements.txt healthcheck: @@ -135,7 +133,6 @@ services: volumes: - ./airflow/dags:/opt/airflow/dags - ./sql:/sql:ro - - airflow_data:/opt/airflow/data - airflow_logs:/opt/airflow/logs - ./airflow/requirements.txt:/opt/airflow/requirements.txt depends_on: @@ -154,12 +151,11 @@ services: volumes: - ./airflow/dags:/opt/airflow/dags - ./sql:/sql:ro - - airflow_data:/opt/airflow/data - airflow_logs:/opt/airflow/logs command: > bash -lc " set -e; - mkdir -p /opt/airflow/data /opt/airflow/logs && chown -R airflow:root /opt/airflow/data /opt/airflow/logs; + mkdir -p /opt/airflow/logs && chown -R airflow:root /opt/airflow/logs; # Дожидаемся готовности БД ретрая миграции for i in {1..30}; do su -s /bin/bash airflow -c \"PATH='/home/airflow/.local/bin:$${PATH}' airflow db migrate\" && break || echo 'waiting for pgmeta' && sleep 3; @@ -175,5 +171,4 @@ volumes: pgmeta: bookings_data: greenplum_data: - airflow_data: airflow_logs: diff --git a/docs/internal/architecture_review.md b/docs/internal/architecture_review.md index fc7f105..0a5276c 100644 --- a/docs/internal/architecture_review.md +++ b/docs/internal/architecture_review.md @@ -91,9 +91,7 @@ GP-специфичная best practice, которую забывают даж - `_resolve_stg_batch_id` с INTERSECT по 4 таблицам — нет комментария **зачем** нужна согласованность - **Нужно**: комментарий в `airflow/dags/bookings_to_gp_ods.py` перед SQL-запросом -- [x] **`helpers/greenplum.py` без пометки «legacy»** - - Использует прямой psycopg2 + ENV — противоречит PostgresOperator-подходу - - **Нужно**: docstring «LEGACY: только для CSV-пайплайна» в `airflow/dags/helpers/greenplum.py` +- [x] **`helpers/greenplum.py`** — удалён вместе с CSV-пайплайном (перенесён в airflow-manual) ### P2: Средние усилия, заметное улучшение качества diff --git a/docs/stack.md b/docs/stack.md index 4a706c6..e446a3a 100644 --- a/docs/stack.md +++ b/docs/stack.md @@ -54,7 +54,7 @@ make gp-psql PXF (Platform Extension Framework) — компонент Greenplum для работы с внешними источниками. В этом стенде PXF используется для чтения таблицы `bookings.bookings` из Postgres прямо из Greenplum через внешнюю таблицу `stg.bookings_ext`. Поэтому загрузка в `stg.bookings` выглядит как обычный -`INSERT ... SELECT` без промежуточных CSV. +`INSERT ... SELECT` без промежуточных файлов. ## Greenplum + PXF: свой образ diff --git a/educational-tasks.md b/educational-tasks.md index 6d1bd48..b418890 100644 --- a/educational-tasks.md +++ b/educational-tasks.md @@ -1,63 +1,18 @@ # Учебные задания по стенду Этот документ собирает в одном месте задания для менти. -Он разбит на блоки: от базовой работы с CSV‑pipeline до более продвинутого сценария с демо‑БД bookings и слоем STG в Greenplum. +Он разбит на блоки: от архитектуры Greenplum и демо‑БД bookings до реализации аналитических слоев DWH. -Если вы только начинаете, выполняйте задания по порядку. К разделу про bookings можно вернуться позже. +Если вы только начинаете, выполняйте задания по порядку. --- -## 1. Базовый CSV‑pipeline (csv_to_greenplum) - -Основная цель этого блока — понять, как устроен простой ETL: генерация данных через pandas, сохранение в CSV и загрузка в Greenplum. - -### 1.1. Разбор готового pipeline - -1. Найдите DAG `csv_to_greenplum` в `airflow/dags/csv_to_greenplum.py`. -2. Ответьте себе на вопросы (можно коротко в отдельном файле/блокноте): - - какие задачи (tasks) входят в DAG и что делает каждая из них; - - какие таблицы создаются в Greenplum; - - где физически лежат CSV‑файлы; - - какие параметры управляют размером датасета. -3. Поднимите стенд и запустите DAG: - - `make up` (Airflow инициализируется автоматически при первом старте) - - включите и запустите DAG `csv_to_greenplum` в Airflow UI. -4. Проверьте результат в Greenplum: - - `make gp-psql` - - `SELECT COUNT(*) FROM public.orders;` - - `SELECT * FROM public.orders LIMIT 5;` - -### 1.2. Изменение параметров генерации - -1. Найдите, где задаётся количество строк для генерации (`CSV_ROWS` в `.env` и параметр в DAG). -2. Поставьте другое значение и перезапустите DAG: - - оцените, как изменилось количество строк в `public.orders`; - - убедитесь, что пайплайн по‑прежнему работает без ошибок. -3. Попробуйте изменить схему данных (добавить колонку в CSV и таблицу в Greenplum): - - добавьте новую колонку в генерацию pandas; - - обновите DDL/SQL, чтобы колонка появилась в таблице `public.orders`; - - перезапустите DAG и убедитесь, что новая колонка заполняется. - -### 1.3. Собственные проверки качества данных - -1. Найдите DAG `csv_to_greenplum_dq` в `airflow/dags/csv_to_greenplum_dq.py`. -2. Посмотрите, какие проверки уже реализованы (наличие таблицы, схема, дубликаты). -3. Добавьте ещё одну простую проверку, например: - - проверка, что в таблице `public.orders` не больше N строк; - - проверка, что поле (например, `order_price`) не содержит отрицательных значений; - - проверка, что нет строк с `NULL` в ключевых колонках. -4. Запустите DAG `csv_to_greenplum_dq` и убедитесь, что: - - новая проверка проходит на «хороших» данных; - - при нарушении условия DAG падает с понятной ошибкой. - ---- - -## 2. Greenplum и модель данных (введение) +## 1. Greenplum и модель данных (введение) В следующих заданиях мы будем опираться на демо‑БД bookings (Postgres) и слой STG в Greenplum. На этом этапе достаточно бегло посмотреть на структуру и понять общую идею, детальная проработка пойдёт позже. -### 2.1. Знакомство с демо‑БД bookings +### 1.1. Знакомство с демо‑БД bookings 1. Прочитайте `bookings/README.md` — какие сервисы и команды относятся к демобазе. 2. Поднимите стенд и выполните: @@ -70,7 +25,7 @@ - какие типы колонок используются; - какие поля выглядят как ключи, даты, суммы. -### 2.2. Знакомство с STG в Greenplum +### 1.2. Знакомство с STG в Greenplum 1. Прочитайте `sql/stg/bookings_ddl.sql` и краткое описание потока `docs/bookings_to_gp_stage.md` (если интересно — `docs/internal/bookings_stg_design.md`). 2. Ответьте себе на вопросы: @@ -81,7 +36,7 @@ - `\dn` и `\dt stg.*` - `SELECT * FROM stg.bookings LIMIT 5;` (после запуска соответствующего DAG). -### 2.3. Как генерируются учебные данные bookings +### 1.3. Как генерируются учебные данные bookings 1. Откройте файл `bookings/generate_next_day.sql` и ответьте себе на вопросы: - с какой даты начинается генерация данных (посмотрите на GUC `bookings.start_date` и переменную `v_start_cfg`); @@ -99,12 +54,12 @@ --- -## 3. DAG bookings_to_gp_stage (заготовка заданий) +## 2. DAG bookings_to_gp_stage (заготовка заданий) Этот DAG показывает путь данных от демо‑БД bookings в Postgres до сырого слоя STG в Greenplum. Сейчас он уже реализован как учебный пример, а в будущем вокруг него появятся отдельные задания по моделированию DWH. -### 3.1. Что есть сейчас +### 2.1. Что есть сейчас 1. Откройте `airflow/dags/bookings_to_gp_stage.py`. 2. Найдите в коде ссылки на SQL‑файлы: @@ -121,9 +76,9 @@ На этом этапе достаточно понять общую цепочку. Детальные задания по переработке модели данных и построению ODS/DDS/DM слоёв будут добавлены позже. -### 3.2. Идеи для будущих заданий (черновик) +### 2.2. Идеи для будущих заданий (черновик) -> Ниже — набросок задач, к которым мы вернёмся, когда базовые темы по Airflow и CSV‑pipeline будут освоены. +> Ниже — набросок задач, к которым мы вернёмся, когда базовые темы по Airflow будут освоены. Планируемые направления: diff --git a/plans/dockerfile-improvements.md b/plans/dockerfile-improvements.md index f117053..5903725 100644 --- a/plans/dockerfile-improvements.md +++ b/plans/dockerfile-improvements.md @@ -284,7 +284,7 @@ airflow-init: - Проверить наличие DAG'ов в списке 2. **Проверить запуск DAG'ов** - - Выбрать любой DAG (например, `csv_to_greenplum`) + - Выбрать любой DAG (например, `bookings_to_gp_stage`) - Запустить вручную через UI - Проверить успешность выполнения задач diff --git a/sql/base/orders_ddl.sql b/sql/base/orders_ddl.sql deleted file mode 100644 index f6b4caf..0000000 --- a/sql/base/orders_ddl.sql +++ /dev/null @@ -1,11 +0,0 @@ --- 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/sql/ddl_gp.sql b/sql/ddl_gp.sql index 893d895..722d48a 100644 --- a/sql/ddl_gp.sql +++ b/sql/ddl_gp.sql @@ -1,6 +1,6 @@ -- Главный входной DDL-скрипт для Greenplum в учебном стенде. -- Выполняется из контейнера командой `make ddl-gp` и создаёт/обновляет --- все объекты, которые нужны базовым DAG (csv_to_greenplum, bookings_to_gp_stage); +-- все объекты, которые нужны DAG'ам bookings ETL-пайплайна; -- подключает файловые DDL через \i, чтобы сохранять единый входной скрипт. -- -- Чтобы не ломать задания, новые объекты лучше добавлять в отдельные файлы @@ -8,9 +8,6 @@ -- -- Подробнее про STG/bookings: см. docs/bookings_to_gp_stage.md. --- Таблица для CSV‑пайплайна (csv_to_greenplum). -\i base/orders_ddl.sql - -- Внешняя таблица для чтения данных из демо-БД bookings через PXF (JDBC). -- Источник: таблица bookings.bookings в базе demo (Postgres, сервис bookings-db). DROP EXTERNAL TABLE IF EXISTS public.ext_bookings_bookings; diff --git a/tests/conftest.py b/tests/conftest.py index 1d5d066..24854e1 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -55,14 +55,3 @@ def _ensure_stub_module(full_name: str) -> ModuleType: module = sys.modules[path] assert isinstance(module, ModuleType) return module - - -def patch_postgres_hook(monkeypatch, hook_cls: Type) -> None: - """ - Patch PostgresHook so that helpers.greenplum can be exercised without real Airflow. - """ - try: - module = importlib.import_module("airflow.providers.postgres.hooks.postgres") - except ModuleNotFoundError: - module = _ensure_stub_module("airflow.providers.postgres.hooks.postgres") - monkeypatch.setattr(module, "PostgresHook", hook_cls, raising=False) diff --git a/tests/test_dags_smoke.py b/tests/test_dags_smoke.py index a1ca87e..f94a597 100644 --- a/tests/test_dags_smoke.py +++ b/tests/test_dags_smoke.py @@ -41,53 +41,6 @@ def _assert_reachable(dag, upstream_task_id: str, downstream_task_id: str) -> No ), f"Expected {downstream_task_id} to be downstream of {upstream_task_id}" -def test_csv_to_greenplum_dag_structure(): - dag = _load_dag("airflow.dags.csv_to_greenplum") - - # tasks - expected_tasks = { - "create_orders_table", - "generate_csv", - "preview_csv", - "load_csv_to_greenplum", - } - assert expected_tasks.issubset(dag.task_dict.keys()) - - # linear dependencies - t1 = dag.get_task("create_orders_table") - t2 = dag.get_task("generate_csv") - t3 = dag.get_task("preview_csv") - t4 = dag.get_task("load_csv_to_greenplum") - - assert t2 in t1.get_direct_relatives(upstream=False) - assert t3 in t2.get_direct_relatives(upstream=False) - assert t4 in t3.get_direct_relatives(upstream=False) - - -def test_csv_to_greenplum_dq_dag_structure(): - dag = _load_dag("airflow.dags.csv_to_greenplum_dq") - - expected_tasks = { - "check_orders_table_exists", - "check_orders_schema", - "check_orders_has_rows", - "check_order_duplicates", - "data_quality_summary", - } - assert expected_tasks.issubset(dag.task_dict.keys()) - - e = dag.get_task("check_orders_table_exists") - s = dag.get_task("check_orders_schema") - h = dag.get_task("check_orders_has_rows") - d = dag.get_task("check_order_duplicates") - q = dag.get_task("data_quality_summary") - - assert s in e.get_direct_relatives(upstream=False) - assert h in s.get_direct_relatives(upstream=False) - assert d in h.get_direct_relatives(upstream=False) - assert q in d.get_direct_relatives(upstream=False) - - def test_bookings_stg_ddl_dag_structure(): """Проверка структуры DAG bookings_stg_ddl.""" dag = _load_dag("airflow.dags.bookings_stg_ddl") diff --git a/tests/test_greenplum_helpers.py b/tests/test_greenplum_helpers.py deleted file mode 100644 index a60876c..0000000 --- a/tests/test_greenplum_helpers.py +++ /dev/null @@ -1,178 +0,0 @@ -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)