Merge branch 'feature/bookings_db'
This commit is contained in:
+16
-1
@@ -1,3 +1,6 @@
|
||||
# Общие настройки
|
||||
TZ=Europe/Moscow
|
||||
|
||||
# Airflow Configuration
|
||||
AIRFLOW_USER=admin
|
||||
AIRFLOW_PASSWORD=admin
|
||||
@@ -8,10 +11,22 @@ PG_USER=airflow
|
||||
PG_PASSWORD=airflow
|
||||
PG_DB=airflow
|
||||
|
||||
# Bookings demo database (Postgres)
|
||||
BOOKINGS_DB_USER=bookings
|
||||
BOOKINGS_DB_PASSWORD=bookings
|
||||
BOOKINGS_DB_NAME=bookings
|
||||
BOOKINGS_DB_PORT=5434
|
||||
# Начальная дата модельного времени для генерации демобазы
|
||||
BOOKINGS_START_DATE=2017-01-01
|
||||
# Количество дней для первой генерации (держим малым, чтобы быстрее увидеть данные)
|
||||
BOOKINGS_INIT_DAYS=1
|
||||
# Количество параллельных джобов генератора (по умолчанию 1; при 1 работаем синхронно без dblink)
|
||||
BOOKINGS_JOBS=1
|
||||
|
||||
# Greenplum Configuration
|
||||
GP_USER=gpadmin
|
||||
GP_PASSWORD=gpadmin
|
||||
GP_DB=gpadmin
|
||||
GP_DB=gp_dwh
|
||||
GP_HOST=greenplum
|
||||
GP_PORT=5432
|
||||
GP_CONN_ID=greenplum_conn
|
||||
|
||||
@@ -7,3 +7,4 @@ __pycache__/
|
||||
*.pyc
|
||||
.venv/
|
||||
data/
|
||||
bookings/demodb/
|
||||
|
||||
@@ -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 (мета-БД).
|
||||
@@ -11,14 +11,16 @@
|
||||
- `.env(.example)` — настройки окружения (реальные секреты не коммитим).
|
||||
|
||||
## Команды (основные)
|
||||
- `make up` — поднять весь стек.
|
||||
- `make airflow-init` — инициализировать мета-БД Airflow и создать пользователя.
|
||||
- `make up` — поднять весь стек (Airflow инициализируется автоматически при первом старте).
|
||||
- `make stop` — остановить контейнеры, не трогая данные.
|
||||
- `make down` — остановить и удалить контейнеры/сети (тома сохраняются).
|
||||
- `make clean` — полный reset: остановить и удалить контейнеры/сети и тома (все данные будут потеряны).
|
||||
- `make airflow-init` — вручную переинициализировать мета-БД Airflow и создать пользователя (обычно не нужно).
|
||||
- `make logs` — логи webserver и scheduler.
|
||||
- `make ddl-gp` — применить DDL к Greenplum.
|
||||
- `make gp-psql` — открыть `psql` в контейнере Greenplum от `gpadmin`.
|
||||
- `make down` — остановить и удалить тома (данные будут потеряны).
|
||||
|
||||
Пример: `make up && make airflow-init`, затем открыть `http://localhost:8080`.
|
||||
Пример: `make up`, затем открыть `http://localhost:8080`.
|
||||
|
||||
## Локальное Python‑окружение
|
||||
- Используем `uv`: достаточно `uv sync` (или `make dev-sync`) — подтянет Python, создаст `.venv`, установит зависимости.
|
||||
@@ -34,6 +36,22 @@
|
||||
- Форматирование: `black` (88 cols) и `isort`. Если не уверены — запустите `make fmt`.
|
||||
- Язык: комментарии, docstring и документацию — на русском; имена идентификаторов — на английском.
|
||||
|
||||
### Структура и нейминг SQL (слои DWH)
|
||||
- В каталоге `sql/` придерживаемся слоёв DWH:
|
||||
- `sql/src/` — скрипты, работающие с исходными системами (например, `bookings_generate_day_if_missing.sql`);
|
||||
- `sql/stg/` — скрипты для стейджинга (`bookings_ddl.sql`, `bookings_load.sql`, `bookings_dq.sql`);
|
||||
- в будущем можно добавить `sql/ods/`, `sql/dds/`, `sql/dm/` по мере роста стенда.
|
||||
- Именование файлов: `{объект}_{роль}.sql`, где:
|
||||
- `объект` — логическое имя сущности (`bookings`, `orders`, и т.п.);
|
||||
- `роль` — `ddl` (создание/изменение объектов), `load` (загрузка/инкремент), `dq` (проверки качества данных) и т.п.
|
||||
- Общие DDL-скрипты (например, `sql/ddl_gp.sql`) могут подключать файловые DDL через `\i`, но сами определения таблиц живут рядом с объектом (`sql/stg/bookings_ddl.sql` и т.п.).
|
||||
|
||||
### Airflow + SQL
|
||||
- В учебных DAG’ах, где основная логика — в SQL, по умолчанию используем `PostgresOperator` + Airflow Connections:
|
||||
- DAG оркестрирует шаги и подключение к БД;
|
||||
- SQL-скрипты лежат в `sql/...` и подключаются по пути (`sql='sql/stg/bookings_load.sql'`).
|
||||
- Сложную ручную работу с подключениями (`psycopg2`, ENV-фоллбеки) используем только там, где реально много Python-логики и это помогает учебной цели.
|
||||
|
||||
## Тестирование
|
||||
- Тесты лежат в `tests/` (pytest). Запуск: `make test`.
|
||||
- Есть юнит‑тесты для `helpers/greenplum.py` и smoke‑тесты DAG‑структуры (`tests/test_dags_smoke.py`).
|
||||
@@ -54,4 +72,4 @@
|
||||
- Пишите простыми словами. Добавляйте короткие комментарии к нетривиальной логике.
|
||||
- Избегайте больших рефакторингов и сложных паттернов — студенты только начинают.
|
||||
- Ошибки и логи — дружелюбные и понятные (лучше с подсказкой «что сделать дальше»).
|
||||
- Перед релевантными правками валидируйте локально: `make up && make airflow-init`, затем откройте DAG в UI и/или прогоните `make test`.
|
||||
- Перед релевантными правками валидируйте локально: `make up`, затем откройте DAG в UI и/или прогоните `make test`.
|
||||
|
||||
@@ -1,15 +1,27 @@
|
||||
SHELL := /bin/bash
|
||||
UV := uv
|
||||
PYTHON_VERSION := 3.11
|
||||
DEMODB_REPO := https://github.com/postgrespro/demodb.git
|
||||
DEMODB_COMMIT := d68de192850237719f09b47688d5f3fc94653ca6
|
||||
BOOKINGS_JOBS ?= 1
|
||||
BOOKINGS_START_DATE ?= 2017-01-01
|
||||
BOOKINGS_INIT_DAYS ?= 1
|
||||
|
||||
.PHONY: up down airflow-init logs gp-psql ddl-gp dev-setup dev-sync dev-lock test lint fmt clean-venv
|
||||
.PHONY: up stop down clean airflow-init logs gp-psql ddl-gp \
|
||||
bookings-clone-demodb bookings-init bookings-psql bookings-generate-day \
|
||||
dev-setup dev-sync dev-lock test lint fmt clean-venv
|
||||
|
||||
up:
|
||||
docker compose -f docker-compose.yml up -d
|
||||
|
||||
stop:
|
||||
docker compose -f docker-compose.yml stop
|
||||
|
||||
down:
|
||||
docker compose -f docker-compose.yml down -v
|
||||
|
||||
clean: down
|
||||
|
||||
airflow-init:
|
||||
docker compose -f docker-compose.yml run --rm airflow-init
|
||||
|
||||
@@ -17,10 +29,65 @@ logs:
|
||||
docker compose -f docker-compose.yml logs -f airflow-webserver airflow-scheduler
|
||||
|
||||
gp-psql:
|
||||
docker compose -f docker-compose.yml exec greenplum bash -c "su - gpadmin -c '/usr/local/greenplum-db/bin/psql -p 5432 -d gpadmin'"
|
||||
docker compose -f docker-compose.yml exec greenplum bash -c "su - gpadmin -c '/usr/local/greenplum-db/bin/psql -p 5432 -d gp_dwh'"
|
||||
|
||||
ddl-gp:
|
||||
docker compose -f docker-compose.yml exec greenplum bash -c "su - gpadmin -c '/usr/local/greenplum-db/bin/psql -d gpadmin -f /sql/ddl_gp.sql'"
|
||||
docker compose -f docker-compose.yml exec greenplum bash -c "su - gpadmin -c 'cd /sql && /usr/local/greenplum-db/bin/psql -d gp_dwh -f ddl_gp.sql'"
|
||||
|
||||
bookings-clone-demodb:
|
||||
mkdir -p bookings
|
||||
if [ ! -d bookings/demodb ]; then \
|
||||
git clone --depth 1 $(DEMODB_REPO) bookings/demodb; \
|
||||
git -C bookings/demodb fetch --depth 1 origin $(DEMODB_COMMIT); \
|
||||
git -C bookings/demodb checkout $(DEMODB_COMMIT); \
|
||||
fi
|
||||
# Патчим generate/continue: при jobs=1 запускаем process_queue синхронно, без dblink
|
||||
if ! grep -q "CALL process_queue(end_date);" bookings/demodb/engine.sql; then \
|
||||
patch -d bookings/demodb -p1 --forward < bookings/patches/engine_jobs1_sync.patch || true; \
|
||||
fi
|
||||
# Делаем установку идемпотентной: DROP DATABASE IF EXISTS demo
|
||||
if ! grep -q "DROP DATABASE IF EXISTS demo;" bookings/demodb/install.sql; then \
|
||||
patch -d bookings/demodb -p1 --forward < bookings/patches/install_drop_if_exists.patch || true; \
|
||||
fi
|
||||
|
||||
bookings-init: bookings-clone-demodb
|
||||
docker compose -f docker-compose.yml up -d bookings-db
|
||||
docker compose -f docker-compose.yml exec bookings-db bash -lc '\
|
||||
until PGPASSWORD="$$POSTGRES_PASSWORD" pg_isready -U "$$POSTGRES_USER" -d "$$POSTGRES_DB" -h localhost; do \
|
||||
echo "Waiting for bookings-db to become ready..."; \
|
||||
sleep 1; \
|
||||
done \
|
||||
'
|
||||
docker compose -f docker-compose.yml exec bookings-db bash -lc 'cd /bookings/demodb && PGPASSWORD="$$POSTGRES_PASSWORD" psql -v ON_ERROR_STOP=1 -U "$$POSTGRES_USER" -d "$$POSTGRES_DB" -f install.sql'
|
||||
docker compose -f docker-compose.yml exec bookings-db bash -lc '\
|
||||
CONNSTR="dbname=demo user=$$POSTGRES_USER password=$$POSTGRES_PASSWORD"; \
|
||||
PGPASSWORD="$$POSTGRES_PASSWORD" psql -v ON_ERROR_STOP=1 -U "$$POSTGRES_USER" -d "$$POSTGRES_DB" -c "ALTER DATABASE demo SET gen.connstr='\''$$CONNSTR'\'';" \
|
||||
'
|
||||
docker compose -f docker-compose.yml exec bookings-db bash -lc '\
|
||||
PGPASSWORD="$$POSTGRES_PASSWORD" psql -v ON_ERROR_STOP=1 -U "$$POSTGRES_USER" -d demo -c "\
|
||||
ALTER DATABASE demo SET bookings.start_date = '\''$(BOOKINGS_START_DATE)'\''; \
|
||||
ALTER DATABASE demo SET bookings.init_days = '\''$(BOOKINGS_INIT_DAYS)'\''; \
|
||||
ALTER DATABASE demo SET bookings.jobs = '\''$(BOOKINGS_JOBS)'\'';" \
|
||||
'
|
||||
# Генерируем первый день данных, чтобы база не оставалась пустой
|
||||
docker compose -f docker-compose.yml exec bookings-db bash -lc '\
|
||||
PGPASSWORD="$$POSTGRES_PASSWORD" psql -v ON_ERROR_STOP=1 -U "$$POSTGRES_USER" -d demo -f /bookings/generate_next_day.sql \
|
||||
'
|
||||
|
||||
bookings-psql:
|
||||
docker compose -f docker-compose.yml exec bookings-db bash -lc 'PGPASSWORD="$$POSTGRES_PASSWORD" psql -U "$$POSTGRES_USER" -d demo'
|
||||
|
||||
bookings-generate-day:
|
||||
docker compose -f docker-compose.yml up -d bookings-db
|
||||
docker compose -f docker-compose.yml exec bookings-db bash -lc '\
|
||||
until PGPASSWORD="$$POSTGRES_PASSWORD" pg_isready -U "$$POSTGRES_USER" -d "$$POSTGRES_DB" -h localhost; do \
|
||||
echo "Waiting for bookings-db to become ready..."; \
|
||||
sleep 1; \
|
||||
done \
|
||||
'
|
||||
docker compose -f docker-compose.yml exec bookings-db bash -lc '\
|
||||
PGPASSWORD="$$POSTGRES_PASSWORD" psql -v ON_ERROR_STOP=1 -U "$$POSTGRES_USER" -d demo -f /bookings/generate_next_day.sql \
|
||||
'
|
||||
|
||||
dev-setup:
|
||||
$(UV) python install $(PYTHON_VERSION)
|
||||
|
||||
@@ -2,7 +2,9 @@
|
||||
|
||||
Учебный стенд: ETL из Postgres в Greenplum с оркестрацией в Airflow.
|
||||
|
||||
Добро пожаловать в учебный стенд для изучения основ Data Engineering! Этот проект поможет вам освоить ключевые инструменты современных data pipeline: **Airflow** для оркестрации, **pandas/CSV** для подготовки данных и **Greenplum** как аналитическую базу данных.
|
||||
Добро пожаловать в учебный стенд для изучения основ Data Engineering! Этот проект поможет вам освоить ключевые инструменты современных data pipeline: **Airflow** для оркестрации, **pandas/CSV** для подготовки данных, **Postgres** с демобазой **bookings** как источник и **Greenplum** как аналитическую базу данных.
|
||||
|
||||
Если вы проходите стенд как серию лабораторных, смотрите также файл с заданиями: `educational-tasks.md`.
|
||||
|
||||
## 🎯 Что вы узнаете
|
||||
|
||||
@@ -12,16 +14,17 @@
|
||||
- Как загружать данные в Greenplum пакетами и избегать дублей
|
||||
- Как проверять качество данных в автоматизированных pipeline
|
||||
- Основы проектирования ETL/ELT процессов
|
||||
- Как работать с демо-БД bookings в Postgres как источником для будущего DWH
|
||||
|
||||
## 👩🎓 Для студентов (10‑минутный чек‑лист)
|
||||
|
||||
- Установите Docker Desktop и Git.
|
||||
- Скопируйте настройки: `cp .env.example .env`.
|
||||
- Поднимите стенд: `docker compose up -d` и инициализируйте Airflow: `docker compose run --rm airflow-init`.
|
||||
- Поднимите стенд: `make up` (или `docker compose up -d` — сервис `airflow-init` запустится автоматически при первом старте).
|
||||
- Откройте 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» ниже.
|
||||
|
||||
@@ -62,7 +65,7 @@ docker compose up -d
|
||||
**Проверка вручную:**
|
||||
```bash
|
||||
# Подключитесь к Greenplum и проверьте данные
|
||||
docker compose exec greenplum bash -c "su - gpadmin -c 'psql -p 5432 -d gpadmin'"
|
||||
docker compose exec greenplum bash -c "su - gpadmin -c 'psql -p 5432 -d gp_dwh'"
|
||||
|
||||
# Внутри psql выполните:
|
||||
\dt # Показать таблицы
|
||||
@@ -89,6 +92,8 @@ make up
|
||||
|
||||
## 🛠️ Подробная настройка (для уверенных пользователей)
|
||||
|
||||
> Если вы впервые запускаете стенд, этот раздел можно пролистать и вернуться к нему позже.
|
||||
|
||||
### Установка Make (опционально)
|
||||
|
||||
Для удобства работы с проектом рекомендуем установить `make`:
|
||||
@@ -102,21 +107,30 @@ make up
|
||||
|
||||
С `make` команды становятся короче:
|
||||
```bash
|
||||
make up && make airflow-init # Запуск стека
|
||||
make up # Запуск стека (включая airflow-init при первом старте)
|
||||
make logs # Просмотр логов
|
||||
make gp-psql # Подключение к Greenplum
|
||||
```
|
||||
|
||||
### Настройка подключения к Greenplum в Airflow
|
||||
|
||||
По умолчанию DAG использует переменные окружения, но вы можете создать Airflow Connection:
|
||||
По умолчанию готовые DAG используют Airflow Connections. В docker-compose они
|
||||
заводятся автоматически через переменные окружения:
|
||||
|
||||
- `AIRFLOW_CONN_GREENPLUM_CONN` — подключение к Greenplum с `conn_id=greenplum_conn`;
|
||||
- `AIRFLOW_CONN_BOOKINGS_DB` — подключение к демо-БД bookings с `conn_id=bookings_db`.
|
||||
|
||||
Такие подключения подхватываются из окружения и могут не отображаться в UI,
|
||||
но для DAG это нормально — `PostgresOperator` найдёт их по `conn_id`.
|
||||
|
||||
При желании вы можете создать или отредактировать подключение вручную в UI:
|
||||
|
||||
1. Airflow UI → **Admin → Connections → Add a new record**
|
||||
2. Заполните поля:
|
||||
- **Conn Id:** `greenplum_conn`
|
||||
- **Conn Type:** `Postgres`
|
||||
- **Host:** `greenplum`
|
||||
- **Schema:** `gpadmin`
|
||||
- **Schema:** `gp_dwh`
|
||||
- **Login:** `gpadmin`
|
||||
- **Password:** `gpadmin`
|
||||
- **Port:** `5432`
|
||||
@@ -161,25 +175,56 @@ uv run black --check airflow tests
|
||||
- **Greenplum** — аналитическая база данных для хранения и анализа данных
|
||||
- **Airflow** — оркестратор workflow и задач
|
||||
- **Postgres** — база метаданных для Airflow
|
||||
- **Postgres (bookings)** — отдельная демо-БД bookings (источник данных для будущего DWH в Greenplum)
|
||||
- **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
|
||||
# Основные команды
|
||||
make up # Запустить весь стенд
|
||||
make down # Остановить и удалить данные
|
||||
make airflow-init # Инициализировать Airflow
|
||||
make ddl-gp # Применить DDL к Greenplum
|
||||
make up # Запустить весь стенд (Airflow инициализируется автоматически при первом старте)
|
||||
make stop # Остановить контейнеры, не трогая данные
|
||||
make down # Остановить и удалить контейнеры/сети (volumes сохраняются)
|
||||
make clean # Полный reset: остановить и удалить контейнеры/сети и тома (данные будут потеряны)
|
||||
make airflow-init # Ручной запуск инициализации Airflow (обычно не нужен)
|
||||
make ddl-gp # Применить DDL к Greenplum вручную
|
||||
make gp-psql # Подключиться к Greenplum через psql
|
||||
make bookings-init # Установить демобазу bookings в Postgres (по умолчанию генерирует 1 день)
|
||||
make bookings-generate-day # Добавить ещё один день данных в bookings (можно вызвать несколько раз)
|
||||
make bookings-psql # Подключиться к демобазе bookings (БД demo)
|
||||
|
||||
# Проверка данных
|
||||
make logs # Следить за логами Airflow
|
||||
|
||||
# Контроль генерации bookings
|
||||
docker compose -f docker-compose.yml exec bookings-db bash -lc 'PGPASSWORD="$POSTGRES_PASSWORD" psql -U "$POSTGRES_USER" -d demo -c "SELECT busy();"'
|
||||
# busy() = t — генерация ещё идёт; f — завершена. При необходимости можно вызвать CALL abort(); и запустить генерацию заново.
|
||||
```
|
||||
|
||||
### Генерация следующего дня в bookings
|
||||
|
||||
- При `make bookings-init` автоматически генерируется `BOOKINGS_INIT_DAYS` суток, начиная с даты `BOOKINGS_START_DATE` (по умолчанию один день с `2017-01-01`).
|
||||
- Дальше каждый вызов `make bookings-generate-day` или генерации через DAG добавляет ровно **один** следующий день после `max(book_date)` в `bookings.bookings` — генератор сам смотрит последнюю дату.
|
||||
- Рекомендуемый учебный сценарий для DAG `bookings_to_gp_stage`: запускать DAG по одному дню вперёд, выбирая в форме Trigger логическую дату `Execution Date (ds)`, совпадающую с тем днём, который вы хотите загрузить (например, `2017-01-01`, затем `2017-01-02` и т.д.).
|
||||
|
||||
- Быстрее всего: `make bookings-generate-day` — читает GUC и сам вызывает `continue`.
|
||||
- Вручную из psql/DBeaver:
|
||||
```sql
|
||||
CALL continue(
|
||||
(SELECT date_trunc('day', max(book_date)) + interval '1 day' FROM bookings.bookings)
|
||||
);
|
||||
-- или с параллельностью: CALL continue((SELECT ...), 4);
|
||||
```
|
||||
- Не вызывайте `CALL generate(...)` поверх существующих данных: она делает TRUNCATE и создаёт демобазу заново.
|
||||
|
||||
---
|
||||
|
||||
## ⚙️ Настройка через переменные окружения
|
||||
@@ -189,8 +234,17 @@ make logs # Следить за логами Airflow
|
||||
### Greenplum
|
||||
- `GP_USER` — пользователь (по умолчанию: gpadmin)
|
||||
- `GP_PASSWORD` — пароль (по умолчанию: gpadmin)
|
||||
- `GP_DB` — база данных (по умолчанию: gpadmin)
|
||||
- `GP_PORT` — порт (по умолчанию: 5432)
|
||||
- `GP_DB` — база данных (по умолчанию: gp_dwh)
|
||||
- `GP_PORT` — порт Greenplum внутри Docker-сети (по умолчанию: 5432, менять обычно не нужно; внешний порт на хосте для подключения клиентов — 5435).
|
||||
|
||||
### Демо-БД bookings (Postgres)
|
||||
- `BOOKINGS_DB_USER` — пользователь Postgres для демобазы (по умолчанию: bookings)
|
||||
- `BOOKINGS_DB_PASSWORD` — пароль пользователя (по умолчанию: bookings)
|
||||
- `BOOKINGS_DB_NAME` — база данных, из которой запускается установка генератора (по умолчанию: bookings)
|
||||
- `BOOKINGS_DB_PORT` — внешний порт для подключения к контейнеру bookings-db (по умолчанию: 5434)
|
||||
- `BOOKINGS_START_DATE` — начальная дата модельного времени (по умолчанию: 2017-01-01)
|
||||
- `BOOKINGS_INIT_DAYS` — сколько дней сгенерировать при первой инициализации (по умолчанию: 1, чтобы увидеть данные без долгого ожидания; можно менять при вызове `BOOKINGS_INIT_DAYS=... make bookings-init`)
|
||||
- `BOOKINGS_JOBS` — число параллельных джобов генератора bookings (по умолчанию: 1; при 1 генерация идёт синхронно без dblink)
|
||||
|
||||
### CSV pipeline
|
||||
- `CSV_DIR` — путь к каталогу с CSV внутри контейнеров Airflow (по умолчанию: `/opt/airflow/data`)
|
||||
@@ -203,6 +257,8 @@ make logs # Следить за логами Airflow
|
||||
|
||||
## 🔍 Продвинутые темы
|
||||
|
||||
> Этот раздел не обязателен при первом прохождении стенда; к нему удобно вернуться, когда базовый CSV‑pipeline уже понятен.
|
||||
|
||||
### Архитектура pipeline
|
||||
|
||||
**Поток данных в DAG `csv_to_greenplum`:**
|
||||
@@ -215,12 +271,38 @@ make logs # Следить за логами Airflow
|
||||
|
||||
### Проверка качества данных
|
||||
|
||||
Запустите DAG `greenplum_data_quality` для автоматической проверки:
|
||||
Запустите DAG `csv_to_greenplum_dq` для автоматической проверки:
|
||||
- Наличие таблицы в базе
|
||||
- Соответствие схемы ожидаемой структуре
|
||||
- Объем загруженных данных
|
||||
- Отсутствие дубликатов записей
|
||||
|
||||
### Пример DAG с SQL-скриптами (bookings → stg)
|
||||
|
||||
> Если вы ещё не дошли до части про bookings и слои DWH, этот подраздел можно пропустить на первом чтении.
|
||||
|
||||
В репозитории есть учебный DAG `bookings_to_gp_stage`, который показывает «канонический» способ работы с SQL в Airflow:
|
||||
|
||||
- подключение к БД через Airflow Connections (`bookings_db`, `greenplum_conn`);
|
||||
- бизнес-логика инкрементальной загрузки и DQ вынесена в SQL-файлы в каталоге `sql/`:
|
||||
- `sql/src/bookings_generate_day_if_missing.sql` — генерация следующего учебного дня в демо-БД bookings (или нескольких стартовых дней, если база пуста);
|
||||
- `sql/stg/bookings_ddl.sql` — DDL для схемы `stg` и таблиц `stg.bookings_ext` / `stg.bookings`;
|
||||
- `sql/stg/bookings_load.sql` — загрузка инкремента из `stg.bookings_ext` в `stg.bookings` на основе «хвоста» после предыдущих батчей;
|
||||
- `sql/stg/bookings_dq.sql` — проверка количества строк между источником и stg за то же окно.
|
||||
|
||||
Фрагмент DAG:
|
||||
|
||||
```python
|
||||
load_bookings_to_stg = PostgresOperator(
|
||||
task_id="load_bookings_to_stg",
|
||||
postgres_conn_id="greenplum_conn",
|
||||
sql="stg/bookings_load.sql",
|
||||
params={"batch_id": "{{ ds_nodash }}"},
|
||||
)
|
||||
```
|
||||
|
||||
Такой подход помогает держать оркестрацию (DAG) и SQL-логику в отдельных файлах и легче сравнивать её с теорией из статьи про моделирование DWH.
|
||||
|
||||
### Ограничения учебного стенда
|
||||
|
||||
- **Greenplum** запущен в single-node режиме (для обучения)
|
||||
@@ -229,15 +311,71 @@ make logs # Следить за логами Airflow
|
||||
|
||||
---
|
||||
|
||||
## 🧩 Подключение к базам через DBeaver
|
||||
|
||||
> Необязательный раздел: нужен только если вы хотите смотреть данные через DBeaver. Для базовых заданий достаточно `make gp-psql`.
|
||||
|
||||
Ниже — краткая инструкция, как подключиться к Greenplum и демо-БД bookings из DBeaver. Перед этим убедитесь, что стенд запущен:
|
||||
|
||||
- `cp .env.example .env` (если ещё не делали)
|
||||
- `make up`
|
||||
- `make bookings-init`
|
||||
- `make ddl-gp`
|
||||
|
||||
### Greenplum (аналитическая БД)
|
||||
|
||||
1. Откройте DBeaver → **New Database Connection**.
|
||||
2. Выберите драйвер **PostgreSQL** (или **Greenplum**, если он есть в вашей версии DBeaver).
|
||||
3. На вкладке **Main** заполните поля (по умолчанию):
|
||||
- `Host`: `localhost`
|
||||
- `Port`: `5435` (внешний порт Greenplum на хосте)
|
||||
- `Database`: значение `GP_DB` (по умолчанию `gp_dwh`)
|
||||
- `Username`: значение `GP_USER` (по умолчанию `gpadmin`)
|
||||
- `Password`: значение `GP_PASSWORD` (по умолчанию `gpadmin`)
|
||||
4. Нажмите **Test Connection** → **OK**, затем **Finish**.
|
||||
|
||||
После подключения:
|
||||
|
||||
- Основные таблицы лаба — в схеме `public` базы `gp_dwh` (например, `public.orders`).
|
||||
- После настройки PXF и выполнения `make ddl-gp` станет доступна внешняя таблица `public.ext_bookings_bookings` — чтение из демо-БД bookings через PXF.
|
||||
|
||||
### bookings-db (демо-БД источника)
|
||||
|
||||
Для работы с исходными данными (демо-БД `demo`) достаточно стандартного PostgreSQL-подключения.
|
||||
|
||||
1. Откройте DBeaver → **New Database Connection** → драйвер **PostgreSQL**.
|
||||
2. На вкладке **Main** укажите (значения по умолчанию из `.env.example`):
|
||||
- `Host`: `localhost`
|
||||
- `Port`: значение `BOOKINGS_DB_PORT` из `.env` (по умолчанию `5434`)
|
||||
- `Database`: `demo`
|
||||
- `Username`: `BOOKINGS_DB_USER` (по умолчанию `bookings`)
|
||||
- `Password`: `BOOKINGS_DB_PASSWORD` (по умолчанию `bookings`)
|
||||
3. Нажмите **Test Connection** → **OK**, затем **Finish**.
|
||||
|
||||
После подключения:
|
||||
|
||||
- Основные таблицы находятся в схеме `bookings` базы `demo` (например, `bookings.bookings`, `bookings.tickets`, `bookings.flights` и т.д.).
|
||||
- Можно сравнивать данные:
|
||||
- между `bookings.bookings` в Postgres и `public.ext_bookings_bookings` в Greenplum;
|
||||
- между временем (`book_date`) в UTC в `demo` и локальным временем в Greenplum (учитывая `TZ=Europe/Moscow`).
|
||||
|
||||
> Если вы меняли порты или креды в `.env`, не забудьте подставить те же значения в настройках соединений в DBeaver.
|
||||
|
||||
---
|
||||
|
||||
## 🆘 Типичные проблемы и решения
|
||||
|
||||
| Проблема | Решение |
|
||||
|----------|---------|
|
||||
| Airflow UI не открывается | Дождитесь сообщения `Listening at: http://0.0.0.0:8080` в логах (`make logs`) |
|
||||
| Ошибка подключения к Greenplum | Убедитесь, что контейнер `greenplum` стал статусом `healthy` (проверьте `docker compose ps`) |
|
||||
| Не открывается порт 8080/5433/5434/5435 | Проверьте, что эти порты не заняты локальными сервисами; при необходимости остановите их или измените порты в `.env`/`docker-compose.yml` |
|
||||
| Нет файла в `./data` после запуска DAG | Проверьте логи задачи `generate_csv`, убедитесь, что `CSV_DIR` смонтирован в docker-compose |
|
||||
| Команда `make` не найдена | Используйте полные команды `docker compose` или установите make |
|
||||
| Greenplum не стартует/падает при старте | Выполните `make down`, затем `make up && make airflow-init` (очищает тома и поднимает заново) |
|
||||
| Greenplum не стартует/падает при старте | Выполните `make down`, затем `make up` (очищает тома и поднимает заново, включая авто‑инициализацию Airflow) |
|
||||
| DAG `bookings_to_gp_stage` падает на внешней таблице/подключении к bookings | Убедитесь, что запущен контейнер `bookings-db` (`docker compose ps`, при необходимости `docker compose start bookings-db`), и выполнены `make bookings-init` и `make ddl-gp` или DAG `bookings_stg_ddl` |
|
||||
| DAG `bookings_to_gp_stage` ругается на отсутствующие таблицы stg | Запустите DAG `bookings_stg_ddl` (или выполните `make ddl-gp`), затем повторите запуск |
|
||||
| DAG не видит Greenplum/DEMObase по Airflow Connections | Убедитесь, что контейнеры `greenplum` и `bookings-db` запущены (`docker compose ps`). Подключения `greenplum_conn` и `bookings_db` задаются через переменные окружения `AIRFLOW_CONN_...` и могут не отображаться в UI, но `PostgresOperator` всё равно найдёт их по `conn_id`. При необходимости вы можете создать/отредактировать их вручную в разделе Connections. |
|
||||
|
||||
---
|
||||
|
||||
@@ -247,12 +385,20 @@ make logs # Следить за логами Airflow
|
||||
├── docker-compose.yml # Описание всех сервисов
|
||||
├── .env.example # Шаблон настроек
|
||||
├── Makefile # Удобные команды для работы
|
||||
├── README.md # Обзор стенда
|
||||
├── TESTING.md # Пошаговый план проверки
|
||||
├── educational-tasks.md # Учебные задания для менти
|
||||
├── airflow/
|
||||
│ └── dags/ # Файлы workflow (DAG)
|
||||
│ ├── csv_to_greenplum.py
|
||||
│ └── data_quality_greenplum.py
|
||||
└── sql/
|
||||
└── ddl_gp.sql # Создание таблицы в Greenplum
|
||||
│ ├── csv_to_greenplum_dq.py
|
||||
│ └── bookings_to_gp_stage.py
|
||||
├── bookings/ # Скрипты и файлы для демобазы bookings в Postgres
|
||||
├── sql/
|
||||
│ └── ddl_gp.sql # Общий DDL для Greenplum (подключает stg/src-скрипты)
|
||||
├── docs/ # Дополнительные документы (архитектура, bookings/STG, PXF)
|
||||
├── tests/ # Автоматические тесты (pytest)
|
||||
└── pxf/ # Конфигурация и файлы для PXF
|
||||
```
|
||||
|
||||
---
|
||||
@@ -260,7 +406,7 @@ make logs # Следить за логами Airflow
|
||||
## 💡 Советы для дальнейшего обучения
|
||||
|
||||
1. **Поэкспериментируйте с DAG** — измените параметры генерации данных или размер батча
|
||||
2. **Добавьте свои проверки** — расширьте DAG `data_quality_greenplum.py`
|
||||
2. **Добавьте свои проверки** — расширьте DAG `csv_to_greenplum_dq.py`
|
||||
3. **Попробуйте другие источники** — замените генератор данных на чтение из файла или API
|
||||
4. **Изучите Airflow deeper** — добавьте зависимости между задачами, настройте расписания
|
||||
|
||||
|
||||
+28
-14
@@ -12,29 +12,40 @@
|
||||
## 2. Локальные автоматические проверки (без Docker)
|
||||
- `make test` — короткие unit-тесты (`tests/test_greenplum_helpers.py`, `tests/test_dags_smoke.py`).
|
||||
- Smoke-тесты DAG автоматически `skip`, если Airflow не установлен в venv, поэтому прогонится за миллисекунды.
|
||||
- `make lint` — black/isort в режиме проверки. Сейчас упадёт из‑за форматирования DAG-файлов.
|
||||
- `make fmt` — автоисправление форматирования; после этого `make lint` должен пройти.
|
||||
- `make lint` — black/isort в режиме проверки (после `make fmt` должен проходить без ошибок).
|
||||
- `make fmt` — автоисправление форматирования; полезно запускать перед пушем.
|
||||
- (опционально) `uv run pytest -q -k dags_smoke` — только DAG smoke.
|
||||
|
||||
## 3. Подготовка Docker-стенда
|
||||
- `cp .env.example .env` (если файла ещё нет) и проверьте переменные:
|
||||
- `GP_PORT` не конфликтует с локальным PostgreSQL.
|
||||
- `GP_PORT` — внутренний порт Greenplum в Docker-сети (по умолчанию 5432, менять не нужно); внешний порт для подключения с хоста фиксирован на `5435`, поэтому локальный PostgreSQL на 5432 не помешает.
|
||||
- `GP_USE_AIRFLOW_CONN=true` при желании использовать Airflow Connection; `false` — fallback на ENV.
|
||||
- `make up` — поднимаем все сервисы. Важно дождаться статуса `healthy` у `pgmeta` и `greenplum` (`docker compose ps`).
|
||||
- `make airflow-init` — миграции мета-БД и создание пользователя Airflow; занимает ~1–2 минуты.
|
||||
- `make up` — поднимаем все сервисы. Важно дождаться статуса `healthy` у `pgmeta` и `greenplum` (`docker compose ps`); инициализация Airflow (`airflow-init`) произойдёт автоматически при первом старте.
|
||||
- `make logs` — следим, пока webserver и scheduler не перейдут в рабочее состояние (`Listening at: http://0.0.0.0:8080`).
|
||||
|
||||
## 4. Smoke тесты DAG в Airflow UI
|
||||
1. Открыть http://localhost:8080 (admin/admin).
|
||||
2. DAG `csv_to_greenplum`:
|
||||
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(*) ...)`.
|
||||
3. DAG `greenplum_data_quality`:
|
||||
4. DAG `csv_to_greenplum_dq`:
|
||||
- Запустить вручную после первого DAG.
|
||||
- Проверить, что все 5 задач Success и логи содержат `Проверка пройдена`.
|
||||
|
||||
- DAG `bookings_to_gp_stage` (полная проверка цепочки bookings → Greenplum STG):
|
||||
- предварительно выполнить один раз: `make bookings-init` (инициализация демо‑БД bookings) и `make ddl-gp` (создаёт `stg.bookings_ext` и `stg.bookings` в Greenplum);
|
||||
- включить DAG `bookings_to_gp_stage` и запустить `Trigger DAG`;
|
||||
- убедиться, что все задачи (`generate_bookings_day`, `load_bookings_to_stg`, `check_row_counts`, `finish_summary`) завершились со статусом Success;
|
||||
- при желании проверить данные: в `bookings-db` появился новый день, а в Greenplum в `stg.bookings` — строки с актуальным `batch_id` (см. пример запросов в разделе 5).
|
||||
|
||||
- (опционально, для менторов/разработчиков) Smoke-тест DAG через Airflow CLI без UI:
|
||||
- `docker compose -f docker-compose.yml exec gp_airflow_web airflow dags test bookings_to_gp_stage 2024-01-01` — прогоняет `bookings_to_gp_stage` целиком в «off-line» режиме;
|
||||
- `docker compose -f docker-compose.yml exec gp_airflow_web airflow dags trigger bookings_to_gp_stage` — создаёт реальный запуск DAG (логи и статус можно смотреть либо через UI, либо командой `airflow tasks list`/`airflow tasks logs` внутри контейнера).
|
||||
|
||||
## 5. Проверка данных в Greenplum
|
||||
- `make gp-psql` — запустить psql в контейнере от имени `gpadmin`.
|
||||
- Команды внутри psql:
|
||||
@@ -42,18 +53,21 @@
|
||||
- `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
|
||||
- **Пустая таблица**: запустить `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` не найдёт дублей.
|
||||
- **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 (если «что-то сломалось»)
|
||||
- Перезапустить стенд с очисткой данных:
|
||||
- `make down` — остановит контейнеры и удалит тома.
|
||||
- `make up && make airflow-init` — заново поднимет всё и проинициализирует Airflow.
|
||||
- Перезапустить стенд:
|
||||
- Мягкий вариант (сохранить данные): `make stop`, затем `make up`.
|
||||
- Полный reset (очистить данные в Docker-томах): `make clean`, затем `make up` (Greenplum/Airflow/bookings будут подняты и инициализированы с нуля).
|
||||
- Иногда Greenplum не стартует после «грязных» остановок (из‑за старых внутренних файлов). Лечение: всегда делайте `make down` перед повторным `make up`.
|
||||
|
||||
## 8. Снятие метрик и мониторинг
|
||||
@@ -67,5 +81,5 @@
|
||||
|
||||
## Текущий статус (пример успешного прогона)
|
||||
- `uv run pytest -q` — 11 passed, 2 smoke-теста DAG пропущены (Airflow не установлен в venv).
|
||||
- `make lint` — падает, потому что `airflow/dags/*.py` не отформатированы black/isort. После `make fmt` проблема уйдёт.
|
||||
- `make lint` — проходит (DAG‑файлы отформатированы black/isort).
|
||||
- Docker-стенд не запускался в рамках этой сессии; ожидается, что инструкции выше обеспечат полноценную проверку.
|
||||
|
||||
@@ -0,0 +1,26 @@
|
||||
# TODO (maintainers / mentors)
|
||||
|
||||
Этот файл собирает идеи по доработке стенда, которые не критичны для текущих задач менти,
|
||||
но улучшат стабильность и удобство сопровождения.
|
||||
|
||||
- Собрать свой образ Airflow поверх `apache/airflow:2.9.2`:
|
||||
- вынести установку Python‑зависимостей из runtime (`pip install ...` при старте контейнеров)
|
||||
в отдельный `Dockerfile`;
|
||||
- переключить `docker-compose.yml` на использование этого образа для `airflow-webserver`,
|
||||
`airflow-scheduler` и `airflow-init`;
|
||||
- обновить документацию (README/TESTING) под новую схему сборки.
|
||||
|
||||
- Переключить Airflow с `SequentialExecutor` (SequentialScheduler) на `LocalExecutor`
|
||||
для docker‑стенда:
|
||||
- проверить, какие параметры достаточно поменять в env/конфиге (`AIRFLOW__CORE__EXECUTOR`)
|
||||
для образа `apache/airflow:2.9.2`;
|
||||
- убедиться, что примерные DAG’и (`csv_to_greenplum`, `bookings_to_gp_stage`) ведут себя
|
||||
предсказуемо в режиме параллельного исполнения;
|
||||
- при необходимости скорректировать тесты и документацию (README/TESTING) с учётом нового executor’а.
|
||||
|
||||
- Разобрать и стабилизировать интеграцию с Greenplum/PXF:
|
||||
- убедиться, что PXF в контейнере `greenplum` всегда корректно инициализируется
|
||||
(нет ошибок вида `protocol "pxf" does not exist` при первом запуске `make ddl-gp`);
|
||||
- при необходимости доработать init‑скрипты в `pxf/init/` и/или документацию,
|
||||
чтобы порядок действий для ментей был однозначным и воспроизводимым;
|
||||
- добавить краткий раздел в README/TESTING о типичных ошибках PXF/Greenplum и шагах по их устранению.
|
||||
@@ -0,0 +1,32 @@
|
||||
from __future__ import annotations
|
||||
|
||||
"""
|
||||
Учебный DAG: создаёт схему stg и таблицы bookings_ext/bookings в Greenplum.
|
||||
Запускается вручную перед DAG загрузки bookings_to_gp_stage или после изменения DDL.
|
||||
"""
|
||||
|
||||
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="bookings_stg_ddl",
|
||||
start_date=datetime(2024, 1, 1),
|
||||
schedule=None,
|
||||
catchup=False,
|
||||
template_searchpath="/sql",
|
||||
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="stg/bookings_ddl.sql",
|
||||
)
|
||||
@@ -0,0 +1,87 @@
|
||||
from __future__ import annotations
|
||||
|
||||
"""
|
||||
Учебный DAG для менти: показывает, как устроен поток
|
||||
от источника bookings-db до слоя stg в Greenplum.
|
||||
|
||||
В этом примере мы сознательно используем:
|
||||
- Airflow Connections для подключения к БД;
|
||||
- PostgresOperator, который берёт SQL-скрипты с диска;
|
||||
чтобы студент увидел «канонический» способ работы с SQL в DAG.
|
||||
|
||||
Каждый запуск DAG работает как «шаг по времени вперёд»:
|
||||
- генератор в демо-БД bookings добавляет следующий учебный день после max(book_date);
|
||||
- загрузка в Greenplum берёт все строки, появившиеся после предыдущих батчей;
|
||||
- логическая дата запуска (`ds`) используется как удобная метка запуска (через `ds_nodash` в `batch_id`, в логах и DQ).
|
||||
"""
|
||||
|
||||
from datetime import datetime, timedelta
|
||||
|
||||
from airflow.operators.python import PythonOperator
|
||||
from airflow.providers.postgres.operators.postgres import PostgresOperator
|
||||
|
||||
from airflow import DAG
|
||||
|
||||
default_args = {
|
||||
"owner": "airflow",
|
||||
"retries": 1,
|
||||
"retry_delay": timedelta(seconds=30),
|
||||
}
|
||||
|
||||
|
||||
BOOKINGS_CONN_ID = "bookings_db"
|
||||
GREENPLUM_CONN_ID = "greenplum_conn"
|
||||
|
||||
|
||||
def _finish_summary() -> None:
|
||||
"""
|
||||
Логирует краткий итог выполнения DAG за один запуск.
|
||||
|
||||
Здесь можно добавить дополнительную агрегацию/логирование,
|
||||
но для учебного примера достаточно простого сообщения.
|
||||
"""
|
||||
from logging import getLogger
|
||||
|
||||
log = getLogger(__name__)
|
||||
log.info("DAG bookings_to_gp_stage завершён. Подробности смотрите в логах задач.")
|
||||
|
||||
|
||||
with DAG(
|
||||
dag_id="bookings_to_gp_stage",
|
||||
start_date=datetime(2024, 1, 1),
|
||||
schedule=None,
|
||||
catchup=False,
|
||||
template_searchpath="/sql",
|
||||
default_args=default_args,
|
||||
tags=["demo", "bookings", "greenplum", "stg"],
|
||||
description="Учебный DAG: загрузка из bookings-db в stg.bookings (Greenplum)",
|
||||
) as dag:
|
||||
# 1. Генерируем один (или несколько стартовых) учебный день в демо-БД bookings
|
||||
generate_bookings_day = PostgresOperator(
|
||||
task_id="generate_bookings_day",
|
||||
postgres_conn_id=BOOKINGS_CONN_ID,
|
||||
sql="src/bookings_generate_day_if_missing.sql",
|
||||
autocommit=True,
|
||||
)
|
||||
|
||||
# 2. Загружаем инкремент из stg.bookings_ext в stg.bookings
|
||||
load_bookings_to_stg = PostgresOperator(
|
||||
task_id="load_bookings_to_stg",
|
||||
postgres_conn_id=GREENPLUM_CONN_ID,
|
||||
sql="stg/bookings_load.sql",
|
||||
)
|
||||
|
||||
# 3. Проверяем количество строк между источником и stg.bookings
|
||||
check_row_counts = PostgresOperator(
|
||||
task_id="check_row_counts",
|
||||
postgres_conn_id=GREENPLUM_CONN_ID,
|
||||
sql="stg/bookings_dq.sql",
|
||||
)
|
||||
|
||||
# 4. Финальный лог/сводка
|
||||
finish_summary = PythonOperator(
|
||||
task_id="finish_summary",
|
||||
python_callable=_finish_summary,
|
||||
)
|
||||
|
||||
generate_bookings_day >> load_bookings_to_stg >> check_row_counts >> finish_summary
|
||||
@@ -4,10 +4,13 @@ import logging
|
||||
from datetime import datetime, timedelta
|
||||
|
||||
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_schema,
|
||||
assert_orders_table_exists, get_gp_conn)
|
||||
assert_orders_table_exists,
|
||||
get_gp_conn,
|
||||
)
|
||||
|
||||
from airflow import DAG
|
||||
|
||||
@@ -16,7 +19,8 @@ def _run_check(check_callable):
|
||||
"""
|
||||
Оборачивает проверку качества данных в контекст подключения к Greenplum.
|
||||
|
||||
Этот DAG предназначен для автоматической проверки качества данных в таблице orders:
|
||||
Этот DAG предназначен для автоматической проверки качества данных
|
||||
после CSV-пайплайна в таблице public.orders:
|
||||
1. Проверяет существование таблицы
|
||||
2. Проверяет соответствие схемы
|
||||
3. Проверяет наличие данных
|
||||
@@ -47,13 +51,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(
|
||||
@@ -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(2024, 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",
|
||||
)
|
||||
@@ -50,14 +50,17 @@ def get_gp_conn():
|
||||
|
||||
# Прямое подключение по переменным окружения
|
||||
conn_params = {
|
||||
"dbname": os.getenv("GP_DB", "gpadmin"),
|
||||
"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", conn_params["host"], conn_params["port"]
|
||||
"🔗 Подключение к Greenplum: %s:%s/%s",
|
||||
conn_params["host"],
|
||||
conn_params["port"],
|
||||
conn_params["dbname"],
|
||||
)
|
||||
return psycopg2.connect(**conn_params)
|
||||
|
||||
|
||||
@@ -0,0 +1,11 @@
|
||||
# Демобаза bookings в Postgres
|
||||
|
||||
Этот каталог используется для работы с демобазой [bookings](https://postgrespro.ru/education/demodb), которая будет источником данных для будущего DWH в Greenplum.
|
||||
|
||||
На первом этапе мы:
|
||||
- поднимаем отдельный контейнер `bookings-db` с Postgres;
|
||||
- устанавливаем в нём генератор демобазы `demodb` (репозиторий `postgrespro/demodb`);
|
||||
- генерируем данные «день за днём» с помощью `make`‑команд.
|
||||
|
||||
Основные команды см. в корневом `Makefile` (`bookings-init`, `bookings-generate-day`, `bookings-psql`) и в `README.md` проекта.
|
||||
|
||||
@@ -0,0 +1,41 @@
|
||||
-- Если база пуста, берём стартовую дату из GUC bookings.start_date (по умолчанию: 2017-01-01).
|
||||
DO $$
|
||||
DECLARE
|
||||
v_max_book_date timestamptz;
|
||||
v_start_date timestamptz;
|
||||
v_end_date timestamptz;
|
||||
v_jobs integer := COALESCE(current_setting('bookings.jobs', true), '1')::integer;
|
||||
v_init_days integer := COALESCE(current_setting('bookings.init_days', true), '1')::integer;
|
||||
v_start_cfg text := COALESCE(current_setting('bookings.start_date', true), '2017-01-01');
|
||||
BEGIN
|
||||
-- Проверяем, что демобаза установлена
|
||||
IF to_regclass('bookings.bookings') IS NULL THEN
|
||||
RAISE EXCEPTION 'Таблица bookings.bookings не найдена. Сначала выполните make bookings-init.';
|
||||
END IF;
|
||||
|
||||
-- Ищем последнюю сгенерированную дату
|
||||
SELECT max(book_date) INTO v_max_book_date FROM bookings.bookings;
|
||||
|
||||
IF v_max_book_date IS NULL THEN
|
||||
-- База пустая: берём стартовую дату из конфигурации (или дефолтную)
|
||||
v_start_date := date_trunc('day', v_start_cfg::timestamptz);
|
||||
ELSE
|
||||
-- Продолжаем с дня, следующего за максимальной датой
|
||||
v_start_date := date_trunc('day', v_max_book_date) + interval '1 day';
|
||||
END IF;
|
||||
|
||||
-- Первая генерация вызывает generate(), последующие — continue()
|
||||
IF v_max_book_date IS NULL THEN
|
||||
v_end_date := v_start_date + (v_init_days || ' days')::interval;
|
||||
CALL generate(v_start_date, v_end_date, v_jobs);
|
||||
ELSE
|
||||
v_end_date := v_start_date + interval '1 day';
|
||||
CALL continue(v_end_date, v_jobs);
|
||||
END IF;
|
||||
|
||||
-- Ждём завершения фоновых джобов генератора, чтобы данные успели записаться
|
||||
WHILE busy() LOOP
|
||||
PERFORM pg_sleep(1);
|
||||
END LOOP;
|
||||
PERFORM dblink_disconnect(unnest(dblink_get_connections()));
|
||||
END $$;
|
||||
@@ -0,0 +1,11 @@
|
||||
--- a/engine.sql
|
||||
+++ b/engine.sql
|
||||
@@ -230,7 +230,8 @@
|
||||
SELECT count(*) > 0
|
||||
FROM pg_stat_activity
|
||||
WHERE application_name = 'Airlines processor'
|
||||
- AND state != 'idle';
|
||||
+ AND state != 'idle'
|
||||
+ AND pid <> pg_backend_pid();
|
||||
END;
|
||||
|
||||
@@ -0,0 +1,11 @@
|
||||
--- a/install.sql
|
||||
+++ b/install.sql
|
||||
@@ -10,7 +10,7 @@
|
||||
THE SOFTWARE IS PROVIDED “AS IS”, WITHOUT WARRANTY OF ANY KIND, EXPRESS OR IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
|
||||
*/
|
||||
|
||||
-DROP DATABASE demo;
|
||||
+DROP DATABASE IF EXISTS demo;
|
||||
CREATE DATABASE demo;
|
||||
\c demo
|
||||
CREATE EXTENSION btree_gist;
|
||||
+44
-6
@@ -5,6 +5,7 @@ services:
|
||||
# container_name: gp_pgmeta
|
||||
env_file: .env
|
||||
environment:
|
||||
TZ: ${TZ:-Europe/Moscow}
|
||||
POSTGRES_USER: ${PG_USER}
|
||||
POSTGRES_PASSWORD: ${PG_PASSWORD}
|
||||
POSTGRES_DB: ${PG_DB}
|
||||
@@ -18,25 +19,50 @@ services:
|
||||
timeout: 5s
|
||||
retries: 20
|
||||
|
||||
# Postgres с демо-БД bookings (источник для будущего DWH)
|
||||
bookings-db:
|
||||
image: postgres:16
|
||||
env_file: .env
|
||||
environment:
|
||||
TZ: ${TZ:-Europe/Moscow}
|
||||
POSTGRES_USER: ${BOOKINGS_DB_USER}
|
||||
POSTGRES_PASSWORD: ${BOOKINGS_DB_PASSWORD}
|
||||
POSTGRES_DB: ${BOOKINGS_DB_NAME}
|
||||
ports:
|
||||
- "${BOOKINGS_DB_PORT:-5434}:5432"
|
||||
volumes:
|
||||
- bookings_data:/var/lib/postgresql/data
|
||||
# В этот каталог будет монтироваться генератор demodb
|
||||
- ./bookings:/bookings:ro
|
||||
healthcheck:
|
||||
test: ["CMD-SHELL", "pg_isready -U ${BOOKINGS_DB_USER} -d ${BOOKINGS_DB_NAME}"]
|
||||
interval: 5s
|
||||
timeout: 5s
|
||||
retries: 20
|
||||
|
||||
greenplum:
|
||||
image: woblerr/greenplum:6.27.1 # based on https://github.com/woblerr/docker-greenplum
|
||||
# container_name: gp_single
|
||||
# hostname: gpdbsne
|
||||
environment:
|
||||
TZ: ${TZ:-Europe/Moscow}
|
||||
GREENPLUM_USER: ${GP_USER:-gpadmin}
|
||||
GREENPLUM_PASSWORD: ${GP_PASSWORD:-gpadmin}
|
||||
GREENPLUM_DATABASE_NAME: ${GP_DB:-gpadmin}
|
||||
# GP_PORT: ${GP_PORT:-5432}
|
||||
# Порты: внешний 5432
|
||||
GREENPLUM_DATABASE_NAME: ${GP_DB:-gp_dwh}
|
||||
GREENPLUM_PXF_ENABLE: "true"
|
||||
# Порты: внешний 5435 (на хосте)
|
||||
ports:
|
||||
- "${GP_PORT}:5432"
|
||||
- "5435:5432"
|
||||
volumes:
|
||||
- ./sql:/sql:ro
|
||||
- ./pxf/postgresql-42.7.3.jar:/pxf-local/postgresql-42.7.3.jar:ro
|
||||
- ./pxf/servers/bookings-db:/pxf-local/servers/bookings-db:ro
|
||||
- ./pxf/init/10_pxf_bookings.sh:/docker-entrypoint-initdb.d/10_pxf_bookings.sh:ro
|
||||
- greenplum_data:/data
|
||||
- airflow_data:/opt/airflow/data
|
||||
# Простая проверка доступности: psql откликается
|
||||
healthcheck:
|
||||
test: ["CMD-SHELL", "pg_isready -h 127.0.0.1 -p 5432 -U ${GP_USER:-gpadmin} -d ${GP_DB:-gpadmin} || echo 1"]
|
||||
test: ["CMD-SHELL", "pg_isready -h 127.0.0.1 -p 5432 -U ${GP_USER:-gpadmin} -d ${GP_DB:-gp_dwh} || echo 1"]
|
||||
interval: 10s
|
||||
timeout: 5s
|
||||
retries: 30
|
||||
@@ -46,9 +72,12 @@ services:
|
||||
container_name: gp_airflow_web
|
||||
env_file: .env
|
||||
environment:
|
||||
TZ: ${TZ:-Europe/Moscow}
|
||||
AIRFLOW__CORE__LOAD_EXAMPLES: "False"
|
||||
AIRFLOW__DATABASE__SQL_ALCHEMY_CONN: postgresql+psycopg2://${PG_USER}:${PG_PASSWORD}@pgmeta:5432/${PG_DB}
|
||||
AIRFLOW__WEBSERVER__SECRET_KEY: ${AIRFLOW__WEBSERVER__SECRET_KEY}
|
||||
AIRFLOW_CONN_GREENPLUM_CONN: postgresql://${GP_USER}:${GP_PASSWORD}@greenplum:${GP_PORT:-5432}/${GP_DB}
|
||||
AIRFLOW_CONN_BOOKINGS_DB: postgresql://${BOOKINGS_DB_USER}:${BOOKINGS_DB_PASSWORD}@bookings-db:5432/demo
|
||||
command: >
|
||||
bash -lc "pip install --no-cache-dir -r /opt/airflow/requirements.txt &&
|
||||
airflow webserver"
|
||||
@@ -57,6 +86,7 @@ services:
|
||||
volumes:
|
||||
- ./airflow/dags:/opt/airflow/dags
|
||||
- ./airflow/requirements.txt:/opt/airflow/requirements.txt
|
||||
- ./sql:/sql:ro
|
||||
- airflow_data:/opt/airflow/data
|
||||
depends_on:
|
||||
pgmeta:
|
||||
@@ -69,15 +99,19 @@ services:
|
||||
container_name: gp_airflow_sch
|
||||
env_file: .env
|
||||
environment:
|
||||
TZ: ${TZ:-Europe/Moscow}
|
||||
AIRFLOW__CORE__LOAD_EXAMPLES: "False"
|
||||
AIRFLOW__DATABASE__SQL_ALCHEMY_CONN: postgresql+psycopg2://${PG_USER}:${PG_PASSWORD}@pgmeta:5432/${PG_DB}
|
||||
AIRFLOW__WEBSERVER__SECRET_KEY: ${AIRFLOW__WEBSERVER__SECRET_KEY}
|
||||
AIRFLOW_CONN_GREENPLUM_CONN: postgresql://${GP_USER}:${GP_PASSWORD}@greenplum:${GP_PORT:-5432}/${GP_DB}
|
||||
AIRFLOW_CONN_BOOKINGS_DB: postgresql://${BOOKINGS_DB_USER}:${BOOKINGS_DB_PASSWORD}@bookings-db:5432/demo
|
||||
command: >
|
||||
bash -lc "pip install --no-cache-dir -r /opt/airflow/requirements.txt &&
|
||||
airflow scheduler"
|
||||
volumes:
|
||||
- ./airflow/dags:/opt/airflow/dags
|
||||
- ./airflow/requirements.txt:/opt/airflow/requirements.txt
|
||||
- ./sql:/sql:ro
|
||||
- airflow_data:/opt/airflow/data
|
||||
depends_on:
|
||||
pgmeta:
|
||||
@@ -90,13 +124,16 @@ services:
|
||||
# container_name: gp_airflow_init
|
||||
env_file: .env
|
||||
environment:
|
||||
TZ: ${TZ:-Europe/Moscow}
|
||||
AIRFLOW__CORE__LOAD_EXAMPLES: "False"
|
||||
AIRFLOW__DATABASE__SQL_ALCHEMY_CONN: postgresql+psycopg2://${PG_USER}:${PG_PASSWORD}@pgmeta:5432/${PG_DB}
|
||||
AIRFLOW__WEBSERVER__SECRET_KEY: ${AIRFLOW__WEBSERVER__SECRET_KEY}
|
||||
user: "0:0"
|
||||
AIRFLOW_CONN_GREENPLUM_CONN: postgresql://${GP_USER}:${GP_PASSWORD}@greenplum:${GP_PORT:-5432}/${GP_DB}
|
||||
AIRFLOW_CONN_BOOKINGS_DB: postgresql://${BOOKINGS_DB_USER}:${BOOKINGS_DB_PASSWORD}@bookings-db:5432/demo
|
||||
volumes:
|
||||
- ./airflow/dags:/opt/airflow/dags
|
||||
- ./airflow/requirements.txt:/opt/airflow/requirements.txt
|
||||
- ./sql:/sql:ro
|
||||
- airflow_data:/opt/airflow/data
|
||||
command: >
|
||||
bash -lc "
|
||||
@@ -117,5 +154,6 @@ services:
|
||||
|
||||
volumes:
|
||||
pgmeta:
|
||||
bookings_data:
|
||||
greenplum_data:
|
||||
airflow_data:
|
||||
|
||||
@@ -0,0 +1,138 @@
|
||||
# Дизайн STG для bookings в Greenplum (черновик)
|
||||
|
||||
_Внутренний документ для учебного стенда. Перед итоговой сдачей можно объединить с основной документацией._
|
||||
|
||||
## 1. Цель и общий контур
|
||||
|
||||
- Источник: Postgres в контейнере `bookings-db`, база `demo`, таблица `bookings.bookings` (см. `docs/internal/bookings_tz.md`).
|
||||
- Цель: показываем путь данных от операционной БД до сырого слоя DWH в Greenplum.
|
||||
- В этом документе описываем только часть `src (bookings-db) → STG (Greenplum)`. Слои ODS/DDS/DM студент проектирует сам по статье про моделирование DWH.
|
||||
|
||||
Логика на уровне слоёв (по статье):
|
||||
|
||||
- `src`: оперативная система (`bookings-db`, схема `bookings`).
|
||||
- `stg`: сырой слой в Greenplum, максимально близкий к источнику, без бизнес‑логики.
|
||||
- дальше по заданию менти можно строить `ods`/`dds`/`dm` поверх STG.
|
||||
|
||||
## 2. Схема и таблицы в Greenplum
|
||||
|
||||
### 2.1. Схема
|
||||
|
||||
- Используем одну схему `stg` в Greenplum.
|
||||
- В этой схеме будут:
|
||||
- внешняя таблица PXF для чтения из `bookings-db`;
|
||||
- внутренняя таблица STG для долговременного хранения «сырых» данных.
|
||||
|
||||
### 2.2. Внешняя таблица (PXF)
|
||||
|
||||
- Имя таблицы: `stg.bookings_ext`.
|
||||
- Назначение: «окно» в исходную таблицу `bookings.bookings` в `bookings-db` через PXF (JDBC).
|
||||
- Типы колонок:
|
||||
- можем использовать «родные» типы из `bookings.bookings` (включая даты/числа);
|
||||
- задача внешней таблицы — корректно читать данные из источника, не заниматься приведением типов.
|
||||
|
||||
DDL будет добавлен в `sql/ddl_gp.sql` в блоке DDL для Greenplum (примерно по шаблону из `docs/internal/pxf_bookings.md`), с `LOCATION ('pxf://bookings.bookings?PROFILE=JDBC&SERVER=bookings-db')`.
|
||||
|
||||
### 2.3. Внутренняя таблица STG
|
||||
|
||||
- Имя таблицы: `stg.bookings`.
|
||||
- Назначение: хранить сырые данные из источника для последующей обработки (ODS/DDS/витрины).
|
||||
- Принципы моделирования:
|
||||
- все бизнес‑колонки из `bookings.bookings` храним как `TEXT` (как в примерах STG из статьи);
|
||||
- не делаем `UPDATE/DELETE`, только `INSERT` новых записей;
|
||||
- бизнес‑колонки по названию совпадают с источником (чтобы проще было маппить).
|
||||
|
||||
Технологические колонки:
|
||||
|
||||
- `src_created_at_ts TIMESTAMP` — дата/время из источника, приведённая к TIMESTAMP:
|
||||
- используется как опорная колонка для инкрементальной загрузки;
|
||||
- заполняется из исходной даты/времени (`created_at` или аналог).
|
||||
- `load_dttm TIMESTAMP NOT NULL DEFAULT now()` — когда запись была загружена в STG.
|
||||
- `batch_id TEXT NOT NULL` — идентификатор «пачки» (например, `{{ ds_nodash }}` или `run_id` Airflow).
|
||||
- при необходимости позже можно добавить `src_system TEXT`, если появятся другие источники.
|
||||
|
||||
Колонки‑бизнес‑ключи (`booking_id` и т.п.) храним как `TEXT`. В слое DDS позже можно будет ввести суррогатные ключи и нормализовать модель под витрины.
|
||||
|
||||
## 3. Инкрементальная загрузка
|
||||
|
||||
### 3.1. Опорное поле для инкремента
|
||||
|
||||
- Опорная колонка: `src_created_at_ts` (внутреннее имя в STG).
|
||||
- Источник значения:
|
||||
- берём из соответствующей колонки в `bookings.bookings` (например, `book_date`/`created_at` — будет уточнено при реализации);
|
||||
- при чтении через `stg.bookings_ext` приводим к `TIMESTAMP`.
|
||||
|
||||
### 3.2. Правила определения full/delta
|
||||
|
||||
- При первом запуске, если таблица `stg.bookings` пуста:
|
||||
- считаем режим `full` — загружаем все строки из `stg.bookings_ext`.
|
||||
- При последующих запусках:
|
||||
- читаем `max(src_created_at_ts)` из `stg.bookings` за все предыдущие загрузки;
|
||||
- считаем, что нужно загрузить только строки, где `src_created_at_ts` больше этой максимальной метки и не позже конца текущего учебного дня.
|
||||
|
||||
Таким образом, вся логика инкремента «замкнута» на один техно‑столбец `src_created_at_ts`, который студент потом сможет использовать и на следующих слоях (например, в CDC‑логике).
|
||||
|
||||
## 4. DAG’и Airflow (логика на уровне задач)
|
||||
|
||||
### 4.1. DAG для DDL
|
||||
|
||||
- `dag_id`: `bookings_stg_ddl` (реализован в `airflow/dags/bookings_stg_ddl.py`).
|
||||
- Назначение: один раз (или при изменении схемы) создать необходимые объекты в Greenplum:
|
||||
- схему `stg` (если её ещё нет);
|
||||
- внешнюю таблицу `stg.bookings_ext` (PXF → `bookings-db`);
|
||||
- внутреннюю таблицу `stg.bookings` с текстовыми колонками и тех.полями.
|
||||
- Этот DAG не загружает данные, только подготавливает структуру.
|
||||
- Вся DDL‑логика (CREATE/ALTER/DROP) сосредоточена здесь; рабочие DAG’и занимаются только DML (INSERT/SELECT).
|
||||
|
||||
### 4.2. DAG для пошаговой загрузки
|
||||
|
||||
- `dag_id`: `bookings_to_gp_stage`.
|
||||
- Основные параметры:
|
||||
- `batch_id` (по умолчанию `{{ ds_nodash }}`) — метка батча, которая попадает в `stg.bookings.batch_id`;
|
||||
- подключения:
|
||||
- `bookings_db_conn_id` — Airflow connection к `bookings-db` (в коде DAG — `BOOKINGS_CONN_ID = "bookings_db"`);
|
||||
- `greenplum_conn_id` — Airflow connection к Greenplum (`GREENPLUM_CONN_ID = "greenplum_conn"`).
|
||||
|
||||
Последовательность задач (упрощённая, но отражающая те же шаги):
|
||||
|
||||
1. `generate_bookings_day`
|
||||
- PostgresOperator к `bookings-db`;
|
||||
- выполняет скрипт `/sql/src/bookings_generate_day_if_missing.sql`;
|
||||
- скрипт смотрит на `max(book_date)` и:
|
||||
- если база пуста — берёт стартовую дату из конфигурации (`bookings.start_date`) и генерирует `bookings.init_days` суток;
|
||||
- если данные уже есть — добавляет один следующий учебный день после `max(book_date)` (логика как в `bookings/generate_next_day.sql`).
|
||||
2. `load_bookings_to_stg`
|
||||
- PostgresOperator к Greenplum;
|
||||
- выполняет скрипт `/sql/stg/bookings_load.sql`;
|
||||
- внутри SQL считается `max(src_created_at_ts)` по «старым» батчам и по нему строится окно инкремента:
|
||||
- первая загрузка (full) — берём все строки из `stg.bookings_ext`;
|
||||
- последующие загрузки — берём только записи, где `book_date` больше предыдущего максимума (верхняя граница по дате не задаётся явно);
|
||||
- при вставке заполняются тех.колонки `src_created_at_ts`, `load_dttm`, `batch_id`.
|
||||
3. `check_row_counts`
|
||||
- PostgresOperator к Greenplum;
|
||||
- выполняет скрипт `/sql/stg/bookings_dq.sql`;
|
||||
- скрипт заново считает окно инкремента по тем же правилам, что и загрузка, и сравнивает:
|
||||
- количество строк в `stg.bookings_ext` с `book_date` позже «старого» максимума,
|
||||
- количество строк в `stg.bookings` для текущего `batch_id`;
|
||||
- при расхождении выполняет `RAISE EXCEPTION` с понятным текстом ошибки.
|
||||
4. `finish_summary`
|
||||
- PythonOperator, который логирует итог выполнения DAG и напоминает, где смотреть детальные логи.
|
||||
|
||||
Таким образом, вся бизнес‑логика инкремента и проверок живёт в SQL‑скриптах, а DAG отвечает за оркестрацию и подключение к нужным БД. Для менти это хороший пример разделения ответственности между SQL и Python.
|
||||
|
||||
## 5. Связь с остальными документами
|
||||
|
||||
- `docs/internal/bookings_tz.md` — как готовится и генерируется источник `bookings-db`.
|
||||
- `docs/internal/pxf_bookings.md` — детали настройки PXF и внешней таблицы для чтения из `bookings-db`.
|
||||
- `sql/stg/bookings_ddl.sql` — DDL для схемы `stg` и таблиц `stg.bookings_ext` / `stg.bookings` (подключается из `sql/ddl_gp.sql` и применяется через `make ddl-gp`).
|
||||
|
||||
Дальнейшая модель DWH (слои ODS/DDS/DM, факт/измерения, SCD) должна быть спроектирована студентом по статье о моделировании данных, используя `stg.bookings` как входной слой.
|
||||
|
||||
## 6. Требования к читаемости и комментариям
|
||||
|
||||
- DAG’и `bookings_stg_ddl` и `bookings_to_gp_stage` — это учебный материал для менти.
|
||||
- В коде DAG’ов должны быть:
|
||||
- понятные docstring на русском у всех функций и DAG;
|
||||
- короткие комментарии рядом с нетривиальной логикой (особенно вокруг инкремента и идемпотентности);
|
||||
- говорящие `task_id` и названия функций, отражающие их роль в процессе.
|
||||
- Цель: чтобы по одному только коду DAG студент мог восстановить архитектуру процесса и сопоставить её с теорией из статьи про моделирование DWH.
|
||||
@@ -0,0 +1,80 @@
|
||||
# Мини‑README по учебному DAG bookings_to_gp_stage (черновик)
|
||||
|
||||
_Внутренний файл, чтобы не забыть договорённости. Перед итоговой сдачей документацию по блоку bookings/STG нужно будет аккуратно собрать и переписать._
|
||||
|
||||
## 1. Что делает DAG
|
||||
|
||||
- DAG `bookings_to_gp_stage` показывает учебный поток:
|
||||
- источник: демо‑БД `bookings-db` (Postgres, схема `bookings`, таблица `bookings.bookings`);
|
||||
- при каждом запуске генерируется один учебный день данных (идемпотентно);
|
||||
- данные из `bookings.bookings` переливаются в сырой слой `stg.bookings` в Greenplum через PXF‑внешнюю таблицу `stg.bookings_ext`.
|
||||
- Слой `stg` задуман как «сырой»:
|
||||
- все бизнес‑колонки (`book_ref`, `book_date`, `total_amount`) хранятся как `TEXT`;
|
||||
- есть тех.колонки `src_created_at_ts`, `load_dttm`, `batch_id`.
|
||||
|
||||
Подробный дизайн описан в `docs/internal/bookings_stg_design.md`.
|
||||
|
||||
## 2. Что нужно, чтобы DAG завёлся
|
||||
|
||||
Минимальные предпосылки:
|
||||
|
||||
- Стенд поднят: `make up`.
|
||||
- Демо‑БД bookings инициализирована: `make bookings-init`.
|
||||
- В Greenplum применён DDL (созданы схема `stg` и таблицы `stg.bookings_ext` / `stg.bookings`):
|
||||
- учебный вариант: запустить DAG `bookings_stg_ddl` (он использует `sql/stg/bookings_ddl.sql`);
|
||||
- технический шорткат: `make ddl-gp` применяет все DDL разом вручную. Команда сама не вызывается при старте контейнеров, её нужно запустить явно.
|
||||
- В Airflow есть коннекты:
|
||||
- `greenplum_conn` — к Greenplum;
|
||||
- `bookings_db` — к сервису `bookings-db`.
|
||||
По умолчанию они задаются через переменные окружения `AIRFLOW_CONN_...` в docker-compose,
|
||||
поэтому могут не отображаться в UI, но `PostgresOperator` найдёт их по `conn_id`. При желании
|
||||
их можно создать/отредактировать вручную через Admin → Connections.
|
||||
|
||||
## 3. Последовательность задач в DAG
|
||||
|
||||
- `generate_bookings_day`:
|
||||
- PostgresOperator к `bookings-db`;
|
||||
- выполняет SQL `/sql/src/bookings_generate_day_if_missing.sql`;
|
||||
- скрипт смотрит на `max(book_date)` в `bookings.bookings`:
|
||||
- если база пустая — берёт стартовую дату из GUC и генерирует `bookings.init_days` суток;
|
||||
- если данные уже есть — добавляет один следующий учебный день после `max(book_date)` и пишет NOTICE с интервалом генерации.
|
||||
- `load_bookings_to_stg`:
|
||||
- PostgresOperator к Greenplum;
|
||||
- выполняет SQL `/sql/stg/bookings_load.sql`;
|
||||
- считает «старый» максимум `src_created_at_ts` (по предыдущим батчам) и грузит только новые строки из `stg.bookings_ext`, заполняя `src_created_at_ts`, `load_dttm`, `batch_id={{ ds_nodash }}`.
|
||||
- `check_row_counts`:
|
||||
- PostgresOperator к Greenplum;
|
||||
- выполняет SQL `/sql/stg/bookings_dq.sql`;
|
||||
- за то же окно инкремента считает количество строк в источнике и в `stg.bookings` (по текущему `batch_id`);
|
||||
- при расхождении делает `RAISE EXCEPTION` с понятным текстом ошибки.
|
||||
- `finish_summary`:
|
||||
- логирует итог выполнения DAG за одно срабатывание.
|
||||
|
||||
## 4. Как этим пользоваться студенту (черновой сценарий)
|
||||
|
||||
1. Поднять стенд и подготовить источники:
|
||||
- `make up` (Airflow инициализируется автоматически при первом старте)
|
||||
- `make bookings-init`
|
||||
- `make ddl-gp`
|
||||
2. Открыть Airflow UI (`http://localhost:8080`) и включить DAG `bookings_to_gp_stage`.
|
||||
3. Вызвать `Trigger` DAG (дату логического запуска можно оставить по умолчанию — она используется только как метка `batch_id`).
|
||||
4. Посмотреть:
|
||||
- в `bookings-db` ― что появился день с бронированиями;
|
||||
- в Greenplum (`make gp-psql`) — данные в `stg.bookings`:
|
||||
- `SELECT * FROM stg.bookings LIMIT 10;`
|
||||
- `SELECT src_created_at_ts, load_dttm, batch_id FROM stg.bookings ORDER BY src_created_at_ts DESC LIMIT 10;`
|
||||
5. Перезапустить DAG ещё несколько раз и увидеть, что:
|
||||
- генерация в `bookings.bookings` идёт по одному дню вперёд от текущего `max(book_date)`;
|
||||
- в `stg.bookings` появляются только новые записи (delta), помеченные разными `batch_id`.
|
||||
|
||||
## 5. Примечания «на потом»
|
||||
|
||||
- Текущая документация по блоку bookings/STG разбросана:
|
||||
- `README.md` (общий обзор стенда),
|
||||
- `docs/internal/bookings_tz.md` (источник bookings),
|
||||
- `docs/internal/pxf_bookings.md` (PXF),
|
||||
- `docs/internal/bookings_stg_design.md` (дизайн STG),
|
||||
- этот файл (мини‑README по DAG).
|
||||
- В будущем всё это нужно будет собрать в одну понятную историю для студента:
|
||||
- отдельный раздел «Учебный пример: bookings → stg → dwh»;
|
||||
- скриншоты DAG, примеры запросов и типичные ошибки.
|
||||
@@ -0,0 +1,25 @@
|
||||
# Временное ТЗ по блоку bookings (для текущей разработки)
|
||||
|
||||
_Этот файл внутренний, удалить перед итоговой сдачей._
|
||||
|
||||
- Контейнер `bookings-db` — отдельный сервис Postgres из `docker-compose.yml`, база по умолчанию `demo` (из upstream demodb), без переименований.
|
||||
- Доступ снаружи не блокируем (порт `5434` по умолчанию), чтобы позже читать через PXF и подключаться из Greenplum.
|
||||
- Инициализация (`make bookings-init`): поднимает контейнер, клонирует demodb с закреплённым коммитом, накладывает патчи (`engine`: `jobs=1` синхронно + `busy()` игнорирует свой pid; `install.sql`: `DROP DATABASE IF EXISTS`), ждёт `pg_isready`, ставит `gen.connstr` и GUC `bookings.start_date/init_days/jobs`, затем запускает `/bookings/generate_next_day.sql` через `psql -f`. Значения по умолчанию: стартовая дата 2017-01-01, `init_days=1`, `jobs=1`.
|
||||
- Генерация следующего дня: `make bookings-generate-day` прогоняет тот же SQL (читает GUC, вызывает `generate/continue`, ждёт `busy()`, закрывает dblink). При `jobs=1` всё синхронно, без dblink.
|
||||
- Исходники demodb: клонируем по требованию с фиксированным хешем, кладём в `bookings/demodb/` (в `.gitignore`), патчи лежат в `bookings/patches/` и применяются автоматически.
|
||||
- Документация: в README описаны команды (`bookings-init`, проверка данных, генерация дня), параметры `.env`; настройка PXF/ETL — следующий этап.
|
||||
|
||||
## Текущее состояние
|
||||
- `Makefile` теперь автоматически применяет патчи (`engine_jobs1_sync.patch`, `install_drop_if_exists.patch`), ждёт готовности Postgres через `pg_isready`, запускает `install.sql`, выставляет `gen.connstr`/GUC и вызывает `generate_next_day.sql` через `psql -f`.
|
||||
- Дефолты: `BOOKINGS_START_DATE=2017-01-01`, `BOOKINGS_INIT_DAYS=1`, `BOOKINGS_JOBS=1`. При `jobs=1` генерация идёт синхронно без dblink, `busy()` не учитывает текущую сессию.
|
||||
- `.env.example`/README обновлены под новые дефолты; каталог `bookings/demodb/` в `.gitignore`.
|
||||
- Патчи лежат в `bookings/patches/` и накладываются при `bookings-clone-demodb`.
|
||||
|
||||
## Текущее состояние тестов/проблем
|
||||
- Чистый прогон `make bookings-init` (после `docker compose down -v` и удаления `bookings/demodb`) проходит за ~1,5 минуты: база ставится, `busy()` → `f`, `bookings.bookings` от `2017-01-01 00:00:18` до `2017-01-01 23:59:59`.
|
||||
- Ранее зависание на `busy()` при `jobs=1` лечится патчем: `process_queue` теперь синхронный, а `busy()` игнорирует текущий backend.
|
||||
- Данных пока только на 1 день по умолчанию, чтобы генерация не занимала много времени.
|
||||
|
||||
## Идеи/следующие шаги
|
||||
- Если понадобится больше дней — увеличивать `BOOKINGS_INIT_DAYS`, но помнить, что генерация может идти долго; контролировать через `SELECT busy();`.
|
||||
- Следующий этап — PXF/ETL в Greenplum; текущая задача — лишь подготовить источник bookings.
|
||||
@@ -0,0 +1,228 @@
|
||||
# Временное ТЗ по PXF для чтения данных из bookings (черновик)
|
||||
|
||||
_Этот файл внутренний, удалить перед итоговой сдачей._
|
||||
|
||||
## 1. Цель и границы
|
||||
|
||||
- Минимальная цель: настроить PXF в контейнере Greenplum так, чтобы из базы `bookings` (Postgres в сервисе `bookings-db`) можно было делать `SELECT` по одной внешней таблице в Greenplum.
|
||||
- На этом этапе **не** делаем загрузку в постоянные таблицы Greenplum, только чтение и ручные smoke‑проверки.
|
||||
- Изменения в коде/конфигурации пока планируем «на бумаге»; реализацию и правки `docker-compose.yml`/SQL/DAG делаем отдельным шагом.
|
||||
|
||||
## 2. Архитектура на уровне контейнеров
|
||||
|
||||
- `bookings-db` — Postgres 16, демо‑БД `demo` из репозитория `demodb` (источник). Доступен внутри сети Docker по имени `bookings-db` и порту `5432`.
|
||||
- `greenplum` — контейнер `woblerr/greenplum:6.27.1` (GPDB 6, Ubuntu 22.04). В нём уже есть:
|
||||
- Greenplum в режиме singlenode;
|
||||
- установленный PXF (`/usr/local/pxf`, `pxf version release-6.10.1`);
|
||||
- стартовый скрипт `/start_gpdb.sh`, который умеет включать PXF по флагу `GREENPLUM_PXF_ENABLE=true`.
|
||||
- Внешний мир (IDE/pytest) подключается к Greenplum по порту `5435` на хосте (см. `docker-compose.yml`), а к `bookings-db` — по порту `${BOOKINGS_DB_PORT}` (см. `.env`).
|
||||
|
||||
## 3. Включение PXF в нашем стенде (дизайн)
|
||||
|
||||
Планируемые изменения (позже будут внесены в `docker-compose.yml`):
|
||||
|
||||
- В сервисе `greenplum` в секцию `environment` добавить:
|
||||
- `GREENPLUM_PXF_ENABLE: "true"`.
|
||||
- При первом старте с этим флагом скрипт `/start_gpdb.sh` сделает за нас:
|
||||
- инициализацию PXF (`pxf cluster prepare`, `pxf cluster register`, `pxf cluster sync`);
|
||||
- создание расширения `pxf` в базе `${GREENPLUM_DATABASE_NAME}` (у нас это `${GP_DB}`, по умолчанию `gp_dwh`);
|
||||
- запуск `pxf cluster start` и привязку остановки/старта PXF к жизненному циклу Greenplum.
|
||||
- База конфигов PXF (`PXF_BASE`) будет располагаться в `${GREENPLUM_DATA_DIRECTORY}/pxf`, в нашем compose — это `/data/pxf` на томе `greenplum_data`.
|
||||
- Важно: `make down` сейчас делает `docker compose down -v`, поэтому при полном сбросе томов будут теряться и данные GP, и конфиги PXF (включая JDBC‑драйвер и `servers/*`).
|
||||
|
||||
## 4. JDBC‑драйвер для Postgres: где и как хранить
|
||||
|
||||
Задача: PXF должен уметь ходить по JDBC в `bookings-db` (Postgres). Для этого нужен PostgreSQL JDBC драйвер (`postgresql-*.jar`).
|
||||
|
||||
Варианты хранения драйвера:
|
||||
|
||||
1. **Коммитить JAR в репозиторий** и монтировать в контейнер.
|
||||
- Плюсы: стенд самодостаточен, не зависит от внешних скачиваний, повторяемость выше (особенно на офлайн‑машинах или при падении зеркал).
|
||||
- Минусы: лишний бинарник в учебном репо, периодически нужно обновлять версию.
|
||||
2. **Скачивать JAR внутрь контейнера один раз вручную** и хранить его в томе `greenplum_data` внутри `PXF_BASE/lib`.
|
||||
- Плюсы: нет бинарников в Git.
|
||||
- Минусы: дополнительный шаг для студентов, зависимость от сети, нужно повторять после полного сброса томов.
|
||||
|
||||
Для учебного стенда окончательно выбираем вариант **(1) — JAR в репозитории**:
|
||||
|
||||
- В репозитории заводим каталог, например `pxf/` или `pxf/jdbc/`, и кладём туда файл `postgresql-42.7.3.jar` (фиксируем версию 42.7.3 как актуальную на момент разработки).
|
||||
- В `docker-compose.yml` (на этапе реализации) смонтируем этот JAR внутрь контейнера `greenplum` в каталог `$PXF_BASE/lib`, например:
|
||||
- `./pxf/postgresql-42.7.3.jar:/data/pxf/lib/postgresql-jdbc.jar:ro`.
|
||||
- PXF по документации поддерживает размещение JDBC‑драйвера в `$PXF_BASE/lib` (общий для всех серверов) или в `$PXF_BASE/servers/<server>/lib` (локальный для сервера). Для простоты используем общий каталог `$PXF_BASE/lib`.
|
||||
- Для студентов не будет лишних подготовительных шагов: после `make up` и инициализации конфигов PXF драйвер уже на месте.
|
||||
|
||||
Договорённость для реализации:
|
||||
|
||||
- Путь в репозитории: условно `pxf/postgresql-42.7.3.jar`.
|
||||
- Путь внутри контейнера: `/data/pxf/lib/postgresql-jdbc.jar` (через bind‑mount, read‑only).
|
||||
- Обновление драйвера в будущем — ручная операция (заменить JAR в `pxf/` и скорректировать путь в `docker-compose.yml` при необходимости).
|
||||
|
||||
## 5. Сервер PXF для bookings-db (jdbc-site.xml)
|
||||
|
||||
PXF использует концепцию «серверов» (`servers/<имя>`), где для каждого сервера хранится свой конфиг подключения (в т.ч. JDBC).
|
||||
|
||||
План:
|
||||
|
||||
- Создать сервер с именем `bookings-db` (название привязываем к сервису Docker, чтобы не путаться).
|
||||
- Конфиг храним в репозитории, например в файле:
|
||||
- `pxf/servers/bookings-db/jdbc-site.xml`.
|
||||
- В контейнере этот файл будет доступен как:
|
||||
- `/data/pxf/servers/bookings-db/jdbc-site.xml` (bind‑mount read‑only).
|
||||
- Внутри прописываем параметры подключения к демо‑БД `demo` в Postgres `bookings-db`.
|
||||
|
||||
Черновой шаблон `jdbc-site.xml` (значения логина/пароля берём из `.env.example`, блок `BOOKINGS_DB_*` — для учебного стенда допускаем хардкод тех же дефолтных значений):
|
||||
|
||||
```xml
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<configuration>
|
||||
|
||||
<property>
|
||||
<name>jdbc.driver</name>
|
||||
<value>org.postgresql.Driver</value>
|
||||
</property>
|
||||
|
||||
<property>
|
||||
<name>jdbc.url</name>
|
||||
<value>jdbc:postgresql://bookings-db:5432/demo</value>
|
||||
</property>
|
||||
|
||||
<property>
|
||||
<name>jdbc.user</name>
|
||||
<value>${BOOKINGS_DB_USER}</value>
|
||||
</property>
|
||||
|
||||
<property>
|
||||
<name>jdbc.password</name>
|
||||
<value>${BOOKINGS_DB_PASSWORD}</value>
|
||||
</property>
|
||||
|
||||
</configuration>
|
||||
```
|
||||
|
||||
Замечания:
|
||||
|
||||
- В реальном файле `pxf/servers/bookings-db/jdbc-site.xml` логин и пароль будут прописаны строками, совпадающими с дефолтами из `.env.example` (`BOOKINGS_DB_USER=bookings`, `BOOKINGS_DB_PASSWORD=bookings`). Это упрощает старт стенда для студентов.
|
||||
- Если студент поменяет креды `BOOKINGS_DB_USER`/`BOOKINGS_DB_PASSWORD` в своём `.env`, он должен **также** поменять их в `pxf/servers/bookings-db/jdbc-site.xml`, иначе PXF не сможет подключиться к источнику.
|
||||
- Адрес `bookings-db:5432` — имя сервиса и внутренний порт Postgres внутри сети `docker compose`. Внешний порт (`${BOOKINGS_DB_PORT:-5434}`) здесь не используется.
|
||||
- При стандартном сценарии (конфиг монтируется read‑only) PXF подхватывает `jdbc-site.xml` при первой инициализации/старте. Если конфиг внутри контейнера всё‑таки меняли вручную, для надёжности можно выполнить:
|
||||
```bash
|
||||
pxf cluster sync
|
||||
pxf cluster restart
|
||||
```
|
||||
чтобы PXF подхватил новые настройки.
|
||||
|
||||
## 6. Внешняя таблица в Greenplum (только для чтения)
|
||||
|
||||
Задача: завести одну external‑таблицу в Greenplum, которая читает данные из демо‑БД `demo` через PXF/JDBC.
|
||||
|
||||
Дизайн:
|
||||
|
||||
- Имя таблицы в Greenplum: `public.ext_bookings_bookings` (подчёркиваем, что это внешнее представление таблицы `bookings.bookings` из Postgres).
|
||||
- Схема колонок должна совпадать со схемой исходной таблицы в `demo` (её нужно будет аккуратно выписать отдельным шагом, через `\d bookings.bookings` в `bookings-db`).
|
||||
- Локация PXF:
|
||||
- `PROFILE=JDBC` — используем JDBC‑профиль.
|
||||
- `SERVER=bookings-db` — имя сервера из `jdbc-site.xml`.
|
||||
|
||||
Черновой шаблон DDL (без конкретных типов, заполним позже по реальной схеме; предполагается, что финальный DDL ляжет в `sql/ddl_gp.sql`, чтобы применяться через `make ddl-gp`):
|
||||
|
||||
```sql
|
||||
CREATE EXTERNAL TABLE public.ext_bookings_bookings (
|
||||
-- TODO: колонки как в bookings.bookings (будет уточнено)
|
||||
)
|
||||
LOCATION ('pxf://bookings.bookings?PROFILE=JDBC&SERVER=bookings-db')
|
||||
FORMAT 'CUSTOM' (formatter='pxfwritable_import');
|
||||
```
|
||||
|
||||
Комментарии:
|
||||
|
||||
- На этапе реализации нужно будет:
|
||||
- в `bookings-db` посмотреть структуру основного факт‑табличного объекта (скорее всего `bookings.bookings`) и перенести DDL;
|
||||
- проверить типы дат/чисел, чтобы избежать сюрпризов на стороне GP.
|
||||
- Для MVP достаточно одной таблицы; позже можно добавить ещё 1–2 внешние таблицы для примеров (например, справочники).
|
||||
|
||||
## 7. План тестирования (без дополнительных make‑таргетов)
|
||||
|
||||
Цель тестов: показать студентам, что PXF настроен корректно и позволяет читать данные; при этом не плодить отдельные `make`‑цели, а использовать уже существующие (`make up`, `make bookings-init`, `make gp-psql`).
|
||||
|
||||
### 7.1. Позитивный сценарий (smoke)
|
||||
|
||||
Предварительные условия:
|
||||
|
||||
- `.env` скопирован из `.env.example` и не изменял дефолтные креды для `bookings-db` (`BOOKINGS_DB_USER=bookings`, `BOOKINGS_DB_PASSWORD=bookings`).
|
||||
- В `docker-compose.yml` включён PXF (`GREENPLUM_PXF_ENABLE=true` в сервисе `greenplum`).
|
||||
- Для `bookings-db` уже выполнен `make bookings-init` (есть данные в `demo`).
|
||||
- JAR драйвера (`pxf/postgresql-42.7.3.jar`) и файл `jdbc-site.xml` (`pxf/servers/bookings-db/jdbc-site.xml`) присутствуют в репозитории (они будут автоматически смонтированы в `/data/pxf/lib` и `/data/pxf/servers/bookings-db`).
|
||||
- DDL внешней таблицы `public.ext_bookings_bookings` добавлен в `sql/ddl_gp.sql` и применяется через `make ddl-gp`.
|
||||
|
||||
Шаги (в будущем попадут в `TESTING.md`):
|
||||
|
||||
1. Поднять стенд:
|
||||
- `make up`
|
||||
- дождаться healthcheck‑ов `pgmeta` и `greenplum`.
|
||||
3. Подготовить демо‑БД bookings:
|
||||
- `make bookings-init`.
|
||||
4. Применить DDL в Greenplum (создать таблицы, включая внешнюю `public.ext_bookings_bookings`):
|
||||
- `make ddl-gp`.
|
||||
5. Зайти в Greenplum:
|
||||
- `make gp-psql`.
|
||||
6. Проверить, что расширение PXF присутствует:
|
||||
- `\dx pxf`.
|
||||
7. Выполнить простые запросы:
|
||||
- `SELECT COUNT(*) FROM public.ext_bookings_bookings;`
|
||||
- `SELECT * FROM public.ext_bookings_bookings LIMIT 5;`
|
||||
|
||||
Ожидаемый результат:
|
||||
|
||||
- Запросы выполняются без ошибок, возвращают ненулевое количество строк.
|
||||
- Структура данных визуально совпадает с данными в `bookings-db` (можно дополнительно открыть `bookings-psql` и сравнить).
|
||||
|
||||
### 7.2. Негативный сценарий (отказ источника)
|
||||
|
||||
Цель: показать, как выглядит ошибка, если источник недоступен, и что с этим делать.
|
||||
|
||||
Шаги:
|
||||
|
||||
1. При работающем стенде остановить только `bookings-db`:
|
||||
- `docker compose stop bookings-db`.
|
||||
2. В `make gp-psql` попробовать снова:
|
||||
- `SELECT 1 FROM public.ext_bookings_bookings LIMIT 1;`
|
||||
|
||||
Ожидаемый результат:
|
||||
|
||||
- Запрос падает с ошибкой подключения к Postgres (через JDBC/pxf).
|
||||
- В `TESTING.md` планируем добавить короткую подсказку: «если видите ошибку подключения — убедитесь, что запущен сервис `bookings-db` (`docker compose start bookings-db`) и повторите запрос».
|
||||
|
||||
### 7.3. Идея для автоматического smoke‑теста (на будущее)
|
||||
|
||||
На будущее (не в рамках текущего этапа) можно добавить простой e2e‑тест в `tests/`, который:
|
||||
|
||||
- с помощью `psycopg2` коннектится к Greenplum (`GP_*` из `.env`);
|
||||
- выполняет `SELECT 1 FROM public.ext_bookings_bookings LIMIT 1`;
|
||||
- помечен как «integration» и запускается только по явному желанию (например, через отдельный маркер или переменную окружения).
|
||||
|
||||
Пока это остаётся идеей: сначала реализуем базовую конфигурацию PXF и ручной smoke‑чек‑лист.
|
||||
|
||||
## 8. Открытые вопросы / TODO
|
||||
|
||||
- Уточнить целевую таблицу(ы) в `demo` для внешнего представления (скорее всего `bookings.bookings`), аккуратно выписать DDL и обновить шаблон из раздела 6.
|
||||
- При переносе DDL проверить типы дат/времени, чтобы не получить неожиданный сдвиг по часовому поясу (см. `docs/internal/bookings_tz.md`).
|
||||
- При необходимости добавить интеграционный тест по мотивам раздела 7.3 (по отдельному маркеру/флагу).
|
||||
|
||||
## 9. Практические детали и нюансы
|
||||
|
||||
- **Таймзона**:
|
||||
- Для единообразия логов и данных задаём `TZ=Europe/Moscow` (GMT+3) в `.env.example` и пробрасываем эту переменную в контейнеры `pgmeta`, `bookings-db`, `greenplum`, `airflow-webserver`, `airflow-scheduler`.
|
||||
- При проверке данных через PXF имеет смысл сравнивать выборки по времени между `bookings-db` и Greenplum, опираясь на договорённости из `docs/internal/bookings_tz.md`.
|
||||
- **Поведение при `make down`**:
|
||||
- `make down` вызывает `docker compose down -v`, что удаляет все тома, включая `greenplum_data` (`/data` в контейнере).
|
||||
- При следующем `make up` Greenplum и PXF будут инициализироваться с нуля, но:
|
||||
- JAR и `jdbc-site.xml` возьмутся из репозитория и снова смонтируются в `/data/pxf/...`;
|
||||
- `make ddl-gp` снова создаст внешнюю таблицу `public.ext_bookings_bookings`.
|
||||
- То есть после полного ресета студенту достаточно повторить цепочку `make up` → `make bookings-init` → `make ddl-gp`.
|
||||
- **Где искать логи при проблемах с PXF**:
|
||||
- Логи PXF: в контейнере `greenplum` под пользователем `gpadmin` в каталоге `${PXF_BASE}/logs` (по умолчанию `/data/pxf/logs`).
|
||||
- Логи Greenplum: в `${GREENPLUM_DATA_DIRECTORY}/master/.../pg_log` (например, `/data/master/gpseg-1/pg_log` для GP6).
|
||||
- При ошибках подключения к `bookings-db` полезно:
|
||||
- проверить, что контейнер `bookings-db` работает (`docker compose ps`);
|
||||
- сверить креды в `.env` и `pxf/servers/bookings-db/jdbc-site.xml`;
|
||||
- посмотреть сообщения в `/data/pxf/logs`.
|
||||
@@ -0,0 +1,136 @@
|
||||
# Учебные задания по стенду
|
||||
|
||||
Этот документ собирает в одном месте задания для менти.
|
||||
Он разбит на блоки: от базовой работы с CSV‑pipeline до более продвинутого сценария с демо‑БД bookings и слоем STG в Greenplum.
|
||||
|
||||
Если вы только начинаете, выполняйте задания по порядку. К разделу про 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 и модель данных (введение)
|
||||
|
||||
В следующих заданиях мы будем опираться на демо‑БД bookings (Postgres) и слой STG в Greenplum.
|
||||
На этом этапе достаточно бегло посмотреть на структуру и понять общую идею, детальная проработка пойдёт позже.
|
||||
|
||||
### 2.1. Знакомство с демо‑БД bookings
|
||||
|
||||
1. Прочитайте `bookings/README.md` — какие сервисы и команды относятся к демобазе.
|
||||
2. Поднимите стенд и выполните:
|
||||
- `make up`
|
||||
- `make bookings-init`
|
||||
3. Подключитесь к демобазе:
|
||||
- `make bookings-psql`
|
||||
- посмотрите таблицы в схеме `bookings` (например, `\dt bookings.*`).
|
||||
4. Найдите таблицу `bookings.bookings` и посмотрите на её структуру:
|
||||
- какие типы колонок используются;
|
||||
- какие поля выглядят как ключи, даты, суммы.
|
||||
|
||||
### 2.2. Знакомство с STG в Greenplum
|
||||
|
||||
1. Прочитайте `sql/stg/bookings_ddl.sql` и мини‑README `docs/internal/bookings_stg_readme.md` (если интересно — `docs/internal/bookings_stg_design.md`).
|
||||
2. Ответьте себе на вопросы:
|
||||
- чем внешняя таблица `stg.bookings_ext` отличается от внутренней `stg.bookings`;
|
||||
- зачем нужны тех.колонки `src_created_at_ts`, `load_dttm`, `batch_id`;
|
||||
- чем слой STG отличается от итоговых витрин (DDS/DM) с точки зрения моделирования.
|
||||
3. Выполните `make ddl-gp`, затем зайдите в Greenplum (`make gp-psql`) и проверьте наличие схемы и таблиц:
|
||||
- `\dn` и `\dt stg.*`
|
||||
- `SELECT * FROM stg.bookings LIMIT 5;` (после запуска соответствующего DAG).
|
||||
|
||||
### 2.3. Как генерируются учебные данные bookings
|
||||
|
||||
1. Откройте файл `bookings/generate_next_day.sql` и ответьте себе на вопросы:
|
||||
- с какой даты начинается генерация данных (посмотрите на GUC `bookings.start_date` и переменную `v_start_cfg`);
|
||||
- сколько дней генерируется при первой установке (переменная `bookings.init_days`);
|
||||
- что происходит, если таблица `bookings.bookings` уже не пустая.
|
||||
2. В демобазе (`make bookings-psql`) выполните:
|
||||
- `SELECT min(book_date), max(book_date) FROM bookings.bookings;`
|
||||
- затем запустите `make bookings-generate-day` и повторите запрос — как изменился максимальный день?
|
||||
3. Откройте `sql/src/bookings_generate_day_if_missing.sql` и обратите внимание, что:
|
||||
- логическая дата запуска DAG (`{{ ds }}`) не влияет на выбор дня генерации;
|
||||
- скрипт всегда смотрит на `max(book_date)` и добавляет **следующий** день (или несколько стартовых дней, если база пуста).
|
||||
4. Сделайте вывод: генератор всегда «шагает» по датам вперёд от максимальной даты, поэтому:
|
||||
- при `make bookings-init` вы получаете `BOOKINGS_INIT_DAYS` дней начиная с `BOOKINGS_START_DATE`;
|
||||
- при последующих вызовах (`make bookings-generate-day` или DAG) добавляется ровно один новый день.
|
||||
|
||||
---
|
||||
|
||||
## 3. DAG bookings_to_gp_stage (заготовка заданий)
|
||||
|
||||
Этот DAG показывает путь данных от демо‑БД bookings в Postgres до сырого слоя STG в Greenplum.
|
||||
Сейчас он уже реализован как учебный пример, а в будущем вокруг него появятся отдельные задания по моделированию DWH.
|
||||
|
||||
### 3.1. Что есть сейчас
|
||||
|
||||
1. Откройте `airflow/dags/bookings_to_gp_stage.py`.
|
||||
2. Найдите в коде ссылки на SQL‑файлы:
|
||||
- `sql/src/bookings_generate_day_if_missing.sql`
|
||||
- `sql/stg/bookings_load.sql`
|
||||
- `sql/stg/bookings_dq.sql`
|
||||
3. Соотнесите шаги DAG с документом `docs/internal/bookings_stg_readme.md`:
|
||||
- генерация учебного дня в `bookings.bookings`;
|
||||
- загрузка инкремента в `stg.bookings`;
|
||||
- проверка количества строк между источником и STG.
|
||||
4. Обратите внимание, как в DAG используется логическая дата запуска:
|
||||
- `{{ ds_nodash }}` используется как `batch_id` — метка загрузки в таблице `stg.bookings` для конкретного запуска;
|
||||
- сами даты данных (какие дни есть в `bookings.bookings`) определяются генератором по `max(book_date)`, а не по `ds`.
|
||||
|
||||
На этом этапе достаточно понять общую цепочку. Детальные задания по переработке модели данных и построению ODS/DDS/DM слоёв будут добавлены позже.
|
||||
|
||||
### 3.2. Идеи для будущих заданий (черновик)
|
||||
|
||||
> Ниже — набросок задач, к которым мы вернёмся, когда базовые темы по Airflow и CSV‑pipeline будут освоены.
|
||||
|
||||
Планируемые направления:
|
||||
|
||||
- Спроектировать модель данных для основных сущностей демобазы bookings (рейсы, билеты, перелёты) в слоях ODS/DDS/DM.
|
||||
- Реализовать слой ODS поверх STG, аккуратно работая с временными атрибутами и ключами.
|
||||
- Построить витрины (DM) для типичных аналитических вопросов: загрузка рейсов, выручка по направлениям, динамика бронирований.
|
||||
- Добавить DAG’и, которые используют `stg.bookings` как источник и строят следующие слои DWH.
|
||||
- Расширить проверки качества данных для потоков bookings → STG → витрины.
|
||||
|
||||
Когда будете готовы к этим темам, вернитесь к этому разделу — он станет основой для следующего «модуля» лабораторных заданий.
|
||||
Executable
+38
@@ -0,0 +1,38 @@
|
||||
#!/usr/bin/env bash
|
||||
|
||||
# Простой init-скрипт для настройки PXF:
|
||||
# - кладёт JDBC-драйвер PostgreSQL в каталог PXF;
|
||||
# - копирует jdbc-site.xml для сервера bookings-db;
|
||||
# - выполняет pxf cluster sync.
|
||||
# Скрипт выполняется только при первой инициализации Greenplum (пустой /data).
|
||||
|
||||
set -euo pipefail
|
||||
|
||||
GREENPLUM_USER=${GREENPLUM_USER:-gpadmin}
|
||||
PXF_BASE_DEFAULT="${GREENPLUM_DATA_DIRECTORY:-/data}/pxf"
|
||||
|
||||
# Пробуем подтянуть PXF_BASE из .bashrc пользователя Greenplum
|
||||
if [ -f "/home/${GREENPLUM_USER}/.bashrc" ]; then
|
||||
# shellcheck disable=SC1090
|
||||
source "/home/${GREENPLUM_USER}/.bashrc"
|
||||
fi
|
||||
|
||||
PXF_BASE=${PXF_BASE:-$PXF_BASE_DEFAULT}
|
||||
|
||||
mkdir -p "${PXF_BASE}/lib" "${PXF_BASE}/servers/bookings-db"
|
||||
|
||||
# Копируем JAR-драйвер, если он ещё не установлен
|
||||
if [ -f /pxf-local/postgresql-42.7.3.jar ] && [ ! -f "${PXF_BASE}/lib/postgresql-jdbc.jar" ]; then
|
||||
cp /pxf-local/postgresql-42.7.3.jar "${PXF_BASE}/lib/postgresql-jdbc.jar"
|
||||
fi
|
||||
|
||||
# Копируем jdbc-site.xml для сервера bookings-db, если его нет
|
||||
if [ -f /pxf-local/servers/bookings-db/jdbc-site.xml ] && [ ! -f "${PXF_BASE}/servers/bookings-db/jdbc-site.xml" ]; then
|
||||
cp /pxf-local/servers/bookings-db/jdbc-site.xml "${PXF_BASE}/servers/bookings-db/jdbc-site.xml"
|
||||
fi
|
||||
|
||||
# Синхронизируем конфигурацию PXF (на всякий случай)
|
||||
if command -v pxf >/dev/null 2>&1; then
|
||||
pxf cluster sync || true
|
||||
fi
|
||||
|
||||
Binary file not shown.
@@ -0,0 +1,31 @@
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<!--
|
||||
JDBC-конфиг для PXF-сервера bookings-db.
|
||||
Для учебного стенда используем те же дефолтные креды, что и в .env.example:
|
||||
BOOKINGS_DB_USER=bookings, BOOKINGS_DB_PASSWORD=bookings.
|
||||
При изменении этих параметров в .env не забывайте обновить их здесь.
|
||||
-->
|
||||
<configuration>
|
||||
|
||||
<property>
|
||||
<name>jdbc.driver</name>
|
||||
<value>org.postgresql.Driver</value>
|
||||
</property>
|
||||
|
||||
<property>
|
||||
<name>jdbc.url</name>
|
||||
<value>jdbc:postgresql://bookings-db:5432/demo</value>
|
||||
</property>
|
||||
|
||||
<property>
|
||||
<name>jdbc.user</name>
|
||||
<value>bookings</value>
|
||||
</property>
|
||||
|
||||
<property>
|
||||
<name>jdbc.password</name>
|
||||
<value>bookings</value>
|
||||
</property>
|
||||
|
||||
</configuration>
|
||||
|
||||
@@ -14,3 +14,7 @@ dev = [
|
||||
"psycopg2-binary==2.9.9",
|
||||
"pytest==7.4.4",
|
||||
]
|
||||
|
||||
[tool.isort]
|
||||
profile = "black"
|
||||
src_paths = ["airflow", "tests"]
|
||||
|
||||
@@ -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);
|
||||
+26
-11
@@ -1,12 +1,27 @@
|
||||
-- Greenplum DDL (GPDB 6 совместимо)
|
||||
-- Колонночная таблица (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
|
||||
-- Главный входной DDL-скрипт для Greenplum в учебном стенде.
|
||||
-- Выполняется из контейнера командой `make ddl-gp` и создаёт/обновляет
|
||||
-- все объекты, которые нужны базовым DAG (csv_to_greenplum, bookings_to_gp_stage);
|
||||
-- подключает файловые DDL через \i, чтобы сохранять единый входной скрипт.
|
||||
--
|
||||
-- Чтобы не ломать задания, новые объекты лучше добавлять в отдельные файлы
|
||||
-- и подключать их отсюда, а существующие определения не удалять.
|
||||
--
|
||||
-- Подробнее про STG/bookings: см. docs/internal/bookings_stg_readme.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;
|
||||
CREATE EXTERNAL TABLE public.ext_bookings_bookings (
|
||||
book_ref CHAR(6),
|
||||
book_date TIMESTAMP,
|
||||
total_amount NUMERIC(10,2)
|
||||
)
|
||||
WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1)
|
||||
DISTRIBUTED BY (order_id);
|
||||
LOCATION ('pxf://bookings.bookings?PROFILE=JDBC&SERVER=bookings-db')
|
||||
FORMAT 'CUSTOM' (formatter='pxfwritable_import');
|
||||
|
||||
-- DDL для слоя stg по таблице bookings вынесен в отдельный файл.
|
||||
-- Здесь подключаем его через psql \i, чтобы сохранить единый входной скрипт.
|
||||
\i stg/bookings_ddl.sql
|
||||
|
||||
@@ -0,0 +1,49 @@
|
||||
-- Генерация учебного дня в демо-БД bookings.
|
||||
-- Этот скрипт всегда добавляет следующий учебный день после максимальной даты
|
||||
-- в таблице bookings.bookings (или несколько дней от стартовой даты, если база пуста).
|
||||
-- Логическая дата запуска DAG ({{ ds }}) используется только как метка батча в стейдже.
|
||||
|
||||
DO $$
|
||||
DECLARE
|
||||
v_max_book_date timestamptz;
|
||||
v_start_date timestamptz;
|
||||
v_end_date timestamptz;
|
||||
v_jobs integer := COALESCE(current_setting('bookings.jobs', true), '1')::integer;
|
||||
v_init_days integer := COALESCE(current_setting('bookings.init_days', true), '1')::integer;
|
||||
v_start_cfg text := COALESCE(current_setting('bookings.start_date', true), '2017-01-01');
|
||||
BEGIN
|
||||
-- Проверяем, что демобаза установлена
|
||||
IF to_regclass('bookings.bookings') IS NULL THEN
|
||||
RAISE EXCEPTION 'Таблица bookings.bookings не найдена. Сначала выполните make bookings-init.';
|
||||
END IF;
|
||||
|
||||
-- Ищем последнюю сгенерированную дату
|
||||
SELECT max(book_date) INTO v_max_book_date FROM bookings.bookings;
|
||||
|
||||
IF v_max_book_date IS NULL THEN
|
||||
-- База пустая: берём стартовую дату из конфигурации (или дефолтную)
|
||||
v_start_date := date_trunc('day', v_start_cfg::timestamptz);
|
||||
ELSE
|
||||
-- Продолжаем с дня, следующего за максимальной датой
|
||||
v_start_date := date_trunc('day', v_max_book_date) + interval '1 day';
|
||||
END IF;
|
||||
|
||||
-- Первая генерация вызывает generate(), последующие — continue()
|
||||
IF v_max_book_date IS NULL THEN
|
||||
v_end_date := v_start_date + (v_init_days || ' days')::interval;
|
||||
CALL generate(v_start_date, v_end_date, v_jobs);
|
||||
ELSE
|
||||
v_end_date := v_start_date + interval '1 day';
|
||||
CALL continue(v_end_date, v_jobs);
|
||||
END IF;
|
||||
|
||||
RAISE NOTICE 'Сгенерированы данные в bookings.bookings за интервал [% - %).',
|
||||
date_trunc('day', v_start_date),
|
||||
date_trunc('day', v_end_date);
|
||||
|
||||
-- Ждём завершения фоновых джобов генератора, чтобы данные успели записаться
|
||||
WHILE busy() LOOP
|
||||
PERFORM pg_sleep(1);
|
||||
END LOOP;
|
||||
PERFORM dblink_disconnect(unnest(dblink_get_connections()));
|
||||
END $$;
|
||||
@@ -0,0 +1,29 @@
|
||||
-- DDL для слоя STG по таблице bookings.
|
||||
-- Используется как из общего скрипта ddl_gp.sql (через \i),
|
||||
-- так и может выполняться отдельно при изменении схемы.
|
||||
|
||||
-- Схема stg для сырого слоя DWH.
|
||||
CREATE SCHEMA IF NOT EXISTS stg;
|
||||
|
||||
-- Внешняя таблица в схеме stg для чтения данных из bookings.bookings через PXF.
|
||||
DROP EXTERNAL TABLE IF EXISTS stg.bookings_ext;
|
||||
CREATE EXTERNAL TABLE stg.bookings_ext (
|
||||
book_ref CHAR(6),
|
||||
book_date TIMESTAMP,
|
||||
total_amount NUMERIC(10,2)
|
||||
)
|
||||
LOCATION ('pxf://bookings.bookings?PROFILE=JDBC&SERVER=bookings-db')
|
||||
FORMAT 'CUSTOM' (formatter='pxfwritable_import');
|
||||
|
||||
-- Внутренняя таблица stg.bookings — сырой слой, все бизнес-колонки как TEXT.
|
||||
CREATE TABLE IF NOT EXISTS stg.bookings (
|
||||
book_ref TEXT,
|
||||
book_date TEXT,
|
||||
total_amount TEXT,
|
||||
src_created_at_ts TIMESTAMP,
|
||||
load_dttm TIMESTAMP NOT NULL DEFAULT now(),
|
||||
batch_id TEXT NOT NULL
|
||||
)
|
||||
WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1)
|
||||
DISTRIBUTED BY (book_ref);
|
||||
|
||||
@@ -0,0 +1,43 @@
|
||||
-- Проверка количества строк между источником stg.bookings_ext и стейджем stg.bookings.
|
||||
-- Считаем строки за то же окно инкремента, что и при загрузке:
|
||||
-- все записи во внешней таблице с book_date больше максимального src_created_at_ts
|
||||
-- из предыдущих батчей должны совпасть по количеству со строками текущего batch_id.
|
||||
|
||||
DO $$
|
||||
DECLARE
|
||||
v_batch_id text := '{{ ds_nodash }}'::text;
|
||||
v_prev_ts timestamp;
|
||||
v_src_count bigint;
|
||||
v_stg_count bigint;
|
||||
BEGIN
|
||||
-- Опорная метка: максимум src_created_at_ts среди предыдущих батчей
|
||||
SELECT max(src_created_at_ts)
|
||||
INTO v_prev_ts
|
||||
FROM stg.bookings
|
||||
WHERE batch_id <> v_batch_id
|
||||
OR batch_id IS NULL;
|
||||
|
||||
-- Источник: считаем строки во внешней таблице, которые вошли в новое окно
|
||||
SELECT COUNT(*)
|
||||
INTO v_src_count
|
||||
FROM stg.bookings_ext
|
||||
WHERE book_date > COALESCE(v_prev_ts, TIMESTAMP '1900-01-01 00:00:00');
|
||||
|
||||
-- Считаем строки, реально вставленные в stg.bookings в этом батче
|
||||
SELECT COUNT(*)
|
||||
INTO v_stg_count
|
||||
FROM stg.bookings
|
||||
WHERE batch_id = v_batch_id;
|
||||
|
||||
IF v_src_count <> v_stg_count THEN
|
||||
RAISE EXCEPTION
|
||||
'Несовпадение количества строк при загрузке bookings: источник=%, stg=%. Проверьте окно инкремента и логи задач загрузки.',
|
||||
v_src_count,
|
||||
v_stg_count;
|
||||
END IF;
|
||||
|
||||
RAISE NOTICE
|
||||
'Проверка количества строк пройдена: источник=%, stg=%',
|
||||
v_src_count,
|
||||
v_stg_count;
|
||||
END $$;
|
||||
@@ -0,0 +1,30 @@
|
||||
-- Загрузка инкремента из stg.bookings_ext в stg.bookings.
|
||||
-- Окно инкремента определяется по src_created_at_ts:
|
||||
-- берём строки, где book_date больше максимального src_created_at_ts
|
||||
-- среди "старых" батчей; верхняя граница по дате не используется.
|
||||
|
||||
INSERT INTO stg.bookings (
|
||||
book_ref,
|
||||
book_date,
|
||||
total_amount,
|
||||
src_created_at_ts,
|
||||
load_dttm,
|
||||
batch_id
|
||||
)
|
||||
SELECT
|
||||
book_ref::text,
|
||||
book_date::text,
|
||||
total_amount::text,
|
||||
book_date::timestamp,
|
||||
now(),
|
||||
'{{ ds_nodash }}'::text
|
||||
FROM stg.bookings_ext
|
||||
WHERE book_date > COALESCE(
|
||||
(
|
||||
SELECT max(src_created_at_ts)
|
||||
FROM stg.bookings
|
||||
WHERE batch_id <> '{{ ds_nodash }}'::text
|
||||
OR batch_id IS NULL
|
||||
),
|
||||
TIMESTAMP '1900-01-01 00:00:00'
|
||||
);
|
||||
@@ -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",
|
||||
|
||||
Reference in New Issue
Block a user