Compare commits
54
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
272013774b | ||
|
|
c3b2050cc4 | ||
|
|
5668edeb7b | ||
|
|
810d3d4727 | ||
|
|
573c853f15 | ||
|
|
1efeb58449 | ||
|
|
f81a831c0c | ||
|
|
4751674ef3 | ||
|
|
dfcc456120 | ||
|
|
5778ad6601 | ||
|
|
33fc42f883 | ||
|
|
1173e0e53f | ||
|
|
e27476fa5f | ||
|
|
4cb11b4769 | ||
|
|
9f028d0728 | ||
|
|
0d81096ded | ||
|
|
31171f84bc | ||
|
|
90aafd7b1d | ||
|
|
c27693b07f | ||
|
|
8f92f2e890 | ||
|
|
98e0c90da0 | ||
|
|
342f9bdb62 | ||
|
|
359245366a | ||
|
|
d936f204bd | ||
|
|
a8ac8f62dc | ||
|
|
3b24d13908 | ||
|
|
c0bb24764e | ||
|
|
409a9f6f31 | ||
|
|
dfd6eca760 | ||
|
|
daed2313e7 | ||
|
|
8fb78f9086 | ||
|
|
fbf9b118e6 | ||
|
|
6e2a43a900 | ||
|
|
250af00b2f | ||
|
|
29490828f4 | ||
|
|
7a9b2d182e | ||
|
|
c865fa282e | ||
|
|
cf95cf95d2 | ||
|
|
c63084be09 | ||
|
|
fe4bb4606f | ||
|
|
cacf989a9c | ||
|
|
2c8f5c1e23 | ||
|
|
196171e3d6 | ||
|
|
ae7f909f5d | ||
|
|
c3d7759e6b | ||
|
|
4515577b4a | ||
|
|
2bdeb17cee | ||
|
|
e3b0fd9ef9 | ||
|
|
86435c282d | ||
|
|
dbd2406765 | ||
|
|
8bc8e5cb3a | ||
|
|
b56ea952e7 | ||
|
|
c17baf447d | ||
|
|
c95dab6478 |
+14
-1
@@ -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
|
||||
|
||||
@@ -7,3 +7,4 @@ __pycache__/
|
||||
*.pyc
|
||||
.venv/
|
||||
data/
|
||||
bookings/demodb/
|
||||
|
||||
@@ -3,7 +3,7 @@
|
||||
Эта репа — учебный стенд для студентов (менти), которые только начинают с Airflow/Greenplum и Python. Пожалуйста, держите решения простыми, стабильными и хорошо объяснёнными.
|
||||
|
||||
## Структура проекта
|
||||
- `airflow/dags/` — DAG-файлы (например, `csv_to_greenplum.py`, `data_quality_greenplum.py`).
|
||||
- `airflow/dags/` — DAG-файлы (например, `csv_to_greenplum.py`, `csv_to_greenplum_dq.py`).
|
||||
- `airflow/requirements.txt` — зависимости, которые ставятся внутри контейнеров Airflow.
|
||||
- `sql/` — DDL и вспомогательные SQL (например, `sql/ddl_gp.sql`).
|
||||
- `docker-compose.yml` — Greenplum, Airflow, Postgres (мета-БД).
|
||||
@@ -11,14 +11,16 @@
|
||||
- `.env(.example)` — настройки окружения (реальные секреты не коммитим).
|
||||
|
||||
## Команды (основные)
|
||||
- `make up` — поднять весь стек.
|
||||
- `make airflow-init` — инициализировать мета-БД Airflow и создать пользователя.
|
||||
- `make up` — поднять весь стек (Airflow инициализируется автоматически при первом старте).
|
||||
- `make stop` — остановить контейнеры, не трогая данные.
|
||||
- `make down` — остановить и удалить контейнеры/сети (тома сохраняются).
|
||||
- `make clean` — полный reset: остановить и удалить контейнеры/сети и тома (все данные будут потеряны).
|
||||
- `make airflow-init` — вручную переинициализировать мета-БД Airflow и создать пользователя (обычно не нужно).
|
||||
- `make logs` — логи webserver и scheduler.
|
||||
- `make ddl-gp` — применить DDL к Greenplum.
|
||||
- `make gp-psql` — открыть `psql` в контейнере Greenplum от `gpadmin`.
|
||||
- `make down` — остановить и удалить тома (данные будут потеряны).
|
||||
|
||||
Пример: `make up && make airflow-init`, затем открыть `http://localhost:8080`.
|
||||
Пример: `make up`, затем открыть `http://localhost:8080`.
|
||||
|
||||
## Локальное Python‑окружение
|
||||
- Используем `uv`: достаточно `uv sync` (или `make dev-sync`) — подтянет Python, создаст `.venv`, установит зависимости.
|
||||
@@ -34,6 +36,22 @@
|
||||
- Форматирование: `black` (88 cols) и `isort`. Если не уверены — запустите `make fmt`.
|
||||
- Язык: комментарии, docstring и документацию — на русском; имена идентификаторов — на английском.
|
||||
|
||||
### Структура и нейминг SQL (слои DWH)
|
||||
- В каталоге `sql/` придерживаемся слоёв DWH:
|
||||
- `sql/src/` — скрипты, работающие с исходными системами (например, `bookings_generate_day_if_missing.sql`);
|
||||
- `sql/stg/` — скрипты для стейджинга (`bookings_ddl.sql`, `bookings_load.sql`, `bookings_dq.sql`);
|
||||
- в будущем можно добавить `sql/ods/`, `sql/dds/`, `sql/dm/` по мере роста стенда.
|
||||
- Именование файлов: `{объект}_{роль}.sql`, где:
|
||||
- `объект` — логическое имя сущности (`bookings`, `orders`, и т.п.);
|
||||
- `роль` — `ddl` (создание/изменение объектов), `load` (загрузка/инкремент), `dq` (проверки качества данных) и т.п.
|
||||
- Общие DDL-скрипты (например, `sql/ddl_gp.sql`) могут подключать файловые DDL через `\i`, но сами определения таблиц живут рядом с объектом (`sql/stg/bookings_ddl.sql` и т.п.).
|
||||
|
||||
### Airflow + SQL
|
||||
- В учебных DAG’ах, где основная логика — в SQL, по умолчанию используем `PostgresOperator` + Airflow Connections:
|
||||
- DAG оркестрирует шаги и подключение к БД;
|
||||
- SQL-скрипты лежат в `sql/...` и подключаются по пути (`sql='sql/stg/bookings_load.sql'`).
|
||||
- Сложную ручную работу с подключениями (`psycopg2`, ENV-фоллбеки) используем только там, где реально много Python-логики и это помогает учебной цели.
|
||||
|
||||
## Тестирование
|
||||
- Тесты лежат в `tests/` (pytest). Запуск: `make test`.
|
||||
- Есть юнит‑тесты для `helpers/greenplum.py` и smoke‑тесты DAG‑структуры (`tests/test_dags_smoke.py`).
|
||||
@@ -54,4 +72,4 @@
|
||||
- Пишите простыми словами. Добавляйте короткие комментарии к нетривиальной логике.
|
||||
- Избегайте больших рефакторингов и сложных паттернов — студенты только начинают.
|
||||
- Ошибки и логи — дружелюбные и понятные (лучше с подсказкой «что сделать дальше»).
|
||||
- Перед релевантными правками валидируйте локально: `make up && make airflow-init`, затем откройте DAG в UI и/или прогоните `make test`.
|
||||
- Перед релевантными правками валидируйте локально: `make up`, затем откройте DAG в UI и/или прогоните `make test`.
|
||||
|
||||
@@ -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
|
||||
@@ -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"]
|
||||
@@ -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
|
||||
|
||||
@@ -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
@@ -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-стенд не запускался в рамках этой сессии; ожидается, что инструкции выше обеспечат полноценную проверку.
|
||||
|
||||
@@ -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 и дополнительных зависимостей на хосте).
|
||||
@@ -0,0 +1,32 @@
|
||||
from __future__ import annotations
|
||||
|
||||
"""
|
||||
Учебный DAG: создаёт схему stg и таблицы bookings_ext/bookings в Greenplum.
|
||||
Запускается вручную перед DAG загрузки bookings_to_gp_stage или после изменения DDL.
|
||||
"""
|
||||
|
||||
from datetime import datetime, timedelta
|
||||
|
||||
from airflow.providers.postgres.operators.postgres import PostgresOperator
|
||||
|
||||
from airflow import DAG
|
||||
|
||||
GREENPLUM_CONN_ID = "greenplum_conn"
|
||||
|
||||
default_args = {"owner": "airflow", "retries": 1, "retry_delay": timedelta(seconds=30)}
|
||||
|
||||
with DAG(
|
||||
dag_id="bookings_stg_ddl",
|
||||
start_date=datetime(2024, 1, 1),
|
||||
schedule=None,
|
||||
catchup=False,
|
||||
template_searchpath="/sql",
|
||||
default_args=default_args,
|
||||
tags=["demo", "greenplum", "ddl", "bookings", "stg"],
|
||||
description="Создаёт/обновляет stg.bookings_ext и stg.bookings для учебного DAG",
|
||||
) as dag:
|
||||
apply_stg_bookings_ddl = PostgresOperator(
|
||||
task_id="apply_stg_bookings_ddl",
|
||||
postgres_conn_id=GREENPLUM_CONN_ID,
|
||||
sql="stg/bookings_ddl.sql",
|
||||
)
|
||||
@@ -0,0 +1,87 @@
|
||||
from __future__ import annotations
|
||||
|
||||
"""
|
||||
Учебный DAG для менти: показывает, как устроен поток
|
||||
от источника bookings-db до слоя stg в Greenplum.
|
||||
|
||||
В этом примере мы сознательно используем:
|
||||
- Airflow Connections для подключения к БД;
|
||||
- PostgresOperator, который берёт SQL-скрипты с диска;
|
||||
чтобы студент увидел «канонический» способ работы с SQL в DAG.
|
||||
|
||||
Каждый запуск DAG работает как «шаг по времени вперёд»:
|
||||
- генератор в демо-БД bookings добавляет следующий учебный день после max(book_date);
|
||||
- загрузка в Greenplum берёт все строки, появившиеся после предыдущих батчей;
|
||||
- логическая дата запуска (`ds`) используется как удобная метка запуска (через `ds_nodash` в `batch_id`, в логах и DQ).
|
||||
"""
|
||||
|
||||
from datetime import datetime, timedelta
|
||||
|
||||
from airflow.operators.python import PythonOperator
|
||||
from airflow.providers.postgres.operators.postgres import PostgresOperator
|
||||
|
||||
from airflow import DAG
|
||||
|
||||
default_args = {
|
||||
"owner": "airflow",
|
||||
"retries": 1,
|
||||
"retry_delay": timedelta(seconds=30),
|
||||
}
|
||||
|
||||
|
||||
BOOKINGS_CONN_ID = "bookings_db"
|
||||
GREENPLUM_CONN_ID = "greenplum_conn"
|
||||
|
||||
|
||||
def _finish_summary() -> None:
|
||||
"""
|
||||
Логирует краткий итог выполнения DAG за один запуск.
|
||||
|
||||
Здесь можно добавить дополнительную агрегацию/логирование,
|
||||
но для учебного примера достаточно простого сообщения.
|
||||
"""
|
||||
from logging import getLogger
|
||||
|
||||
log = getLogger(__name__)
|
||||
log.info("DAG bookings_to_gp_stage завершён. Подробности смотрите в логах задач.")
|
||||
|
||||
|
||||
with DAG(
|
||||
dag_id="bookings_to_gp_stage",
|
||||
start_date=datetime(2024, 1, 1),
|
||||
schedule=None,
|
||||
catchup=False,
|
||||
template_searchpath="/sql",
|
||||
default_args=default_args,
|
||||
tags=["demo", "bookings", "greenplum", "stg"],
|
||||
description="Учебный DAG: загрузка из bookings-db в stg.bookings (Greenplum)",
|
||||
) as dag:
|
||||
# 1. Генерируем один (или несколько стартовых) учебный день в демо-БД bookings
|
||||
generate_bookings_day = PostgresOperator(
|
||||
task_id="generate_bookings_day",
|
||||
postgres_conn_id=BOOKINGS_CONN_ID,
|
||||
sql="src/bookings_generate_day_if_missing.sql",
|
||||
autocommit=True,
|
||||
)
|
||||
|
||||
# 2. Загружаем инкремент из stg.bookings_ext в stg.bookings
|
||||
load_bookings_to_stg = PostgresOperator(
|
||||
task_id="load_bookings_to_stg",
|
||||
postgres_conn_id=GREENPLUM_CONN_ID,
|
||||
sql="stg/bookings_load.sql",
|
||||
)
|
||||
|
||||
# 3. Проверяем количество строк между источником и stg.bookings
|
||||
check_row_counts = PostgresOperator(
|
||||
task_id="check_row_counts",
|
||||
postgres_conn_id=GREENPLUM_CONN_ID,
|
||||
sql="stg/bookings_dq.sql",
|
||||
)
|
||||
|
||||
# 4. Финальный лог/сводка
|
||||
finish_summary = PythonOperator(
|
||||
task_id="finish_summary",
|
||||
python_callable=_finish_summary,
|
||||
)
|
||||
|
||||
generate_bookings_day >> load_bookings_to_stg >> check_row_counts >> finish_summary
|
||||
@@ -4,10 +4,13 @@ import logging
|
||||
from datetime import datetime, timedelta
|
||||
|
||||
from airflow.operators.python import PythonOperator
|
||||
from helpers.greenplum import (assert_orders_have_rows,
|
||||
from helpers.greenplum import (
|
||||
assert_orders_have_rows,
|
||||
assert_orders_no_duplicates,
|
||||
assert_orders_schema,
|
||||
assert_orders_table_exists, get_gp_conn)
|
||||
assert_orders_table_exists,
|
||||
get_gp_conn,
|
||||
)
|
||||
|
||||
from airflow import DAG
|
||||
|
||||
@@ -16,7 +19,8 @@ def _run_check(check_callable):
|
||||
"""
|
||||
Оборачивает проверку качества данных в контекст подключения к Greenplum.
|
||||
|
||||
Этот DAG предназначен для автоматической проверки качества данных в таблице orders:
|
||||
Этот DAG предназначен для автоматической проверки качества данных
|
||||
после CSV-пайплайна в таблице public.orders:
|
||||
1. Проверяет существование таблицы
|
||||
2. Проверяет соответствие схемы
|
||||
3. Проверяет наличие данных
|
||||
@@ -47,13 +51,13 @@ def _log_dq_summary():
|
||||
default_args = {"owner": "airflow", "retries": 1, "retry_delay": timedelta(seconds=30)}
|
||||
|
||||
with DAG(
|
||||
dag_id="greenplum_data_quality",
|
||||
dag_id="csv_to_greenplum_dq",
|
||||
start_date=datetime(2024, 1, 1),
|
||||
schedule=None,
|
||||
catchup=False,
|
||||
default_args=default_args,
|
||||
tags=["demo", "greenplum", "quality"],
|
||||
description="Автоматизированные проверки качества данных в Greenplum",
|
||||
tags=["demo", "greenplum", "quality", "csv", "dq"],
|
||||
description="Проверки качества данных после CSV → public.orders в Greenplum",
|
||||
) as dag:
|
||||
# Задача 1: Проверка существования таблицы
|
||||
check_exists = PythonOperator(
|
||||
@@ -0,0 +1,32 @@
|
||||
from __future__ import annotations
|
||||
|
||||
"""
|
||||
Учебный DAG: применяет DDL для базовой таблицы orders в Greenplum.
|
||||
Запускается вручную перед CSV‑пайплайном или после изменения схемы.
|
||||
"""
|
||||
|
||||
from datetime import datetime, timedelta
|
||||
|
||||
from airflow.providers.postgres.operators.postgres import PostgresOperator
|
||||
|
||||
from airflow import DAG
|
||||
|
||||
GREENPLUM_CONN_ID = "greenplum_conn"
|
||||
|
||||
default_args = {"owner": "airflow", "retries": 1, "retry_delay": timedelta(seconds=30)}
|
||||
|
||||
with DAG(
|
||||
dag_id="orders_base_ddl",
|
||||
start_date=datetime(2024, 1, 1),
|
||||
schedule=None,
|
||||
catchup=False,
|
||||
template_searchpath="/sql",
|
||||
default_args=default_args,
|
||||
tags=["demo", "greenplum", "ddl", "orders"],
|
||||
description="Создаёт/обновляет базовую таблицу orders в схеме public",
|
||||
) as dag:
|
||||
apply_orders_ddl = PostgresOperator(
|
||||
task_id="apply_orders_ddl",
|
||||
postgres_conn_id=GREENPLUM_CONN_ID,
|
||||
sql="base/orders_ddl.sql",
|
||||
)
|
||||
@@ -50,14 +50,17 @@ def get_gp_conn():
|
||||
|
||||
# Прямое подключение по переменным окружения
|
||||
conn_params = {
|
||||
"dbname": os.getenv("GP_DB", "gpadmin"),
|
||||
"dbname": os.getenv("GP_DB", "gp_dwh"),
|
||||
"user": os.getenv("GP_USER", "gpadmin"),
|
||||
"password": os.getenv("GP_PASSWORD", ""),
|
||||
"host": os.getenv("GP_HOST", "greenplum"),
|
||||
"port": int(os.getenv("GP_PORT", "5432")),
|
||||
}
|
||||
logging.info(
|
||||
"🔗 Подключение к Greenplum: %s:%s", conn_params["host"], conn_params["port"]
|
||||
"🔗 Подключение к Greenplum: %s:%s/%s",
|
||||
conn_params["host"],
|
||||
conn_params["port"],
|
||||
conn_params["dbname"],
|
||||
)
|
||||
return psycopg2.connect(**conn_params)
|
||||
|
||||
|
||||
@@ -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 $$;
|
||||
```
|
||||
|
||||
@@ -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 $$;
|
||||
|
||||
|
||||
@@ -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
@@ -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:
|
||||
|
||||
@@ -0,0 +1,138 @@
|
||||
# Дизайн STG для bookings в Greenplum (черновик)
|
||||
|
||||
_Внутренний документ для учебного стенда. Перед итоговой сдачей можно объединить с основной документацией._
|
||||
|
||||
## 1. Цель и общий контур
|
||||
|
||||
- Источник: Postgres в контейнере `bookings-db`, база `demo`, таблица `bookings.bookings` (см. `docs/internal/bookings_tz.md`).
|
||||
- Цель: показываем путь данных от операционной БД до сырого слоя DWH в Greenplum.
|
||||
- В этом документе описываем только часть `src (bookings-db) → STG (Greenplum)`. Слои ODS/DDS/DM студент проектирует сам по статье про моделирование DWH.
|
||||
|
||||
Логика на уровне слоёв (по статье):
|
||||
|
||||
- `src`: оперативная система (`bookings-db`, схема `bookings`).
|
||||
- `stg`: сырой слой в Greenplum, максимально близкий к источнику, без бизнес‑логики.
|
||||
- дальше по заданию менти можно строить `ods`/`dds`/`dm` поверх STG.
|
||||
|
||||
## 2. Схема и таблицы в Greenplum
|
||||
|
||||
### 2.1. Схема
|
||||
|
||||
- Используем одну схему `stg` в Greenplum.
|
||||
- В этой схеме будут:
|
||||
- внешняя таблица PXF для чтения из `bookings-db`;
|
||||
- внутренняя таблица STG для долговременного хранения «сырых» данных.
|
||||
|
||||
### 2.2. Внешняя таблица (PXF)
|
||||
|
||||
- Имя таблицы: `stg.bookings_ext`.
|
||||
- Назначение: «окно» в исходную таблицу `bookings.bookings` в `bookings-db` через PXF (JDBC).
|
||||
- Типы колонок:
|
||||
- можем использовать «родные» типы из `bookings.bookings` (включая даты/числа);
|
||||
- задача внешней таблицы — корректно читать данные из источника, не заниматься приведением типов.
|
||||
|
||||
DDL будет добавлен в `sql/ddl_gp.sql` в блоке DDL для Greenplum (примерно по шаблону из `docs/internal/pxf_bookings.md`), с `LOCATION ('pxf://bookings.bookings?PROFILE=JDBC&SERVER=bookings-db')`.
|
||||
|
||||
### 2.3. Внутренняя таблица STG
|
||||
|
||||
- Имя таблицы: `stg.bookings`.
|
||||
- Назначение: хранить сырые данные из источника для последующей обработки (ODS/DDS/витрины).
|
||||
- Принципы моделирования:
|
||||
- все бизнес‑колонки из `bookings.bookings` храним как `TEXT` (как в примерах STG из статьи);
|
||||
- не делаем `UPDATE/DELETE`, только `INSERT` новых записей;
|
||||
- бизнес‑колонки по названию совпадают с источником (чтобы проще было маппить).
|
||||
|
||||
Технологические колонки:
|
||||
|
||||
- `src_created_at_ts TIMESTAMP` — дата/время из источника, приведённая к TIMESTAMP:
|
||||
- используется как опорная колонка для инкрементальной загрузки;
|
||||
- заполняется из исходной даты/времени (`created_at` или аналог).
|
||||
- `load_dttm TIMESTAMP NOT NULL DEFAULT now()` — когда запись была загружена в STG.
|
||||
- `batch_id TEXT NOT NULL` — идентификатор «пачки» (например, `{{ ds_nodash }}` или `run_id` Airflow).
|
||||
- при необходимости позже можно добавить `src_system TEXT`, если появятся другие источники.
|
||||
|
||||
Колонки‑бизнес‑ключи (`booking_id` и т.п.) храним как `TEXT`. В слое DDS позже можно будет ввести суррогатные ключи и нормализовать модель под витрины.
|
||||
|
||||
## 3. Инкрементальная загрузка
|
||||
|
||||
### 3.1. Опорное поле для инкремента
|
||||
|
||||
- Опорная колонка: `src_created_at_ts` (внутреннее имя в STG).
|
||||
- Источник значения:
|
||||
- берём из соответствующей колонки в `bookings.bookings` (например, `book_date`/`created_at` — будет уточнено при реализации);
|
||||
- при чтении через `stg.bookings_ext` приводим к `TIMESTAMP`.
|
||||
|
||||
### 3.2. Правила определения full/delta
|
||||
|
||||
- При первом запуске, если таблица `stg.bookings` пуста:
|
||||
- считаем режим `full` — загружаем все строки из `stg.bookings_ext`.
|
||||
- При последующих запусках:
|
||||
- читаем `max(src_created_at_ts)` из `stg.bookings` за все предыдущие загрузки;
|
||||
- считаем, что нужно загрузить только строки, где `src_created_at_ts` больше этой максимальной метки и не позже конца текущего учебного дня.
|
||||
|
||||
Таким образом, вся логика инкремента «замкнута» на один техно‑столбец `src_created_at_ts`, который студент потом сможет использовать и на следующих слоях (например, в CDC‑логике).
|
||||
|
||||
## 4. DAG’и Airflow (логика на уровне задач)
|
||||
|
||||
### 4.1. DAG для DDL
|
||||
|
||||
- `dag_id`: `bookings_stg_ddl` (реализован в `airflow/dags/bookings_stg_ddl.py`).
|
||||
- Назначение: один раз (или при изменении схемы) создать необходимые объекты в Greenplum:
|
||||
- схему `stg` (если её ещё нет);
|
||||
- внешнюю таблицу `stg.bookings_ext` (PXF → `bookings-db`);
|
||||
- внутреннюю таблицу `stg.bookings` с текстовыми колонками и тех.полями.
|
||||
- Этот DAG не загружает данные, только подготавливает структуру.
|
||||
- Вся DDL‑логика (CREATE/ALTER/DROP) сосредоточена здесь; рабочие DAG’и занимаются только DML (INSERT/SELECT).
|
||||
|
||||
### 4.2. DAG для пошаговой загрузки
|
||||
|
||||
- `dag_id`: `bookings_to_gp_stage`.
|
||||
- Основные параметры:
|
||||
- `batch_id` (по умолчанию `{{ ds_nodash }}`) — метка батча, которая попадает в `stg.bookings.batch_id`;
|
||||
- подключения:
|
||||
- `bookings_db_conn_id` — Airflow connection к `bookings-db` (в коде DAG — `BOOKINGS_CONN_ID = "bookings_db"`);
|
||||
- `greenplum_conn_id` — Airflow connection к Greenplum (`GREENPLUM_CONN_ID = "greenplum_conn"`).
|
||||
|
||||
Последовательность задач (упрощённая, но отражающая те же шаги):
|
||||
|
||||
1. `generate_bookings_day`
|
||||
- PostgresOperator к `bookings-db`;
|
||||
- выполняет скрипт `/sql/src/bookings_generate_day_if_missing.sql`;
|
||||
- скрипт смотрит на `max(book_date)` и:
|
||||
- если база пуста — берёт стартовую дату из конфигурации (`bookings.start_date`) и генерирует `bookings.init_days` суток;
|
||||
- если данные уже есть — добавляет один следующий учебный день после `max(book_date)` (логика как в `bookings/generate_next_day.sql`).
|
||||
2. `load_bookings_to_stg`
|
||||
- PostgresOperator к Greenplum;
|
||||
- выполняет скрипт `/sql/stg/bookings_load.sql`;
|
||||
- внутри SQL считается `max(src_created_at_ts)` по «старым» батчам и по нему строится окно инкремента:
|
||||
- первая загрузка (full) — берём все строки из `stg.bookings_ext`;
|
||||
- последующие загрузки — берём только записи, где `book_date` больше предыдущего максимума (верхняя граница по дате не задаётся явно);
|
||||
- при вставке заполняются тех.колонки `src_created_at_ts`, `load_dttm`, `batch_id`.
|
||||
3. `check_row_counts`
|
||||
- PostgresOperator к Greenplum;
|
||||
- выполняет скрипт `/sql/stg/bookings_dq.sql`;
|
||||
- скрипт заново считает окно инкремента по тем же правилам, что и загрузка, и сравнивает:
|
||||
- количество строк в `stg.bookings_ext` с `book_date` позже «старого» максимума,
|
||||
- количество строк в `stg.bookings` для текущего `batch_id`;
|
||||
- при расхождении выполняет `RAISE EXCEPTION` с понятным текстом ошибки.
|
||||
4. `finish_summary`
|
||||
- PythonOperator, который логирует итог выполнения DAG и напоминает, где смотреть детальные логи.
|
||||
|
||||
Таким образом, вся бизнес‑логика инкремента и проверок живёт в SQL‑скриптах, а DAG отвечает за оркестрацию и подключение к нужным БД. Для менти это хороший пример разделения ответственности между SQL и Python.
|
||||
|
||||
## 5. Связь с остальными документами
|
||||
|
||||
- `docs/internal/bookings_tz.md` — как готовится и генерируется источник `bookings-db`.
|
||||
- `docs/internal/pxf_bookings.md` — детали настройки PXF и внешней таблицы для чтения из `bookings-db`.
|
||||
- `sql/stg/bookings_ddl.sql` — DDL для схемы `stg` и таблиц `stg.bookings_ext` / `stg.bookings` (подключается из `sql/ddl_gp.sql` и применяется через `make ddl-gp`).
|
||||
|
||||
Дальнейшая модель DWH (слои ODS/DDS/DM, факт/измерения, SCD) должна быть спроектирована студентом по статье о моделировании данных, используя `stg.bookings` как входной слой.
|
||||
|
||||
## 6. Требования к читаемости и комментариям
|
||||
|
||||
- DAG’и `bookings_stg_ddl` и `bookings_to_gp_stage` — это учебный материал для менти.
|
||||
- В коде DAG’ов должны быть:
|
||||
- понятные docstring на русском у всех функций и DAG;
|
||||
- короткие комментарии рядом с нетривиальной логикой (особенно вокруг инкремента и идемпотентности);
|
||||
- говорящие `task_id` и названия функций, отражающие их роль в процессе.
|
||||
- Цель: чтобы по одному только коду DAG студент мог восстановить архитектуру процесса и сопоставить её с теорией из статьи про моделирование DWH.
|
||||
@@ -0,0 +1,80 @@
|
||||
# Мини‑README по учебному DAG bookings_to_gp_stage (черновик)
|
||||
|
||||
_Внутренний файл, чтобы не забыть договорённости. Перед итоговой сдачей документацию по блоку bookings/STG нужно будет аккуратно собрать и переписать._
|
||||
|
||||
## 1. Что делает DAG
|
||||
|
||||
- DAG `bookings_to_gp_stage` показывает учебный поток:
|
||||
- источник: демо‑БД `bookings-db` (Postgres, схема `bookings`, таблица `bookings.bookings`);
|
||||
- при каждом запуске генерируется один учебный день данных (идемпотентно);
|
||||
- данные из `bookings.bookings` переливаются в сырой слой `stg.bookings` в Greenplum через PXF‑внешнюю таблицу `stg.bookings_ext`.
|
||||
- Слой `stg` задуман как «сырой»:
|
||||
- все бизнес‑колонки (`book_ref`, `book_date`, `total_amount`) хранятся как `TEXT`;
|
||||
- есть тех.колонки `src_created_at_ts`, `load_dttm`, `batch_id`.
|
||||
|
||||
Подробный дизайн описан в `docs/internal/bookings_stg_design.md`.
|
||||
|
||||
## 2. Что нужно, чтобы DAG завёлся
|
||||
|
||||
Минимальные предпосылки:
|
||||
|
||||
- Стенд поднят: `make up`.
|
||||
- Демо‑БД bookings инициализирована: `make bookings-init`.
|
||||
- В Greenplum применён DDL (созданы схема `stg` и таблицы `stg.bookings_ext` / `stg.bookings`):
|
||||
- учебный вариант: запустить DAG `bookings_stg_ddl` (он использует `sql/stg/bookings_ddl.sql`);
|
||||
- технический шорткат: `make ddl-gp` применяет все DDL разом вручную. Команда сама не вызывается при старте контейнеров, её нужно запустить явно.
|
||||
- В Airflow есть коннекты:
|
||||
- `greenplum_conn` — к Greenplum;
|
||||
- `bookings_db` — к сервису `bookings-db`.
|
||||
По умолчанию они задаются через переменные окружения `AIRFLOW_CONN_...` в docker-compose,
|
||||
поэтому могут не отображаться в UI, но `PostgresOperator` найдёт их по `conn_id`. При желании
|
||||
их можно создать/отредактировать вручную через Admin → Connections.
|
||||
|
||||
## 3. Последовательность задач в DAG
|
||||
|
||||
- `generate_bookings_day`:
|
||||
- PostgresOperator к `bookings-db`;
|
||||
- выполняет SQL `/sql/src/bookings_generate_day_if_missing.sql`;
|
||||
- скрипт смотрит на `max(book_date)` в `bookings.bookings`:
|
||||
- если база пустая — берёт стартовую дату из GUC и генерирует `bookings.init_days` суток;
|
||||
- если данные уже есть — добавляет один следующий учебный день после `max(book_date)` и пишет NOTICE с интервалом генерации.
|
||||
- `load_bookings_to_stg`:
|
||||
- PostgresOperator к Greenplum;
|
||||
- выполняет SQL `/sql/stg/bookings_load.sql`;
|
||||
- считает «старый» максимум `src_created_at_ts` (по предыдущим батчам) и грузит только новые строки из `stg.bookings_ext`, заполняя `src_created_at_ts`, `load_dttm`, `batch_id={{ ds_nodash }}`.
|
||||
- `check_row_counts`:
|
||||
- PostgresOperator к Greenplum;
|
||||
- выполняет SQL `/sql/stg/bookings_dq.sql`;
|
||||
- за то же окно инкремента считает количество строк в источнике и в `stg.bookings` (по текущему `batch_id`);
|
||||
- при расхождении делает `RAISE EXCEPTION` с понятным текстом ошибки.
|
||||
- `finish_summary`:
|
||||
- логирует итог выполнения DAG за одно срабатывание.
|
||||
|
||||
## 4. Как этим пользоваться студенту (черновой сценарий)
|
||||
|
||||
1. Поднять стенд и подготовить источники:
|
||||
- `make up` (Airflow инициализируется автоматически при первом старте)
|
||||
- `make bookings-init`
|
||||
- `make ddl-gp`
|
||||
2. Открыть Airflow UI (`http://localhost:8080`) и включить DAG `bookings_to_gp_stage`.
|
||||
3. Вызвать `Trigger` DAG (дату логического запуска можно оставить по умолчанию — она используется только как метка `batch_id`).
|
||||
4. Посмотреть:
|
||||
- в `bookings-db` ― что появился день с бронированиями;
|
||||
- в Greenplum (`make gp-psql`) — данные в `stg.bookings`:
|
||||
- `SELECT * FROM stg.bookings LIMIT 10;`
|
||||
- `SELECT src_created_at_ts, load_dttm, batch_id FROM stg.bookings ORDER BY src_created_at_ts DESC LIMIT 10;`
|
||||
5. Перезапустить DAG ещё несколько раз и увидеть, что:
|
||||
- генерация в `bookings.bookings` идёт по одному дню вперёд от текущего `max(book_date)`;
|
||||
- в `stg.bookings` появляются только новые записи (delta), помеченные разными `batch_id`.
|
||||
|
||||
## 5. Примечания «на потом»
|
||||
|
||||
- Текущая документация по блоку bookings/STG разбросана:
|
||||
- `README.md` (общий обзор стенда),
|
||||
- `docs/internal/bookings_tz.md` (источник bookings),
|
||||
- `docs/internal/pxf_bookings.md` (PXF),
|
||||
- `docs/internal/bookings_stg_design.md` (дизайн STG),
|
||||
- этот файл (мини‑README по DAG).
|
||||
- В будущем всё это нужно будет собрать в одну понятную историю для студента:
|
||||
- отдельный раздел «Учебный пример: bookings → stg → dwh»;
|
||||
- скриншоты DAG, примеры запросов и типичные ошибки.
|
||||
@@ -0,0 +1,25 @@
|
||||
# Временное ТЗ по блоку bookings (для текущей разработки)
|
||||
|
||||
_Этот файл внутренний, удалить перед итоговой сдачей._
|
||||
|
||||
- Контейнер `bookings-db` — отдельный сервис Postgres из `docker-compose.yml`, база по умолчанию `demo` (из upstream demodb), без переименований.
|
||||
- Доступ снаружи не блокируем (порт `5434` по умолчанию), чтобы позже читать через PXF и подключаться из Greenplum.
|
||||
- Инициализация (`make bookings-init`): поднимает контейнер, клонирует demodb с закреплённым коммитом, накладывает патчи (`engine`: `jobs=1` синхронно + `busy()` игнорирует свой pid; `install.sql`: `DROP DATABASE IF EXISTS`), ждёт `pg_isready`, ставит `gen.connstr` и GUC `bookings.start_date/init_days/jobs`, затем запускает `/bookings/generate_next_day.sql` через `psql -f`. Значения по умолчанию: стартовая дата 2017-01-01, `init_days=1`, `jobs=1`.
|
||||
- Генерация следующего дня: `make bookings-generate-day` прогоняет тот же SQL (читает GUC, вызывает `generate/continue`, ждёт `busy()`, закрывает dblink). При `jobs=1` всё синхронно, без dblink.
|
||||
- Исходники demodb: клонируем по требованию с фиксированным хешем, кладём в `bookings/demodb/` (в `.gitignore`), патчи лежат в `bookings/patches/` и применяются автоматически.
|
||||
- Документация: в README описаны команды (`bookings-init`, проверка данных, генерация дня), параметры `.env`; настройка PXF/ETL — следующий этап.
|
||||
|
||||
## Текущее состояние
|
||||
- `Makefile` теперь автоматически применяет патчи (`engine_jobs1_sync.patch`, `install_drop_if_exists.patch`), ждёт готовности Postgres через `pg_isready`, запускает `install.sql`, выставляет `gen.connstr`/GUC и вызывает `generate_next_day.sql` через `psql -f`.
|
||||
- Дефолты: `BOOKINGS_START_DATE=2017-01-01`, `BOOKINGS_INIT_DAYS=1`, `BOOKINGS_JOBS=1`. При `jobs=1` генерация идёт синхронно без dblink, `busy()` не учитывает текущую сессию.
|
||||
- `.env.example`/README обновлены под новые дефолты; каталог `bookings/demodb/` в `.gitignore`.
|
||||
- Патчи лежат в `bookings/patches/` и накладываются при `bookings-clone-demodb`.
|
||||
|
||||
## Текущее состояние тестов/проблем
|
||||
- Чистый прогон `make bookings-init` (после `docker compose down -v` и удаления `bookings/demodb`) проходит за ~1,5 минуты: база ставится, `busy()` → `f`, `bookings.bookings` от `2017-01-01 00:00:18` до `2017-01-01 23:59:59`.
|
||||
- Ранее зависание на `busy()` при `jobs=1` лечится патчем: `process_queue` теперь синхронный, а `busy()` игнорирует текущий backend.
|
||||
- Данных пока только на 1 день по умолчанию, чтобы генерация не занимала много времени.
|
||||
|
||||
## Идеи/следующие шаги
|
||||
- Если понадобится больше дней — увеличивать `BOOKINGS_INIT_DAYS`, но помнить, что генерация может идти долго; контролировать через `SELECT busy();`.
|
||||
- Следующий этап — PXF/ETL в Greenplum; текущая задача — лишь подготовить источник bookings.
|
||||
@@ -0,0 +1,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/` (JDBC‑JAR и `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` (а не в local‑auth).
|
||||
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 уже запущен.
|
||||
@@ -0,0 +1,136 @@
|
||||
# Учебные задания по стенду
|
||||
|
||||
Этот документ собирает в одном месте задания для менти.
|
||||
Он разбит на блоки: от базовой работы с CSV‑pipeline до более продвинутого сценария с демо‑БД bookings и слоем STG в Greenplum.
|
||||
|
||||
Если вы только начинаете, выполняйте задания по порядку. К разделу про bookings можно вернуться позже.
|
||||
|
||||
---
|
||||
|
||||
## 1. Базовый CSV‑pipeline (csv_to_greenplum)
|
||||
|
||||
Основная цель этого блока — понять, как устроен простой ETL: генерация данных через pandas, сохранение в CSV и загрузка в Greenplum.
|
||||
|
||||
### 1.1. Разбор готового pipeline
|
||||
|
||||
1. Найдите DAG `csv_to_greenplum` в `airflow/dags/csv_to_greenplum.py`.
|
||||
2. Ответьте себе на вопросы (можно коротко в отдельном файле/блокноте):
|
||||
- какие задачи (tasks) входят в DAG и что делает каждая из них;
|
||||
- какие таблицы создаются в Greenplum;
|
||||
- где физически лежат CSV‑файлы;
|
||||
- какие параметры управляют размером датасета.
|
||||
3. Поднимите стенд и запустите DAG:
|
||||
- `make up` (Airflow инициализируется автоматически при первом старте)
|
||||
- включите и запустите DAG `csv_to_greenplum` в Airflow UI.
|
||||
4. Проверьте результат в Greenplum:
|
||||
- `make gp-psql`
|
||||
- `SELECT COUNT(*) FROM public.orders;`
|
||||
- `SELECT * FROM public.orders LIMIT 5;`
|
||||
|
||||
### 1.2. Изменение параметров генерации
|
||||
|
||||
1. Найдите, где задаётся количество строк для генерации (`CSV_ROWS` в `.env` и параметр в DAG).
|
||||
2. Поставьте другое значение и перезапустите DAG:
|
||||
- оцените, как изменилось количество строк в `public.orders`;
|
||||
- убедитесь, что пайплайн по‑прежнему работает без ошибок.
|
||||
3. Попробуйте изменить схему данных (добавить колонку в CSV и таблицу в Greenplum):
|
||||
- добавьте новую колонку в генерацию pandas;
|
||||
- обновите DDL/SQL, чтобы колонка появилась в таблице `public.orders`;
|
||||
- перезапустите DAG и убедитесь, что новая колонка заполняется.
|
||||
|
||||
### 1.3. Собственные проверки качества данных
|
||||
|
||||
1. Найдите DAG `csv_to_greenplum_dq` в `airflow/dags/csv_to_greenplum_dq.py`.
|
||||
2. Посмотрите, какие проверки уже реализованы (наличие таблицы, схема, дубликаты).
|
||||
3. Добавьте ещё одну простую проверку, например:
|
||||
- проверка, что в таблице `public.orders` не больше N строк;
|
||||
- проверка, что поле (например, `order_price`) не содержит отрицательных значений;
|
||||
- проверка, что нет строк с `NULL` в ключевых колонках.
|
||||
4. Запустите DAG `csv_to_greenplum_dq` и убедитесь, что:
|
||||
- новая проверка проходит на «хороших» данных;
|
||||
- при нарушении условия DAG падает с понятной ошибкой.
|
||||
|
||||
---
|
||||
|
||||
## 2. Greenplum и модель данных (введение)
|
||||
|
||||
В следующих заданиях мы будем опираться на демо‑БД bookings (Postgres) и слой STG в Greenplum.
|
||||
На этом этапе достаточно бегло посмотреть на структуру и понять общую идею, детальная проработка пойдёт позже.
|
||||
|
||||
### 2.1. Знакомство с демо‑БД bookings
|
||||
|
||||
1. Прочитайте `bookings/README.md` — какие сервисы и команды относятся к демобазе.
|
||||
2. Поднимите стенд и выполните:
|
||||
- `make up`
|
||||
- `make bookings-init`
|
||||
3. Подключитесь к демобазе:
|
||||
- `make bookings-psql`
|
||||
- посмотрите таблицы в схеме `bookings` (например, `\dt bookings.*`).
|
||||
4. Найдите таблицу `bookings.bookings` и посмотрите на её структуру:
|
||||
- какие типы колонок используются;
|
||||
- какие поля выглядят как ключи, даты, суммы.
|
||||
|
||||
### 2.2. Знакомство с STG в Greenplum
|
||||
|
||||
1. Прочитайте `sql/stg/bookings_ddl.sql` и мини‑README `docs/internal/bookings_stg_readme.md` (если интересно — `docs/internal/bookings_stg_design.md`).
|
||||
2. Ответьте себе на вопросы:
|
||||
- чем внешняя таблица `stg.bookings_ext` отличается от внутренней `stg.bookings`;
|
||||
- зачем нужны тех.колонки `src_created_at_ts`, `load_dttm`, `batch_id`;
|
||||
- чем слой STG отличается от итоговых витрин (DDS/DM) с точки зрения моделирования.
|
||||
3. Выполните `make ddl-gp`, затем зайдите в Greenplum (`make gp-psql`) и проверьте наличие схемы и таблиц:
|
||||
- `\dn` и `\dt stg.*`
|
||||
- `SELECT * FROM stg.bookings LIMIT 5;` (после запуска соответствующего DAG).
|
||||
|
||||
### 2.3. Как генерируются учебные данные bookings
|
||||
|
||||
1. Откройте файл `bookings/generate_next_day.sql` и ответьте себе на вопросы:
|
||||
- с какой даты начинается генерация данных (посмотрите на GUC `bookings.start_date` и переменную `v_start_cfg`);
|
||||
- сколько дней генерируется при первой установке (переменная `bookings.init_days`);
|
||||
- что происходит, если таблица `bookings.bookings` уже не пустая.
|
||||
2. В демобазе (`make bookings-psql`) выполните:
|
||||
- `SELECT min(book_date), max(book_date) FROM bookings.bookings;`
|
||||
- затем запустите `make bookings-generate-day` и повторите запрос — как изменился максимальный день?
|
||||
3. Откройте `sql/src/bookings_generate_day_if_missing.sql` и обратите внимание, что:
|
||||
- логическая дата запуска DAG (`{{ ds }}`) не влияет на выбор дня генерации;
|
||||
- скрипт всегда смотрит на `max(book_date)` и добавляет **следующий** день (или несколько стартовых дней, если база пуста).
|
||||
4. Сделайте вывод: генератор всегда «шагает» по датам вперёд от максимальной даты, поэтому:
|
||||
- при `make bookings-init` вы получаете `BOOKINGS_INIT_DAYS` дней начиная с `BOOKINGS_START_DATE`;
|
||||
- при последующих вызовах (`make bookings-generate-day` или DAG) добавляется ровно один новый день.
|
||||
|
||||
---
|
||||
|
||||
## 3. DAG bookings_to_gp_stage (заготовка заданий)
|
||||
|
||||
Этот DAG показывает путь данных от демо‑БД bookings в Postgres до сырого слоя STG в Greenplum.
|
||||
Сейчас он уже реализован как учебный пример, а в будущем вокруг него появятся отдельные задания по моделированию DWH.
|
||||
|
||||
### 3.1. Что есть сейчас
|
||||
|
||||
1. Откройте `airflow/dags/bookings_to_gp_stage.py`.
|
||||
2. Найдите в коде ссылки на SQL‑файлы:
|
||||
- `sql/src/bookings_generate_day_if_missing.sql`
|
||||
- `sql/stg/bookings_load.sql`
|
||||
- `sql/stg/bookings_dq.sql`
|
||||
3. Соотнесите шаги DAG с документом `docs/internal/bookings_stg_readme.md`:
|
||||
- генерация учебного дня в `bookings.bookings`;
|
||||
- загрузка инкремента в `stg.bookings`;
|
||||
- проверка количества строк между источником и STG.
|
||||
4. Обратите внимание, как в DAG используется логическая дата запуска:
|
||||
- `{{ ds_nodash }}` используется как `batch_id` — метка загрузки в таблице `stg.bookings` для конкретного запуска;
|
||||
- сами даты данных (какие дни есть в `bookings.bookings`) определяются генератором по `max(book_date)`, а не по `ds`.
|
||||
|
||||
На этом этапе достаточно понять общую цепочку. Детальные задания по переработке модели данных и построению ODS/DDS/DM слоёв будут добавлены позже.
|
||||
|
||||
### 3.2. Идеи для будущих заданий (черновик)
|
||||
|
||||
> Ниже — набросок задач, к которым мы вернёмся, когда базовые темы по Airflow и CSV‑pipeline будут освоены.
|
||||
|
||||
Планируемые направления:
|
||||
|
||||
- Спроектировать модель данных для основных сущностей демобазы bookings (рейсы, билеты, перелёты) в слоях ODS/DDS/DM.
|
||||
- Реализовать слой ODS поверх STG, аккуратно работая с временными атрибутами и ключами.
|
||||
- Построить витрины (DM) для типичных аналитических вопросов: загрузка рейсов, выручка по направлениям, динамика бронирований.
|
||||
- Добавить DAG’и, которые используют `stg.bookings` как источник и строят следующие слои DWH.
|
||||
- Расширить проверки качества данных для потоков bookings → STG → витрины.
|
||||
|
||||
Когда будете готовы к этим темам, вернитесь к этому разделу — он станет основой для следующего «модуля» лабораторных заданий.
|
||||
@@ -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`, под который уже готовы патчи (делать только вместе с обновлением документации и проверкой, что генерация стабильна).
|
||||
|
||||
@@ -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'ах.
|
||||
@@ -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`
|
||||
- 3–10 рестартов:
|
||||
- `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`).
|
||||
Executable
+100
@@ -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
|
||||
Executable
+161
@@ -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.
@@ -0,0 +1,31 @@
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<!--
|
||||
JDBC-конфиг для PXF-сервера bookings-db.
|
||||
Для учебного стенда используем те же дефолтные креды, что и в .env.example:
|
||||
BOOKINGS_DB_USER=bookings, BOOKINGS_DB_PASSWORD=bookings.
|
||||
При изменении этих параметров в .env не забывайте обновить их здесь.
|
||||
-->
|
||||
<configuration>
|
||||
|
||||
<property>
|
||||
<name>jdbc.driver</name>
|
||||
<value>org.postgresql.Driver</value>
|
||||
</property>
|
||||
|
||||
<property>
|
||||
<name>jdbc.url</name>
|
||||
<value>jdbc:postgresql://bookings-db:5432/demo</value>
|
||||
</property>
|
||||
|
||||
<property>
|
||||
<name>jdbc.user</name>
|
||||
<value>bookings</value>
|
||||
</property>
|
||||
|
||||
<property>
|
||||
<name>jdbc.password</name>
|
||||
<value>bookings</value>
|
||||
</property>
|
||||
|
||||
</configuration>
|
||||
|
||||
@@ -14,3 +14,7 @@ dev = [
|
||||
"psycopg2-binary==2.9.9",
|
||||
"pytest==7.4.4",
|
||||
]
|
||||
|
||||
[tool.isort]
|
||||
profile = "black"
|
||||
src_paths = ["airflow", "tests"]
|
||||
|
||||
Executable
+67
@@ -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"
|
||||
@@ -0,0 +1,11 @@
|
||||
-- DDL для базовой таблицы orders, которую использует CSV‑pipeline.
|
||||
-- Выполняется идемпотентно: таблица создаётся, если ещё не существует.
|
||||
|
||||
CREATE TABLE IF NOT EXISTS public.orders (
|
||||
order_id BIGINT,
|
||||
order_ts TIMESTAMP NOT NULL,
|
||||
customer_id BIGINT NOT NULL,
|
||||
amount NUMERIC(12,2) NOT NULL
|
||||
)
|
||||
WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1)
|
||||
DISTRIBUTED BY (order_id);
|
||||
+26
-11
@@ -1,12 +1,27 @@
|
||||
-- Greenplum DDL (GPDB 6 совместимо)
|
||||
-- Колонночная таблица (append-optimized) и распределение по ключу.
|
||||
-- Внимание: append-optimized таблицы не поддерживают UNIQUE/PRIMARY KEY,
|
||||
-- поэтому контроль дублей выполняем в DAG при загрузке.
|
||||
CREATE TABLE IF NOT EXISTS public.orders (
|
||||
order_id BIGINT,
|
||||
order_ts TIMESTAMP NOT NULL,
|
||||
customer_id BIGINT NOT NULL,
|
||||
amount NUMERIC(12,2) NOT NULL
|
||||
-- Главный входной DDL-скрипт для Greenplum в учебном стенде.
|
||||
-- Выполняется из контейнера командой `make ddl-gp` и создаёт/обновляет
|
||||
-- все объекты, которые нужны базовым DAG (csv_to_greenplum, bookings_to_gp_stage);
|
||||
-- подключает файловые DDL через \i, чтобы сохранять единый входной скрипт.
|
||||
--
|
||||
-- Чтобы не ломать задания, новые объекты лучше добавлять в отдельные файлы
|
||||
-- и подключать их отсюда, а существующие определения не удалять.
|
||||
--
|
||||
-- Подробнее про STG/bookings: см. docs/internal/bookings_stg_readme.md.
|
||||
|
||||
-- Таблица для CSV‑пайплайна (csv_to_greenplum).
|
||||
\i base/orders_ddl.sql
|
||||
|
||||
-- Внешняя таблица для чтения данных из демо-БД bookings через PXF (JDBC).
|
||||
-- Источник: таблица bookings.bookings в базе demo (Postgres, сервис bookings-db).
|
||||
DROP EXTERNAL TABLE IF EXISTS public.ext_bookings_bookings;
|
||||
CREATE EXTERNAL TABLE public.ext_bookings_bookings (
|
||||
book_ref CHAR(6),
|
||||
book_date TIMESTAMP,
|
||||
total_amount NUMERIC(10,2)
|
||||
)
|
||||
WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1)
|
||||
DISTRIBUTED BY (order_id);
|
||||
LOCATION ('pxf://bookings.bookings?PROFILE=JDBC&SERVER=bookings-db')
|
||||
FORMAT 'CUSTOM' (formatter='pxfwritable_import');
|
||||
|
||||
-- DDL для слоя stg по таблице bookings вынесен в отдельный файл.
|
||||
-- Здесь подключаем его через psql \i, чтобы сохранить единый входной скрипт.
|
||||
\i stg/bookings_ddl.sql
|
||||
|
||||
@@ -0,0 +1,49 @@
|
||||
-- Генерация учебного дня в демо-БД bookings.
|
||||
-- Этот скрипт всегда добавляет следующий учебный день после максимальной даты
|
||||
-- в таблице bookings.bookings (или несколько дней от стартовой даты, если база пуста).
|
||||
-- Логическая дата запуска DAG ({{ ds }}) используется только как метка батча в стейдже.
|
||||
|
||||
DO $$
|
||||
DECLARE
|
||||
v_max_book_date timestamptz;
|
||||
v_start_date timestamptz;
|
||||
v_end_date timestamptz;
|
||||
v_jobs integer := COALESCE(current_setting('bookings.jobs', true), '1')::integer;
|
||||
v_init_days integer := COALESCE(current_setting('bookings.init_days', true), '1')::integer;
|
||||
v_start_cfg text := COALESCE(current_setting('bookings.start_date', true), '2017-01-01');
|
||||
BEGIN
|
||||
-- Проверяем, что демобаза установлена
|
||||
IF to_regclass('bookings.bookings') IS NULL THEN
|
||||
RAISE EXCEPTION 'Таблица bookings.bookings не найдена. Сначала выполните make bookings-init.';
|
||||
END IF;
|
||||
|
||||
-- Ищем последнюю сгенерированную дату
|
||||
SELECT max(book_date) INTO v_max_book_date FROM bookings.bookings;
|
||||
|
||||
IF v_max_book_date IS NULL THEN
|
||||
-- База пустая: берём стартовую дату из конфигурации (или дефолтную)
|
||||
v_start_date := date_trunc('day', v_start_cfg::timestamptz);
|
||||
ELSE
|
||||
-- Продолжаем с дня, следующего за максимальной датой
|
||||
v_start_date := date_trunc('day', v_max_book_date) + interval '1 day';
|
||||
END IF;
|
||||
|
||||
-- Первая генерация вызывает generate(), последующие — continue()
|
||||
IF v_max_book_date IS NULL THEN
|
||||
v_end_date := v_start_date + (v_init_days || ' days')::interval;
|
||||
CALL generate(v_start_date, v_end_date, v_jobs);
|
||||
ELSE
|
||||
v_end_date := v_start_date + interval '1 day';
|
||||
CALL continue(v_end_date, v_jobs);
|
||||
END IF;
|
||||
|
||||
RAISE NOTICE 'Сгенерированы данные в bookings.bookings за интервал [% - %).',
|
||||
date_trunc('day', v_start_date),
|
||||
date_trunc('day', v_end_date);
|
||||
|
||||
-- Ждём завершения фоновых джобов генератора, чтобы данные успели записаться
|
||||
WHILE busy() LOOP
|
||||
PERFORM pg_sleep(1);
|
||||
END LOOP;
|
||||
PERFORM dblink_disconnect(unnest(dblink_get_connections()));
|
||||
END $$;
|
||||
@@ -0,0 +1,29 @@
|
||||
-- DDL для слоя STG по таблице bookings.
|
||||
-- Используется как из общего скрипта ddl_gp.sql (через \i),
|
||||
-- так и может выполняться отдельно при изменении схемы.
|
||||
|
||||
-- Схема stg для сырого слоя DWH.
|
||||
CREATE SCHEMA IF NOT EXISTS stg;
|
||||
|
||||
-- Внешняя таблица в схеме stg для чтения данных из bookings.bookings через PXF.
|
||||
DROP EXTERNAL TABLE IF EXISTS stg.bookings_ext;
|
||||
CREATE EXTERNAL TABLE stg.bookings_ext (
|
||||
book_ref CHAR(6),
|
||||
book_date TIMESTAMP,
|
||||
total_amount NUMERIC(10,2)
|
||||
)
|
||||
LOCATION ('pxf://bookings.bookings?PROFILE=JDBC&SERVER=bookings-db')
|
||||
FORMAT 'CUSTOM' (formatter='pxfwritable_import');
|
||||
|
||||
-- Внутренняя таблица stg.bookings — сырой слой, все бизнес-колонки как TEXT.
|
||||
CREATE TABLE IF NOT EXISTS stg.bookings (
|
||||
book_ref TEXT,
|
||||
book_date TEXT,
|
||||
total_amount TEXT,
|
||||
src_created_at_ts TIMESTAMP,
|
||||
load_dttm TIMESTAMP NOT NULL DEFAULT now(),
|
||||
batch_id TEXT NOT NULL
|
||||
)
|
||||
WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1)
|
||||
DISTRIBUTED BY (book_ref);
|
||||
|
||||
@@ -0,0 +1,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 $$;
|
||||
@@ -0,0 +1,30 @@
|
||||
-- Загрузка инкремента из stg.bookings_ext в stg.bookings.
|
||||
-- Окно инкремента определяется по src_created_at_ts:
|
||||
-- берём строки, где book_date больше максимального src_created_at_ts
|
||||
-- среди "старых" батчей; верхняя граница по дате не используется.
|
||||
|
||||
INSERT INTO stg.bookings (
|
||||
book_ref,
|
||||
book_date,
|
||||
total_amount,
|
||||
src_created_at_ts,
|
||||
load_dttm,
|
||||
batch_id
|
||||
)
|
||||
SELECT
|
||||
book_ref::text,
|
||||
book_date::text,
|
||||
total_amount::text,
|
||||
book_date::timestamp,
|
||||
now(),
|
||||
'{{ ds_nodash }}'::text
|
||||
FROM stg.bookings_ext
|
||||
WHERE book_date > COALESCE(
|
||||
(
|
||||
SELECT max(src_created_at_ts)
|
||||
FROM stg.bookings
|
||||
WHERE batch_id <> '{{ ds_nodash }}'::text
|
||||
OR batch_id IS NULL
|
||||
),
|
||||
TIMESTAMP '1900-01-01 00:00:00'
|
||||
);
|
||||
@@ -48,8 +48,8 @@ def test_csv_to_greenplum_dag_structure():
|
||||
assert t4 in t3.get_direct_relatives("downstream")
|
||||
|
||||
|
||||
def test_data_quality_greenplum_dag_structure():
|
||||
dag = _load_dag("airflow.dags.data_quality_greenplum")
|
||||
def test_csv_to_greenplum_dq_dag_structure():
|
||||
dag = _load_dag("airflow.dags.csv_to_greenplum_dq")
|
||||
|
||||
expected_tasks = {
|
||||
"check_orders_table_exists",
|
||||
|
||||
Reference in New Issue
Block a user