Первая публикация кода
This commit is contained in:
+10
@@ -0,0 +1,10 @@
|
||||
.vscode/settings.json
|
||||
|
||||
# Do not commit secrets
|
||||
.env
|
||||
.env.*
|
||||
__pycache__/
|
||||
*/__pycache__/
|
||||
*.pyc
|
||||
.venv/
|
||||
data/
|
||||
@@ -0,0 +1 @@
|
||||
3.11
|
||||
@@ -0,0 +1,57 @@
|
||||
# Repository Guidelines (для агентa и контрибьюторов)
|
||||
|
||||
Эта репа — учебный стенд для студентов (менти), которые только начинают с Airflow/Greenplum и Python. Пожалуйста, держите решения простыми, стабильными и хорошо объяснёнными.
|
||||
|
||||
## Структура проекта
|
||||
- `airflow/dags/` — DAG-файлы (например, `csv_to_greenplum.py`, `data_quality_greenplum.py`).
|
||||
- `airflow/requirements.txt` — зависимости, которые ставятся внутри контейнеров Airflow.
|
||||
- `sql/` — DDL и вспомогательные SQL (например, `sql/ddl_gp.sql`).
|
||||
- `docker-compose.yml` — Greenplum, Airflow, Postgres (мета-БД).
|
||||
- `Makefile` — удобные команды для локальной работы.
|
||||
- `.env(.example)` — настройки окружения (реальные секреты не коммитим).
|
||||
|
||||
## Команды (основные)
|
||||
- `make up` — поднять весь стек.
|
||||
- `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`.
|
||||
|
||||
## Локальное Python‑окружение
|
||||
- Используем `uv`: достаточно `uv sync` (или `make dev-sync`) — подтянет Python, создаст `.venv`, установит зависимости.
|
||||
- `make dev-setup` полезен при смене версии Python (выполнит `uv python install` + `uv python pin` перед `uv sync`).
|
||||
- Проверки: `make test`, `make lint`, `make fmt` (выполняются через `uv run`).
|
||||
- Не используем `pip install --user`; если что‑то попало в user‑site — удалить `pip uninstall <package>` и проверить `pip list --user`.
|
||||
- В IDE выбираем интерпретатор из `.venv`.
|
||||
|
||||
## Стиль кода
|
||||
- Python: PEP 8, 4 пробела, `snake_case`; `dag_id` — `lower_snake_case`.
|
||||
- Импорты: stdlib → third‑party → local, по одному модулю в строке.
|
||||
- SQL: ключевые слова UPPERCASE, идентификаторы `snake_case`, завершаем `;`.
|
||||
- Форматирование: `black` (88 cols) и `isort`. Если не уверены — запустите `make fmt`.
|
||||
- Язык: комментарии, docstring и документацию — на русском; имена идентификаторов — на английском.
|
||||
|
||||
## Тестирование
|
||||
- Тесты лежат в `tests/` (pytest). Запуск: `make test`.
|
||||
- Есть юнит‑тесты для `helpers/greenplum.py` и smoke‑тесты DAG‑структуры (`tests/test_dags_smoke.py`).
|
||||
- Smoke‑тесты DAG автоматически пропускаются, если Airflow не установлен в venv.
|
||||
- Для ручного прогона стенда см. `TESTING.md` (пошаговый чек‑лист для студентов).
|
||||
|
||||
## Pull Requests
|
||||
- Conventional Commits: `feat:`, `fix:`, `docs:`, `chore:`, `refactor:`. Пример: `feat(dags): load orders to Greenplum`.
|
||||
- Держите изменения минимальными и локальными. Не переименовывайте Make‑таргеты без обновления документации.
|
||||
- В описании PR добавляйте скрин DAG‑графа или логи задач, если менялась логика.
|
||||
- При изменении схемы/поведения — обновляйте `README.md` и `sql/ddl_gp.sql`.
|
||||
|
||||
## Безопасность и конфигурация
|
||||
- Все настройки — через `.env`; креды в коде не хардкодим. Частые переменные: `GP_*`, `PG_*`, `AIRFLOW_*`, `CSV_*`.
|
||||
- `make down` удаляет тома — предупреждайте студентов, что данные пропадут.
|
||||
|
||||
## Для агента (особенности аудитории)
|
||||
- Пишите простыми словами. Добавляйте короткие комментарии к нетривиальной логике.
|
||||
- Избегайте больших рефакторингов и сложных паттернов — студенты только начинают.
|
||||
- Ошибки и логи — дружелюбные и понятные (лучше с подсказкой «что сделать дальше»).
|
||||
- Перед релевантными правками валидируйте локально: `make up && make airflow-init`, затем откройте DAG в UI и/или прогоните `make test`.
|
||||
@@ -0,0 +1,48 @@
|
||||
SHELL := /bin/bash
|
||||
UV := uv
|
||||
PYTHON_VERSION := 3.11
|
||||
|
||||
.PHONY: up down airflow-init logs gp-psql ddl-gp dev-setup dev-sync dev-lock test lint fmt clean-venv
|
||||
|
||||
up:
|
||||
docker compose -f docker-compose.yml up -d
|
||||
|
||||
down:
|
||||
docker compose -f docker-compose.yml down -v
|
||||
|
||||
airflow-init:
|
||||
docker compose -f docker-compose.yml run --rm airflow-init
|
||||
|
||||
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'"
|
||||
|
||||
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'"
|
||||
|
||||
dev-setup:
|
||||
$(UV) python install $(PYTHON_VERSION)
|
||||
$(UV) python pin $(PYTHON_VERSION)
|
||||
$(UV) sync
|
||||
|
||||
dev-sync:
|
||||
$(UV) sync
|
||||
|
||||
dev-lock:
|
||||
$(UV) lock --upgrade
|
||||
|
||||
test:
|
||||
$(UV) run pytest -q
|
||||
|
||||
lint:
|
||||
$(UV) run black --check airflow tests
|
||||
$(UV) run isort --check-only airflow tests
|
||||
|
||||
fmt:
|
||||
$(UV) run black airflow tests
|
||||
$(UV) run isort airflow tests
|
||||
|
||||
clean-venv:
|
||||
python -c "import shutil; shutil.rmtree('.venv', ignore_errors=True)"
|
||||
@@ -0,0 +1,275 @@
|
||||
# DE Starter Kit — Airflow + Greenplum + CSV
|
||||
|
||||
Добро пожаловать в учебный стенд для изучения основ Data Engineering! Этот проект поможет вам освоить ключевые инструменты современных data pipeline: **Airflow** для оркестрации, **pandas/CSV** для подготовки данных и **Greenplum** как аналитическую базу данных.
|
||||
|
||||
## 🎯 Что вы узнаете
|
||||
|
||||
- Как настроить локальный стек данных с помощью Docker
|
||||
- Как Airflow управляет workflow и координирует задачи
|
||||
- Как генерировать датасеты через pandas и сохранять их в CSV
|
||||
- Как загружать данные в Greenplum пакетами и избегать дублей
|
||||
- Как проверять качество данных в автоматизированных pipeline
|
||||
- Основы проектирования ETL/ELT процессов
|
||||
|
||||
## 👩🎓 Для студентов (10‑минутный чек‑лист)
|
||||
|
||||
- Установите Docker Desktop и Git.
|
||||
- Скопируйте настройки: `cp .env.example .env`.
|
||||
- Поднимите стенд: `docker compose up -d` и инициализируйте Airflow: `docker compose run --rm airflow-init`.
|
||||
- Откройте UI: http://localhost:8080 (admin/admin).
|
||||
- Включите и запустите DAG `csv_to_greenplum`. Дождитесь Success.
|
||||
- Проверьте данные: `make gp-psql` → `SELECT COUNT(*) FROM public.orders;`.
|
||||
- Дополнительно: запустите `greenplum_data_quality` — все проверки должны быть зелёные.
|
||||
|
||||
Если что‑то не работает — смотрите «Типичные проблемы» и «Быстрый reset» ниже.
|
||||
|
||||
## Локальное окружение разработчика
|
||||
|
||||
Локальным окружением управляет [uv](https://docs.astral.sh/uv/) — он скачивает нужный Python и создаёт `.venv` на основе `pyproject.toml` / `uv.lock`.
|
||||
|
||||
```bash
|
||||
uv sync
|
||||
```
|
||||
|
||||
`uv sync` сам подтянет версию Python из `.python-version`/`pyproject.toml`, создаст `.venv` и установит зависимости. Для тех же действий можно использовать `make dev-sync`. Цель `make dev-setup` (или вручную `uv python install` + `uv python pin`) нужна только когда вы меняете версию Python или прогреваете кэш.
|
||||
|
||||
> Если требуется «классическое» активированное окружение, после `uv sync` выполните `.\.venv\Scripts\Activate.ps1` в PowerShell или `source .venv/bin/activate` в Unix-терминале.
|
||||
|
||||
Проверки и форматирование выполняем через uv:
|
||||
|
||||
```bash
|
||||
make test # uv run pytest -q
|
||||
make lint # black/isort в режиме проверки
|
||||
make fmt # автоформатирование black + isort
|
||||
```
|
||||
|
||||
### Быстрый старт с uv
|
||||
|
||||
```bash
|
||||
uv sync
|
||||
uv run pytest -q
|
||||
uv run black --check airflow tests
|
||||
```
|
||||
|
||||
> Не устанавливайте пакеты напрямую через `pip install --user ...`. Если что-то уже попало в user-site, удалите `pip uninstall <package>` и проверьте `pip list --user`.
|
||||
---
|
||||
|
||||
## 🚀 Быстрый старт (для новичков)
|
||||
|
||||
### Шаг 1: Подготовка окружения
|
||||
|
||||
**Требования:**
|
||||
- Docker Desktop (Windows/Mac) или Docker Engine 24+ (Linux)
|
||||
- Git для клонирования репозитория
|
||||
|
||||
> 💡 **Совет:** Если у вас Windows, рекомендуем использовать WSL (Windows Subsystem for Linux) для лучшей совместимости.
|
||||
|
||||
### Шаг 2: Настройка проекта
|
||||
|
||||
```bash
|
||||
# Скопируйте файл настроек
|
||||
cp .env.example .env
|
||||
|
||||
# Запустите стек (это может занять 2-3 минуты при первом запуске)
|
||||
docker compose up -d
|
||||
|
||||
# Инициализируйте Airflow
|
||||
docker compose run --rm airflow-init
|
||||
```
|
||||
|
||||
### Шаг 3: Первый запуск pipeline
|
||||
|
||||
1. Откройте Airflow UI: **http://localhost:8080** (логин/пароль: admin/admin)
|
||||
2. Найдите DAG с названием **csv_to_greenplum**
|
||||
3. Нажмите на переключатель слева от названия DAG, чтобы включить его
|
||||
4. Нажмите кнопку **Trigger** (значок воспроизведения ▶️)
|
||||
|
||||
🎉 **Поздравляем!** Вы только что запустили свой первый data pipeline:
|
||||
- Система сгенерировала 1000 тестовых заказов при помощи pandas
|
||||
- Датасет сохранился в CSV-файл в каталоге `./data`
|
||||
- Airflow загрузил данные из CSV в Greenplum без дублей по `order_id`
|
||||
|
||||
### Шаг 4: Проверка результатов
|
||||
|
||||
**Проверка вручную:**
|
||||
```bash
|
||||
# Подключитесь к Greenplum и проверьте данные
|
||||
docker compose exec greenplum bash -c "su - gpadmin -c 'psql -p 5432 -d gpadmin'"
|
||||
|
||||
# Внутри psql выполните:
|
||||
\dt # Показать таблицы
|
||||
SELECT count(*) FROM public.orders; # Посчитать записи
|
||||
|
||||
# Посмотреть несколько строк
|
||||
SELECT * FROM public.orders LIMIT 5;
|
||||
```
|
||||
|
||||
CSV-файлы после выполнения DAG остаются в директории `./data`. Их можно открыть любым редактором или изучить через pandas.
|
||||
|
||||
### Быстрый reset
|
||||
|
||||
Если после изменений что‑то «сломалось»:
|
||||
|
||||
```bash
|
||||
make down # Остановить и стереть данные в контейнерах
|
||||
make up && make airflow-init
|
||||
```
|
||||
|
||||
Это помогает, когда Greenplum не стартует из‑за «грязной» остановки и внутренних файлов.
|
||||
|
||||
---
|
||||
|
||||
## 🛠️ Подробная настройка (для уверенных пользователей)
|
||||
|
||||
### Установка Make (опционально)
|
||||
|
||||
Для удобства работы с проектом рекомендуем установить `make`:
|
||||
|
||||
- **Linux (Debian/Ubuntu):** `sudo apt install -y make`
|
||||
- **macOS:** `brew install make`
|
||||
- **Windows:**
|
||||
- WSL: `sudo apt install -y make`
|
||||
- Chocolatey: `choco install make`
|
||||
- Scoop: `scoop install make`
|
||||
|
||||
С `make` команды становятся короче:
|
||||
```bash
|
||||
make up && make airflow-init # Запуск стека
|
||||
make logs # Просмотр логов
|
||||
make gp-psql # Подключение к Greenplum
|
||||
```
|
||||
|
||||
### Настройка подключения к Greenplum в Airflow
|
||||
|
||||
По умолчанию DAG использует переменные окружения, но вы можете создать Airflow Connection:
|
||||
|
||||
1. Airflow UI → **Admin → Connections → Add a new record**
|
||||
2. Заполните поля:
|
||||
- **Conn Id:** `greenplum_conn`
|
||||
- **Conn Type:** `Postgres`
|
||||
- **Host:** `greenplum`
|
||||
- **Schema:** `gpadmin`
|
||||
- **Login:** `gpadmin`
|
||||
- **Password:** `gpadmin`
|
||||
- **Port:** `5432`
|
||||
|
||||
---
|
||||
|
||||
## 📋 Что входит в стенд
|
||||
|
||||
### Основные компоненты
|
||||
- **Greenplum** — аналитическая база данных для хранения и анализа данных
|
||||
- **Airflow** — оркестратор workflow и задач
|
||||
- **Postgres** — база метаданных для Airflow
|
||||
- **pandas** — библиотека для генерации и анализа данных в формате CSV
|
||||
|
||||
### Готовые DAG (workflow)
|
||||
- **csv_to_greenplum** — базовый pipeline: pandas → CSV → Greenplum
|
||||
- **greenplum_data_quality** — проверки качества данных (наличие таблицы, схема, дубликаты)
|
||||
|
||||
### Полезные команды
|
||||
```bash
|
||||
# Основные команды
|
||||
make up # Запустить весь стенд
|
||||
make down # Остановить и удалить данные
|
||||
make airflow-init # Инициализировать Airflow
|
||||
make ddl-gp # Применить DDL к Greenplum
|
||||
make gp-psql # Подключиться к Greenplum через psql
|
||||
|
||||
# Проверка данных
|
||||
make logs # Следить за логами Airflow
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## ⚙️ Настройка через переменные окружения
|
||||
|
||||
Все настройки находятся в файле `.env`. Основные параметры:
|
||||
|
||||
### Greenplum
|
||||
- `GP_USER` — пользователь (по умолчанию: gpadmin)
|
||||
- `GP_PASSWORD` — пароль (по умолчанию: gpadmin)
|
||||
- `GP_DB` — база данных (по умолчанию: gpadmin)
|
||||
- `GP_PORT` — порт (по умолчанию: 5432)
|
||||
|
||||
### CSV pipeline
|
||||
- `CSV_DIR` — путь к каталогу с CSV внутри контейнеров Airflow (по умолчанию: `/opt/airflow/data`)
|
||||
- `CSV_ROWS` — количество строк, генерируемых DAG (по умолчанию: 1000)
|
||||
|
||||
### Airflow
|
||||
- `GP_CONN_ID` — ID подключения (по умолчанию: greenplum_conn)
|
||||
|
||||
---
|
||||
|
||||
## 🔍 Продвинутые темы
|
||||
|
||||
### Архитектура pipeline
|
||||
|
||||
**Поток данных в DAG `csv_to_greenplum`:**
|
||||
1. `create_orders_table` — создаёт таблицу `public.orders` в Greenplum
|
||||
2. `generate_csv` — генерирует датасет при помощи pandas и сохраняет CSV в `CSV_DIR`
|
||||
3. `preview_csv` — выводит предпросмотр и статистику по данным
|
||||
4. `load_csv_to_greenplum` — загружает CSV во временную таблицу и переносит новые строки в `public.orders`
|
||||
|
||||
> 💡 **Безопасность повторного запуска:** Pipeline защищен от дубликатов, поэтому его можно запускать многократно.
|
||||
|
||||
### Проверка качества данных
|
||||
|
||||
Запустите DAG `greenplum_data_quality` для автоматической проверки:
|
||||
- Наличие таблицы в базе
|
||||
- Соответствие схемы ожидаемой структуре
|
||||
- Объем загруженных данных
|
||||
- Отсутствие дубликатов записей
|
||||
|
||||
### Ограничения учебного стенда
|
||||
|
||||
- **Greenplum** запущен в single-node режиме (для обучения)
|
||||
- В продакшене Greenplum обычно разворачивают кластером на нескольких серверах
|
||||
- Используется Greenplum 6 (широко доступная версия), хотя Greenplum 7 предлагает больше возможностей
|
||||
|
||||
---
|
||||
|
||||
## 🆘 Типичные проблемы и решения
|
||||
|
||||
| Проблема | Решение |
|
||||
|----------|---------|
|
||||
| Airflow UI не открывается | Дождитесь сообщения `Listening at: http://0.0.0.0:8080` в логах (`make logs`) |
|
||||
| Ошибка подключения к Greenplum | Убедитесь, что контейнер `greenplum` стал статусом `healthy` (проверьте `docker compose ps`) |
|
||||
| Нет файла в `./data` после запуска DAG | Проверьте логи задачи `generate_csv`, убедитесь, что `CSV_DIR` смонтирован в docker-compose |
|
||||
| Команда `make` не найдена | Используйте полные команды `docker compose` или установите make |
|
||||
| Greenplum не стартует/падает при старте | Выполните `make down`, затем `make up && make airflow-init` (очищает тома и поднимает заново) |
|
||||
|
||||
---
|
||||
|
||||
## 📁 Структура проекта
|
||||
|
||||
```
|
||||
├── docker-compose.yml # Описание всех сервисов
|
||||
├── .env.example # Шаблон настроек
|
||||
├── Makefile # Удобные команды для работы
|
||||
├── airflow/
|
||||
│ └── dags/ # Файлы workflow (DAG)
|
||||
│ ├── csv_to_greenplum.py
|
||||
│ └── data_quality_greenplum.py
|
||||
└── sql/
|
||||
└── ddl_gp.sql # Создание таблицы в Greenplum
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## 💡 Советы для дальнейшего обучения
|
||||
|
||||
1. **Поэкспериментируйте с DAG** — измените параметры генерации данных или размер батча
|
||||
2. **Добавьте свои проверки** — расширьте DAG `data_quality_greenplum.py`
|
||||
3. **Попробуйте другие источники** — замените генератор данных на чтение из файла или API
|
||||
4. **Изучите Airflow deeper** — добавьте зависимости между задачами, настройте расписания
|
||||
|
||||
---
|
||||
|
||||
## ✅ Тестирование
|
||||
|
||||
- Локальные проверки: `make test` (pytest). Для форматирования — `make fmt`, для проверки — `make lint`.
|
||||
- Пошаговый сценарий с Docker (включая негативные кейсы и reset) — см. `TESTING.md`.
|
||||
|
||||
Удачи в изучении Data Engineering! 🚀
|
||||
|
||||
+71
@@ -0,0 +1,71 @@
|
||||
# План тестирования (для студентов)
|
||||
|
||||
Этот документ — пошаговый чек‑лист, как проверить, что всё работает: от «быстрых локальных проверок» до запуска стенда в Docker и просмотра данных в Greenplum. Подходит начинающим: просто выполняйте шаги по порядку.
|
||||
|
||||
Если что‑то пошло не так, смотрите раздел «Быстрый reset» ниже.
|
||||
|
||||
## 1. Быстрая проверка окружения
|
||||
- `uv sync` — подтягиваем Python и зависимости из `pyproject.toml`/`uv.lock`.
|
||||
- Проверяем версию uv: `uv --version` (ожидаем ≥ 0.9).
|
||||
- Убедитесь, что `docker compose version` доступна и Docker запущен.
|
||||
|
||||
## 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` должен пройти.
|
||||
- (опционально) `uv run pytest -q -k dags_smoke` — только DAG smoke.
|
||||
|
||||
## 3. Подготовка Docker-стенда
|
||||
- `cp .env.example .env` (если файла ещё нет) и проверьте переменные:
|
||||
- `GP_PORT` не конфликтует с локальным PostgreSQL.
|
||||
- `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 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`:
|
||||
- Включить переключатель.
|
||||
- Нажать «Trigger DAG».
|
||||
- Контроль: все таски Success, в `data/` появился CSV, в логах `load_csv_to_greenplum` видно `INSERT`.
|
||||
- В Greenplum (см. п.5) убедиться в наличии строк `(SELECT COUNT(*) ...)`.
|
||||
3. DAG `greenplum_data_quality`:
|
||||
- Запустить вручную после первого DAG.
|
||||
- Проверить, что все 5 задач Success и логи содержат `Проверка пройдена`.
|
||||
|
||||
## 5. Проверка данных в Greenplum
|
||||
- `make gp-psql` — запустить psql в контейнере от имени `gpadmin`.
|
||||
- Команды внутри 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;` — поиск дублей.
|
||||
- Завершить `\q`.
|
||||
|
||||
## 6. Негативные сценарии и fallback
|
||||
- **Пустая таблица**: запустить `greenplum_data_quality` до `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` не найдёт дублей.
|
||||
|
||||
## 7. Быстрый reset (если «что-то сломалось»)
|
||||
- Перезапустить стенд с очисткой данных:
|
||||
- `make down` — остановит контейнеры и удалит тома.
|
||||
- `make up && make airflow-init` — заново поднимет всё и проинициализирует Airflow.
|
||||
- Иногда Greenplum не стартует после «грязных» остановок (из‑за старых внутренних файлов). Лечение: всегда делайте `make down` перед повторным `make up`.
|
||||
|
||||
## 8. Снятие метрик и мониторинг
|
||||
- Контейнеры: `docker compose ps`, `docker stats` (по желанию).
|
||||
- Логи задач: в Airflow UI → конкретный таск → Log.
|
||||
- Хостовые CSV: каталог `data/` (можно открыть любой файл и убедиться в структуре).
|
||||
|
||||
## 9. Завершение работы
|
||||
- `make down` — выключает сервисы и удаляет тома (перезапишет данные в Greenplum!).
|
||||
- При необходимости сохранить данные: скопировать CSV из `data/` и дампы из контейнера до `make down`.
|
||||
|
||||
## Текущий статус (пример успешного прогона)
|
||||
- `uv run pytest -q` — 11 passed, 2 smoke-теста DAG пропущены (Airflow не установлен в venv).
|
||||
- `make lint` — падает, потому что `airflow/dags/*.py` не отформатированы black/isort. После `make fmt` проблема уйдёт.
|
||||
- Docker-стенд не запускался в рамках этой сессии; ожидается, что инструкции выше обеспечат полноценную проверку.
|
||||
@@ -0,0 +1,160 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import os
|
||||
import random
|
||||
from datetime import datetime, timedelta
|
||||
from pathlib import Path
|
||||
from typing import List
|
||||
|
||||
import pandas as pd
|
||||
from airflow import DAG
|
||||
from airflow.operators.python import PythonOperator
|
||||
from helpers.greenplum import get_gp_conn
|
||||
|
||||
CSV_DIR = Path(os.getenv("CSV_DIR", "/opt/airflow/data"))
|
||||
CSV_ROWS = int(os.getenv("CSV_ROWS", "1000"))
|
||||
|
||||
|
||||
def _create_table() -> None:
|
||||
"""Создаёт таблицу public.orders, если она ещё не существует."""
|
||||
ddl = """
|
||||
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);
|
||||
"""
|
||||
with get_gp_conn() as conn, conn.cursor() as cur:
|
||||
cur.execute(ddl)
|
||||
conn.commit()
|
||||
|
||||
|
||||
def _generate_csv(rows: int, csv_dir: Path) -> str:
|
||||
"""Генерирует CSV c заказами с помощью pandas и сохраняет на диск."""
|
||||
csv_dir.mkdir(parents=True, exist_ok=True)
|
||||
timestamp = datetime.utcnow().strftime("%Y%m%d_%H%M%S")
|
||||
csv_path = csv_dir / f"orders_{timestamp}.csv"
|
||||
|
||||
# Генерируем данные в pandas-стиле
|
||||
base_order_id = int(datetime.utcnow().timestamp() * 1_000)
|
||||
|
||||
# Создаём DataFrame с использованием pandas методов
|
||||
df = pd.DataFrame({
|
||||
# Уникальные order_id начиная с базового значения
|
||||
"order_id": pd.Series(range(base_order_id, base_order_id + rows), dtype="int64"),
|
||||
|
||||
# Временные метки с интервалом в 1 секунду в обратном порядке
|
||||
"order_ts": pd.date_range(
|
||||
end=datetime.utcnow(),
|
||||
periods=rows,
|
||||
freq="1S"
|
||||
).sort_values(ascending=False),
|
||||
|
||||
# Случайные customer_id от 1 до 1000
|
||||
"customer_id": pd.Series(
|
||||
random.choices(range(1, 1001), k=rows),
|
||||
dtype="int64"
|
||||
),
|
||||
|
||||
# Случайные суммы от 10 до 500 с округлением до 2 знаков
|
||||
"amount": pd.Series(
|
||||
[round(random.uniform(10, 500), 2) for _ in range(rows)],
|
||||
dtype="float64"
|
||||
)
|
||||
})
|
||||
|
||||
# Сохраняем CSV без индекса
|
||||
df.to_csv(csv_path, index=False)
|
||||
logging.info("CSV сохранён: %s (строк: %s)", csv_path, len(df))
|
||||
return str(csv_path)
|
||||
|
||||
|
||||
def _preview_csv(csv_path: str, sample_rows: int = 5) -> None:
|
||||
"""Отображает предпросмотр CSV через pandas (head и describe)."""
|
||||
df = pd.read_csv(csv_path)
|
||||
df["order_ts"] = pd.to_datetime(df["order_ts"], errors="coerce")
|
||||
logging.info("Первые %s строк:\n%s", sample_rows, df.head(sample_rows).to_string(index=False))
|
||||
numeric_summary = df.describe(include="number")
|
||||
logging.info("Числовая статистика:\n%s", numeric_summary.to_string())
|
||||
if df["order_ts"].notna().any():
|
||||
logging.info(
|
||||
"Диапазон order_ts: %s → %s",
|
||||
df["order_ts"].min().isoformat(),
|
||||
df["order_ts"].max().isoformat(),
|
||||
)
|
||||
|
||||
|
||||
def _load_csv(csv_path: str) -> None:
|
||||
"""Загружает CSV в Greenplum через временную таблицу и anti-join."""
|
||||
csv_file = Path(csv_path)
|
||||
if not csv_file.exists():
|
||||
raise FileNotFoundError(f"CSV не найден: {csv_file}")
|
||||
|
||||
with get_gp_conn() as conn, conn.cursor() as cur, csv_file.open("r", encoding="utf-8") as f:
|
||||
cur.execute("CREATE TEMP TABLE tmp_orders (LIKE public.orders INCLUDING DEFAULTS) ON COMMIT DROP;")
|
||||
cur.copy_expert(
|
||||
"COPY tmp_orders (order_id, order_ts, customer_id, amount) FROM STDIN WITH CSV HEADER",
|
||||
f,
|
||||
)
|
||||
|
||||
cur.execute("SELECT COUNT(*) FROM tmp_orders")
|
||||
tmp_rows = cur.fetchone()[0]
|
||||
|
||||
cur.execute(
|
||||
"""
|
||||
INSERT INTO public.orders(order_id, order_ts, customer_id, amount)
|
||||
SELECT t.order_id, t.order_ts, t.customer_id, t.amount
|
||||
FROM tmp_orders t
|
||||
LEFT JOIN public.orders o ON o.order_id = t.order_id
|
||||
WHERE o.order_id IS NULL
|
||||
"""
|
||||
)
|
||||
inserted = cur.rowcount if cur.rowcount != -1 else 0
|
||||
conn.commit()
|
||||
|
||||
logging.info("Загружено строк: %s (прочитано из CSV: %s)", inserted, tmp_rows)
|
||||
|
||||
|
||||
default_args = {"owner": "airflow", "retries": 1, "retry_delay": timedelta(seconds=30)}
|
||||
|
||||
with DAG(
|
||||
dag_id="csv_to_greenplum",
|
||||
start_date=datetime(2024, 1, 1),
|
||||
schedule=None,
|
||||
catchup=False,
|
||||
default_args=default_args,
|
||||
tags=["demo", "greenplum", "csv"],
|
||||
) as dag:
|
||||
create_table = PythonOperator(
|
||||
task_id="create_orders_table",
|
||||
python_callable=_create_table,
|
||||
)
|
||||
|
||||
generate_csv = PythonOperator(
|
||||
task_id="generate_csv",
|
||||
python_callable=_generate_csv,
|
||||
op_kwargs={"rows": CSV_ROWS, "csv_dir": CSV_DIR},
|
||||
)
|
||||
|
||||
preview_csv = PythonOperator(
|
||||
task_id="preview_csv",
|
||||
python_callable=_preview_csv,
|
||||
op_kwargs={
|
||||
"csv_path": "{{ ti.xcom_pull(task_ids='generate_csv') }}",
|
||||
"sample_rows": 5,
|
||||
},
|
||||
)
|
||||
|
||||
load_csv = PythonOperator(
|
||||
task_id="load_csv_to_greenplum",
|
||||
python_callable=_load_csv,
|
||||
op_kwargs={
|
||||
"csv_path": "{{ ti.xcom_pull(task_ids='generate_csv') }}",
|
||||
},
|
||||
)
|
||||
|
||||
create_table >> generate_csv >> preview_csv >> load_csv
|
||||
@@ -0,0 +1,96 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from datetime import datetime, timedelta
|
||||
|
||||
from airflow import DAG
|
||||
from airflow.operators.python import PythonOperator
|
||||
|
||||
from helpers.greenplum import (
|
||||
assert_orders_have_rows,
|
||||
assert_orders_no_duplicates,
|
||||
assert_orders_schema,
|
||||
assert_orders_table_exists,
|
||||
get_gp_conn,
|
||||
)
|
||||
|
||||
|
||||
def _run_check(check_callable):
|
||||
"""
|
||||
Оборачивает проверку качества данных в контекст подключения к Greenplum.
|
||||
|
||||
Этот DAG предназначен для автоматической проверки качества данных в таблице orders:
|
||||
1. Проверяет существование таблицы
|
||||
2. Проверяет соответствие схемы
|
||||
3. Проверяет наличие данных
|
||||
4. Проверяет отсутствие дубликатов
|
||||
|
||||
Args:
|
||||
check_callable: Функция проверки, принимающая подключение к БД
|
||||
"""
|
||||
# Получаем имя функции для логов
|
||||
check_name = check_callable.__name__.replace("assert_", "")
|
||||
logging.info("🚀 Запуск проверки: %s", check_name)
|
||||
|
||||
with get_gp_conn() as conn:
|
||||
check_callable(conn)
|
||||
|
||||
logging.info("✅ Проверка пройдена: %s", check_name)
|
||||
|
||||
|
||||
def _log_dq_summary():
|
||||
"""
|
||||
Логирует итоговую сводку по качеству данных.
|
||||
Эта задача выполняется после всех проверок и показывает общий результат.
|
||||
"""
|
||||
logging.info("🎉 Все проверки качества данных пройдены успешно!")
|
||||
logging.info("📊 Качество данных в таблице orders соответствует требованиям.")
|
||||
|
||||
|
||||
default_args = {"owner": "airflow", "retries": 1, "retry_delay": timedelta(seconds=30)}
|
||||
|
||||
with DAG(
|
||||
dag_id="greenplum_data_quality",
|
||||
start_date=datetime(2024, 1, 1),
|
||||
schedule=None,
|
||||
catchup=False,
|
||||
default_args=default_args,
|
||||
tags=["demo", "greenplum", "quality"],
|
||||
description="Автоматизированные проверки качества данных в Greenplum",
|
||||
) as dag:
|
||||
# Задача 1: Проверка существования таблицы
|
||||
check_exists = PythonOperator(
|
||||
task_id="check_orders_table_exists",
|
||||
python_callable=_run_check,
|
||||
op_args=[assert_orders_table_exists],
|
||||
)
|
||||
|
||||
# Задача 2: Проверка соответствия схемы таблицы
|
||||
check_schema = PythonOperator(
|
||||
task_id="check_orders_schema",
|
||||
python_callable=_run_check,
|
||||
op_args=[assert_orders_schema],
|
||||
)
|
||||
|
||||
# Задача 3: Проверка наличия данных
|
||||
check_has_rows = PythonOperator(
|
||||
task_id="check_orders_has_rows",
|
||||
python_callable=_run_check,
|
||||
op_args=[assert_orders_have_rows],
|
||||
)
|
||||
|
||||
# Задача 4: Проверка отсутствия дубликатов
|
||||
check_no_duplicates = PythonOperator(
|
||||
task_id="check_order_duplicates",
|
||||
python_callable=_run_check,
|
||||
op_args=[assert_orders_no_duplicates],
|
||||
)
|
||||
|
||||
# Задача 5: Итоговая сводка
|
||||
dq_summary = PythonOperator(
|
||||
task_id="data_quality_summary",
|
||||
python_callable=_log_dq_summary,
|
||||
)
|
||||
|
||||
# Определяем последовательность выполнения задач
|
||||
check_exists >> check_schema >> check_has_rows >> check_no_duplicates >> dq_summary
|
||||
@@ -0,0 +1,199 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import os
|
||||
from typing import List, Sequence, Tuple
|
||||
|
||||
import psycopg2
|
||||
|
||||
# Настройки для подключения к Greenplum. По умолчанию используем Airflow Connection,
|
||||
# но при проблемах можно переключиться на ENV-подключение, установив GP_USE_AIRFLOW_CONN=false.
|
||||
GP_CONN_ID = os.getenv("GP_CONN_ID", "greenplum_conn")
|
||||
GP_USE_AIRFLOW_CONN = os.getenv("GP_USE_AIRFLOW_CONN", "true").lower() in ("1", "true", "yes")
|
||||
|
||||
# Ожидаемая схема таблицы orders для проверки качества данных
|
||||
EXPECTED_ORDERS_SCHEMA: List[Tuple[str, str]] = [
|
||||
("order_id", "bigint"),
|
||||
("order_ts", "timestamp without time zone"),
|
||||
("customer_id", "bigint"),
|
||||
("amount", "numeric"),
|
||||
]
|
||||
|
||||
|
||||
def get_gp_conn():
|
||||
"""
|
||||
Возвращает psycopg2 connection к Greenplum.
|
||||
|
||||
Приоритет подключения:
|
||||
1. Через Airflow Connection (если настроено и доступно)
|
||||
2. Прямое подключение по переменным окружения (фоллбек)
|
||||
|
||||
Returns:
|
||||
psycopg2 connection object
|
||||
"""
|
||||
if GP_USE_AIRFLOW_CONN:
|
||||
try:
|
||||
from airflow.providers.postgres.hooks.postgres import PostgresHook
|
||||
|
||||
hook = PostgresHook(postgres_conn_id=GP_CONN_ID)
|
||||
conn = hook.get_conn()
|
||||
logging.info("✅ Подключение через Airflow Connection успешно")
|
||||
return conn
|
||||
except Exception as e:
|
||||
logging.warning("⚠️ Не удалось подключиться через Airflow Connection: %s", e)
|
||||
logging.info("🔄 Переключаемся на прямое подключение по ENV переменным")
|
||||
# Фоллбек на прямое подключение по переменным окружения.
|
||||
|
||||
# Прямое подключение по переменным окружения
|
||||
conn_params = {
|
||||
"dbname": os.getenv("GP_DB", "gpadmin"),
|
||||
"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"])
|
||||
return psycopg2.connect(**conn_params)
|
||||
|
||||
|
||||
def assert_orders_table_exists(conn) -> None:
|
||||
"""
|
||||
Проверяет наличие таблицы orders в схеме public.
|
||||
|
||||
Args:
|
||||
conn: Подключение к Greenplum
|
||||
|
||||
Raises:
|
||||
ValueError: Если таблица не найдена
|
||||
"""
|
||||
logging.info("🔍 Проверяем существование таблицы public.orders...")
|
||||
with conn.cursor() as cur:
|
||||
cur.execute(
|
||||
"""
|
||||
SELECT 1
|
||||
FROM pg_catalog.pg_tables
|
||||
WHERE schemaname = 'public' AND tablename = 'orders'
|
||||
"""
|
||||
)
|
||||
if cur.fetchone() is None:
|
||||
raise ValueError("❌ Таблица public.orders не найдена; запусти DAG csv_to_greenplum.")
|
||||
logging.info("✅ Таблица public.orders существует")
|
||||
|
||||
|
||||
def fetch_orders_schema(conn) -> Sequence[Tuple[str, str]]:
|
||||
"""
|
||||
Получает схему таблицы orders из information_schema.
|
||||
|
||||
Args:
|
||||
conn: Подключение к Greenplum
|
||||
|
||||
Returns:
|
||||
Список кортежей (имя_колонки, тип_данных)
|
||||
"""
|
||||
with conn.cursor() as cur:
|
||||
cur.execute(
|
||||
"""
|
||||
SELECT column_name, data_type
|
||||
FROM information_schema.columns
|
||||
WHERE table_schema = 'public' AND table_name = 'orders'
|
||||
ORDER BY ordinal_position
|
||||
"""
|
||||
)
|
||||
return cur.fetchall()
|
||||
|
||||
|
||||
def assert_orders_schema(conn) -> None:
|
||||
"""
|
||||
Проверяет, что схема таблицы orders соответствует ожидаемой.
|
||||
|
||||
Args:
|
||||
conn: Подключение к Greenplum
|
||||
|
||||
Raises:
|
||||
ValueError: Если схема не соответствует ожидаемой
|
||||
"""
|
||||
logging.info("📋 Проверяем схему таблицы orders...")
|
||||
schema = fetch_orders_schema(conn)
|
||||
logging.info("📊 Фактическая схема: %s", list(schema))
|
||||
logging.info("📊 Ожидаемая схема: %s", EXPECTED_ORDERS_SCHEMA)
|
||||
|
||||
if list(schema) != EXPECTED_ORDERS_SCHEMA:
|
||||
raise ValueError(f"❌ Неожиданная схема orders: {schema}. Ожидали {EXPECTED_ORDERS_SCHEMA}.")
|
||||
logging.info("✅ Схема таблицы orders соответствует ожиданиям")
|
||||
|
||||
|
||||
def fetch_orders_count(conn) -> int:
|
||||
"""
|
||||
Получает количество строк в таблице orders.
|
||||
|
||||
Args:
|
||||
conn: Подключение к Greenplum
|
||||
|
||||
Returns:
|
||||
Количество строк в таблице
|
||||
"""
|
||||
with conn.cursor() as cur:
|
||||
cur.execute("SELECT COUNT(*) FROM public.orders")
|
||||
return cur.fetchone()[0]
|
||||
|
||||
|
||||
def assert_orders_have_rows(conn) -> None:
|
||||
"""
|
||||
Проверяет, что таблица orders не пустая.
|
||||
|
||||
Args:
|
||||
conn: Подключение к Greenplum
|
||||
|
||||
Raises:
|
||||
ValueError: Если таблица пустая
|
||||
"""
|
||||
logging.info("📊 Проверяем наличие данных в таблице orders...")
|
||||
row_count = fetch_orders_count(conn)
|
||||
logging.info("📈 Количество строк в orders: %s", row_count)
|
||||
|
||||
if row_count <= 0:
|
||||
raise ValueError("❌ Таблица public.orders пустая — запусти DAG csv_to_greenplum перед проверкой.")
|
||||
logging.info("✅ Таблица orders содержит данные (%s строк)", row_count)
|
||||
|
||||
|
||||
def fetch_orders_duplicates(conn) -> int:
|
||||
"""
|
||||
Подсчитывает количество дубликатов по order_id.
|
||||
|
||||
Args:
|
||||
conn: Подключение к Greenplum
|
||||
|
||||
Returns:
|
||||
Количество дублирующихся order_id
|
||||
"""
|
||||
with conn.cursor() as cur:
|
||||
cur.execute(
|
||||
"""
|
||||
SELECT COUNT(*) FROM (
|
||||
SELECT order_id
|
||||
FROM public.orders
|
||||
GROUP BY order_id
|
||||
HAVING COUNT(*) > 1
|
||||
) d
|
||||
"""
|
||||
)
|
||||
return cur.fetchone()[0]
|
||||
|
||||
|
||||
def assert_orders_no_duplicates(conn) -> None:
|
||||
"""
|
||||
Проверяет, что в таблице нет дублей по order_id.
|
||||
|
||||
Args:
|
||||
conn: Подключение к Greenplum
|
||||
|
||||
Raises:
|
||||
ValueError: Если обнаружены дубликаты
|
||||
"""
|
||||
logging.info("🔍 Проверяем отсутствие дубликатов по order_id...")
|
||||
duplicates = fetch_orders_duplicates(conn)
|
||||
logging.info("📊 Найдено дубликатов: %s", duplicates)
|
||||
|
||||
if duplicates:
|
||||
raise ValueError(f"❌ Обнаружены дубли по order_id ({duplicates} шт.) — проверь загрузку данных.")
|
||||
logging.info("✅ Дубликаты не обнаружены")
|
||||
@@ -0,0 +1,2 @@
|
||||
psycopg2-binary==2.9.9
|
||||
pandas==2.1.4
|
||||
@@ -0,0 +1,113 @@
|
||||
services:
|
||||
# Postgres только для Airflow метаданных
|
||||
pgmeta:
|
||||
image: postgres:16
|
||||
# container_name: gp_pgmeta
|
||||
env_file: .env
|
||||
environment:
|
||||
POSTGRES_USER: ${PG_USER}
|
||||
POSTGRES_PASSWORD: ${PG_PASSWORD}
|
||||
POSTGRES_DB: ${PG_DB}
|
||||
ports:
|
||||
- "5433:5432"
|
||||
volumes:
|
||||
- pgmeta:/var/lib/postgresql/data
|
||||
healthcheck:
|
||||
test: ["CMD-SHELL", "pg_isready -U ${PG_USER} -d ${PG_DB}"]
|
||||
interval: 5s
|
||||
timeout: 5s
|
||||
retries: 20
|
||||
|
||||
greenplum:
|
||||
image: woblerr/greenplum:6.27.1
|
||||
# container_name: gp_single
|
||||
# hostname: gpdbsne
|
||||
environment:
|
||||
GREENPLUM_USER: ${GP_USER:-gpadmin}
|
||||
GREENPLUM_PASSWORD: ${GP_PASSWORD:-gpadmin}
|
||||
GREENPLUM_DATABASE_NAME: ${GP_DB:-gpadmin}
|
||||
# GP_PORT: ${GP_PORT:-5432}
|
||||
# Порты: внешний 5432
|
||||
ports:
|
||||
- "${GP_PORT}:5432"
|
||||
volumes:
|
||||
- ./sql:/sql:ro
|
||||
- greenplum_data:/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"]
|
||||
interval: 10s
|
||||
timeout: 5s
|
||||
retries: 30
|
||||
|
||||
airflow-webserver:
|
||||
image: apache/airflow:2.9.2
|
||||
container_name: gp_airflow_web
|
||||
env_file: .env
|
||||
environment:
|
||||
AIRFLOW__CORE__LOAD_EXAMPLES: "False"
|
||||
AIRFLOW__DATABASE__SQL_ALCHEMY_CONN: postgresql+psycopg2://${PG_USER}:${PG_PASSWORD}@pgmeta:5432/${PG_DB}
|
||||
command: >
|
||||
bash -lc "pip install --no-cache-dir -r /opt/airflow/requirements.txt &&
|
||||
airflow webserver"
|
||||
ports:
|
||||
- "8080:8080"
|
||||
volumes:
|
||||
- ./airflow/dags:/opt/airflow/dags
|
||||
- ./airflow/requirements.txt:/opt/airflow/requirements.txt
|
||||
- ./data:/opt/airflow/data
|
||||
depends_on:
|
||||
pgmeta:
|
||||
condition: service_healthy
|
||||
greenplum:
|
||||
condition: service_healthy
|
||||
|
||||
airflow-scheduler:
|
||||
image: apache/airflow:2.9.2
|
||||
container_name: gp_airflow_sch
|
||||
env_file: .env
|
||||
environment:
|
||||
AIRFLOW__CORE__LOAD_EXAMPLES: "False"
|
||||
AIRFLOW__DATABASE__SQL_ALCHEMY_CONN: postgresql+psycopg2://${PG_USER}:${PG_PASSWORD}@pgmeta:5432/${PG_DB}
|
||||
command: >
|
||||
bash -lc "pip install --no-cache-dir -r /opt/airflow/requirements.txt &&
|
||||
airflow scheduler"
|
||||
volumes:
|
||||
- ./airflow/dags:/opt/airflow/dags
|
||||
- ./airflow/requirements.txt:/opt/airflow/requirements.txt
|
||||
- ./data:/opt/airflow/data
|
||||
depends_on:
|
||||
pgmeta:
|
||||
condition: service_healthy
|
||||
greenplum:
|
||||
condition: service_healthy
|
||||
|
||||
airflow-init:
|
||||
image: apache/airflow:2.9.2
|
||||
# container_name: gp_airflow_init
|
||||
env_file: .env
|
||||
environment:
|
||||
AIRFLOW__CORE__LOAD_EXAMPLES: "False"
|
||||
AIRFLOW__DATABASE__SQL_ALCHEMY_CONN: postgresql+psycopg2://${PG_USER}:${PG_PASSWORD}@pgmeta:5432/${PG_DB}
|
||||
volumes:
|
||||
- ./airflow/dags:/opt/airflow/dags
|
||||
- ./airflow/requirements.txt:/opt/airflow/requirements.txt
|
||||
- ./data:/opt/airflow/data
|
||||
command: >
|
||||
bash -lc "
|
||||
set -e;
|
||||
pip install --no-cache-dir -r /opt/airflow/requirements.txt;
|
||||
# Дожидаемся готовности БД ретрая миграции
|
||||
for i in {1..30}; do
|
||||
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
|
||||
"
|
||||
depends_on:
|
||||
pgmeta:
|
||||
condition: service_healthy
|
||||
|
||||
volumes:
|
||||
pgmeta:
|
||||
greenplum_data:
|
||||
@@ -0,0 +1,16 @@
|
||||
[project]
|
||||
name = "airflow-greenplum"
|
||||
version = "0.1.0"
|
||||
requires-python = ">=3.11,<3.12"
|
||||
dependencies = []
|
||||
|
||||
[dependency-groups]
|
||||
dev = [
|
||||
"apache-airflow==2.9.2",
|
||||
"apache-airflow-providers-postgres==5.11.1",
|
||||
"black==24.4.2",
|
||||
"isort>=7.0.0",
|
||||
"pandas==2.1.4",
|
||||
"psycopg2-binary==2.9.9",
|
||||
"pytest==7.4.4",
|
||||
]
|
||||
@@ -0,0 +1,12 @@
|
||||
-- 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
|
||||
)
|
||||
WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1)
|
||||
DISTRIBUTED BY (order_id);
|
||||
@@ -0,0 +1,69 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import importlib
|
||||
import sys
|
||||
from pathlib import Path
|
||||
from types import ModuleType
|
||||
from typing import Type
|
||||
|
||||
|
||||
PROJECT_ROOT = Path(__file__).resolve().parents[1]
|
||||
if str(PROJECT_ROOT) not in sys.path:
|
||||
sys.path.append(str(PROJECT_ROOT))
|
||||
|
||||
if "airflow" not in sys.modules:
|
||||
airflow_module = ModuleType("airflow")
|
||||
airflow_module.__path__ = [str(PROJECT_ROOT / "airflow")]
|
||||
sys.modules["airflow"] = airflow_module
|
||||
|
||||
providers_module = ModuleType("airflow.providers")
|
||||
providers_module.__path__ = []
|
||||
sys.modules["airflow.providers"] = providers_module
|
||||
airflow_module.providers = providers_module
|
||||
|
||||
postgres_module = ModuleType("airflow.providers.postgres")
|
||||
postgres_module.__path__ = []
|
||||
sys.modules["airflow.providers.postgres"] = postgres_module
|
||||
providers_module.postgres = postgres_module
|
||||
|
||||
hooks_module = ModuleType("airflow.providers.postgres.hooks")
|
||||
hooks_module.__path__ = []
|
||||
sys.modules["airflow.providers.postgres.hooks"] = hooks_module
|
||||
postgres_module.hooks = hooks_module
|
||||
|
||||
if "psycopg2" not in sys.modules:
|
||||
psycopg2_stub = ModuleType("psycopg2")
|
||||
psycopg2_stub.connect = lambda **_: None # type: ignore[assignment]
|
||||
sys.modules["psycopg2"] = psycopg2_stub
|
||||
|
||||
|
||||
def _ensure_stub_module(full_name: str) -> ModuleType:
|
||||
"""
|
||||
Ensure that module placeholders exist for a dotted path and return leaf module.
|
||||
"""
|
||||
parts = full_name.split(".")
|
||||
module: ModuleType | None = None
|
||||
path = ""
|
||||
for part in parts:
|
||||
path = f"{path}.{part}" if path else part
|
||||
if path not in sys.modules:
|
||||
new_module = ModuleType(path)
|
||||
if module is not None:
|
||||
setattr(module, part, new_module)
|
||||
sys.modules[path] = new_module
|
||||
module = new_module
|
||||
else:
|
||||
module = sys.modules[path]
|
||||
assert isinstance(module, ModuleType)
|
||||
return module
|
||||
|
||||
|
||||
def patch_postgres_hook(monkeypatch, hook_cls: Type) -> None:
|
||||
"""
|
||||
Patch PostgresHook so that helpers.greenplum can be exercised without real Airflow.
|
||||
"""
|
||||
try:
|
||||
module = importlib.import_module("airflow.providers.postgres.hooks.postgres")
|
||||
except ModuleNotFoundError:
|
||||
module = _ensure_stub_module("airflow.providers.postgres.hooks.postgres")
|
||||
monkeypatch.setattr(module, "PostgresHook", hook_cls, raising=False)
|
||||
@@ -0,0 +1,70 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import importlib
|
||||
|
||||
import pytest
|
||||
|
||||
|
||||
def _airflow_available() -> bool:
|
||||
try:
|
||||
af = importlib.import_module("airflow")
|
||||
except Exception:
|
||||
return False
|
||||
# Real Airflow exposes DAG at top-level
|
||||
return hasattr(af, "DAG")
|
||||
|
||||
|
||||
pytestmark = pytest.mark.skipif(not _airflow_available(), reason="Airflow is not installed for DAG smoke tests")
|
||||
|
||||
|
||||
def _load_dag(module_name: str):
|
||||
mod = importlib.import_module(module_name)
|
||||
assert hasattr(mod, "dag"), f"{module_name} must expose variable 'dag'"
|
||||
return getattr(mod, "dag")
|
||||
|
||||
|
||||
def test_csv_to_greenplum_dag_structure():
|
||||
dag = _load_dag("airflow.dags.csv_to_greenplum")
|
||||
|
||||
# tasks
|
||||
expected_tasks = {
|
||||
"create_orders_table",
|
||||
"generate_csv",
|
||||
"preview_csv",
|
||||
"load_csv_to_greenplum",
|
||||
}
|
||||
assert expected_tasks.issubset(dag.task_dict.keys())
|
||||
|
||||
# linear dependencies
|
||||
t1 = dag.get_task("create_orders_table")
|
||||
t2 = dag.get_task("generate_csv")
|
||||
t3 = dag.get_task("preview_csv")
|
||||
t4 = dag.get_task("load_csv_to_greenplum")
|
||||
|
||||
assert t2 in t1.get_direct_relatives("downstream")
|
||||
assert t3 in t2.get_direct_relatives("downstream")
|
||||
assert t4 in t3.get_direct_relatives("downstream")
|
||||
|
||||
|
||||
def test_data_quality_greenplum_dag_structure():
|
||||
dag = _load_dag("airflow.dags.data_quality_greenplum")
|
||||
|
||||
expected_tasks = {
|
||||
"check_orders_table_exists",
|
||||
"check_orders_schema",
|
||||
"check_orders_has_rows",
|
||||
"check_order_duplicates",
|
||||
"data_quality_summary",
|
||||
}
|
||||
assert expected_tasks.issubset(dag.task_dict.keys())
|
||||
|
||||
e = dag.get_task("check_orders_table_exists")
|
||||
s = dag.get_task("check_orders_schema")
|
||||
h = dag.get_task("check_orders_has_rows")
|
||||
d = dag.get_task("check_order_duplicates")
|
||||
q = dag.get_task("data_quality_summary")
|
||||
|
||||
assert s in e.get_direct_relatives("downstream")
|
||||
assert h in s.get_direct_relatives("downstream")
|
||||
assert d in h.get_direct_relatives("downstream")
|
||||
assert q in d.get_direct_relatives("downstream")
|
||||
@@ -0,0 +1,178 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from dataclasses import dataclass
|
||||
from typing import Any, List, Sequence
|
||||
|
||||
import pytest
|
||||
|
||||
import airflow.dags.helpers.greenplum as greenplum
|
||||
from tests.conftest import patch_postgres_hook
|
||||
|
||||
|
||||
@dataclass
|
||||
class FakeCursor:
|
||||
fetchone_value: Any = None
|
||||
fetchall_value: Sequence[Any] | None = None
|
||||
rowcount: int | None = None
|
||||
|
||||
def __post_init__(self) -> None:
|
||||
self.queries: List[Any] = []
|
||||
|
||||
def execute(self, query: str, params: Any | None = None) -> None:
|
||||
self.queries.append((query, params))
|
||||
|
||||
def fetchone(self) -> Any:
|
||||
return self.fetchone_value
|
||||
|
||||
def fetchall(self) -> Sequence[Any] | None:
|
||||
return self.fetchall_value
|
||||
|
||||
def __enter__(self) -> FakeCursor:
|
||||
return self
|
||||
|
||||
def __exit__(self, exc_type, exc, tb) -> None:
|
||||
return None
|
||||
|
||||
|
||||
class FakeConn:
|
||||
def __init__(self, cursors: Sequence[FakeCursor]) -> None:
|
||||
self._cursors = list(cursors)
|
||||
self._index = 0
|
||||
self.commits = 0
|
||||
|
||||
def cursor(self) -> FakeCursor:
|
||||
cursor = self._cursors[self._index]
|
||||
self._index += 1
|
||||
return cursor
|
||||
|
||||
def commit(self) -> None:
|
||||
self.commits += 1
|
||||
|
||||
|
||||
def test_get_gp_conn_uses_airflow_hook(monkeypatch) -> None:
|
||||
class FakeHook:
|
||||
def __init__(self, postgres_conn_id: str) -> None:
|
||||
self.postgres_conn_id = postgres_conn_id
|
||||
|
||||
def get_conn(self) -> str:
|
||||
return "hook_connection"
|
||||
|
||||
patch_postgres_hook(monkeypatch, FakeHook)
|
||||
monkeypatch.setattr(greenplum, "GP_CONN_ID", "demo_conn", raising=False)
|
||||
monkeypatch.setattr(greenplum, "GP_USE_AIRFLOW_CONN", True, raising=False)
|
||||
|
||||
conn = greenplum.get_gp_conn()
|
||||
|
||||
assert conn == "hook_connection"
|
||||
|
||||
|
||||
def test_get_gp_conn_fallback_to_psycopg(monkeypatch) -> None:
|
||||
class BrokenHook:
|
||||
def __init__(self, postgres_conn_id: str) -> None:
|
||||
self.postgres_conn_id = postgres_conn_id
|
||||
|
||||
def get_conn(self):
|
||||
raise RuntimeError("boom")
|
||||
|
||||
patch_postgres_hook(monkeypatch, BrokenHook)
|
||||
monkeypatch.setattr(greenplum, "GP_USE_AIRFLOW_CONN", True, raising=False)
|
||||
monkeypatch.setattr(greenplum, "GP_CONN_ID", "demo_conn", raising=False)
|
||||
monkeypatch.setenv("GP_DB", "demo_db")
|
||||
monkeypatch.setenv("GP_USER", "demo_user")
|
||||
monkeypatch.setenv("GP_PASSWORD", "secret")
|
||||
monkeypatch.setenv("GP_HOST", "greenplum-host")
|
||||
monkeypatch.setenv("GP_PORT", "5434")
|
||||
|
||||
captured_kwargs = {}
|
||||
|
||||
def fake_connect(**kwargs):
|
||||
captured_kwargs.update(kwargs)
|
||||
return "psycopg_connection"
|
||||
|
||||
monkeypatch.setattr(greenplum.psycopg2, "connect", fake_connect)
|
||||
|
||||
conn = greenplum.get_gp_conn()
|
||||
|
||||
assert conn == "psycopg_connection"
|
||||
assert captured_kwargs == {
|
||||
"dbname": "demo_db",
|
||||
"user": "demo_user",
|
||||
"password": "secret",
|
||||
"host": "greenplum-host",
|
||||
"port": 5434,
|
||||
}
|
||||
|
||||
|
||||
def test_get_gp_conn_without_airflow(monkeypatch) -> None:
|
||||
monkeypatch.setattr(greenplum, "GP_USE_AIRFLOW_CONN", False, raising=False)
|
||||
monkeypatch.setenv("GP_DB", "demo_db")
|
||||
monkeypatch.setenv("GP_USER", "demo_user")
|
||||
monkeypatch.setenv("GP_PASSWORD", "secret")
|
||||
monkeypatch.setenv("GP_HOST", "greenplum-host")
|
||||
monkeypatch.setenv("GP_PORT", "5435")
|
||||
|
||||
captured_kwargs = {}
|
||||
|
||||
def fake_connect(**kwargs):
|
||||
captured_kwargs.update(kwargs)
|
||||
return "direct_psycopg"
|
||||
|
||||
monkeypatch.setattr(greenplum.psycopg2, "connect", fake_connect)
|
||||
|
||||
conn = greenplum.get_gp_conn()
|
||||
|
||||
assert conn == "direct_psycopg"
|
||||
assert captured_kwargs["port"] == 5435
|
||||
|
||||
|
||||
def test_assert_orders_table_exists_ok() -> None:
|
||||
conn = FakeConn([FakeCursor(fetchone_value=(1,))])
|
||||
|
||||
greenplum.assert_orders_table_exists(conn)
|
||||
|
||||
|
||||
def test_assert_orders_table_exists_missing() -> None:
|
||||
conn = FakeConn([FakeCursor(fetchone_value=None)])
|
||||
|
||||
with pytest.raises(ValueError):
|
||||
greenplum.assert_orders_table_exists(conn)
|
||||
|
||||
|
||||
def test_assert_orders_schema_ok() -> None:
|
||||
expected = list(greenplum.EXPECTED_ORDERS_SCHEMA)
|
||||
conn = FakeConn([FakeCursor(fetchall_value=expected)])
|
||||
|
||||
greenplum.assert_orders_schema(conn)
|
||||
|
||||
|
||||
def test_assert_orders_schema_mismatch() -> None:
|
||||
conn = FakeConn([FakeCursor(fetchall_value=[("order_id", "bigint")])])
|
||||
|
||||
with pytest.raises(ValueError):
|
||||
greenplum.assert_orders_schema(conn)
|
||||
|
||||
|
||||
def test_assert_orders_have_rows_ok() -> None:
|
||||
conn = FakeConn([FakeCursor(fetchone_value=(5,))])
|
||||
|
||||
greenplum.assert_orders_have_rows(conn)
|
||||
|
||||
|
||||
def test_assert_orders_have_rows_empty() -> None:
|
||||
conn = FakeConn([FakeCursor(fetchone_value=(0,))])
|
||||
|
||||
with pytest.raises(ValueError):
|
||||
greenplum.assert_orders_have_rows(conn)
|
||||
|
||||
|
||||
def test_assert_orders_no_duplicates_ok() -> None:
|
||||
conn = FakeConn([FakeCursor(fetchone_value=(0,))])
|
||||
|
||||
greenplum.assert_orders_no_duplicates(conn)
|
||||
|
||||
|
||||
def test_assert_orders_no_duplicates_detected() -> None:
|
||||
conn = FakeConn([FakeCursor(fetchone_value=(3,))])
|
||||
|
||||
with pytest.raises(ValueError):
|
||||
greenplum.assert_orders_no_duplicates(conn)
|
||||
Reference in New Issue
Block a user