10 Commits
19 changed files with 3112 additions and 382 deletions
+3 -7
View File
@@ -16,10 +16,6 @@ GP_PORT=5432
GP_CONN_ID=greenplum_conn
GP_USE_AIRFLOW_CONN=true
# Kafka Configuration
KAFKA_BOOTSTRAP=kafka:29092
KAFKA_TOPIC=orders
KAFKA_BATCH_SIZE=500
KAFKA_POLL_TIMEOUT=10
KAFKA_MAX_EMPTY_POLLS=3
KAFKA_UI_CLUSTER_NAME=KafkaCluster
# CSV Pipeline
CSV_DIR=/opt/airflow/data
CSV_ROWS=1000
+5 -1
View File
@@ -3,4 +3,8 @@
# Do not commit secrets
.env
.env.*
/airflow/dags/__pycache__
__pycache__/
*/__pycache__/
*.pyc
.venv/
data/
+1
View File
@@ -0,0 +1 @@
3.11
+50 -37
View File
@@ -1,44 +1,57 @@
# Repository Guidelines
# Repository Guidelines (для агентa и контрибьюторов)
## Project Structure & Module Organization
- `airflow/dags/` — Airflow DAGs (e.g., `airflow/dags/kafka_to_greenplum.py`).
- `airflow/requirements.txt` — Python deps installed inside Airflow containers.
- `sql/` — database DDL and helpers (e.g., `sql/ddl_gp.sql`).
- `docker-compose.yml` — Greenplum, Kafka, Airflow, Postgres (metadata DB).
- `Makefile` — local DX commands; see targets below.
- `.env(.example)` — runtime configuration; never commit real secrets.
Эта репа — учебный стенд для студентов (менти), которые только начинают с Airflow/Greenplum и Python. Пожалуйста, держите решения простыми, стабильными и хорошо объяснёнными.
## Build, Test, and Development Commands
- `make up` — start the full stack.
- `make airflow-init` — migrate metadata DB and create admin user.
- `make logs` — follow webserver and scheduler logs.
- `make ddl-gp` — apply DDL to Greenplum.
- `make gp-psql` — open `psql` in the GP container.
- `make down` — stop stack and remove volumes.
Example: `make up && make airflow-init` then open `http://localhost:8080`.
## Структура проекта
- `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)` — настройки окружения (реальные секреты не коммитим).
## Coding Style & Naming Conventions
- Python: PEP 8, 4-space indents, `snake_case` for functions/vars, DAG IDs lower_snake_case.
- Imports: stdlib → third-party → local; prefer one module per line.
- SQL: uppercase keywords, `snake_case` identifiers, end statements with `;`.
- Filenames: DAGs as `<source>_to_<target>.py` (e.g., `kafka_to_greenplum.py`).
- Formatting: if available, use `black` (88 cols) and `isort`; otherwise keep existing style.
- Language: комментарии, docstrings и документацию (README, описания PR/Issues) пишем на русском; имена идентификаторов и код — на английском.
## Команды (основные)
- `make up` — поднять весь стек.
- `make airflow-init` — инициализировать мета-БД Airflow и создать пользователя.
- `make logs` — логи webserver и scheduler.
- `make ddl-gp` — применить DDL к Greenplum.
- `make gp-psql` — открыть `psql` в контейнере Greenplum от `gpadmin`.
- `make down` — остановить и удалить тома (данные будут потеряны).
## Testing Guidelines
- No test suite yet. If adding tests, use `pytest` under `tests/` with `test_*.py`.
- Prefer unit tests for Python callables used by tasks; mock env vars and external systems.
- Run locally with `pytest -q`.
Пример: `make up && make airflow-init`, затем открыть `http://localhost:8080`.
## Commit & Pull Request Guidelines
- Use Conventional Commits: `feat:`, `fix:`, `docs:`, `chore:`, `refactor:` etc. Example: `feat(dags): load orders to Greenplum`.
- Keep PRs focused; include a description, run steps, and relevant screenshots (e.g., DAG graph or task logs).
- Link issues; update `README.md` and DDL when behavior or schema changes.
## Локальное 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`.
## Security & Configuration Tips
- Configure via `.env`; do not hardcode credentials. Common vars: `GP_USER`, `GP_PASSWORD`, `GP_DB`, `GP_PORT`, `PG_*`, `AIRFLOW_*`.
- Be cautious with `make down` (removes volumes). Pin images/deps; prefer digests for critical images.
## Стиль кода
- Python: PEP 8, 4 пробела, `snake_case`; `dag_id``lower_snake_case`.
- Импорты: stdlib → thirdparty → local, по одному модулю в строке.
- SQL: ключевые слова UPPERCASE, идентификаторы `snake_case`, завершаем `;`.
- Форматирование: `black` (88 cols) и `isort`. Если не уверены — запустите `make fmt`.
- Язык: комментарии, docstring и документацию — на русском; имена идентификаторов — на английском.
## Agent-Specific Notes
- Keep changes minimal and localized; do not rename Make targets without updating docs.
- Validate by running `make up`, `make airflow-init`, and inspecting the DAG in Airflow.
## Тестирование
- Тесты лежат в `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`.
+29
View File
@@ -1,4 +1,8 @@
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
@@ -17,3 +21,28 @@ gp-psql:
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)"
+237 -127
View File
@@ -1,165 +1,275 @@
# DE Starter Kit — Greenplum + Kafka + Airflow
# DE Starter Kit — Airflow + Greenplum + CSV
Обновлено: 2025-09-15 11:12
Добро пожаловать в учебный стенд для изучения основ Data Engineering! Этот проект поможет вам освоить ключевые инструменты современных data pipeline: **Airflow** для оркестрации, **pandas/CSV** для подготовки данных и **Greenplum** как аналитическую базу данных.
Этот вариант повторяет логику Postgres-стенда, но в роли DWH — **Greenplum (single node в Docker)**.
Airflow по‑прежнему использует **Postgres** только как metadata DB (это стандартная и простая схема).
## 🎯 Что вы узнаете
## Требования
- Docker Desktop (Windows/Mac) или Docker Engine 24+ (Linux) с `docker compose v2`.
- Make (опционально). Если `make` нет — используйте приведённые ниже команды `docker compose` напрямую.
- Windows: запускайте команды в Git Bash или WSL; PowerShell тоже подойдёт, но для `make` удобнее Git Bash/WSL.
- Как настроить локальный стек данных с помощью Docker
- Как Airflow управляет workflow и координирует задачи
- Как генерировать датасеты через pandas и сохранять их в CSV
- Как загружать данные в Greenplum пакетами и избегать дублей
- Как проверять качество данных в автоматизированных pipeline
- Основы проектирования ETL/ELT процессов
## Что внутри
- **Greenplum** (single node) — `woblerr/greenplum:6.27.1` (локальный стенд для разработки).
- **Kafka (KRaft, без Zookeeper)** — генерация событий.
- **Airflow (2.9)** — оркестрация пайплайна, metadata в Postgres.
- **DAG** `kafka_to_greenplum.py` — генерирует данные → пишет в Kafka → читает и грузит в Greenplum.
- **DAG** `greenplum_data_quality.py` — выполняет проверки качества данных (наличие таблицы, схема, заполненность, дубли).
- **Kafka UI** — автоматически подключается к стендовому брокеру (bootstrap берётся из `KAFKA_BOOTSTRAP` или значения по умолчанию).
- **DDL** `sql/ddl_gp.sql` — создаёт колонночную таблицу `orders` без PRIMARY KEY (AO-таблицы GP6 не поддерживают его), распределённую по `order_id`; контроль дублей реализован в DAG.
## 👩‍🎓 Для студентов (10‑минутный чек‑лист)
> Примечание по версиям и надёжности: используется образ `woblerr/greenplum:6.27.1` с поддержкой переменных окружения и fallback значениями. Для продакшен‑подобных тестов зафиксируй digest (SHA256) конкретного тега на Docker Hub.
- Установите 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`.
## Быстрый старт
Установим make (Linux/Mac)
```bash
sudo apt install -y make
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
# При необходимости отредактируйте .env для ваших настроек
# Запустите стек (это может занять 2-3 минуты при первом запуске)
docker compose up -d
# Инициализируйте Airflow
docker compose run --rm airflow-init
```
Запустим приложение (вариант с Make)
### Шаг 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
make logs # ждем "Listening at: http://0.0.0.0:8080"
```
Открой Airflow: http://localhost:8080 (логин/пароль см. `.env`, по умолчанию admin/admin).
Включи DAG **kafka_to_greenplum** и нажми **Trigger** — он создаст таблицу и загрузит ~1000 записей в `gpadmin.public.orders`.
Альтернатива без Make (на всех ОС):
Это помогает, когда 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
docker compose -f docker-compose.yml up -d
docker compose -f docker-compose.yml run --rm airflow-init
docker compose -f docker-compose.yml logs -f airflow-webserver airflow-scheduler
make up && make airflow-init # Запуск стека
make logs # Просмотр логов
make gp-psql # Подключение к Greenplum
```
### Создаём Airflow Connection для Greenplum
1. Открой Airflow UI → **Admin → Connections****Add a new record**.
2. Заполни поля:
- `Conn Id`: `greenplum_conn` (или своё значение, тогда пропиши его в переменной `GP_CONN_ID`).
- `Conn Type`: `Postgres`.
- `Host`: `greenplum`.
- `Schema`: значение `GP_DB` (по умолчанию `gpadmin`).
- `Login`: `GP_USER` (по умолчанию `gpadmin`).
- `Password`: `GP_PASSWORD`.
- `Port`: `5432`.
3. Сохрани соединение и перезапусти DAG (если он уже был активирован).
### Настройка подключения к Greenplum в Airflow
CLI-альтернатива (выполняется внутри контейнера 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
docker compose -f docker-compose.yml exec airflow-webserver bash -lc "
airflow connections add 'greenplum_conn' \
--conn-type postgres \
--conn-host greenplum \
--conn-login ${GP_USER:-gpadmin} \
--conn-password ${GP_PASSWORD:-gpadmin} \
--conn-schema ${GP_DB:-gpadmin} \
--conn-port 5432"
# Основные команды
make up # Запустить весь стенд
make down # Остановить и удалить данные
make airflow-init # Инициализировать Airflow
make ddl-gp # Применить DDL к Greenplum
make gp-psql # Подключиться к Greenplum через psql
# Проверка данных
make logs # Следить за логами Airflow
```
### Интерфейсы и порты
- Airflow UI: http://localhost:8080 (admin/admin по умолчанию)
- Kafka UI: http://localhost:8082 (просмотр топиков/сообщений; подключение к кластеру создаётся автоматически)
- Greenplum: `localhost:${GP_PORT:-5432}` (внешний порт проброшен из контейнера)
- Postgres (Airflow metadata): `localhost:5433`
- Kafka (для клиентов на хосте): `localhost:9092`
- Kafka (из контейнеров Docker): `kafka:29092`
---
### Параметры чтения/загрузки
- `KAFKA_BATCH_SIZE` — размер батча при вставке в Greenplum (по умолчанию 500).
- `KAFKA_POLL_TIMEOUT` — таймаут ожидания сообщения в секундах (по умолчанию 10).
- `KAFKA_MAX_EMPTY_POLLS` — сколько подряд пустых `poll` допускается перед выходом из цикла (по умолчанию 3).
## ⚙️ Настройка через переменные окружения
### Проверка загрузки
```bash
# Подключение к Greenplum через внешний psql клиент
psql -h localhost -p 5432 -U gpadmin -d gpadmin
Все настройки находятся в файле `.env`. Основные параметры:
# Или через Docker (используя make команду)
make gp-psql
# Или напрямую через Docker
docker compose -f docker-compose.yml exec greenplum bash -c "su - gpadmin -c '/usr/local/greenplum-db/bin/psql -p 5432 -d gpadmin'"
# Внутри psql:
\dt
SELECT count(*) FROM public.orders;
```
### Проверка Kafka
- Открой Kafka UI: http://localhost:8082 — проверь, что существует топик `orders` и в нём появляются сообщения после запуска DAG.
- Если автосоздание топиков в брокере отключено, создай топик вручную через UI перед запуском DAG.
## Файлы
- `docker-compose.yml` — сервисы Greenplum + Kafka + Airflow + Postgres (metadata).
- `.env.example` — шаблон переменных окружения.
- `.env` — переменные окружения (создается из .env.example).
- `airflow/dags/kafka_to_greenplum.py` — сам DAG.
- `sql/ddl_gp.sql` — DDL таблицы в Greenplum.
- `Makefile` — обёртки команд (`up`, `down`, `airflow-init`, `gp-psql`, `ddl-gp`).
## Конфигурация через переменные окружения
Все настройки Greenplum передаются через переменные окружения в файле `.env` с fallback значениями:
- `GP_USER` — пользователь Greenplum (по умолчанию: gpadmin)
- `GP_PASSWORD` — пароль пользователя (по умолчанию: gpadmin)
### Greenplum
- `GP_USER` — пользователь (по умолчанию: gpadmin)
- `GP_PASSWORD` — пароль (по умолчанию: gpadmin)
- `GP_DB` — база данных (по умолчанию: gpadmin)
- `GP_PORT` — порт для подключения (по умолчанию: 5432)
- `GP_CONN_ID` — ID Airflow Connection (по умолчанию: `greenplum_conn`)
- `GP_USE_AIRFLOW_CONN` — использовать ли Airflow Connection (`true`/`false`). Если `false`, DAG подключается к БД напрямую по ENV.
- `GP_PORT` — порт (по умолчанию: 5432)
Настройки Kafka берутся из тех же переменных окружения:
### CSV pipeline
- `CSV_DIR` — путь к каталогу с CSV внутри контейнеров Airflow (по умолчанию: `/opt/airflow/data`)
- `CSV_ROWS` — количество строк, генерируемых DAG (по умолчанию: 1000)
- `KAFKA_BOOTSTRAP` — bootstrap-адрес брокера для Airflow и Kafka UI (по умолчанию `kafka:29092`).
- `KAFKA_TOPIC` — имя демо-топика (по умолчанию `orders`).
- `KAFKA_BATCH_SIZE`, `KAFKA_POLL_TIMEOUT`, `KAFKA_MAX_EMPTY_POLLS` — параметры чтения/загрузки (см. раздел выше).
- `KAFKA_UI_CLUSTER_NAME` — отображаемое имя кластера в Kafka UI (по умолчанию `KafkaCluster`).
### Airflow
- `GP_CONN_ID` — ID подключения (по умолчанию: greenplum_conn)
Образ `woblerr/greenplum:6.27.1` использует переменные:
- `GREENPLUM_USER` (маппится на `GP_USER`)
- `GREENPLUM_PASSWORD` (маппится на `GP_PASSWORD`)
- `GREENPLUM_DATABASE_NAME` (маппится на `GP_DB`)
---
Внутри контейнеров Airflow хост для подключения к БД — `greenplum` (см. `GP_HOST`), а с вашей машины — `localhost:${GP_PORT}`. Созданный в Airflow Connection реиспользует те же значения, что и `.env`.
## 🔍 Продвинутые темы
## Пинning и альтернативы
- Зафиксируй digest образа `woblerr/greenplum:6.27.1` (Docker Hub → Tag → «Copy digest») и замени тег на `@sha256:...` в `docker-compose.yml`.
- Альтернативы: можно использовать другие образы Greenplum или собрать собственный образ для GPDB 6/7.
- Особенность GPDB 6: в нём **нет** `INSERT ... ON CONFLICT`. В DAG используется безопасная для GP6 конструкция `INSERT ... WHERE NOT EXISTS` внутри транзакции.
### Архитектура pipeline
## Ограничения и заметки
- Этот стенд — учебный. Для высокой надёжности и производительности Greenplum обычно разворачивают кластерами на нескольких узлах, на裸‑железе/VM с отдельными дисками под сегменты.
- Для GP7 (основан на новее PostgreSQL) можно упростить загрузку, включая `ON CONFLICT`. В учебных целях мы остались на широко доступном GP6 образе.
**Поток данных в 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`
## Поток данных (DAG)
Последовательность задач в `kafka_to_greenplum`:
- `create_table` — создаёт таблицу `public.orders` в Greenplum (колоночная, AO/CO, распределение по `order_id`).
- `produce_messages` — генерирует ~1000 сообщений и пишет их в Kafka-топик `orders`.
- `consume_and_load` — читает сообщения из Kafka и вставляет в `public.orders` батчами (по `KAFKA_BATCH_SIZE`) с защитой от дублей для GP6 через anti-join. По умолчанию использует Airflow Connection `greenplum_conn`, но при `GP_USE_AIRFLOW_CONN=false` подключается по ENV (`GP_HOST`, `GP_PORT`, `GP_DB`, `GP_USER`, `GP_PASSWORD`).
> 💡 **Безопасность повторного запуска:** Pipeline защищен от дубликатов, поэтому его можно запускать многократно.
Повторный запуск DAG безопасен: при вставке используется проверка на существование `order_id`.
### Проверка качества данных
## Типичные проблемы и решения
- Airflow UI не открывается: проверь `make logs` и дождись строки `Listening at: http://0.0.0.0:8080`.
- Ошибка подключения к Greenplum: дождись, пока контейнер `greenplum` станет `healthy`; проверь, что порт `GP_PORT` не занят локальными сервисами.
- Нет топика `orders`: создай его через Kafka UI (или перезапусти DAG после включения авто‑создания топиков).
- `make` отсутствует на Windows: используй команды `docker compose` из раздела «Альтернатива без Make» или установи Git Bash/WSL.
Запустите 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! 🚀
## Проверка данных
- Подними стенд (`make up && make airflow-init`) и запусти DAG `kafka_to_greenplum`, чтобы заполнить таблицу `orders`.
- Активируй и запусти DAG `greenplum_data_quality` — он последовательно проверит наличие таблицы, схему, объём данных и отсутствие дублей. Все проверки выполняются внутри Airflow и используют те же настройки подключений.
+71
View File
@@ -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
@@ -1,5 +1,6 @@
from __future__ import annotations
import logging
from datetime import datetime, timedelta
from airflow import DAG
@@ -15,9 +16,35 @@ from helpers.greenplum import (
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)}
@@ -29,26 +56,41 @@ with DAG(
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
# Определяем последовательность выполнения задач
check_exists >> check_schema >> check_has_rows >> check_no_duplicates >> dq_summary
@@ -1,5 +1,6 @@
from __future__ import annotations
import logging
import os
from typing import List, Sequence, Tuple
@@ -10,6 +11,7 @@ import psycopg2
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"),
@@ -19,28 +21,52 @@ EXPECTED_ORDERS_SCHEMA: List[Tuple[str, str]] = [
def get_gp_conn():
"""Возвращает psycopg2 connection к Greenplum (через Airflow Connection или напрямую по ENV)."""
"""
Возвращает 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)
return hook.get_conn()
except Exception:
conn = hook.get_conn()
logging.info("✅ Подключение через Airflow Connection успешно")
return conn
except Exception as e:
logging.warning("⚠️ Не удалось подключиться через Airflow Connection: %s", e)
logging.info("🔄 Переключаемся на прямое подключение по ENV переменным")
# Фоллбек на прямое подключение по переменным окружения.
pass
return psycopg2.connect(
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")),
)
# Прямое подключение по переменным окружения
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."""
"""
Проверяет наличие таблицы orders в схеме public.
Args:
conn: Подключение к Greenplum
Raises:
ValueError: Если таблица не найдена
"""
logging.info("🔍 Проверяем существование таблицы public.orders...")
with conn.cursor() as cur:
cur.execute(
"""
@@ -50,10 +76,20 @@ def assert_orders_table_exists(conn) -> None:
"""
)
if cur.fetchone() is None:
raise ValueError("Таблица public.orders не найдена; запусти DAG kafka_to_greenplum.")
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(
"""
@@ -67,25 +103,69 @@ def fetch_orders_schema(conn) -> Sequence[Tuple[str, str]]:
def assert_orders_schema(conn) -> None:
"""Проверяет, что схема таблицы orders соответствует ожидаемой."""
"""
Проверяет, что схема таблицы 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}.")
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 не пустая."""
if fetch_orders_count(conn) <= 0:
raise ValueError("Таблица public.orders пустая — запусти DAG kafka_to_greenplum перед проверкой.")
"""
Проверяет, что таблица 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(
"""
@@ -101,7 +181,19 @@ def fetch_orders_duplicates(conn) -> int:
def assert_orders_no_duplicates(conn) -> None:
"""Проверяет, что в таблице нет дублей по order_id."""
"""
Проверяет, что в таблице нет дублей по 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} шт.) — проверь загрузку данных.")
raise ValueError(f"Обнаружены дубли по order_id ({duplicates} шт.) — проверь загрузку данных.")
logging.info("✅ Дубликаты не обнаружены")
@@ -1,146 +0,0 @@
from __future__ import annotations
import json
import os
import random
from datetime import datetime, timedelta
from typing import List, Tuple, Optional
from airflow import DAG
from airflow.operators.python import PythonOperator
from confluent_kafka import Consumer, KafkaException, Producer
from psycopg2.extras import execute_values
from helpers.greenplum import get_gp_conn
KAFKA_BOOTSTRAP = os.getenv("KAFKA_BOOTSTRAP", "kafka:29092")
TOPIC = os.getenv("KAFKA_TOPIC", "orders")
BATCH_SIZE = int(os.getenv("KAFKA_BATCH_SIZE", "500"))
POLL_TIMEOUT_S = int(os.getenv("KAFKA_POLL_TIMEOUT", "10"))
MAX_EMPTY_POLLS = int(os.getenv("KAFKA_MAX_EMPTY_POLLS", "3"))
def _create_table():
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=column, compresstype=zlib)
DISTRIBUTED BY (order_id);
"""
with get_gp_conn() as conn, conn.cursor() as cur:
cur.execute(ddl)
conn.commit()
def _produce(n=1000):
producer = Producer({"bootstrap.servers": KAFKA_BOOTSTRAP})
for idx in range(n):
payload = {
"order_id": idx + 1,
"order_ts": datetime.utcnow().isoformat(),
"customer_id": random.randint(1, 100),
"amount": round(random.uniform(10, 500), 2),
}
producer.produce(TOPIC, json.dumps(payload).encode("utf-8"))
producer.flush()
def _flush_batch(cur, rows: List[Tuple]):
"""Insert deduplicated batch of rows into public.orders for GP6 (no PK support)."""
if not rows:
return
# Дедупликация внутри батча по первичному ключу (order_id)
by_id = {int(r[0]): r for r in rows}
unique_rows = list(by_id.values())
# Вставка через VALUES + anti-join для GP6 (без ON CONFLICT)
execute_values(
cur,
"""
INSERT INTO public.orders(order_id, order_ts, customer_id, amount)
SELECT v.order_id, v.order_ts, v.customer_id, v.amount
FROM (VALUES %s) AS v(order_id, order_ts, customer_id, amount)
LEFT JOIN public.orders o ON o.order_id = v.order_id
WHERE o.order_id IS NULL
""",
unique_rows,
template="(%s,%s,%s,%s)",
)
def _consume_and_load(max_messages=1000, timeout_s: Optional[int] = None):
if timeout_s is None:
timeout_s = POLL_TIMEOUT_S
consumer = Consumer(
{
"bootstrap.servers": KAFKA_BOOTSTRAP,
"group.id": "airflow-loader-gp",
"auto.offset.reset": "earliest",
"enable.auto.commit": False,
}
)
consumer.subscribe([TOPIC])
with get_gp_conn() as conn, conn.cursor() as cur:
batch: List[Tuple] = []
consumed = 0
empty_polls = 0
while consumed < max_messages and empty_polls < MAX_EMPTY_POLLS:
msg = consumer.poll(timeout_s)
if msg is None:
empty_polls += 1
continue
empty_polls = 0
if msg.error():
raise KafkaException(msg.error())
data = json.loads(msg.value().decode("utf-8"))
batch.append(
(
int(data["order_id"]),
datetime.fromisoformat(data["order_ts"]),
int(data["customer_id"]),
float(data["amount"]),
)
)
consumed += 1
if len(batch) >= BATCH_SIZE:
_flush_batch(cur, batch)
conn.commit()
batch.clear()
# Финальный сброс, если вышли по лимиту сообщений
if batch:
_flush_batch(cur, batch)
conn.commit()
# Фиксируем оффсеты после успешной загрузки
consumer.commit()
consumer.close()
default_args = {"owner": "airflow", "retries": 1, "retry_delay": timedelta(seconds=30)}
with DAG(
dag_id="kafka_to_greenplum",
start_date=datetime(2024, 1, 1),
schedule=None,
catchup=False,
default_args=default_args,
tags=["demo", "kafka", "greenplum"],
) as dag:
create_table = PythonOperator(task_id="create_table", python_callable=_create_table)
produce = PythonOperator(
task_id="produce_messages", python_callable=_produce, op_kwargs={"n": 1000}
)
consume_and_load = PythonOperator(
task_id="consume_and_load",
python_callable=_consume_and_load,
op_kwargs={"max_messages": 1000},
)
create_table >> produce >> consume_and_load
+1 -1
View File
@@ -1,2 +1,2 @@
confluent-kafka==2.3.0
psycopg2-binary==2.9.9
pandas==2.1.4
+3 -41
View File
@@ -40,42 +40,6 @@ services:
timeout: 5s
retries: 30
kafka:
image: apache/kafka:3.8.0
environment:
KAFKA_BROKER_ID: 1
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT,CONTROLLER:PLAINTEXT
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:29092,PLAINTEXT_HOST://localhost:9092
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_NODE_ID: 1
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka:29093
KAFKA_LISTENERS: PLAINTEXT://kafka:29092,CONTROLLER://kafka:29093,PLAINTEXT_HOST://0.0.0.0:9092
KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
KAFKA_LOG_DIRS: /tmp/kraft-combined-logs
# 22-символьный base64 идентификатор для KRaft-кластера
CLUSTER_ID: aGVsbG93b3JsZGtpdGNoZW4
ports:
- "9092:9092"
kafka-ui:
image: provectuslabs/kafka-ui:v0.7.2
env_file: .env
ports:
- 8082:8080
environment:
DYNAMIC_CONFIG_ENABLED: "true"
KAFKA_CLUSTERS_0_NAME: ${KAFKA_UI_CLUSTER_NAME:-KafkaCluster}
KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS: ${KAFKA_BOOTSTRAP:-kafka:29092}
volumes:
- kafka-ui-data:/app/data
depends_on:
- kafka
airflow-webserver:
image: apache/airflow:2.9.2
container_name: gp_airflow_web
@@ -91,11 +55,10 @@ services:
volumes:
- ./airflow/dags:/opt/airflow/dags
- ./airflow/requirements.txt:/opt/airflow/requirements.txt
- ./data:/opt/airflow/data
depends_on:
pgmeta:
condition: service_healthy
kafka:
condition: service_started
greenplum:
condition: service_healthy
@@ -112,11 +75,10 @@ services:
volumes:
- ./airflow/dags:/opt/airflow/dags
- ./airflow/requirements.txt:/opt/airflow/requirements.txt
- ./data:/opt/airflow/data
depends_on:
pgmeta:
condition: service_healthy
kafka:
condition: service_started
greenplum:
condition: service_healthy
@@ -130,6 +92,7 @@ services:
volumes:
- ./airflow/dags:/opt/airflow/dags
- ./airflow/requirements.txt:/opt/airflow/requirements.txt
- ./data:/opt/airflow/data
command: >
bash -lc "
set -e;
@@ -147,5 +110,4 @@ services:
volumes:
pgmeta:
kafka-ui-data:
greenplum_data:
+16
View File
@@ -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",
]
+69
View File
@@ -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)
+2063
View File
File diff suppressed because it is too large Load Diff