Author SHA1 Message Date
ddadmin 272013774b Сохранение логов airflow между перезапусками 2026-01-09 13:37:14 +03:00
ddadmin c3b2050cc4 еще баги и документирование сделанного 2026-01-09 12:41:36 +03:00
ddadmin 5668edeb7b Диагностика проблемы 2026-01-09 00:33:45 +03:00
ddadmin 810d3d4727 дальшейшая отладка 2026-01-08 23:51:30 +03:00
ddadmin 573c853f15 фикс pxf 2026-01-08 23:34:26 +03:00
ddadmin 1efeb58449 диагностика проблемы с pxf 2026-01-08 23:14:53 +03:00
ddadmin f81a831c0c патч 2026-01-08 22:52:56 +03:00
ddadmin 4751674ef3 планы по доработке 2026-01-08 22:41:11 +03:00
ddadmin dfcc456120 Документирование решения 2026-01-08 22:24:24 +03:00
ddadmin 5778ad6601 фикс бага pxf 2026-01-08 22:02:57 +03:00
ddadmin 33fc42f883 Merge branch 'chore/airflow-dockerfile' 2026-01-06 23:31:27 +03:00
ddadmin 1173e0e53f Отладка airflow dockerfile 2026-01-06 23:20:41 +03:00
ddadmin e27476fa5f Раздел "Благодарности" 2026-01-06 15:01:26 +03:00
ddadmin 4cb11b4769 feat: use custom docker image for airflow to avoid runtime pip install 2025-12-13 22:21:41 +03:00
ddadmin 9f028d0728 Доработки по full smoke тесту 2025-12-11 11:09:05 +03:00
ddadmin 0d81096ded полный smoke тест стенда 2025-12-11 10:51:39 +03:00
ddadmin 31171f84bc Баг генерации bookings 2025-12-11 10:34:16 +03:00
ddadmin 90aafd7b1d Фикс прав доступа у data раздела 2025-12-11 10:05:21 +03:00
ddadmin c27693b07f Планы по интеграции psql 2025-12-11 10:04:58 +03:00
ddadmin 8f92f2e890 раздел с благодарностями в ридми 2025-12-11 00:23:57 +03:00
ddadmin 98e0c90da0 Merge branch 'feature/bookings_db' 2025-12-11 00:01:41 +03:00
ddadmin 342f9bdb62 даг наконец-то работает 2025-12-10 23:49:01 +03:00
ddadmin 359245366a Способ тестирования даг 2025-12-10 23:43:35 +03:00
ddadmin d936f204bd Ошибка с вложенной транзакцией 2025-12-10 23:40:27 +03:00
ddadmin a8ac8f62dc проблемы с запуском GP 2025-12-10 23:32:45 +03:00
ddadmin 3b24d13908 фикс создания таблиц в make 2025-12-10 23:31:34 +03:00
ddadmin c0bb24764e фикс ошибки с etl 2025-12-10 23:23:34 +03:00
ddadmin 409a9f6f31 мысли по todo 2025-12-10 23:15:12 +03:00
ddadmin dfd6eca760 Удобная команда остановки стенда без удаления 2025-12-10 23:04:45 +03:00
ddadmin daed2313e7 Фикс создания connections 2025-12-10 22:58:19 +03:00
ddadmin 8fb78f9086 fix: добавленны airflow connections и обновлена структура запуска 2025-12-10 22:52:08 +03:00
ddadmin fbf9b118e6 фикс вфп инициализации ddl в gp 2025-12-10 22:31:12 +03:00
ddadmin 6e2a43a900 fix root issue for airflow-init 2025-12-10 22:15:49 +03:00
ddadmin 250af00b2f refactor(bookings): упростить работу с датами в DAG и SQL 2025-12-10 21:56:12 +03:00
ddadmin 29490828f4 Уход от запуска DAG за конкретную дату.
Теперь новый запуск генерирует и переливает новый, следующий день
2025-12-10 21:53:17 +03:00
ddadmin 7a9b2d182e Принесен isort 2025-12-10 18:49:49 +03:00
ddadmin c865fa282e feat(dags): add ddl dag for gp and rename csv dq 2025-12-10 18:41:00 +03:00
ddadmin cf95cf95d2 пояснения назначения файла 2025-12-10 17:36:38 +03:00
ddadmin c63084be09 ДОравботка структуры readme для лучшей читабильности 2025-12-10 17:33:55 +03:00
ddadmin fe4bb4606f Как сделать новый день 2025-12-10 17:13:08 +03:00
ddadmin cacf989a9c вынесли sql из dag в папку sql 2025-12-10 16:59:50 +03:00
ddadmin 2c8f5c1e23 Учебный даг первоначально протестирован 2025-12-09 11:22:56 +03:00
ddadmin 196171e3d6 Склелет учебного даг 2025-12-09 10:43:43 +03:00
ddadmin ae7f909f5d Проработка задачи на ETL 2025-12-09 10:39:47 +03:00
ddadmin c3d7759e6b Переименование БД GP в gp_dwh 2025-12-09 10:09:46 +03:00
ddadmin 4515577b4a Смена внешнего порта GreenPlum 2025-12-09 09:57:52 +03:00
ddadmin 2bdeb17cee Тестирование PXF 2025-12-09 09:52:50 +03:00
ddadmin e3b0fd9ef9 Первоначальная реализация PXF (не тестировано) 2025-12-09 09:42:52 +03:00
ddadmin 86435c282d Уточнения планов 2025-12-09 09:35:31 +03:00
ddadmin dbd2406765 Планы по включению PXF для GP 2025-12-09 00:09:33 +03:00
ddadmin 8bc8e5cb3a Улучшение документации по bookings 2025-12-08 23:34:39 +03:00
ddadmin b56ea952e7 Подключение bookings db заработало 2025-12-08 23:31:13 +03:00
ddadmin c17baf447d bookings debug 2025-12-08 22:08:07 +03:00
ddadmin c95dab6478 Доработки по генерации БД - в процессе 2025-12-08 18:26:10 +03:00
40 changed files with 2660 additions and 136 deletions
+14 -1
View File
@@ -1,3 +1,6 @@
# Общие настройки
TZ=Europe/Moscow
# Airflow Configuration
AIRFLOW_USER=admin
AIRFLOW_PASSWORD=admin
@@ -15,16 +18,26 @@ 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
GP_USE_AIRFLOW_CONN=true
# PXF (optional)
# PXF_SEED_OVERWRITE=1 — принудительно перезаписать seed-файлы из образа в PXF_BASE (полезно после правок в pxf/)
PXF_SEED_OVERWRITE=0
# PXF_SYNC_ON_START=1 — выполнять `pxf cluster sync` при старте контейнера (дольше, но гарантирует актуальные конфиги)
PXF_SYNC_ON_START=0
# CSV Pipeline
CSV_DIR=/opt/airflow/data
CSV_ROWS=1000
+1
View File
@@ -7,3 +7,4 @@ __pycache__/
*.pyc
.venv/
data/
bookings/demodb/
+24 -6
View File
@@ -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`.
+19
View File
@@ -0,0 +1,19 @@
# Apache Airflow с дополнительными зависимостями для Greenplum
FROM apache/airflow:2.9.2
LABEL maintainer="your-email@example.com"
LABEL description="Airflow with Greenplum and Pandas dependencies"
LABEL version="1.0"
# Копируем requirements в стандартное расположение
COPY airflow/requirements.txt /opt/airflow/requirements.txt
# Устанавливаем зависимости
RUN pip install --no-cache-dir -r /opt/airflow/requirements.txt \
&& pip check \
&& rm -rf /tmp/pip-*
# Переключаемся на пользователя airflow
USER airflow
WORKDIR /opt/airflow
+16
View File
@@ -0,0 +1,16 @@
FROM woblerr/greenplum:6.27.1
ENV PXF_SEED_DIR=/opt/pxf-seed \
PXF_SEED_OVERWRITE=0 \
PXF_SYNC_ON_START=0
RUN mkdir -p /opt/pxf-seed/servers/bookings-db /opt/pxf-scripts
COPY pxf/postgresql-42.7.3.jar /opt/pxf-seed/postgresql-42.7.3.jar
COPY pxf/servers/bookings-db/jdbc-site.xml /opt/pxf-seed/servers/bookings-db/jdbc-site.xml
COPY pxf/init/10_pxf_bookings.sh /opt/pxf-scripts/ensure_pxf_bookings.sh
COPY pxf/init/start_greenplum_with_pxf.sh /start_greenplum_with_pxf.sh
RUN chmod 755 /opt/pxf-scripts/ensure_pxf_bookings.sh /start_greenplum_with_pxf.sh
CMD ["/start_greenplum_with_pxf.sh"]
+70 -6
View File
@@ -1,15 +1,30 @@
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 \
.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
dev-setup dev-sync dev-lock test lint fmt clean-venv build
SHELL := /bin/bash
up:
docker compose -f docker-compose.yml up -d
build:
docker compose -f docker-compose.yml build
stop:
docker compose -f docker-compose.yml stop
down:
docker compose -f docker-compose.yml down
clean:
docker compose -f docker-compose.yml down -v
airflow-init:
@@ -19,25 +34,71 @@ 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
test -d bookings/demodb || git clone --depth 1 https://github.com/postgrespro/demodb.git bookings/demodb
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 "Job 1 (local): ok" bookings/demodb/engine.sql; then \
if ! patch -d bookings/demodb -p1 --forward < bookings/patches/engine_jobs1_sync.patch; then \
echo "Не удалось применить патч engine_jobs1_sync.patch. Удалите bookings/demodb и повторите make bookings-init." >&2; \
exit 1; \
fi; \
fi
# Делаем установку идемпотентной и принудительной: DROP DATABASE IF EXISTS demo WITH (FORCE)
if ! grep -q "DROP DATABASE IF EXISTS demo WITH (FORCE);" bookings/demodb/install.sql; then \
if ! patch -d bookings/demodb -p1 --forward < bookings/patches/install_drop_if_exists.patch; then \
echo "Не удалось применить патч install_drop_if_exists.patch. Удалите bookings/demodb и повторите make bookings-init." >&2; \
exit 1; \
fi; \
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 'PGPASSWORD="$$POSTGRES_PASSWORD" psql -v ON_ERROR_STOP=1 -U "$$POSTGRES_USER" -d demo -v start_date="$${BOOKINGS_START_DATE:-2017-01-01}" -f /bookings/generate_next_day.sql'
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)
@@ -63,3 +124,6 @@ fmt:
clean-venv:
python -c "import shutil; shutil.rmtree('.venv', ignore_errors=True)"
e2e-smoke:
./scripts/e2e_smoke.sh
+218 -25
View File
@@ -4,6 +4,8 @@
Добро пожаловать в учебный стенд для изучения основ Data Engineering! Этот проект поможет вам освоить ключевые инструменты современных data pipeline: **Airflow** для оркестрации, **pandas/CSV** для подготовки данных, **Postgres** с демобазой **bookings** как источник и **Greenplum** как аналитическую базу данных.
Если вы проходите стенд как серию лабораторных, смотрите также файл с заданиями: `educational-tasks.md`.
## 🎯 Что вы узнаете
- Как настроить локальный стек данных с помощью Docker
@@ -18,11 +20,11 @@
- Установите 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» ниже.
@@ -63,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 # Показать таблицы
@@ -80,16 +82,19 @@ CSV-файлы после выполнения DAG остаются в дире
Если после изменений что‑то «сломалось»:
```bash
make down # Остановить и стереть данные в контейнерах
make down # Остановить и удалить контейнеры/сети (volumes сохраняются)
make up
```
Это помогает, когда Greenplum не стартует из‑за «грязной» остановки и внутренних файлов.
Если проблема связана с «грязной» остановкой и данными в томах (например, Greenplum не стартует),
используйте полный reset: `make clean && make up` (данные в Docker-томах будут потеряны).
---
## 🛠️ Подробная настройка (для уверенных пользователей)
> Если вы впервые запускаете стенд, этот раздел можно пролистать и вернуться к нему позже.
### Установка Make (опционально)
Для удобства работы с проектом рекомендуем установить `make`:
@@ -103,27 +108,72 @@ 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`
---
### Greenplum + PXF: свой образ
Чтобы избежать проблем с правами и нестабильных запусков, Greenplum собирается
из собственного `Dockerfile.greenplum`. В образ вшиты:
- JDBC-драйвер PostgreSQL;
- конфигурация PXF-сервера `bookings-db`;
- ensure-скрипт, который при каждом старте контейнера докладывает файлы в `PXF_BASE`.
Дополнительно при старте контейнера:
- базовые конфиги PXF копируются в `PXF_BASE/conf` (если их ещё нет);
- создаются каталоги `PXF_BASE/run` и `PXF_BASE/logs`;
- `CREATE EXTENSION pxf` выполняется автоматически, когда Greenplum становится доступен (с ретраями).
Healthcheck сервиса `greenplum` учитывает не только готовность Greenplum, но и запуск PXF,
а также наличие `extension pxf` — это нужно, чтобы Airflow не стартовал раньше PXF.
Сборка и запуск:
- `make build` — собрать образ (явно);
- `make up` — поднимет стек и соберёт образ, если он ещё не создан.
Обновление PXF-конфигов:
- изменили файлы в `pxf/` → выполните `make build` и перезапустите контейнер;
- для принудительной перезаписи файлов в `PXF_BASE` используйте `PXF_SEED_OVERWRITE=1`;
- для принудительного `pxf cluster sync` при старте используйте `PXF_SYNC_ON_START=1`.
Проверка PXF:
- статус: `docker compose exec greenplum bash -lc "su - gpadmin -c '/usr/local/pxf/bin/pxf cluster status'"`;
- логи: `greenplum_data:/data/pxf/logs` (внутри контейнера — `/data/pxf/logs`).
---
### Локальное окружение разработчика
Локальным окружением управляет [uv](https://docs.astral.sh/uv/) — он скачивает нужный Python и создаёт `.venv` на основе `pyproject.toml` / `uv.lock`.
@@ -166,25 +216,60 @@ uv run black --check airflow tests
- **pandas** — библиотека для генерации и анализа данных в формате CSV
### Готовые DAG (workflow)
- **orders_base_ddl** — создаёт базовую таблицу `public.orders` для CSV‑пайплайна
- **bookings_stg_ddl** — готовит схему `stg` и таблицы `stg.bookings_ext` / `stg.bookings`
- **csv_to_greenplum** — базовый pipeline: pandas → CSV → Greenplum
- **greenplum_data_quality** — проверки качества данных (наличие таблицы, схема, дубликаты)
- **bookings_to_gp_stage** — пример загрузки из демо‑БД bookings в слой STG
- **csv_to_greenplum_dq** — проверки качества данных (наличие таблицы, схема, дубликаты)
> Учебный путь — триггернуть DDL‑DAG: для CSV `orders_base_ddl`, для bookings `bookings_stg_ddl`. Технический шорткат для быстрой инициализации — `make ddl-gp` (он не вызывается автоматически при старте контейнеров).
### Полезные команды
```bash
# Основные команды
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
make bookings-generate-day # Добавить ещё один день данных в bookings
make bookings-init # Установить демобазу bookings в Postgres (по умолчанию генерирует 1 день)
make bookings-generate-day # Добавить ещё один день данных в bookings (можно вызвать несколько раз)
make bookings-psql # Подключиться к демобазе bookings (БД demo)
# Проверка данных
make logs # Следить за логами Airflow
# Логи задач Airflow сохраняются в Docker-томе `airflow_logs`
# и переживают `docker compose down`/`up` (удаляются при `docker compose down -v` / `make clean`).
# Контроль генерации 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 (подзапросы в аргументах CALL не работают, поэтому через DO-блок):
```sql
DO $$
DECLARE
v_next_day timestamptz;
BEGIN
SELECT date_trunc('day', max(book_date)) + interval '1 day'
INTO v_next_day
FROM bookings.bookings;
CALL continue(v_next_day); -- или CALL continue(v_next_day, 4) для параллельности
END $$;
```
- Не вызывайте `CALL generate(...)` поверх существующих данных: она делает TRUNCATE и создаёт демобазу заново.
---
## ⚙️ Настройка через переменные окружения
@@ -194,15 +279,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_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`)
@@ -215,6 +302,8 @@ make logs # Следить за логами Airflow
## 🔍 Продвинутые темы
> Этот раздел не обязателен при первом прохождении стенда; к нему удобно вернуться, когда базовый CSV‑pipeline уже понятен.
### Архитектура pipeline
**Поток данных в DAG `csv_to_greenplum`:**
@@ -227,12 +316,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 режиме (для обучения)
@@ -241,15 +356,75 @@ 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`) |
| `database "demo" does not exist` в bookings‑DAG | Вы сделали полный reset с удалением томов (`docker compose down -v` / `make clean`), поэтому демобаза bookings не установлена. Запустите `make bookings-init` и повторите DAG. |
| Ошибка подключения к 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`. Если не помогло — полный reset: `make clean && make up` (удалит тома). |
| `protocol "pxf" does not exist` | Перезапустите `greenplum` и повторите `bookings_stg_ddl`/`make ddl-gp` — расширение `pxf` создаётся автоматически при старте контейнера. |
| PXF не отвечает (Connection refused к порту 5888) | Проверьте `pxf cluster status` в контейнере `greenplum` и перезапустите сервис `greenplum`. |
| PXF не подхватывает изменения конфигов | Пересоберите образ (`make build`) и перезапустите `greenplum`. Для принудительной перезаписи файлов задайте `PXF_SEED_OVERWRITE=1`. |
| 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. |
---
@@ -257,15 +432,24 @@ make logs # Следить за логами Airflow
```
├── docker-compose.yml # Описание всех сервисов
├── Dockerfile.airflow # Образ Airflow с зависимостями
├── Dockerfile.greenplum # Образ Greenplum с интегрированным PXF
├── .env.example # Шаблон настроек
├── Makefile # Удобные команды для работы
├── README.md # Обзор стенда
├── TESTING.md # Пошаговый план проверки
├── educational-tasks.md # Учебные задания для менти
├── airflow/
│ └── dags/ # Файлы workflow (DAG)
│ ├── csv_to_greenplum.py
── data_quality_greenplum.py
├── bookings/ # Скрипты и вспомогательные файлы для демобазы bookings в Postgres
── 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
```
---
@@ -273,7 +457,7 @@ make logs # Следить за логами Airflow
## 💡 Советы для дальнейшего обучения
1. **Поэкспериментируйте с DAG** — измените параметры генерации данных или размер батча
2. **Добавьте свои проверки** — расширьте DAG `data_quality_greenplum.py`
2. **Добавьте свои проверки** — расширьте DAG `csv_to_greenplum_dq.py`
3. **Попробуйте другие источники** — замените генератор данных на чтение из файла или API
4. **Изучите Airflow deeper** — добавьте зависимости между задачами, настройте расписания
@@ -283,5 +467,14 @@ make logs # Следить за логами Airflow
- Локальные проверки: `make test` (pytest). Для форматирования — `make fmt`, для проверки — `make lint`.
- Пошаговый сценарий с Docker (включая негативные кейсы и reset) — см. `TESTING.md`.
- Полный smoke-тест стенда (сносит volumes!): `make e2e-smoke` — поднимает стек с нуля, прогоняет `csv_to_greenplum` и `bookings_to_gp_stage` через `airflow dags test` и проверяет, что в `public.orders` и `stg.bookings` появились строки.
---
## Благодарности
- **Postgres Pro** — за демо-БД bookings и генератор данных `demodb`: https://github.com/postgrespro/demodb (лицензия MIT: https://github.com/postgrespro/demodb/blob/main/LICENSE).
- **woblerr** — за Docker-сборку Greenplum: https://github.com/woblerr/docker-greenplum (образ: `woblerr/greenplum`, лицензия MIT: https://github.com/woblerr/docker-greenplum/blob/master/LICENSE).
Удачи в изучении Data Engineering! 🚀
+31 -14
View File
@@ -12,48 +12,65 @@
## 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`); `greenplum` считается `healthy` только когда поднят и Greenplum, и PXF.
- `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` (установка демобазы `demo` в контейнере `bookings-db`) и `make ddl-gp` (создаёт `stg.bookings_ext` и `stg.bookings` в Greenplum);
- важно: DAG `bookings_stg_ddl` **не** создаёт базу `demo` в `bookings-db`; если вы делали `docker compose down -v` / `make clean`, `make bookings-init` обязателен;
- включить 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 airflow-webserver airflow dags test bookings_to_gp_stage 2024-01-01` — прогоняет `bookings_to_gp_stage` целиком в «off-line» режиме;
- `docker compose -f docker-compose.yml exec airflow-webserver airflow dags trigger bookings_to_gp_stage` — создаёт реальный запуск DAG (логи и статус можно смотреть либо через UI, либо командой `airflow tasks list`/`airflow tasks logs` внутри контейнера).
## 5. Проверка данных в Greenplum
- `make gp-psql` — запустить psql в контейнере от имени `gpadmin`.
- (опционально) Проверить, что PXF действительно запущен:
- `docker compose exec greenplum bash -lc "su - gpadmin -c '/usr/local/pxf/bin/pxf cluster status'"`
- Команды внутри psql:
- `\dt public.*` — таблицы схему public.
- `SELECT COUNT(*) FROM public.orders;` — оценка объёма.
- `SELECT * FROM public.orders LIMIT 5;` — визуальная проверка.
- `SELECT order_id FROM public.orders GROUP BY 1 HAVING COUNT(*) > 1;` — поиск дублей.
- (после настройки PXF) `SELECT COUNT(*) FROM public.ext_bookings_bookings;` — проверка чтения из демо-БД bookings через PXF.
- (после настройки PXF) `SELECT * FROM public.ext_bookings_bookings LIMIT 5;` — визуальное сравнение с таблицей `bookings.bookings` в исходной БД.
- Завершить `\q`.
## 6. Негативные сценарии и fallback
- **Пустая таблица**: запустить `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 +84,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-стенд не запускался в рамках этой сессии; ожидается, что инструкции выше обеспечат полноценную проверку.
+41
View File
@@ -0,0 +1,41 @@
# 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 и шагах по их устранению.
- диагностика текущего кейса: `docs/internal/pxf_bookings.md` (раздел «Известная проблема»).
- Разобраться с генератором demodb:
- после `make bookings-init` таблица `bookings.bookings` остаётся пустой;
- патчи `bookings/patches/engine_jobs1_sync.patch` и `bookings/patches/install_drop_if_exists.patch`
падают при применении (hunk failed / garbage in patch);
- из‑за этого DAG `bookings_to_gp_stage` валится на проверках (источник пустой).
- план: `plans/bookings-demodb-bugfix-plan.md`
- Добавить раздел «Благодарности» в `README.md`:
- явно поблагодарить Postgres Pro за демо‑БД bookings (репозиторий `postgrespro/demodb`);
- указать автора Docker‑сборки Greenplum (`woblerr/docker-greenplum`, образ `woblerr/greenplum`);
- при необходимости сослаться на соответствующие лицензии/README исходных проектов.
- Добавить в образ Airflow установку `psql`, чтобы тестировать загрузку CSV из CLI внутри контейнера (без root и дополнительных зависимостей на хосте).
+32
View File
@@ -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",
)
+87
View File
@@ -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(
+32
View File
@@ -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",
)
+5 -2
View File
@@ -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)
+30
View File
@@ -9,3 +9,33 @@
Основные команды см. в корневом `Makefile` (`bookings-init`, `bookings-generate-day`, `bookings-psql`) и в `README.md` проекта.
## Источник и версия
- Репозиторий демобазы: `postgrespro/demodb`.
- Закреплённый коммит: `d68de192850237719f09b47688d5f3fc94653ca6` (см. `DEMODB_COMMIT` в корневом `Makefile`).
## Что мы патчим в demodb
- `install.sql`: `DROP DATABASE IF EXISTS demo WITH (FORCE)` — установка не падает, даже если демобазу держат активные сессии (например, из Airflow).
- `engine.sql`: два изменения в `engine_jobs1_sync.patch`:
- `busy()` игнорирует свой `pid`, чтобы не считать собственное подключение занятым;
- `continue()` при `jobs=1` вызывает `process_queue` синхронно (без `dblink`), иначе генерация обрывается при выходе из `psql` и данных не появляется.
- Патчи применяются автоматически в `make bookings-init`. Если что-то пошло не так, их можно накатить вручную:
```
patch -d bookings/demodb -p1 --forward < bookings/patches/install_drop_if_exists.patch
patch -d bookings/demodb -p1 --forward < bookings/patches/engine_jobs1_sync.patch
```
## Быстрая проверка после init/обновления
- `make bookings-init` должен завершиться без ошибок; в `bookings.bookings` ожидаем >0 строк (примерно 15k).
- `make bookings-generate-day` добавляет следующий день после `max(book_date)`.
- Ручной вызов генерации из psql/DBeaver — только через DO-блок (подзапрос в аргументах `CALL` не работает):
```sql
DO $$
DECLARE
v_next_day timestamptz;
BEGIN
SELECT date_trunc('day', max(book_date)) + interval '1 day'
INTO v_next_day
FROM bookings.bookings;
CALL continue(v_next_day); -- или CALL continue(v_next_day, 4)
END $$;
```
+23 -9
View File
@@ -1,10 +1,13 @@
-- Генерация ещё одного дня данных в демобазе bookings.
-- Если база пуста, берём стартовую дату из параметра :start_date (формат YYYY-MM-DD).
-- Если база пуста, берём стартовую дату из GUC bookings.start_date (по умолчанию: 2017-01-01).
DO $$
DECLARE
v_max_book_date timestamptz;
v_start_date timestamptz;
v_end_date timestamptz;
v_bookings_cnt bigint;
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
@@ -15,20 +18,31 @@ BEGIN
SELECT max(book_date) INTO v_max_book_date FROM bookings.bookings;
IF v_max_book_date IS NULL THEN
-- База пустая: берём стартовую дату из psql-переменной
v_start_date := date_trunc('day', :'start_date'::timestamptz);
-- База пустая: берём стартовую дату из конфигурации (или дефолтную)
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;
v_end_date := v_start_date + interval '1 day';
-- Первая генерация вызывает generate(), последующие — continue()
IF v_max_book_date IS NULL THEN
CALL generate(v_start_date, v_end_date, 1);
v_end_date := v_start_date + (v_init_days || ' days')::interval;
CALL generate(v_start_date, v_end_date, v_jobs);
ELSE
CALL continue(v_end_date, 1);
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()));
-- Если данных нет, останавливаемся с понятной ошибкой
SELECT COUNT(*) INTO v_bookings_cnt FROM bookings.bookings;
IF v_bookings_cnt = 0 THEN
RAISE EXCEPTION 'Генератор demodb завершился, но bookings.bookings пустая. Проверьте применение патчей и логи генератора.';
END IF;
END $$;
+45
View File
@@ -0,0 +1,45 @@
--- a/engine.sql
+++ b/engine.sql
@@ -165,16 +165,22 @@
-- disconnect all previously opened connections
PERFORM dblink_disconnect(unnest(dblink_get_connections()));
- -- start parallel jobs
- FOR i IN 1 .. jobs LOOP
- connname := 'job' || i;
- PERFORM dblink_connect(connname, current_setting('gen.connstr'));
- res := CASE dblink_send_query(connname, format('CALL process_queue(%L)',end_date))
- WHEN 1 THEN 'ok' ELSE 'FAILED'
- END CASE;
- CALL log_message(0, format('Job %s (connname=%s): %s', i, connname, res));
- RAISE NOTICE 'Starting job %: %', i, res;
- END LOOP;
+ IF jobs = 1 THEN
+ -- jobs=1: синхронно, чтобы не убивать dblink-сессию при выходе из psql
+ CALL process_queue(end_date);
+ CALL log_message(0, 'Job 1 (local): ok');
+ ELSE
+ -- start parallel jobs
+ FOR i IN 1 .. jobs LOOP
+ connname := 'job' || i;
+ PERFORM dblink_connect(connname, current_setting('gen.connstr'));
+ res := CASE dblink_send_query(connname, format('CALL process_queue(%L)',end_date))
+ WHEN 1 THEN 'ok' ELSE 'FAILED'
+ END CASE;
+ CALL log_message(0, format('Job %s (connname=%s): %s', i, connname, res));
+ RAISE NOTICE 'Starting job %: %', i, res;
+ END LOOP;
+ END IF;
-- re-create booking.now()
EXECUTE format(
@@ -230,7 +236,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,5 @@
--- a/install.sql
+++ b/install.sql
@@ -13,1 +13,1 @@
-DROP DATABASE demo;
+DROP DATABASE IF EXISTS demo WITH (FORCE);
+83 -48
View File
@@ -1,3 +1,25 @@
x-airflow-common-env: &airflow-env
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
x-airflow-common-volumes: &airflow-volumes
- ./airflow/dags:/opt/airflow/dags
- ./sql:/sql:ro
- airflow_data:/opt/airflow/data
- airflow_logs:/opt/airflow/logs
x-airflow-common-depends: &airflow-depends
pgmeta:
condition: service_healthy
greenplum:
condition: service_healthy
airflow-init:
condition: service_completed_successfully
services:
# Postgres только для Airflow метаданных
pgmeta:
@@ -5,6 +27,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}
@@ -13,7 +36,7 @@ services:
volumes:
- pgmeta:/var/lib/postgresql/data
healthcheck:
test: ["CMD-SHELL", "pg_isready -U ${PG_USER} -d ${PG_DB}"]
test: [ "CMD-SHELL", "pg_isready -U ${PG_USER} -d ${PG_DB}" ]
interval: 5s
timeout: 5s
retries: 20
@@ -23,6 +46,7 @@ services:
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}
@@ -33,103 +57,113 @@ services:
# В этот каталог будет монтироваться генератор demodb
- ./bookings:/bookings:ro
healthcheck:
test: ["CMD-SHELL", "pg_isready -U ${BOOKINGS_DB_USER} -d ${BOOKINGS_DB_NAME}"]
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
build:
context: .
dockerfile: Dockerfile.greenplum
image: greenplum-custom:6.27.1
# container_name: gp_single
# hostname: gpdbsne
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"
PXF_SEED_OVERWRITE: ${PXF_SEED_OVERWRITE:-0}
PXF_SYNC_ON_START: ${PXF_SYNC_ON_START:-0}
# Порты: внешний 5435 (на хосте)
ports:
- "${GP_PORT}:5432"
- "5435:5432"
volumes:
- ./sql:/sql: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"]
# Ждём не только доступность GPDB, но и готовность PXF,
# чтобы Airflow не стартовал раньше PXF.
test: [ "CMD-SHELL", "/usr/local/greenplum-db/bin/pg_isready -h 127.0.0.1 -p 5432 -U ${GP_USER:-gpadmin} -d ${GP_DB:-gp_dwh} && PGPASSWORD=${GP_PASSWORD:-gpadmin} /usr/local/greenplum-db/bin/psql -h 127.0.0.1 -p 5432 -U ${GP_USER:-gpadmin} -d ${GP_DB:-gp_dwh} -t -A -c \"SELECT 1 FROM pg_extension WHERE extname='pxf';\" | grep -q 1 && su - ${GP_USER:-gpadmin} -c '/usr/local/pxf/bin/pxf cluster status' >/dev/null 2>&1" ]
interval: 10s
timeout: 5s
retries: 30
start_period: 40s
airflow-webserver:
image: apache/airflow:2.9.2
container_name: gp_airflow_web
build:
context: .
dockerfile: Dockerfile.airflow
image: airflow-custom:latest
container_name: gp_airflow_webserver
env_file: .env
environment:
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}
command: >
bash -lc "pip install --no-cache-dir -r /opt/airflow/requirements.txt &&
airflow webserver"
<<: *airflow-env
command: bash -lc "rm -f /opt/airflow/airflow-webserver.pid && exec airflow webserver"
ports:
- "8080:8080"
volumes:
- ./airflow/dags:/opt/airflow/dags
- ./airflow/requirements.txt:/opt/airflow/requirements.txt
- ./sql:/sql:ro
- airflow_data:/opt/airflow/data
- airflow_logs:/opt/airflow/logs
- ./airflow/requirements.txt:/opt/airflow/requirements.txt
healthcheck:
test: ["CMD", "curl", "-f", "http://localhost:8080/health"]
interval: 30s
timeout: 10s
retries: 5
start_period: 40s
depends_on:
pgmeta:
condition: service_healthy
greenplum:
condition: service_healthy
<<: *airflow-depends
airflow-scheduler:
image: apache/airflow:2.9.2
container_name: gp_airflow_sch
build:
context: .
dockerfile: Dockerfile.airflow
image: airflow-custom:latest
container_name: gp_airflow_scheduler
env_file: .env
environment:
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}
command: >
bash -lc "pip install --no-cache-dir -r /opt/airflow/requirements.txt &&
airflow scheduler"
<<: *airflow-env
command: airflow scheduler
volumes:
- ./airflow/dags:/opt/airflow/dags
- ./airflow/requirements.txt:/opt/airflow/requirements.txt
- ./sql:/sql:ro
- airflow_data:/opt/airflow/data
- airflow_logs:/opt/airflow/logs
- ./airflow/requirements.txt:/opt/airflow/requirements.txt
depends_on:
pgmeta:
condition: service_healthy
greenplum:
condition: service_healthy
<<: *airflow-depends
airflow-init:
image: apache/airflow:2.9.2
# container_name: gp_airflow_init
build:
context: .
dockerfile: Dockerfile.airflow
image: airflow-custom:latest
# Имя контейнера закомментировано (одноразовый сервис)
user: "0"
env_file: .env
environment:
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-env
volumes:
- ./airflow/dags:/opt/airflow/dags
- ./airflow/requirements.txt:/opt/airflow/requirements.txt
- ./sql:/sql:ro
- airflow_data:/opt/airflow/data
- airflow_logs:/opt/airflow/logs
command: >
bash -lc "
set -e;
# mkdir -p /opt/airflow/data && chmod -R 777 /opt/airflow/data;
# su - airflow;
pip install --no-cache-dir -r /opt/airflow/requirements.txt;
mkdir -p /opt/airflow/data /opt/airflow/logs && chown -R airflow:root /opt/airflow/data /opt/airflow/logs;
# Дожидаемся готовности БД ретрая миграции
for i in {1..30}; do
airflow db migrate && break || echo 'waiting for pgmeta' && sleep 3;
su -s /bin/bash airflow -c \"PATH='/home/airflow/.local/bin:$${PATH}' airflow db migrate\" && break || echo 'waiting for pgmeta' && sleep 3;
done;
# Создаём админа; при повторном запуске не падаем
airflow users create --username ${AIRFLOW_USER} --password ${AIRFLOW_PASSWORD} --firstname Admin --lastname User --role Admin --email admin@example.org || true
su -s /bin/bash airflow -c \"PATH='/home/airflow/.local/bin:$${PATH}' airflow users create --username ${AIRFLOW_USER} --password ${AIRFLOW_PASSWORD} --firstname Admin --lastname User --role Admin --email admin@example.org\" || true
"
depends_on:
pgmeta:
@@ -140,3 +174,4 @@ volumes:
bookings_data:
greenplum_data:
airflow_data:
airflow_logs:
+138
View File
@@ -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.
+80
View File
@@ -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, примеры запросов и типичные ошибки.
+25
View File
@@ -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.
+201
View File
@@ -0,0 +1,201 @@
# PXF для bookings в учебном стенде (актуально)
Этот документ описывает **текущую реализацию** PXF в проекте: чтение данных из демо‑БД
`bookings` (Postgres, сервис `bookings-db`) в Greenplum через JDBC.
## 1. Что должно работать
- В Greenplum доступны внешние таблицы:
- `public.ext_bookings_bookings` (создаётся `make ddl-gp`);
- `stg.bookings_ext` (создаётся DAG `bookings_stg_ddl`).
- PXF должен быть готов **после каждого старта** контейнера `greenplum`.
## 2. Почему мы делаем свой образ Greenplum
Изначально PXF‑скрипты/конфиги монтировались в контейнер как bind‑mount `:ro`.
Базовый entrypoint образа Greenplum пытается делать `chown` файлов в
`/docker-entrypoint-initdb.d/`, из‑за чего контейнер иногда падал с ошибкой:
`chown: changing ownership ... Read-only file system`
Снять `:ro` тоже нежелательно — можно получить проблемы с правами на файлах хоста
(файл становится `root`, IDE перестаёт сохранять, появляются лишние изменения в git).
Решение для учебного стенда:
- собрать **свой образ** Greenplum (`Dockerfile.greenplum`);
- «вшить» в образ seed‑файлы и скрипты PXF;
- на каждом старте контейнера идемпотентно докладывать файлы в `PXF_BASE`,
который живёт на persistent volume.
## 3. Где что хранится
**Внутри образа (immutable):**
- seed для PXF: `/opt/pxf-seed/` (JDBCJAR и `servers/bookings-db/jdbc-site.xml`);
- скрипты:
- `/opt/pxf-scripts/ensure_pxf_bookings.sh` (подготовка `PXF_BASE`);
- `/start_greenplum_with_pxf.sh` (startup wrapper).
**На persistent volume (переживает рестарты):**
- `PXF_BASE` по умолчанию: `${GREENPLUM_DATA_DIRECTORY}/pxf` → в нашем compose это
`/data/pxf` на томе `greenplum_data`.
Важно: так как `PXF_BASE` лежит на томе, обновления seed‑файлов из нового образа
**не перезатирают** файлы в `PXF_BASE` автоматически (это сделано намеренно, чтобы
не ломать ручные правки студентов).
## 4. Что происходит при старте контейнера `greenplum`
1) Docker запускает контейнер с базовым entrypoint образа и командой
`/start_greenplum_with_pxf.sh` (она задана в `Dockerfile.greenplum` как `CMD`).
2) `/start_greenplum_with_pxf.sh` выполняет подготовку PXF:
- запускает ensure‑скрипт `/opt/pxf-scripts/ensure_pxf_bookings.sh`;
- параллельно пытается выполнить `CREATE EXTENSION IF NOT EXISTS pxf`
в базе `${GP_DB}` (по умолчанию `gp_dwh`), когда Greenplum начинает принимать
подключения.
3) Затем управление передаётся оригинальному старту Greenplum: `exec /start_gpdb.sh`.
4) Healthcheck сервиса `greenplum` ждёт и готовность Greenplum, и то, что PXF уже
запущен (`pxf cluster status`). Это нужно, чтобы Airflow не стартовал раньше PXF.
## 5. Управляющие переменные окружения
Все переменные можно задать в `.env` (см. `.env.example`):
- `PXF_SEED_OVERWRITE=1` — принудительно перезаписать seed‑файлы из образа в `PXF_BASE`
(обычно нужно после правок в каталоге `pxf/`).
- `PXF_SYNC_ON_START=1` — выполнять `pxf cluster sync` при старте контейнера
(делает старт чуть дольше, но гарантирует актуальные конфиги на хостах кластера).
## 6. Быстрая ручная проверка
1) Дождаться `healthy` у `greenplum`:
`docker compose ps`
2) Проверить статус PXF (PXF CLI запускается только под пользователем `gpadmin`):
`docker compose exec greenplum bash -lc "su - gpadmin -c '/usr/local/pxf/bin/pxf cluster status'"`
3) После применения DDL (`make ddl-gp`) проверить чтение через PXF:
- `make gp-psql`
- `SELECT COUNT(*) FROM public.ext_bookings_bookings;`
## 7. Типовые ошибки
- `protocol "pxf" does not exist`
- причина: не создано расширение `pxf` в базе Greenplum;
- решение: перезапустить `greenplum` (скрипт сделает `CREATE EXTENSION IF NOT EXISTS pxf`)
или выполнить вручную `CREATE EXTENSION pxf;`.
- `Connection refused` к порту `5888`
- причина: PXF не поднялся/не успел подняться;
- решение: проверить `pxf cluster status`, посмотреть логи PXF в `/data/pxf/logs`,
перезапустить сервис `greenplum`.
- PXF «не подхватывает» изменения конфигов
- причина: файлы уже лежат в `PXF_BASE` на томе, а seed из образа по умолчанию не перетирает их;
- решение: `make build` + restart `greenplum` + (при необходимости) `PXF_SEED_OVERWRITE=1`.
## 9. Известная проблема: `protocol "pxf" does not exist` на «холодном старте» (исправлено)
Раньше (воспроизводилось в `./scripts/e2e_smoke.sh`) при первом `make ddl-gp` можно было получить:
`ERROR: protocol "pxf" does not exist`
### Почему так происходило
В базовом `/start_gpdb.sh` из образа Greenplum создание расширения `pxf` связано с проверкой
файла `${PXF_BASE}/conf/pxf-env.sh`:
- если `pxf-env.sh` **отсутствует**, скрипт выполняет `pxf cluster prepare/register` и затем
`CREATE EXTENSION IF NOT EXISTS pxf`;
- если `pxf-env.sh` **уже существует**, этот блок **пропускается**, и расширение может не появиться.
При этом наш ensure‑скрипт `pxf/init/10_pxf_bookings.sh` копировал `pxf-env.sh` в `${PXF_BASE}`
ещё до запуска Greenplum, из‑за чего базовый скрипт считал PXF “уже настроенным” и
пропускал создание расширения.
### Что изменили
- создание `extension pxf` вынесено в `start_greenplum_with_pxf.sh` и обёрнуто ретраями;
- `pxf-env.sh` по‑прежнему копируется в `PXF_BASE`, чтобы `/start_gpdb.sh` не пытался выполнять
`pxf cluster prepare` на непустом `PXF_BASE`;
- healthcheck `greenplum` ждёт не только PXF, но и наличие `extension pxf`.
- добавлен экспорт `PGPASSWORD` для `pxf cluster start`, чтобы `docker compose stop/start`
не ломал запуск из‑за `password authentication failed` для `gpadmin`.
### Если ошибка всё ещё возникает
1) Пересоберите образ и перезапустите контейнер `greenplum`:
`make build && make down && make up`
2) Проверьте наличие extension:
`docker compose exec greenplum bash -lc "su - gpadmin -c '/usr/local/greenplum-db/bin/psql -d gp_dwh -t -A -c \"SELECT extname FROM pg_extension WHERE extname = ''pxf'';\"'"`
## 8. Связанные файлы
- `Dockerfile.greenplum`
- `docker-compose.yml` (сервис `greenplum`: `build`, `hostname`, env, healthcheck)
- `pxf/init/10_pxf_bookings.sh` (ensure‑логика)
- `pxf/init/start_greenplum_with_pxf.sh` (старт контейнера)
- `README.md` (раздел «Greenplum + PXF: свой образ»)
## 10. Известная проблема: после `docker compose stop/start` Greenplum может упасть (auth для PXF)
### Симптом
После `docker compose stop`, затем `docker compose start` контейнер `greenplum` иногда уходит в `Exited (1)`.
В логах видно, что GPDB поднялся, но упал на старте PXF:
- `INFO - pxf cluster start`
- `ERROR: Could not connect to GPDB`
- `FATAL: password authentication failed for user "gpadmin"`
### Текущее понимание причины (почему это “иногда”)
1) При старте GPDB образ `woblerr/greenplum` генерирует/дописывает `pg_hba.conf` на persistent volume.
2) В `pg_hba.conf` присутствует trust‑правило для **конкретного IP** контейнера в docker‑сети
(пример из диагностики: `host all gpadmin 172.21.0.2/32 trust`).
3) После `docker compose stop/start` Docker может выдать контейнеру **другой IP** (например, `172.21.0.3`).
Тогда trust‑правило больше не подходит, и подключение начинает идти по `md5`.
4) `pxf cluster start` подключается к GPDB по TCP на `host=gpdbsne` (hostname контейнера),
то есть попадает именно в `pg_hba.conf` (а не в localauth).
5) В результате при “не совпавшем IP” получаем `md5` + пароль (возможно пустой/не тот) → падение на `28P01`.
Эта проблема выглядит флапающей, потому что IP после `stop/start` иногда совпадает с захардкоженным trust‑/32,
а иногда нет.
### Как подтвердить при следующем воспроизведении
1) Посмотреть логи `greenplum`:
`docker compose logs --tail=200 greenplum`
2) Найти реальный IP клиента в master‑логах GPDB (на томе):
`Password does not match ...` обычно содержит адрес вида `172.21.0.X`.
3) Сравнить его с trust‑строкой в `pg_hba.conf` на томе:
`/data/master/gpseg-1/pg_hba.conf`
Если IP в ошибке (например, `172.21.0.3`) **не** совпадает с trust‑/32 (например, `172.21.0.2/32`) —
это почти наверняка корень падения.
### Что с этим делать дальше (варианты решения, без реализации здесь)
Основная цель — убрать зависимость от “случайного IP после stop/start”:
- заставить `pxf cluster start` подключаться к GPDB через `127.0.0.1` (тогда работает существующий trust на localhost);
- или перестать добавлять в `pg_hba.conf` trust на конкретный `172.21.0.2/32` и заменить на более стабильное правило
(например, на подсеть docker‑сети или на `samehost`);
- или закрепить IP контейнера в compose (static IP), чтобы он не “плавал”;
- или отказаться от `stop/start` в пользу сценария, который не меняет сетевое окружение (но это хуже для UX студентов).
### Что реализовано
- В `pxf/init/start_greenplum_with_pxf.sh` добавлен шаг, который на каждом старте
обеспечивает в `pg_hba.conf` trust‑правило `host all gpadmin samehost trust`
(вставка перед `host all all 0.0.0.0/0 md5`), и делает `pg_ctl reload`, если GPDB уже запущен.
+136
View File
@@ -0,0 +1,136 @@
# Учебные задания по стенду
Этот документ собирает в одном месте задания для менти.
Он разбит на блоки: от базовой работы с CSV‑pipeline до более продвинутого сценария с демо‑БД bookings и слоем STG в Greenplum.
Если вы только начинаете, выполняйте задания по порядку. К разделу про bookings можно вернуться позже.
---
## 1. Базовый CSVpipeline (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 и CSVpipeline будут освоены.
Планируемые направления:
- Спроектировать модель данных для основных сущностей демобазы bookings (рейсы, билеты, перелёты) в слоях ODS/DDS/DM.
- Реализовать слой ODS поверх STG, аккуратно работая с временными атрибутами и ключами.
- Построить витрины (DM) для типичных аналитических вопросов: загрузка рейсов, выручка по направлениям, динамика бронирований.
- Добавить DAG’и, которые используют `stg.bookings` как источник и строят следующие слои DWH.
- Расширить проверки качества данных для потоков bookings → STG → витрины.
Когда будете готовы к этим темам, вернитесь к этому разделу — он станет основой для следующего «модуля» лабораторных заданий.
+98
View File
@@ -0,0 +1,98 @@
# План исправления: генератор demodb (bookings) остаётся пустым
## Контекст
В стенде используется демобаза bookings из репозитория `postgrespro/demodb`, закреплённая на коммите `d68de192850237719f09b47688d5f3fc94653ca6` (см. `DEMODB_COMMIT` в `Makefile`).
Инициализация источника для ETL выполняется командой `make bookings-init`:
- клонирует demodb в `bookings/demodb/`;
- пытается применить патчи из `bookings/patches/`;
- запускает `install.sql` в контейнере `bookings-db`;
- выставляет GUC-параметры (`gen.connstr`, `bookings.start_date/init_days/jobs`);
- запускает `/bookings/generate_next_day.sql` (должен сгенерировать минимум 1 день данных).
## Симптомы (как в TODO)
- после `make bookings-init` таблица `bookings.bookings` остаётся пустой;
- патчи `bookings/patches/engine_jobs1_sync.patch` и `bookings/patches/install_drop_if_exists.patch` падают при применении;
- из‑за этого DAG `bookings_to_gp_stage` валится на проверках (источник пустой).
## Предварительный диагноз (что уже видно)
1) `engine_jobs1_sync.patch` не является валидным unified diff (в hunk’ах нет номеров строк вида `@@ -N,M +N,M @@`), поэтому `patch` отвечает:
`patch: **** Only garbage was found in the patch input.`
2) `install_drop_if_exists.patch` устарел относительно закреплённого коммита demodb: в `install.sql` уже есть `DROP DATABASE IF EXISTS demo;`, поэтому hunk “не находится” и патч не накатывается.
3) Ошибки патча сейчас замаскированы в `Makefile` через `|| true`, поэтому `make bookings-init` может завершаться “успешно”, хотя критичные правки в demodb не применились.
## Цель фикса
- `make bookings-init` воспроизводимо создаёт и наполняет `demo.bookings.bookings` (>0 строк).
- Если патчи не применяются — процесс останавливается с понятным сообщением, что делать дальше.
- Патчи соответствуют закреплённому коммиту demodb и применяются идемпотентно.
## План диагностики (чтобы быстро подтвердить проблему)
1) Чистое воспроизведение:
- `make clean`
- `rm -rf bookings/demodb`
- `make bookings-init`
2) Проверка данных:
- `make bookings-psql`
- выполнить:
- `SELECT COUNT(*) FROM bookings.bookings;`
- `SELECT min(book_date), max(book_date) FROM bookings.bookings;`
3) Проверка патчей (без изменения файлов):
- `patch -d bookings/demodb -p1 --dry-run < bookings/patches/engine_jobs1_sync.patch`
- `patch -d bookings/demodb -p1 --dry-run < bookings/patches/install_drop_if_exists.patch`
Ожидаемо: сейчас dry-run показывает “garbage in patch” и/или “Hunk FAILED”.
## План решения
### Шаг 1. Пересобрать патчи под закреплённый демо‑коммит
Собираем патчи через `git diff`, чтобы получился корректный unified diff.
1) `bookings/patches/engine_jobs1_sync.patch`:
- Добавить/подтвердить 2 изменения в `engine.sql`:
- `busy()` игнорирует текущий backend: `AND pid <> pg_backend_pid()`.
- `continue()` при `jobs = 1` выполняет `process_queue(end_date)` синхронно и пишет заметный маркер в лог (`Job 1 (local): ok`), иначе — оставляет текущую логику через `dblink`.
2) `bookings/patches/install_drop_if_exists.patch`:
- Поменять строку (в актуальном `install.sql`):
- было: `DROP DATABASE IF EXISTS demo;`
- стало: `DROP DATABASE IF EXISTS demo WITH (FORCE);`
### Шаг 2. Сделать `make bookings-init` fail-fast на проблемах с патчами
В `Makefile`:
- убрать `|| true` у применения патчей;
- при ошибке патча — завершать `make` с ненулевым кодом и короткой подсказкой:
- “удалите `bookings/demodb` и повторите `make bookings-init`”,
- “если не помогло — проверьте, что `DEMODB_COMMIT` не менялся и патчи собраны под него”.
### Шаг 3. Добавить “защиту от тихого пустого результата”
После запуска `/bookings/generate_next_day.sql``Makefile` или внутри SQL):
- выполнить проверку `COUNT(*)` по `bookings.bookings`;
- если 0 — завершаться ошибкой с подсказкой, куда смотреть (патчи/логи генератора).
Цель: чтобы проблема не уезжала дальше в DAG’и и DQ‑проверки, а ловилась сразу при init.
## Проверка (критерии готовности)
- `make clean && rm -rf bookings/demodb && make bookings-init` завершается без ошибок.
- `make bookings-psql``SELECT COUNT(*) FROM bookings.bookings;` возвращает `> 0`.
- `make bookings-generate-day` добавляет следующий день:
- `max(book_date)` сдвигается на +1 сутки.
- `./scripts/e2e_smoke.sh` проходит до проверки `stg.bookings` (или хотя бы DAG `bookings_to_gp_stage` перестаёт падать на “источник пустой”).
## Откат (если нужно быстро вернуть стенд в рабочее состояние)
- Временно отключить применение патчей в `Makefile` и явно предупреждать, что генерация может быть нестабильной (нежелательно для студентов).
- Или зафиксировать альтернативный `DEMODB_COMMIT`, под который уже готовы патчи (делать только вместе с обновлением документации и проверкой, что генерация стабильна).
+489
View File
@@ -0,0 +1,489 @@
# План улучшения Dockerfile и docker-compose.yml
## Обзор
Документ описывает план улучшения Dockerfile для Airflow и его интеграции с docker-compose.yml на основе анализа best practices.
## Согласованные изменения
✅ Переименовать `Dockerfile``Dockerfile.airflow`
✅ Обновить `build: .``build: { context: ., dockerfile: Dockerfile.airflow }`
✅ Добавить YAML anchors для устранения дублирования конфигурации
✅ Добавить healthcheck для airflow-webserver
✅ Исправить расположение requirements.txt в Dockerfile
✅ Добавить LABEL в Dockerfile
✅ Добавить `USER airflow` после установки зависимостей
✅ Добавить проверку `pip check`
✅ Переименовать контейнеры (`gp_airflow_web``gp_airflow_webserver`, `gp_airflow_sch``gp_airflow_scheduler`)
---
## Часть 1: Изменения в Dockerfile
### Текущее состояние (Dockerfile)
```dockerfile
FROM apache/airflow:2.9.2
COPY airflow/requirements.txt /requirements.txt
RUN pip install --no-cache-dir -r /requirements.txt
```
### Новое состояние (Dockerfile.airflow)
```dockerfile
# Apache Airflow с дополнительными зависимостями для Greenplum
FROM apache/airflow:2.9.2
LABEL maintainer="your-email@example.com"
LABEL description="Airflow with Greenplum and Pandas dependencies"
LABEL version="1.0"
# Копируем requirements в стандартное расположение
COPY airflow/requirements.txt /opt/airflow/requirements.txt
# Устанавливаем зависимости
RUN pip install --no-cache-dir -r /opt/airflow/requirements.txt \
&& pip check \
&& rm -rf /tmp/pip-*
# Переключаемся на пользователя airflow
USER airflow
WORKDIR /opt/airflow
```
### Обоснование изменений:
1. **LABEL** - стандартная практика для документирования образов
2. **`/opt/airflow/requirements.txt`** - стандартное расположение для Airflow
3. **`pip check`** - проверка совместимости установленных пакетов
4. **`USER airflow`** - безопасность и соответствие best practices
5. **`WORKDIR /opt/airflow`** - явное указание рабочей директории
---
## Часть 2: Изменения в docker-compose.yml
### Основные изменения:
1. **Добавить YAML anchors** для устранения дублирования
2. **Обновить build** для всех Airflow сервисов
3. **Добавить healthcheck** для airflow-webserver
4. **Добавить комментарии** для пояснения сокращений в именах контейнеров
### Структура YAML anchors:
```yaml
x-airflow-common-env: &airflow-env
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
x-airflow-common-volumes: &airflow-volumes
- ./airflow/dags:/opt/airflow/dags
- ./sql:/sql:ro
- airflow_data:/opt/airflow/data
x-airflow-common-depends: &airflow-depends
pgmeta:
condition: service_healthy
greenplum:
condition: service_healthy
airflow-init:
condition: service_completed_successfully
```
### Изменения для airflow-webserver:
```yaml
airflow-webserver:
build:
context: .
dockerfile: Dockerfile.airflow
image: airflow-custom:latest
container_name: gp_airflow_webserver
env_file: .env
environment:
<<: *airflow-env
command: airflow webserver
ports:
- "8080:8080"
volumes:
<<: *airflow-volumes
# requirements.txt монтируется как volume для удобства разработки
# При изменении зависимостей не требуется пересборка образа
- ./airflow/requirements.txt:/opt/airflow/requirements.txt
healthcheck:
test: ["CMD", "curl", "-f", "http://localhost:8080/health"]
interval: 30s
timeout: 10s
retries: 5
start_period: 40s
depends_on:
<<: *airflow-depends
```
### Изменения для airflow-scheduler:
```yaml
airflow-scheduler:
build:
context: .
dockerfile: Dockerfile.airflow
image: airflow-custom:latest
container_name: gp_airflow_scheduler
env_file: .env
environment:
<<: *airflow-env
command: airflow scheduler
volumes:
<<: *airflow-volumes
- ./airflow/requirements.txt:/opt/airflow/requirements.txt
depends_on:
<<: *airflow-depends
```
### Изменения для airflow-init:
```yaml
airflow-init:
build:
context: .
dockerfile: Dockerfile.airflow
image: airflow-custom:latest
# Имя контейнера закомментировано (одноразовый сервис)
user: "0"
env_file: .env
environment:
<<: *airflow-env
volumes:
<<: *airflow-volumes
command: >
bash -lc "
set -e;
mkdir -p /opt/airflow/data && chown -R airflow:root /opt/airflow/data;
for i in {1..30}; do
su -s /bin/bash airflow -c \"PATH='/home/airflow/.local/bin:$${PATH}' airflow db migrate\" && break || echo 'waiting for pgmeta' && sleep 3;
done;
su -s /bin/bash airflow -c \"PATH='/home/airflow/.local/bin:$${PATH}' airflow users create --username ${AIRFLOW_USER} --password ${AIRFLOW_PASSWORD} --firstname Admin --lastname User --role Admin --email admin@example.org\" || true
"
depends_on:
pgmeta:
condition: service_healthy
```
---
## Часть 3: Влияние переименования контейнеров
### Важно: Переименование контейнеров
При переименовании контейнеров (`gp_airflow_web``gp_airflow_webserver`, `gp_airflow_sch``gp_airflow_scheduler`) необходимо учитывать:
1. **Имена сервисов в docker-compose.yml остаются без изменений**
- `airflow-webserver`, `airflow-scheduler`, `airflow-init` - это имена сервисов
- Они используются в командах `docker-compose logs`, `docker-compose restart` и т.д.
- Эти имена НЕ меняются
2. **Имена контейнеров меняются**
- `gp_airflow_web``gp_airflow_webserver`
- `gp_airflow_sch``gp_airflow_scheduler`
- Они используются в командах `docker exec`, `docker inspect`, `docker logs`
3. **Проверка использования старых имён контейнеров**
```bash
# Поиск упоминаний старых имён в скриптах
grep -r "gp_airflow_web" .
grep -r "gp_airflow_sch" .
```
4. **Обновление Makefile (если используется)**
- Проверить команды, которые используют старые имена контейнеров
- Пример: `make logs` может использовать `docker-compose logs airflow-webserver` (это корректно)
5. **Обновление документации**
- Обновить все упоминания имён контейнеров в README.md, TESTING.md
- Обновить примеры команд в документации
---
## Часть 4: План тестирования изменений
### Этап 1: Подготовка окружения
1. **Остановить текущие контейнеры**
```bash
make down
```
2. **Удалить старые образы (опционально)**
```bash
docker rmi airflow-custom:latest
```
3. **Проверить наличие .env файла**
```bash
ls -la .env
# Если отсутствует, скопировать из .env.example
cp .env.example .env
```
### Этап 2: Сборка новых образов
1. **Собрать образы с новым Dockerfile**
```bash
docker-compose build
```
2. **Проверить успешность сборки**
```bash
docker images | grep airflow-custom
```
3. **Проверить наличие LABEL**
```bash
docker inspect airflow-custom:latest | grep -A 10 "Labels"
```
### Этап 3: Запуск сервисов
1. **Запустить весь стек**
```bash
make up
```
2. **Проверить статус контейнеров**
```bash
docker-compose ps
```
3. **Проверить логи инициализации**
```bash
docker-compose logs airflow-init
```
### Этап 4: Проверка healthcheck
1. **Проверить healthcheck для airflow-webserver**
```bash
docker inspect gp_airflow_webserver | grep -A 20 "Health"
```
2. **Дождаться healthy статуса**
```bash
watch -n 2 'docker inspect --format="{{.State.Health.Status}}" gp_airflow_webserver'
```
### Этап 5: Функциональное тестирование
1. **Проверить доступ к Airflow UI**
- Открыть http://localhost:8080
- Проверить авторизацию (использовать креды из .env)
- Проверить наличие DAG'ов в списке
2. **Проверить запуск DAG'ов**
- Выбрать любой DAG (например, `csv_to_greenplum`)
- Запустить вручную через UI
- Проверить успешность выполнения задач
3. **Проверить подключения к БД**
- Проверить Admin → Connections
- Убедиться, что `greenplum_conn` и `bookings_db` доступны
- Проверить Test Connection для каждого подключения
4. **Проверить логи scheduler**
```bash
docker-compose logs airflow-scheduler | tail -50
```
### Этап 6: Проверка использования старых имён контейнеров
1. **Поиск упоминаний старых имён в проекте**
```bash
grep -r "gp_airflow_web" . --exclude-dir=.git --exclude-dir=__pycache__
grep -r "gp_airflow_sch" . --exclude-dir=.git --exclude-dir=__pycache__
```
2. **Проверка Makefile**
```bash
grep -E "gp_airflow_web|gp_airflow_sch" Makefile
```
3. **Проверка документации**
```bash
grep -r "gp_airflow_web" README.md TESTING.md AGENTS.md docs/
grep -r "gp_airflow_sch" README.md TESTING.md AGENTS.md docs/
```
4. **Обновление найденных упоминаний**
- Заменить `gp_airflow_web` на `gp_airflow_webserver`
- Заменить `gp_airflow_sch` на `gp_airflow_scheduler`
- Обновить примеры команд в документации
### Этап 7: Тестирование требований
1. **Проверить установленные пакеты**
```bash
docker exec gp_airflow_webserver pip list | grep -E "psycopg2|pandas"
```
2. **Проверить совместимость пакетов**
```bash
docker exec gp_airflow_webserver pip check
```
3. **Проверить пользователя внутри контейнера**
```bash
docker exec gp_airflow_webserver whoami
# Ожидаемый результат: airflow
```
### Этап 8: Тестирование изменений requirements.txt
1. **Добавить тестовый пакет в requirements.txt**
```bash
echo "requests==2.31.0" >> airflow/requirements.txt
```
2. **Перезапустить контейнеры**
```bash
docker-compose restart airflow-webserver airflow-scheduler
```
3. **Проверить установку пакета**
```bash
docker exec gp_airflow_webserver pip list | grep requests
```
4. **Удалить тестовый пакет**
```bash
# Вернуть исходный requirements.txt
git checkout airflow/requirements.txt
```
### Этап 9: Регрессионное тестирование
1. **Запустить существующие тесты**
```bash
make test
```
2. **Проверить smoke-тесты DAG**
```bash
uv run pytest tests/test_dags_smoke.py -v
```
3. **Проверить тесты helpers**
```bash
uv run pytest tests/test_greenplum_helpers.py -v
```
### Этап 10: Проверка после перезапуска
1. **Полный перезапуск стека**
```bash
make down
make up
```
2. **Проверить сохранность данных**
- Проверить наличие DAG'ов в UI
- Проверить историю запусков DAG'ов
- Проверить подключения к БД
3. **Проверить логи на наличие ошибок**
```bash
docker-compose logs airflow-webserver | grep -i error
docker-compose logs airflow-scheduler | grep -i error
# Обратите внимание: имена сервисов не изменились, только имена контейнеров
```
---
## Часть 5: Критерии успеха
### Функциональные требования:
- ✅ Все контейнеры успешно запускаются
- ✅ Airflow UI доступен на http://localhost:8080
- ✅ Авторизация работает корректно
- ✅ Все DAG'ы отображаются в UI
- ✅ DAG'и успешно выполняются
- ✅ Подключения к Greenplum и bookings-db работают
- ✅ Healthcheck для airflow-webserver работает корректно
### Технические требования:
- ✅ Образ собирается без ошибок
- ✅ LABEL присутствуют в образе
- ✅ Пакеты устанавливаются от пользователя airflow
- ✅ `pip check` не возвращает ошибок
- ✅ requirements.txt монтируется как volume
- ✅ YAML anchors работают корректно
- ✅ Нет дублирования конфигурации
### Требования к совместимости:
- ✅ Существующие тесты проходят успешно
- ✅ Данные сохраняются после перезапуска
- ✅ История запусков DAG'ов сохраняется
- ✅ Подключения к БД работают как раньше
---
## Часть 6: Откат изменений
Если изменения вызывают проблемы, план отката:
1. **Остановить контейнеры**
```bash
make down
```
2. **Вернуть старые файлы**
```bash
git checkout Dockerfile docker-compose.yml
```
3. **Удалить новый образ**
```bash
docker rmi airflow-custom:latest
```
4. **Перезапустить стек**
```bash
make up
```
---
## Часть 7: Документация
После успешного внедрения изменений необходимо обновить:
1. **README.md** - добавить информацию о новом Dockerfile
2. **AGENTS.md** - обновить инструкции для агентов
3. **Комментарии в docker-compose.yml** - добавить пояснения к YAML anchors
---
## Резюме
План включает:
- 2 файла для изменения: `Dockerfile` → `Dockerfile.airflow`, `docker-compose.yml`
- 10 этапов тестирования
- 12 функциональных критериев успеха
- План отката на случай проблем
- Переименование контейнеров для улучшения читаемости
- Проверка использования старых имён контейнеров в проекте
### Важные замечания:
1. **Имена сервисов НЕ меняются** - `airflow-webserver`, `airflow-scheduler`, `airflow-init`
2. **Имена контейнеров меняются** - `gp_airflow_web` → `gp_airflow_webserver`, `gp_airflow_sch` → `gp_airflow_scheduler`
3. **Необходимо проверить** использование старых имён контейнеров в скриптах и документации
4. **Makefile команды** используют имена сервисов, поэтому они продолжат работать без изменений
Все изменения следуют best practices для Docker и Airflow, сохраняют обратную совместимость и не требуют изменений в DAG'ах.
+142
View File
@@ -0,0 +1,142 @@
# План работ: свой образ Greenplum с интегрированным PXF
## Контекст и проблема
Иногда (и у нас воспроизводится стабильно) контейнер `greenplum` падает при старте с ошибкой:
`chown: changing ownership of '/docker-entrypoint-initdb.d/10_pxf_bookings.sh': Read-only file system`
Причина: init-скрипт `pxf/init/10_pxf_bookings.sh` проброшен в контейнер как bind-mount `:ro`, а entrypoint базового образа пытается сделать `chown` файлов в `/docker-entrypoint-initdb.d/`.
## Цели
- Greenplum стабильно стартует после любого `restart/up` (без “иногда не стартует”).
- PXF готов к работе **после каждого запуска** контейнера.
- Не меняем права/владельца файлов в репозитории на хосте (никаких “файл стал root’ом / не редактируется”).
- Для студентов всё остаётся простым: `make build` (явная сборка) + `make up` (поднимает стенд; если образа нет — Docker Compose соберёт сам).
## Выбранный подход (высокоуровневый дизайн)
1) Делаем свой образ Greenplum на базе `woblerr/greenplum:6.27.1`.
2) Встраиваем в образ “seed” для PXF:
- JDBC-драйвер (JAR),
- конфиг сервера `bookings-db` (`jdbc-site.xml`),
- скрипт “ensure”, который идемпотентно гарантирует, что файлы лежат в `PXF_BASE` (на persistent volume).
3) Запускаем “ensure” **на каждом старте контейнера** через wrapper-entrypoint, а затем передаём управление оригинальному entrypoint базового образа.
4) `PXF_BASE` по умолчанию остаётся на volume (`/data/pxf`), чтобы настройки переживали рестарты.
## План работ (по шагам)
### Шаг 1. Разведка базового образа
- Проверить, где находится оригинальный entrypoint и как он запускается (путь, параметры, пользователь).
- Понять, как `GREENPLUM_PXF_ENABLE=true` влияет на старт (чтобы wrapper не ломал поведение).
Результат: фиксируем в README “как устроен старт” (1–2 абзаца).
### Шаг 2. Новый Dockerfile для Greenplum
- Добавить `Dockerfile.greenplum`:
- `FROM woblerr/greenplum:6.27.1`
- `COPY` seed-артефакты в образ (например, в `/opt/pxf-seed/...`)
- `COPY` wrapper-entrypoint в образ
- настроить права/владельца внутри образа так, чтобы старт был без ошибок
Результат: образ собирается локально через `make build` и/или автоматически через `make up`.
### Шаг 3. Wrapper-entrypoint (каждый старт)
- Добавить скрипт entrypoint-обёртки (например, `greenplum/entrypoint-wrapper.sh` или `pxf/entrypoint-wrapper.sh`):
- на старте вызывает `ensure`-скрипт;
- затем делает `exec` оригинального entrypoint базового образа с теми же аргументами.
Важно: wrapper не должен “перехватывать” логику инициализации кластера — только добавлять шаг подготовки PXF.
### Шаг 4. Переписать текущий init-скрипт в “ensure” (идемпотентный)
- Превратить `pxf/init/10_pxf_bookings.sh` в скрипт, который можно безопасно выполнять на каждом запуске:
- не опираться на `~/.bashrc`;
- `PXF_BASE` вычислять через env (`PXF_BASE`, `GREENPLUM_DATA_DIRECTORY`, fallback `/data/pxf`);
- seed-копирование делать “если файла нет”;
- добавить понятные логи (что сделано / что пропущено);
- `pxf cluster sync`:
- выполнять, только если команда доступна,
- не валить контейнер при ошибке (но писать предупреждение).
Дополнительно (опционально, но полезно для стенда):
- env-переключатель `PXF_SEED_OVERWRITE=1` — принудительно перезаписывать конфиг из образа в volume (для обновлений без удаления volume).
### Шаг 5. Обновить `docker-compose.yml`
- Для сервиса `greenplum` перейти на `build:` (и при желании оставить `image:` как тег).
- Убрать bind-mount’ы PXF (jar/config/init-скрипт), т.к. теперь всё в образе.
- Оставить `greenplum_data:/data` и `./sql:/sql:ro`.
- Исправить healthcheck Greenplum (сейчас конструкция вида `... || echo 1` делает healthcheck “вечно успешным”):
- healthcheck должен возвращать ненулевой код, если БД не готова;
- добавить `start_period`, чтобы не ловить ложные падения на холодном старте.
### Шаг 6. Обновить Makefile и README
- `Makefile`:
- убедиться, что `make build` собирает также Greenplum-образ (если введём build для сервиса);
- оставить `make up` как есть (Compose сам соберёт образ, если его нет).
- `README.md`:
- зачем свой образ (устойчивость, права на хосте, меньше mount’ов),
- как пересобрать образ,
- как обновить PXF-конфиг (через rebuild + `PXF_SEED_OVERWRITE=1` или через очистку volume).
## Проверка и критерии готовности
- `docker compose up -d greenplum` → контейнер остаётся `Up`, не падает.
- `docker compose restart greenplum` повторить 10–20 раз → без падений.
- Healthcheck Greenplum становится `healthy` (не “вечно healthy” и не “вечно starting”).
- Поднятие всего стека (`make up`) приводит к старту Airflow (scheduler/webserver), т.к. `depends_on: condition: service_healthy` начинает работать корректно.
## Риски и как их снизить
- **Стартап может замедлиться**, если `pxf cluster sync` делать каждый раз: поэтому скрипт должен быть быстрым, а sync — не фатальным при ошибках.
- **Обновление конфигов**: так как PXF_BASE на volume, изменения в образе сами не перетрут файлы — поэтому нужен `PXF_SEED_OVERWRITE=1` или понятная инструкция “как обновить”.
## Откат
- Вернуться к использованию `image: woblerr/greenplum:6.27.1` в `docker-compose.yml`.
- Вернуть mount’ы PXF, если нужно (но это вернёт риск с `:ro`).
---
## Статус (реализовано)
- Добавлен кастомный образ Greenplum: `Dockerfile.greenplum` (seed PXF + startup wrapper).
- Ensure‑скрипт перенесён в образ и стал идемпотентным: `pxf/init/10_pxf_bookings.sh`.
- Стартовый скрипт контейнера включает ensure и создаёт `EXTENSION pxf`: `pxf/init/start_greenplum_with_pxf.sh`.
- В `docker-compose.yml`:
- `greenplum` собирается через `build: Dockerfile.greenplum`;
- добавлен `hostname: gpdbsne`;
- убраны PXF bind-mount’ы (jar/config/init);
- healthcheck ждёт не только GPDB, но и готовность PXF.
- Обновлены инструкции: `README.md`, `.env.example`.
- Проблема с генератором demodb (пустая `bookings.bookings`) зафиксирована в `TODO.md`.
## Проверка (как воспроизвести)
Команды для ручной проверки:
- Пересобрать и перезапустить Greenplum:
- `make build`
- `docker compose up -d --force-recreate greenplum`
- 310 рестартов:
- `docker compose restart greenplum`
- дождаться `healthy` в `docker compose ps`
- Проверить PXF:
- `docker compose exec greenplum bash -lc "su - gpadmin -c '/usr/local/pxf/bin/pxf cluster status'"`
- Проверить DAG, который использует PXF:
- `docker exec -i gp_airflow_scheduler airflow dags test bookings_stg_ddl`
Текущий результат:
- `bookings_stg_ddl` проходит (PXF и `protocol pxf` доступны).
- `bookings_to_gp_stage` падает не из-за PXF, а из-за пустого источника
(`demo.bookings.bookings` = 0 строк). Это отдельная задача (см. `TODO.md`).
+100
View File
@@ -0,0 +1,100 @@
#!/usr/bin/env bash
# Идемпотентная подготовка PXF на каждом запуске контейнера:
# - копирует JDBC-драйвер и конфиг сервера из образа в PXF_BASE;
# - не перезаписывает файлы, если не задан PXF_SEED_OVERWRITE=1;
# - синхронизацию pxf выполняет только при PXF_SYNC_ON_START=1.
set -euo pipefail
PXF_SEED_DIR="${PXF_SEED_DIR:-/opt/pxf-seed}"
PXF_CONF_SEED_DIR="${PXF_CONF_SEED_DIR:-/usr/local/pxf/conf}"
PXF_BASE_DEFAULT="${GREENPLUM_DATA_DIRECTORY:-/data}/pxf"
PXF_BASE="${PXF_BASE:-$PXF_BASE_DEFAULT}"
PXF_SEED_OVERWRITE="${PXF_SEED_OVERWRITE:-0}"
PXF_SYNC_ON_START="${PXF_SYNC_ON_START:-0}"
PXF_CLI="${PXF_CLI:-/usr/local/pxf/bin/pxf}"
GP_USER="${GREENPLUM_USER:-gpadmin}"
log_info() {
echo "INFO - $*"
}
log_warn() {
echo "WARN - $*"
}
copy_seed_file() {
local src="$1"
local dst="$2"
local label="$3"
if [ ! -f "${src}" ]; then
log_warn "seed-файл не найден: ${src}"
return 0
fi
if [ "${PXF_SEED_OVERWRITE}" = "1" ] || [ ! -f "${dst}" ]; then
cp -f "${src}" "${dst}"
log_info "${label}: установлено в ${dst}"
return 0
fi
log_info "${label}: уже существует, пропускаем"
}
mkdir -p "${PXF_BASE}/lib" "${PXF_BASE}/servers/bookings-db" "${PXF_BASE}/conf"
mkdir -p "${PXF_BASE}/run" "${PXF_BASE}/logs"
copy_seed_file \
"${PXF_SEED_DIR}/postgresql-42.7.3.jar" \
"${PXF_BASE}/lib/postgresql-jdbc.jar" \
"JDBC драйвер PostgreSQL"
copy_seed_file \
"${PXF_SEED_DIR}/servers/bookings-db/jdbc-site.xml" \
"${PXF_BASE}/servers/bookings-db/jdbc-site.xml" \
"jdbc-site.xml для bookings-db"
copy_seed_file \
"${PXF_CONF_SEED_DIR}/pxf-application.properties" \
"${PXF_BASE}/conf/pxf-application.properties" \
"pxf-application.properties"
copy_seed_file \
"${PXF_CONF_SEED_DIR}/pxf-env.sh" \
"${PXF_BASE}/conf/pxf-env.sh" \
"pxf-env.sh"
# Устанавливаем уменьшенные JVM-опции, если они ещё не заданы явно
if [ -f "${PXF_BASE}/conf/pxf-env.sh" ]; then
if ! grep -Eq '^[[:space:]]*export[[:space:]]+PXF_JVM_OPTS=' "${PXF_BASE}/conf/pxf-env.sh"; then
echo 'export PXF_JVM_OPTS="-Xmx512m -Xms256m"' >> "${PXF_BASE}/conf/pxf-env.sh"
log_info "PXF_JVM_OPTS: установлен уменьшенный профиль памяти"
fi
fi
copy_seed_file \
"${PXF_CONF_SEED_DIR}/pxf-log4j2.xml" \
"${PXF_BASE}/conf/pxf-log4j2.xml" \
"pxf-log4j2.xml"
copy_seed_file \
"${PXF_CONF_SEED_DIR}/pxf-profiles.xml" \
"${PXF_BASE}/conf/pxf-profiles.xml" \
"pxf-profiles.xml"
if [ "${PXF_SYNC_ON_START}" = "1" ]; then
if [ ! -x "${PXF_CLI}" ]; then
log_warn "pxf cli не найден: ${PXF_CLI}, пропускаем pxf cluster sync"
exit 0
fi
# PXF CLI не запускается под root, поэтому выполняем синхронизацию под gpadmin.
# PXF_BASE передаём явно, т.к. `su -` сбрасывает окружение.
if ! su - "${GP_USER}" -c "PXF_BASE='${PXF_BASE}' '${PXF_CLI}' cluster sync"; then
log_warn "pxf cluster sync завершился с ошибкой"
else
log_info "pxf cluster sync выполнен"
fi
fi
+161
View File
@@ -0,0 +1,161 @@
#!/usr/bin/env bash
set -euo pipefail
# Этот скрипт — "обёртка" над стандартным стартом Greenplum (exec /start_gpdb.sh),
# которая добавляет устойчивый старт PXF для учебного стенда.
#
# Задачи скрипта (почему он нужен):
# 1) Подготовить PXF_BASE на persistent volume (через ensure-скрипт).
# 2) Обеспечить стабильную аутентификацию PXF → GPDB после `docker compose stop/start`
# (не зависеть от «плавающего» IP контейнера в docker-сети).
# 3) Гарантировать наличие расширения `pxf` в БД GPDB (иначе внешние таблицы падают с
# `ERROR: protocol "pxf" does not exist`).
#
# Важно: это учебный стенд, поэтому мы сознательно выбираем простые и надёжные решения
# (например, trust для samehost), а не «боевой» hardened security.
ensure_script="/opt/pxf-scripts/ensure_pxf_bookings.sh"
# PXF CLI использует libpq и при `pxf cluster start` подключается к GPDB как к обычному Postgres.
# После `docker compose stop/start` может внезапно потребоваться пароль (см. ensure_pg_hba_trust ниже),
# поэтому экспортируем PGPASSWORD заранее: это уменьшает "флап" и делает поведение воспроизводимым.
if [ -z "${PGPASSWORD:-}" ]; then
export PGPASSWORD="${GREENPLUM_PASSWORD:-gpadmin}"
fi
ensure_pg_hba_trust() {
# После `docker compose stop/start` Docker может выдать контейнеру другой IP.
# У базового образа Greenplum встречается trust-правило на конкретный /32 (старый IP),
# и тогда аутентификация по TCP начинает идти через md5 → PXF падает на `28P01`.
#
# Решение: добавить правило `host all gpadmin samehost trust`, которое срабатывает для
# подключений "с этого же контейнера" (PXF запускается рядом с master).
# Это максимально простое и стабильное правило для учебного стенда.
local gp_user="${GREENPLUM_USER:-gpadmin}"
local data_dir="${GREENPLUM_DATA_DIRECTORY:-/data}"
local attempts=60
local pg_hba=""
local trust_line="host all ${gp_user} samehost trust"
find_pg_hba() {
local candidate
for candidate in "${data_dir}/master"/*/pg_hba.conf; do
if [ -f "${candidate}" ]; then
echo "${candidate}"
return 0
fi
done
return 1
}
for _ in $(seq 1 "${attempts}"); do
pg_hba="$(find_pg_hba || true)"
if [ -n "${pg_hba}" ]; then
break
fi
sleep 2
done
if [ -z "${pg_hba}" ]; then
echo "WARN - pg_hba.conf не найден, пропускаем trust для ${gp_user}"
return 0
fi
if grep -Eq "^[[:space:]]*host[[:space:]]+all[[:space:]]+${gp_user}[[:space:]]+samehost[[:space:]]+trust" "${pg_hba}"; then
echo "INFO - pg_hba.conf уже содержит trust для ${gp_user} samehost"
return 0
fi
# Вставляем trust-правило перед самым "общим" md5-правилом (0.0.0.0/0),
# чтобы samehost гарантированно матчился раньше.
awk -v trust_line="${trust_line}" '
BEGIN { added = 0 }
$0 ~ /^[[:space:]]*host[[:space:]]+all[[:space:]]+all[[:space:]]+0\.0\.0\.0\/0[[:space:]]+md5/ && added == 0 {
print trust_line
added = 1
}
{ print }
END {
if (added == 0) {
print trust_line
}
}
' "${pg_hba}" > "${pg_hba}.tmp" && mv "${pg_hba}.tmp" "${pg_hba}"
echo "INFO - pg_hba.conf: добавлен trust для ${gp_user} samehost"
local master_dir="${pg_hba%/pg_hba.conf}"
# Если GPDB уже поднялся, достаточно reload, чтобы новое правило применилось без рестарта.
if /usr/local/greenplum-db/bin/pg_ctl -D "${master_dir}" status >/dev/null 2>&1; then
if ! /usr/local/greenplum-db/bin/pg_ctl -D "${master_dir}" reload >/dev/null 2>&1; then
echo "WARN - не удалось перезагрузить pg_hba.conf (pg_ctl reload)"
fi
fi
}
ensure_pxf_extension() {
# В базовом /start_gpdb.sh создание `CREATE EXTENSION pxf` зависит от наличия
# `${PXF_BASE}/conf/pxf-env.sh`. В нашем стенде pxf-env.sh может быть уже создан
# ensure-скриптом (PXF_BASE на томе), и тогда /start_gpdb.sh пропускает extension.
#
# Чтобы внешние таблицы через PXF работали после любого рестарта, создаём extension сами,
# но только когда GPDB начнёт принимать подключения (с ретраями).
local gp_user="${GREENPLUM_USER:-gpadmin}"
local gp_db="${GREENPLUM_DATABASE_NAME:-gp_dwh}"
local gp_password="${GREENPLUM_PASSWORD:-gpadmin}"
local attempts=60
if [ "${GREENPLUM_PXF_ENABLE:-false}" != "true" ]; then
return 0
fi
extension_exists() {
local result
# Подключаемся к localhost: это "локальный" путь и на нём обычно уже есть trust.
result=$(PGPASSWORD="${gp_password}" /usr/local/greenplum-db/bin/psql \
-h 127.0.0.1 -p 5432 -U "${gp_user}" -d "${gp_db}" \
-t -A -c "SELECT 1 FROM pg_extension WHERE extname='pxf';" 2>/dev/null || true)
[ "${result}" = "1" ]
}
for attempt in $(seq 1 "${attempts}"); do
# Ждём, пока master начнёт принимать подключения.
if /usr/local/greenplum-db/bin/pg_isready \
-h 127.0.0.1 -p 5432 -U "${gp_user}" -d "${gp_db}" >/dev/null 2>&1; then
if extension_exists; then
echo "INFO - extension pxf уже создано"
return 0
fi
if PGPASSWORD="${gp_password}" /usr/local/greenplum-db/bin/psql \
-h 127.0.0.1 -p 5432 -U "${gp_user}" -d "${gp_db}" \
-v ON_ERROR_STOP=1 -c "CREATE EXTENSION IF NOT EXISTS pxf;" >/dev/null 2>&1; then
if extension_exists; then
echo "INFO - extension pxf готово"
return 0
fi
fi
echo "WARN - попытка ${attempt}/${attempts}: не удалось создать extension pxf"
fi
sleep 2
done
echo "WARN - Greenplum не готов или extension pxf не создано"
}
# 1) Подготовка PXF_BASE (копирование seed, конфигов и т.д.).
if [ -f "${ensure_script}" ]; then
"${ensure_script}"
else
echo "WARN - не найден ensure-скрипт PXF: ${ensure_script}"
fi
# 2) Дальше запускаем две "подстраховки" параллельно, чтобы не замедлять старт контейнера:
# - правка pg_hba.conf (как только он появится на томе);
# - создание extension pxf (как только GPDB начнёт отвечать).
ensure_pg_hba_trust &
ensure_pxf_extension &
# 3) Стартуем GPDB "как обычно". Важно использовать exec, чтобы сигналы Docker
# (stop/restart) корректно приходили в основной процесс entrypoint.
exec /start_gpdb.sh
Binary file not shown.
+31
View File
@@ -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>
+4
View File
@@ -14,3 +14,7 @@ dev = [
"psycopg2-binary==2.9.9",
"pytest==7.4.4",
]
[tool.isort]
profile = "black"
src_paths = ["airflow", "tests"]
+67
View File
@@ -0,0 +1,67 @@
#!/usr/bin/env bash
set -euo pipefail
# Полный smoke-тест стенда: сносит volumes, поднимает стек,
# и проверяет оба учебных DAG через airflow dags test.
ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)"
cd "$ROOT"
warn() { echo ":: ${1}" >&2; }
wait_up() {
local svc="$1"
local attempts=60
while true; do
if docker compose -f docker-compose.yml ps "$svc" 2>/dev/null | grep -q "Up"; then
break
fi
attempts=$((attempts - 1))
if [ "$attempts" -le 0 ]; then
echo "Service $svc is not up after waiting" >&2
exit 1
fi
sleep 2
done
}
warn "Reset stack (containers + volumes)"
make clean
warn "Starting stack (docker compose up)"
make up
warn "Waiting for Airflow services"
wait_up airflow-webserver
wait_up airflow-scheduler
warn "Init demo DB bookings"
make bookings-init
warn "Apply DDL to Greenplum"
make ddl-gp
warn "Run local pytest suite"
make test
warn "Airflow DAG test: csv_to_greenplum"
docker compose -f docker-compose.yml exec airflow-webserver airflow dags test csv_to_greenplum 2024-01-01
warn "Check orders count in Greenplum"
ORDERS_COUNT=$(docker compose -f docker-compose.yml exec greenplum bash -lc "su - gpadmin -c \"/usr/local/greenplum-db/bin/psql -t -A -d gp_dwh -c 'SELECT COUNT(*) FROM public.orders;'\"")
if [ "${ORDERS_COUNT:-0}" -le 0 ]; then
echo "orders table is empty after csv_to_greenplum (COUNT=${ORDERS_COUNT:-0})" >&2
exit 1
fi
warn "Airflow DAG test: bookings_to_gp_stage"
docker compose -f docker-compose.yml exec airflow-webserver airflow dags test bookings_to_gp_stage 2024-01-01
warn "Check stg.bookings count in Greenplum"
BOOKINGS_COUNT=$(docker compose -f docker-compose.yml exec greenplum bash -lc "su - gpadmin -c \"/usr/local/greenplum-db/bin/psql -t -A -d gp_dwh -c 'SELECT COUNT(*) FROM stg.bookings;'\"")
if [ "${BOOKINGS_COUNT:-0}" -le 0 ]; then
echo "stg.bookings is empty after bookings_to_gp_stage (COUNT=${BOOKINGS_COUNT:-0})" >&2
exit 1
fi
warn "Smoke test completed successfully"
+11
View File
@@ -0,0 +1,11 @@
-- DDL для базовой таблицы orders, которую использует CSV‑pipeline.
-- Выполняется идемпотентно: таблица создаётся, если ещё не существует.
CREATE TABLE IF NOT EXISTS public.orders (
order_id BIGINT,
order_ts TIMESTAMP NOT NULL,
customer_id BIGINT NOT NULL,
amount NUMERIC(12,2) NOT NULL
)
WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1)
DISTRIBUTED BY (order_id);
+26 -11
View File
@@ -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 $$;
+29
View File
@@ -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);
+49
View File
@@ -0,0 +1,49 @@
-- Проверка количества строк между источником 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');
IF v_src_count = 0 THEN
RAISE EXCEPTION
'В источнике bookings_ext нет строк для окна инкремента (book_date > %). Проверьте генерацию данных (make bookings-init / make bookings-generate-day или таск generate_bookings_day).',
COALESCE(v_prev_ts, TIMESTAMP '1900-01-01 00:00:00');
END IF;
-- Считаем строки, реально вставленные в 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 $$;
+30
View File
@@ -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'
);
+2 -2
View File
@@ -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",