diff --git a/AGENTS.md b/AGENTS.md index 1cb3e3e..0fb94f8 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -3,7 +3,7 @@ Эта репа — учебный стенд для студентов (менти), которые только начинают с Airflow/Greenplum и Python. Пожалуйста, держите решения простыми, стабильными и хорошо объяснёнными. ## Структура проекта -- `airflow/dags/` — DAG-файлы (например, `csv_to_greenplum.py`, `data_quality_greenplum.py`). +- `airflow/dags/` — DAG-файлы (например, `csv_to_greenplum.py`, `csv_to_greenplum_dq.py`). - `airflow/requirements.txt` — зависимости, которые ставятся внутри контейнеров Airflow. - `sql/` — DDL и вспомогательные SQL (например, `sql/ddl_gp.sql`). - `docker-compose.yml` — Greenplum, Airflow, Postgres (мета-БД). diff --git a/README.md b/README.md index eacb37f..ad30171 100644 --- a/README.md +++ b/README.md @@ -24,7 +24,7 @@ - Откройте UI: http://localhost:8080 (admin/admin). - Включите и запустите DAG `csv_to_greenplum`. Дождитесь Success. - Проверьте данные: `make gp-psql` → `SELECT COUNT(*) FROM public.orders;`. -- Дополнительно: запустите `greenplum_data_quality` — все проверки должны быть зелёные. +- Дополнительно: запустите `csv_to_greenplum_dq` — все проверки должны быть зелёные. Если что‑то не работает — смотрите «Типичные проблемы» и «Быстрый reset» ниже. @@ -170,8 +170,13 @@ uv run black --check airflow tests - **pandas** — библиотека для генерации и анализа данных в формате CSV ### Готовые DAG (workflow) +- **orders_base_ddl** — создаёт базовую таблицу `public.orders` для CSV‑пайплайна +- **bookings_stg_ddl** — готовит схему `stg` и таблицы `stg.bookings_ext` / `stg.bookings` - **csv_to_greenplum** — базовый pipeline: pandas → CSV → Greenplum -- **greenplum_data_quality** — проверки качества данных (наличие таблицы, схема, дубликаты) +- **bookings_to_gp_stage** — пример загрузки из демо‑БД bookings в слой STG +- **csv_to_greenplum_dq** — проверки качества данных (наличие таблицы, схема, дубликаты) + +> Учебный путь — триггернуть DDL‑DAG: для CSV `orders_base_ddl`, для bookings `bookings_stg_ddl`. Технический шорткат для быстрой инициализации — `make ddl-gp` (он не вызывается автоматически при старте контейнеров). ### Полезные команды ```bash @@ -179,7 +184,7 @@ uv run black --check airflow tests make up # Запустить весь стенд make down # Остановить и удалить данные make airflow-init # Инициализировать Airflow -make ddl-gp # Применить DDL к Greenplum +make ddl-gp # Применить DDL к Greenplum вручную make gp-psql # Подключиться к Greenplum через psql make bookings-init # Установить демобазу bookings в Postgres (по умолчанию генерирует 1 день) make bookings-generate-day # Добавить ещё один день данных в bookings (можно вызвать несколько раз) @@ -250,7 +255,7 @@ docker compose -f docker-compose.yml exec bookings-db bash -lc 'PGPASSWORD="$POS ### Проверка качества данных -Запустите DAG `greenplum_data_quality` для автоматической проверки: +Запустите DAG `csv_to_greenplum_dq` для автоматической проверки: - Наличие таблицы в базе - Соответствие схемы ожидаемой структуре - Объем загруженных данных @@ -366,7 +371,7 @@ load_bookings_to_stg = PostgresOperator( ├── airflow/ │ └── dags/ # Файлы workflow (DAG) │ ├── csv_to_greenplum.py -│ ├── data_quality_greenplum.py +│ ├── csv_to_greenplum_dq.py │ └── bookings_to_gp_stage.py ├── bookings/ # Скрипты и файлы для демобазы bookings в Postgres ├── sql/ @@ -381,7 +386,7 @@ load_bookings_to_stg = PostgresOperator( ## 💡 Советы для дальнейшего обучения 1. **Поэкспериментируйте с DAG** — измените параметры генерации данных или размер батча -2. **Добавьте свои проверки** — расширьте DAG `data_quality_greenplum.py` +2. **Добавьте свои проверки** — расширьте DAG `csv_to_greenplum_dq.py` 3. **Попробуйте другие источники** — замените генератор данных на чтение из файла или API 4. **Изучите Airflow deeper** — добавьте зависимости между задачами, настройте расписания diff --git a/TESTING.md b/TESTING.md index ff6fcb5..56986e3 100644 --- a/TESTING.md +++ b/TESTING.md @@ -31,7 +31,7 @@ - Нажать «Trigger DAG». - Контроль: все таски Success, в `data/` появился CSV, в логах `load_csv_to_greenplum` видно `INSERT`. - В Greenplum (см. п.5) убедиться в наличии строк `(SELECT COUNT(*) ...)`. -3. DAG `greenplum_data_quality`: +3. DAG `csv_to_greenplum_dq`: - Запустить вручную после первого DAG. - Проверить, что все 5 задач Success и логи содержат `Проверка пройдена`. @@ -47,10 +47,10 @@ - Завершить `\q`. ## 6. Негативные сценарии и fallback -- **Пустая таблица**: запустить `greenplum_data_quality` до `csv_to_greenplum`. Ожидается ошибка на таске `check_orders_has_rows`. +- **Пустая таблица**: запустить `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 && make airflow-init`), удостовериться, что загрузка и DQ работают через ENV. -- **Дубликаты**: дважды вызвать `csv_to_greenplum` — ожидаем, что количество строк в `public.orders` не увеличится на размер CSV, а DAG `greenplum_data_quality` не найдёт дублей. +- **Дубликаты**: дважды вызвать `csv_to_greenplum` — ожидаем, что количество строк в `public.orders` не увеличится на размер CSV, а DAG `csv_to_greenplum_dq` не найдёт дублей. ## 7. Быстрый reset (если «что-то сломалось») - Перезапустить стенд с очисткой данных: diff --git a/airflow/dags/bookings_stg_ddl.py b/airflow/dags/bookings_stg_ddl.py new file mode 100644 index 0000000..fa7dd68 --- /dev/null +++ b/airflow/dags/bookings_stg_ddl.py @@ -0,0 +1,30 @@ +from __future__ import annotations + +""" +Учебный DAG: создаёт схему stg и таблицы bookings_ext/bookings в Greenplum. +Запускается вручную перед DAG загрузки bookings_to_gp_stage или после изменения DDL. +""" + +from datetime import datetime, timedelta + +from airflow import DAG +from airflow.providers.postgres.operators.postgres import PostgresOperator + +GREENPLUM_CONN_ID = "greenplum_conn" + +default_args = {"owner": "airflow", "retries": 1, "retry_delay": timedelta(seconds=30)} + +with DAG( + dag_id="bookings_stg_ddl", + start_date=datetime(2024, 1, 1), + schedule=None, + catchup=False, + default_args=default_args, + tags=["demo", "greenplum", "ddl", "bookings", "stg"], + description="Создаёт/обновляет stg.bookings_ext и stg.bookings для учебного DAG", +) as dag: + apply_stg_bookings_ddl = PostgresOperator( + task_id="apply_stg_bookings_ddl", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="/sql/stg/bookings_ddl.sql", + ) diff --git a/airflow/dags/data_quality_greenplum.py b/airflow/dags/csv_to_greenplum_dq.py similarity index 85% rename from airflow/dags/data_quality_greenplum.py rename to airflow/dags/csv_to_greenplum_dq.py index 8fab0e7..b127005 100644 --- a/airflow/dags/data_quality_greenplum.py +++ b/airflow/dags/csv_to_greenplum_dq.py @@ -3,20 +3,23 @@ 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 +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, +) def _run_check(check_callable): """ Оборачивает проверку качества данных в контекст подключения к Greenplum. - Этот DAG предназначен для автоматической проверки качества данных в таблице orders: + Этот DAG предназначен для автоматической проверки качества данных + после CSV-пайплайна в таблице public.orders: 1. Проверяет существование таблицы 2. Проверяет соответствие схемы 3. Проверяет наличие данных @@ -47,13 +50,13 @@ def _log_dq_summary(): default_args = {"owner": "airflow", "retries": 1, "retry_delay": timedelta(seconds=30)} with DAG( - dag_id="greenplum_data_quality", + dag_id="csv_to_greenplum_dq", start_date=datetime(2024, 1, 1), schedule=None, catchup=False, default_args=default_args, - tags=["demo", "greenplum", "quality"], - description="Автоматизированные проверки качества данных в Greenplum", + tags=["demo", "greenplum", "quality", "csv", "dq"], + description="Проверки качества данных после CSV → public.orders в Greenplum", ) as dag: # Задача 1: Проверка существования таблицы check_exists = PythonOperator( diff --git a/airflow/dags/ddl_greenplum_base.py b/airflow/dags/ddl_greenplum_base.py new file mode 100644 index 0000000..7c2d4d7 --- /dev/null +++ b/airflow/dags/ddl_greenplum_base.py @@ -0,0 +1,30 @@ +from __future__ import annotations + +""" +Учебный DAG: применяет DDL для базовой таблицы orders в Greenplum. +Запускается вручную перед CSV‑пайплайном или после изменения схемы. +""" + +from datetime import datetime, timedelta + +from airflow import DAG +from airflow.providers.postgres.operators.postgres import PostgresOperator + +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(2024, 1, 1), + schedule=None, + catchup=False, + 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="/sql/base/orders_ddl.sql", + ) diff --git a/docs/internal/bookings_stg_design.md b/docs/internal/bookings_stg_design.md index 75725b0..3dee531 100644 --- a/docs/internal/bookings_stg_design.md +++ b/docs/internal/bookings_stg_design.md @@ -76,7 +76,7 @@ DDL будет добавлен в `sql/ddl_gp.sql` в блоке DDL для Gre ### 4.1. DAG для DDL -- `dag_id`: `bookings_stg_ddl` (рабочее имя). +- `dag_id`: `bookings_stg_ddl` (реализован в `airflow/dags/bookings_stg_ddl.py`). - Назначение: один раз (или при изменении схемы) создать необходимые объекты в Greenplum: - схему `stg` (если её ещё нет); - внешнюю таблицу `stg.bookings_ext` (PXF → `bookings-db`); diff --git a/docs/internal/bookings_stg_readme.md b/docs/internal/bookings_stg_readme.md index a7fe356..8ab8dc4 100644 --- a/docs/internal/bookings_stg_readme.md +++ b/docs/internal/bookings_stg_readme.md @@ -21,7 +21,8 @@ _Внутренний файл, чтобы не забыть договорён - Стенд поднят: `make up`. - Демо‑БД bookings инициализирована: `make bookings-init`. - В Greenplum применён DDL (созданы схема `stg` и таблицы `stg.bookings_ext` / `stg.bookings`): - - `make ddl-gp` (использует `sql/ddl_gp.sql`, который подтягивает `sql/stg/bookings_ddl.sql`). + - учебный вариант: запустить DAG `bookings_stg_ddl` (он использует `sql/stg/bookings_ddl.sql`); + - технический шорткат: `make ddl-gp` применяет все DDL разом вручную. Команда сама не вызывается при старте контейнеров, её нужно запустить явно. - В Airflow есть коннекты: - `greenplum_conn` (по умолчанию уже используется в helpers/greenplum.py); - `bookings_db` (Postgres к сервису `bookings-db`, если не хочется полагаться на ENV). diff --git a/educational-tasks.md b/educational-tasks.md index 2125892..4e9c28b 100644 --- a/educational-tasks.md +++ b/educational-tasks.md @@ -40,13 +40,13 @@ ### 1.3. Собственные проверки качества данных -1. Найдите DAG `greenplum_data_quality` в `airflow/dags/data_quality_greenplum.py`. +1. Найдите DAG `csv_to_greenplum_dq` в `airflow/dags/csv_to_greenplum_dq.py`. 2. Посмотрите, какие проверки уже реализованы (наличие таблицы, схема, дубликаты). 3. Добавьте ещё одну простую проверку, например: - проверка, что в таблице `public.orders` не больше N строк; - проверка, что поле (например, `order_price`) не содержит отрицательных значений; - проверка, что нет строк с `NULL` в ключевых колонках. -4. Запустите DAG `greenplum_data_quality` и убедитесь, что: +4. Запустите DAG `csv_to_greenplum_dq` и убедитесь, что: - новая проверка проходит на «хороших» данных; - при нарушении условия DAG падает с понятной ошибкой. @@ -115,4 +115,3 @@ - Расширить проверки качества данных для потоков bookings → STG → витрины. Когда будете готовы к этим темам, вернитесь к этому разделу — он станет основой для следующего «модуля» лабораторных заданий. - diff --git a/sql/base/orders_ddl.sql b/sql/base/orders_ddl.sql new file mode 100644 index 0000000..c1bfac6 --- /dev/null +++ b/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=zlib, compresslevel=1) +DISTRIBUTED BY (order_id); diff --git a/sql/ddl_gp.sql b/sql/ddl_gp.sql index a33111d..ae63093 100644 --- a/sql/ddl_gp.sql +++ b/sql/ddl_gp.sql @@ -1,28 +1,15 @@ -- Главный входной DDL-скрипт для Greenplum в учебном стенде. -- Выполняется из контейнера командой `make ddl-gp` и создаёт/обновляет --- все объекты, которые нужны базовым DAG (csv_to_greenplum, bookings_to_gp_stage). +-- все объекты, которые нужны базовым DAG (csv_to_greenplum, bookings_to_gp_stage); +-- подключает файловые DDL через \i, чтобы сохранять единый входной скрипт. -- --- Идея такая: --- - здесь описаны только верхнеуровневые объекты (orders, внешняя таблица bookings); --- - более подробный DDL для отдельных слоёв (stg, src и т.п.) лежит в соседних файлах --- в каталоге sql/ и подключается через psql-команду \i; --- - чтобы не ломать задания, новые объекты лучше добавлять в отдельные файлы и --- подключать их отсюда, а существующие определения не удалять. +-- Чтобы не ломать задания, новые объекты лучше добавлять в отдельные файлы +-- и подключать их отсюда, а существующие определения не удалять. -- -- Подробнее про STG/bookings: см. docs/internal/bookings_stg_readme.md. -- Таблица для CSV‑пайплайна (csv_to_greenplum). --- Колонночная таблица (append-optimized) и распределение по ключу. --- Внимание: append-optimized таблицы не поддерживают UNIQUE/PRIMARY KEY, --- поэтому контроль дублей выполняем в DAG при загрузке. -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=zlib, compresslevel=1) -DISTRIBUTED BY (order_id); +\i base/orders_ddl.sql -- Внешняя таблица для чтения данных из демо-БД bookings через PXF (JDBC). -- Источник: таблица bookings.bookings в базе demo (Postgres, сервис bookings-db). diff --git a/tests/test_dags_smoke.py b/tests/test_dags_smoke.py index e1ec388..2d2f727 100644 --- a/tests/test_dags_smoke.py +++ b/tests/test_dags_smoke.py @@ -48,8 +48,8 @@ def test_csv_to_greenplum_dag_structure(): assert t4 in t3.get_direct_relatives("downstream") -def test_data_quality_greenplum_dag_structure(): - dag = _load_dag("airflow.dags.data_quality_greenplum") +def test_csv_to_greenplum_dq_dag_structure(): + dag = _load_dag("airflow.dags.csv_to_greenplum_dq") expected_tasks = { "check_orders_table_exists",