Compare commits
10
Commits
main
...
ba3e681645
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ba3e681645 | ||
|
|
919dee8e8b | ||
|
|
dfc1f2ca84 | ||
|
|
bf6b424a5e | ||
|
|
9341983f26 | ||
|
|
fe8c29cbe1 | ||
|
|
9ba57c1696 | ||
|
|
cea33ea8d8 | ||
|
|
806e1dc3e4 | ||
|
|
e7e40d18f7 |
@@ -0,0 +1,25 @@
|
|||||||
|
# Airflow Configuration
|
||||||
|
AIRFLOW_USER=admin
|
||||||
|
AIRFLOW_PASSWORD=admin
|
||||||
|
|
||||||
|
# PostgreSQL (Airflow metadata)
|
||||||
|
PG_USER=airflow
|
||||||
|
PG_PASSWORD=airflow
|
||||||
|
PG_DB=airflow
|
||||||
|
|
||||||
|
# Greenplum Configuration
|
||||||
|
GP_USER=gpadmin
|
||||||
|
GP_PASSWORD=gpadmin
|
||||||
|
GP_DB=gpadmin
|
||||||
|
GP_HOST=greenplum
|
||||||
|
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
|
||||||
@@ -0,0 +1,6 @@
|
|||||||
|
.vscode/settings.json
|
||||||
|
|
||||||
|
# Do not commit secrets
|
||||||
|
.env
|
||||||
|
.env.*
|
||||||
|
/airflow/dags/__pycache__
|
||||||
@@ -0,0 +1,44 @@
|
|||||||
|
# Repository Guidelines
|
||||||
|
|
||||||
|
## 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.
|
||||||
|
|
||||||
|
## 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`.
|
||||||
|
|
||||||
|
## 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) пишем на русском; имена идентификаторов и код — на английском.
|
||||||
|
|
||||||
|
## 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`.
|
||||||
|
|
||||||
|
## 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.
|
||||||
|
|
||||||
|
## 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.
|
||||||
|
|
||||||
|
## 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.
|
||||||
@@ -0,0 +1,19 @@
|
|||||||
|
SHELL := /bin/bash
|
||||||
|
|
||||||
|
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'"
|
||||||
@@ -0,0 +1,165 @@
|
|||||||
|
# DE Starter Kit — Greenplum + Kafka + Airflow
|
||||||
|
|
||||||
|
Обновлено: 2025-09-15 11:12
|
||||||
|
|
||||||
|
Этот вариант повторяет логику 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.
|
||||||
|
|
||||||
|
## Что внутри
|
||||||
|
- **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.
|
||||||
|
|
||||||
|
> Примечание по версиям и надёжности: используется образ `woblerr/greenplum:6.27.1` с поддержкой переменных окружения и fallback значениями. Для продакшен‑подобных тестов зафиксируй digest (SHA256) конкретного тега на Docker Hub.
|
||||||
|
|
||||||
|
## Быстрый старт
|
||||||
|
Установим make (Linux/Mac)
|
||||||
|
```bash
|
||||||
|
sudo apt install -y make
|
||||||
|
```
|
||||||
|
|
||||||
|
Настроим переменные окружения
|
||||||
|
```bash
|
||||||
|
cp .env.example .env
|
||||||
|
# При необходимости отредактируйте .env для ваших настроек
|
||||||
|
```
|
||||||
|
|
||||||
|
Запустим приложение (вариант с 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`.
|
||||||
|
|
||||||
|
Альтернатива без 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
|
||||||
|
```
|
||||||
|
|
||||||
|
### Создаём 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`
|
||||||
|
|
||||||
|
### Параметры чтения/загрузки
|
||||||
|
- `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.
|
||||||
|
|
||||||
|
## Файлы
|
||||||
|
- `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)
|
||||||
|
- `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.
|
||||||
|
|
||||||
|
Настройки Kafka берутся из тех же переменных окружения:
|
||||||
|
|
||||||
|
- `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`).
|
||||||
|
|
||||||
|
Образ `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` внутри транзакции.
|
||||||
|
|
||||||
|
## Ограничения и заметки
|
||||||
|
- Этот стенд — учебный. Для высокой надёжности и производительности Greenplum обычно разворачивают кластерами на нескольких узлах, на裸‑железе/VM с отдельными дисками под сегменты.
|
||||||
|
- Для GP7 (основан на новее PostgreSQL) можно упростить загрузку, включая `ON CONFLICT`. В учебных целях мы остались на широко доступном GP6 образе.
|
||||||
|
|
||||||
|
## Поток данных (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`).
|
||||||
|
|
||||||
|
Повторный запуск 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.
|
||||||
|
|
||||||
|
## Проверка данных
|
||||||
|
- Подними стенд (`make up && make airflow-init`) и запусти DAG `kafka_to_greenplum`, чтобы заполнить таблицу `orders`.
|
||||||
|
- Активируй и запусти DAG `greenplum_data_quality` — он последовательно проверит наличие таблицы, схему, объём данных и отсутствие дублей. Все проверки выполняются внутри Airflow и используют те же настройки подключений.
|
||||||
@@ -0,0 +1,54 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
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):
|
||||||
|
"""Оборачиваем проверку в контекст подключения."""
|
||||||
|
with get_gp_conn() as conn:
|
||||||
|
check_callable(conn)
|
||||||
|
|
||||||
|
|
||||||
|
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"],
|
||||||
|
) as dag:
|
||||||
|
check_exists = PythonOperator(
|
||||||
|
task_id="check_orders_table_exists",
|
||||||
|
python_callable=_run_check,
|
||||||
|
op_args=[assert_orders_table_exists],
|
||||||
|
)
|
||||||
|
check_schema = PythonOperator(
|
||||||
|
task_id="check_orders_schema",
|
||||||
|
python_callable=_run_check,
|
||||||
|
op_args=[assert_orders_schema],
|
||||||
|
)
|
||||||
|
check_has_rows = PythonOperator(
|
||||||
|
task_id="check_orders_has_rows",
|
||||||
|
python_callable=_run_check,
|
||||||
|
op_args=[assert_orders_have_rows],
|
||||||
|
)
|
||||||
|
check_no_duplicates = PythonOperator(
|
||||||
|
task_id="check_order_duplicates",
|
||||||
|
python_callable=_run_check,
|
||||||
|
op_args=[assert_orders_no_duplicates],
|
||||||
|
)
|
||||||
|
|
||||||
|
check_exists >> check_schema >> check_has_rows >> check_no_duplicates
|
||||||
Binary file not shown.
@@ -0,0 +1,107 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
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")
|
||||||
|
|
||||||
|
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 (через Airflow Connection или напрямую по ENV)."""
|
||||||
|
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:
|
||||||
|
# Фоллбек на прямое подключение по переменным окружения.
|
||||||
|
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")),
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def assert_orders_table_exists(conn) -> None:
|
||||||
|
"""Проверяет наличие таблицы orders в схеме public."""
|
||||||
|
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 kafka_to_greenplum.")
|
||||||
|
|
||||||
|
|
||||||
|
def fetch_orders_schema(conn) -> Sequence[Tuple[str, str]]:
|
||||||
|
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 соответствует ожидаемой."""
|
||||||
|
schema = fetch_orders_schema(conn)
|
||||||
|
if list(schema) != EXPECTED_ORDERS_SCHEMA:
|
||||||
|
raise ValueError(f"Неожиданная схема orders: {schema}. Ожидали {EXPECTED_ORDERS_SCHEMA}.")
|
||||||
|
|
||||||
|
|
||||||
|
def fetch_orders_count(conn) -> int:
|
||||||
|
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 перед проверкой.")
|
||||||
|
|
||||||
|
|
||||||
|
def fetch_orders_duplicates(conn) -> int:
|
||||||
|
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."""
|
||||||
|
duplicates = fetch_orders_duplicates(conn)
|
||||||
|
if duplicates:
|
||||||
|
raise ValueError(f"Обнаружены дубли по order_id ({duplicates} шт.) — проверь загрузку данных.")
|
||||||
@@ -0,0 +1,146 @@
|
|||||||
|
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
|
||||||
@@ -0,0 +1,2 @@
|
|||||||
|
confluent-kafka==2.3.0
|
||||||
|
psycopg2-binary==2.9.9
|
||||||
@@ -0,0 +1,151 @@
|
|||||||
|
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
|
||||||
|
|
||||||
|
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
|
||||||
|
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
|
||||||
|
depends_on:
|
||||||
|
pgmeta:
|
||||||
|
condition: service_healthy
|
||||||
|
kafka:
|
||||||
|
condition: service_started
|
||||||
|
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
|
||||||
|
depends_on:
|
||||||
|
pgmeta:
|
||||||
|
condition: service_healthy
|
||||||
|
kafka:
|
||||||
|
condition: service_started
|
||||||
|
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
|
||||||
|
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:
|
||||||
|
kafka-ui-data:
|
||||||
|
greenplum_data:
|
||||||
@@ -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);
|
||||||
Reference in New Issue
Block a user