diff --git a/airflow-greenplum/README.md b/airflow-greenplum/README.md index fc75b6e..e8cc850 100644 --- a/airflow-greenplum/README.md +++ b/airflow-greenplum/README.md @@ -1,198 +1,211 @@ # DE Starter Kit — Greenplum + Kafka + Airflow -Обновлено: 2025-09-15 11:12 +Добро пожаловать в учебный стенд для изучения основ Data Engineering! Этот проект поможет вам освоить ключевые инструменты современных data pipeline: **Airflow** для оркестрации, **Kafka** для потоковой передачи данных и **Greenplum** как аналитическую базу данных. -Это учебный стенд для знакомства с Airflow и практикой загрузки данных в **Greenplum (single node в Docker)**. -На его базе можно собирать и запускать демонстрационные пайплайны, исследовать конфигурацию Airflow и тренировать разработку ETL/ELT‑оркестраций. +## 🎯 Что вы узнаете -Базовая поставка включает пример интеграции с Kafka: генерация данных, запись в топик и загрузка в Greenplum. -При желании вы можете заменить Kafka на любой другой источник или приёмник — стенд предназначен для экспериментов с разными интеграциями. +- Как настроить локальный стек данных с помощью Docker +- Как Airflow управляет workflow и координирует задачи +- Как данные перемещаются из Kafka в аналитическую базу +- Как проверять качество данных в автоматизированных pipeline +- Основы проектирования ETL/ELT процессов -Airflow по‑прежнему использует **Postgres** только как metadata DB (это стандартная и простая схема). +--- -## Требования -- Docker Desktop (Windows/Mac) или Docker Engine 24+ (Linux) с `docker compose v2`. -- Рекомендуемые среды: Linux, macOS, или Windows через WSL. Чистый Windows работает, но с ограничениями (пути/права, отсутствие `make` по умолчанию). -- Make (опционально). Если `make` нет — используйте команды `docker compose` из примеров ниже. -- Windows: в стандартный Git Bash `make` отсутствует. Для `make` используйте WSL/Chocolatey/Scoop/MSYS2; иначе работайте через `docker compose`. +## 🚀 Быстрый старт (для новичков) -## Что внутри -- **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. +### Шаг 1: Подготовка окружения -> Примечание по версиям и надёжности: используется образ `woblerr/greenplum:6.27.1` с поддержкой переменных окружения и fallback значениями. +**Требования:** +- Docker Desktop (Windows/Mac) или Docker Engine 24+ (Linux) +- Git для клонирования репозитория -## Быстрый старт -Опционально установим `make` (если нужен): +> 💡 **Совет:** Если у вас Windows, рекомендуем использовать WSL (Windows Subsystem for Linux) для лучшей совместимости. -- Linux (Debian/Ubuntu): - ```bash - sudo apt update && sudo apt install -y make - ``` -- macOS (Homebrew): - ```bash - brew install make - ``` -- Windows варианты: - - WSL (Ubuntu): `sudo apt install -y make` +### Шаг 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 с названием **kafka_to_greenplum** +3. Нажмите на переключатель слева от названия DAG, чтобы включить его +4. Нажмите кнопку **Trigger** (значок воспроизведения ▶️) + +🎉 **Поздравляем!** Вы только что запустили свой первый data pipeline: +- Система сгенерировала 1000 тестовых заказов +- Данные отправились в Kafka (систему потоковой передачи сообщений) +- Airflow прочитал данные из Kafka и загрузил их в Greenplum + +### Шаг 4: Проверка результатов + +**Через веб-интерфейс:** +- Kafka UI: **http://localhost:8082** — посмотрите топик `orders` и сообщения +- Airflow UI: **http://localhost:8080** — отслеживайте выполнение задач + +**Через командную строку:** +```bash +# Подключитесь к Greenplum и проверьте данные +docker compose exec greenplum bash -c "su - gpadmin -c 'psql -p 5432 -d gpadmin'" + +# Внутри psql выполните: +\dt # Показать таблицы +SELECT count(*) FROM public.orders; # Посчитать записи +``` + +--- + +## 🛠️ Подробная настройка (для уверенных пользователей) + +### Установка 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` - - MSYS2: `pacman -S make` -Настроим переменные окружения +С `make` команды становятся короче: ```bash -cp .env.example .env -# При необходимости отредактируйте .env для ваших настроек +make up && make airflow-init # Запуск стека +make logs # Просмотр логов +make gp-psql # Подключение к Greenplum ``` -Запустим приложение (вариант с Make) -```bash -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`. +### Настройка подключения к Greenplum в Airflow -Альтернатива без Make (на всех ОС): +По умолчанию 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** — аналитическая база данных для хранения и анализа данных +- **Kafka** — система потоковой передачи данных (без Zookeeper) +- **Airflow** — оркестратор workflow и задач +- **Postgres** — база метаданных для Airflow + +### Готовые DAG (workflow) +- **kafka_to_greenplum** — базовый pipeline: генерация → Kafka → Greenplum +- **greenplum_data_quality** — проверки качества данных (наличие таблицы, схема, дубликаты) + +### Полезные команды ```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 down # Остановить и удалить данные +make airflow-init # Инициализировать Airflow +make ddl-gp # Применить DDL к Greenplum +make gp-psql # Подключиться к Greenplum через psql + +# Проверка данных +make logs # Следить за логами Airflow ``` -### Создаём 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 (если он уже был активирован). +--- -CLI-альтернатива (выполняется внутри контейнера Airflow): -```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" -``` +## ⚙️ Настройка через переменные окружения -### Интерфейсы и порты -- 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` +Все настройки находятся в файле `.env`. Основные параметры: -### Параметры чтения/загрузки -- `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 - -# Или через 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. - -**Подключение К Greenplum** -- DBeaver (PostgreSQL/Greenplum драйвер): - - Запусти DBeaver → Database → New Connection → выбери `Greenplum`. - - Host: `localhost` - - Port: `5432` (или значение из `GP_PORT`) - - Database: `gpadmin` (или `GP_DB`) - - Username: `gpadmin` (или `GP_USER`) - - Password: `gpadmin` (или `GP_PASSWORD`) - - SSL: Disabled (для локального стенда) - - Нажми Test Connection → Finish. -- psql внутри контейнера (rootless доступ к встроенному psql): - - Команда Make: `make 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'"` - - Примеры команд внутри psql: -`\dt`, `SELECT version()`, `SELECT count(*) FROM public.orders;` - -## Файлы -- `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 берутся из тех же переменных окружения: +### Kafka +- `KAFKA_TOPIC` — имя топика (по умолчанию: orders) +- `KAFKA_BATCH_SIZE` — размер пакета при загрузке (по умолчанию: 500) -- `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`. +## 🔍 Продвинутые темы -## Фиксация образов (pinning) и альтернативы -- Зафиксируй 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 `kafka_to_greenplum`:** +1. `create_table` — создает таблицу `public.orders` в Greenplum +2. `produce_messages` — генерирует и отправляет сообщения в Kafka +3. `consume_and_load` — читает из Kafka и загружает в Greenplum батчами -## Поток данных (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` для автоматической проверки: +- Наличие таблицы в базе +- Соответствие схемы ожидаемой структуре +- Объем загруженных данных +- Отсутствие дубликатов записей -## Проверка данных -- Подними стенд (`make up && make airflow-init`) и запусти DAG `kafka_to_greenplum`, чтобы заполнить таблицу `orders`. -- Активируй и запусти DAG `greenplum_data_quality` — он последовательно проверит наличие таблицы, схему, объём данных и отсутствие дублей. Все проверки выполняются внутри Airflow и используют те же настройки подключений. +### Ограничения учебного стенда + +- **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`) | +| Нет топика `orders` в Kafka | Создайте топик через Kafka UI или дождитесь авто-создания при первом запуске DAG | +| Команда `make` не найдена | Используйте полные команды `docker compose` или установите make | + +--- + +## 📁 Структура проекта + +``` +├── docker-compose.yml # Описание всех сервисов +├── .env.example # Шаблон настроек +├── Makefile # Удобные команды для работы +├── airflow/ +│ └── dags/ # Файлы workflow (DAG) +│ ├── kafka_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** — добавьте зависимости между задачами, настройте расписания + +Удачи в изучении Data Engineering! 🚀