feat(dags): add ddl dag for gp and rename csv dq

This commit is contained in:
2025-12-10 18:41:00 +03:00
parent cf95cf95d2
commit c865fa282e
12 changed files with 111 additions and 45 deletions
+1 -1
View File
@@ -3,7 +3,7 @@
Эта репа — учебный стенд для студентов (менти), которые только начинают с Airflow/Greenplum и Python. Пожалуйста, держите решения простыми, стабильными и хорошо объяснёнными. Эта репа — учебный стенд для студентов (менти), которые только начинают с 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. - `airflow/requirements.txt` — зависимости, которые ставятся внутри контейнеров Airflow.
- `sql/` — DDL и вспомогательные SQL (например, `sql/ddl_gp.sql`). - `sql/` — DDL и вспомогательные SQL (например, `sql/ddl_gp.sql`).
- `docker-compose.yml` — Greenplum, Airflow, Postgres (мета-БД). - `docker-compose.yml` — Greenplum, Airflow, Postgres (мета-БД).
+11 -6
View File
@@ -24,7 +24,7 @@
- Откройте UI: http://localhost:8080 (admin/admin). - Откройте UI: http://localhost:8080 (admin/admin).
- Включите и запустите DAG `csv_to_greenplum`. Дождитесь Success. - Включите и запустите DAG `csv_to_greenplum`. Дождитесь Success.
- Проверьте данные: `make gp-psql``SELECT COUNT(*) FROM public.orders;`. - Проверьте данные: `make gp-psql``SELECT COUNT(*) FROM public.orders;`.
- Дополнительно: запустите `greenplum_data_quality` — все проверки должны быть зелёные. - Дополнительно: запустите `csv_to_greenplum_dq` — все проверки должны быть зелёные.
Если что‑то не работает — смотрите «Типичные проблемы» и «Быстрый reset» ниже. Если что‑то не работает — смотрите «Типичные проблемы» и «Быстрый reset» ниже.
@@ -170,8 +170,13 @@ uv run black --check airflow tests
- **pandas** — библиотека для генерации и анализа данных в формате CSV - **pandas** — библиотека для генерации и анализа данных в формате CSV
### Готовые DAG (workflow) ### Готовые DAG (workflow)
- **orders_base_ddl** — создаёт базовую таблицу `public.orders` для CSV‑пайплайна
- **bookings_stg_ddl** — готовит схему `stg` и таблицы `stg.bookings_ext` / `stg.bookings`
- **csv_to_greenplum** — базовый pipeline: pandas → CSV → Greenplum - **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 ```bash
@@ -179,7 +184,7 @@ uv run black --check airflow tests
make up # Запустить весь стенд make up # Запустить весь стенд
make down # Остановить и удалить данные make down # Остановить и удалить данные
make airflow-init # Инициализировать Airflow make airflow-init # Инициализировать Airflow
make ddl-gp # Применить DDL к Greenplum make ddl-gp # Применить DDL к Greenplum вручную
make gp-psql # Подключиться к Greenplum через psql make gp-psql # Подключиться к Greenplum через psql
make bookings-init # Установить демобазу bookings в Postgres (по умолчанию генерирует 1 день) make bookings-init # Установить демобазу bookings в Postgres (по умолчанию генерирует 1 день)
make bookings-generate-day # Добавить ещё один день данных в bookings (можно вызвать несколько раз) 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/ ├── airflow/
│ └── dags/ # Файлы workflow (DAG) │ └── dags/ # Файлы workflow (DAG)
│ ├── csv_to_greenplum.py │ ├── csv_to_greenplum.py
│ ├── data_quality_greenplum.py │ ├── csv_to_greenplum_dq.py
│ └── bookings_to_gp_stage.py │ └── bookings_to_gp_stage.py
├── bookings/ # Скрипты и файлы для демобазы bookings в Postgres ├── bookings/ # Скрипты и файлы для демобазы bookings в Postgres
├── sql/ ├── sql/
@@ -381,7 +386,7 @@ load_bookings_to_stg = PostgresOperator(
## 💡 Советы для дальнейшего обучения ## 💡 Советы для дальнейшего обучения
1. **Поэкспериментируйте с DAG** — измените параметры генерации данных или размер батча 1. **Поэкспериментируйте с DAG** — измените параметры генерации данных или размер батча
2. **Добавьте свои проверки** — расширьте DAG `data_quality_greenplum.py` 2. **Добавьте свои проверки** — расширьте DAG `csv_to_greenplum_dq.py`
3. **Попробуйте другие источники** — замените генератор данных на чтение из файла или API 3. **Попробуйте другие источники** — замените генератор данных на чтение из файла или API
4. **Изучите Airflow deeper** — добавьте зависимости между задачами, настройте расписания 4. **Изучите Airflow deeper** — добавьте зависимости между задачами, настройте расписания
+3 -3
View File
@@ -31,7 +31,7 @@
- Нажать «Trigger DAG». - Нажать «Trigger DAG».
- Контроль: все таски Success, в `data/` появился CSV, в логах `load_csv_to_greenplum` видно `INSERT`. - Контроль: все таски Success, в `data/` появился CSV, в логах `load_csv_to_greenplum` видно `INSERT`.
- В Greenplum (см. п.5) убедиться в наличии строк `(SELECT COUNT(*) ...)`. - В Greenplum (см. п.5) убедиться в наличии строк `(SELECT COUNT(*) ...)`.
3. DAG `greenplum_data_quality`: 3. DAG `csv_to_greenplum_dq`:
- Запустить вручную после первого DAG. - Запустить вручную после первого DAG.
- Проверить, что все 5 задач Success и логи содержат `Проверка пройдена`. - Проверить, что все 5 задач Success и логи содержат `Проверка пройдена`.
@@ -47,10 +47,10 @@
- Завершить `\q`. - Завершить `\q`.
## 6. Негативные сценарии и fallback ## 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`). - **Проблемы с подключением**: временно изменить `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. - **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 (если «что-то сломалось») ## 7. Быстрый reset (если «что-то сломалось»)
- Перезапустить стенд с очисткой данных: - Перезапустить стенд с очисткой данных:
+30
View File
@@ -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",
)
@@ -3,20 +3,23 @@ from __future__ import annotations
import logging import logging
from datetime import datetime, timedelta from datetime import datetime, timedelta
from airflow import DAG
from airflow.operators.python import PythonOperator from airflow.operators.python import PythonOperator
from helpers.greenplum import (assert_orders_have_rows, from helpers.greenplum import (
assert_orders_have_rows,
assert_orders_no_duplicates, assert_orders_no_duplicates,
assert_orders_schema, assert_orders_schema,
assert_orders_table_exists, get_gp_conn) assert_orders_table_exists,
get_gp_conn,
from airflow import DAG )
def _run_check(check_callable): def _run_check(check_callable):
""" """
Оборачивает проверку качества данных в контекст подключения к Greenplum. Оборачивает проверку качества данных в контекст подключения к Greenplum.
Этот DAG предназначен для автоматической проверки качества данных в таблице orders: Этот DAG предназначен для автоматической проверки качества данных
после CSV-пайплайна в таблице public.orders:
1. Проверяет существование таблицы 1. Проверяет существование таблицы
2. Проверяет соответствие схемы 2. Проверяет соответствие схемы
3. Проверяет наличие данных 3. Проверяет наличие данных
@@ -47,13 +50,13 @@ def _log_dq_summary():
default_args = {"owner": "airflow", "retries": 1, "retry_delay": timedelta(seconds=30)} default_args = {"owner": "airflow", "retries": 1, "retry_delay": timedelta(seconds=30)}
with DAG( with DAG(
dag_id="greenplum_data_quality", dag_id="csv_to_greenplum_dq",
start_date=datetime(2024, 1, 1), start_date=datetime(2024, 1, 1),
schedule=None, schedule=None,
catchup=False, catchup=False,
default_args=default_args, default_args=default_args,
tags=["demo", "greenplum", "quality"], tags=["demo", "greenplum", "quality", "csv", "dq"],
description="Автоматизированные проверки качества данных в Greenplum", description="Проверки качества данных после CSV → public.orders в Greenplum",
) as dag: ) as dag:
# Задача 1: Проверка существования таблицы # Задача 1: Проверка существования таблицы
check_exists = PythonOperator( check_exists = PythonOperator(
+30
View File
@@ -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",
)
+1 -1
View File
@@ -76,7 +76,7 @@ DDL будет добавлен в `sql/ddl_gp.sql` в блоке DDL для Gre
### 4.1. DAG для DDL ### 4.1. DAG для DDL
- `dag_id`: `bookings_stg_ddl` (рабочее имя). - `dag_id`: `bookings_stg_ddl` (реализован в `airflow/dags/bookings_stg_ddl.py`).
- Назначение: один раз (или при изменении схемы) создать необходимые объекты в Greenplum: - Назначение: один раз (или при изменении схемы) создать необходимые объекты в Greenplum:
- схему `stg` (если её ещё нет); - схему `stg` (если её ещё нет);
- внешнюю таблицу `stg.bookings_ext` (PXF → `bookings-db`); - внешнюю таблицу `stg.bookings_ext` (PXF → `bookings-db`);
+2 -1
View File
@@ -21,7 +21,8 @@ _Внутренний файл, чтобы не забыть договорён
- Стенд поднят: `make up`. - Стенд поднят: `make up`.
- Демо‑БД bookings инициализирована: `make bookings-init`. - Демо‑БД bookings инициализирована: `make bookings-init`.
- В Greenplum применён DDL (созданы схема `stg` и таблицы `stg.bookings_ext` / `stg.bookings`): - В 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 есть коннекты: - В Airflow есть коннекты:
- `greenplum_conn` (по умолчанию уже используется в helpers/greenplum.py); - `greenplum_conn` (по умолчанию уже используется в helpers/greenplum.py);
- `bookings_db` (Postgres к сервису `bookings-db`, если не хочется полагаться на ENV). - `bookings_db` (Postgres к сервису `bookings-db`, если не хочется полагаться на ENV).
+2 -3
View File
@@ -40,13 +40,13 @@
### 1.3. Собственные проверки качества данных ### 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. Посмотрите, какие проверки уже реализованы (наличие таблицы, схема, дубликаты). 2. Посмотрите, какие проверки уже реализованы (наличие таблицы, схема, дубликаты).
3. Добавьте ещё одну простую проверку, например: 3. Добавьте ещё одну простую проверку, например:
- проверка, что в таблице `public.orders` не больше N строк; - проверка, что в таблице `public.orders` не больше N строк;
- проверка, что поле (например, `order_price`) не содержит отрицательных значений; - проверка, что поле (например, `order_price`) не содержит отрицательных значений;
- проверка, что нет строк с `NULL` в ключевых колонках. - проверка, что нет строк с `NULL` в ключевых колонках.
4. Запустите DAG `greenplum_data_quality` и убедитесь, что: 4. Запустите DAG `csv_to_greenplum_dq` и убедитесь, что:
- новая проверка проходит на «хороших» данных; - новая проверка проходит на «хороших» данных;
- при нарушении условия DAG падает с понятной ошибкой. - при нарушении условия DAG падает с понятной ошибкой.
@@ -115,4 +115,3 @@
- Расширить проверки качества данных для потоков bookings → STG → витрины. - Расширить проверки качества данных для потоков bookings → STG → витрины.
Когда будете готовы к этим темам, вернитесь к этому разделу — он станет основой для следующего «модуля» лабораторных заданий. Когда будете готовы к этим темам, вернитесь к этому разделу — он станет основой для следующего «модуля» лабораторных заданий.
+11
View File
@@ -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);
+5 -18
View File
@@ -1,28 +1,15 @@
-- Главный входной DDL-скрипт для Greenplum в учебном стенде. -- Главный входной DDL-скрипт для Greenplum в учебном стенде.
-- Выполняется из контейнера командой `make ddl-gp` и создаёт/обновляет -- Выполняется из контейнера командой `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. -- Подробнее про STG/bookings: см. docs/internal/bookings_stg_readme.md.
-- Таблица для CSV‑пайплайна (csv_to_greenplum). -- Таблица для CSV‑пайплайна (csv_to_greenplum).
-- Колонночная таблица (append-optimized) и распределение по ключу. \i base/orders_ddl.sql
-- Внимание: 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);
-- Внешняя таблица для чтения данных из демо-БД bookings через PXF (JDBC). -- Внешняя таблица для чтения данных из демо-БД bookings через PXF (JDBC).
-- Источник: таблица bookings.bookings в базе demo (Postgres, сервис bookings-db). -- Источник: таблица bookings.bookings в базе demo (Postgres, сервис bookings-db).
+2 -2
View File
@@ -48,8 +48,8 @@ def test_csv_to_greenplum_dag_structure():
assert t4 in t3.get_direct_relatives("downstream") assert t4 in t3.get_direct_relatives("downstream")
def test_data_quality_greenplum_dag_structure(): def test_csv_to_greenplum_dq_dag_structure():
dag = _load_dag("airflow.dags.data_quality_greenplum") dag = _load_dag("airflow.dags.csv_to_greenplum_dq")
expected_tasks = { expected_tasks = {
"check_orders_table_exists", "check_orders_table_exists",