diff --git a/airflow-greenplum/.env.example b/airflow-greenplum/.env.example index ad2f0e0..fb6257a 100644 --- a/airflow-greenplum/.env.example +++ b/airflow-greenplum/.env.example @@ -13,7 +13,11 @@ 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:9092 KAFKA_TOPIC=orders +KAFKA_BATCH_SIZE=500 +KAFKA_POLL_TIMEOUT=10 diff --git a/airflow-greenplum/.gitignore b/airflow-greenplum/.gitignore index c6f9a44..1c94086 100644 --- a/airflow-greenplum/.gitignore +++ b/airflow-greenplum/.gitignore @@ -1 +1,5 @@ .vscode/settings.json + +# Do not commit secrets +.env +.env.* diff --git a/airflow-greenplum/README.md b/airflow-greenplum/README.md index ed8b10b..c2f25af 100644 --- a/airflow-greenplum/README.md +++ b/airflow-greenplum/README.md @@ -5,6 +5,11 @@ Этот вариант повторяет логику 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)** — генерация событий. @@ -15,7 +20,7 @@ Airflow по‑прежнему использует **Postgres** только > Примечание по версиям и надёжности: используется образ `woblerr/greenplum:6.27.1` с поддержкой переменных окружения и fallback значениями. Для продакшен‑подобных тестов зафиксируй digest (SHA256) конкретного тега на Docker Hub. ## Быстрый старт -Установим make под Вашу операционную систему +Установим make (Linux/Mac) ```bash sudo apt install -y make ``` @@ -26,7 +31,7 @@ cp .env.example .env # При необходимости отредактируйте .env для ваших настроек ``` -Запустим приложение +Запустим приложение (вариант с Make) ```bash make up && make airflow-init make logs # ждем "Listening at: http://0.0.0.0:8080" @@ -34,6 +39,48 @@ 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_BATCH_SIZE` — размер батча при вставке в Greenplum (по умолчанию 500). +- `KAFKA_POLL_TIMEOUT` — таймаут ожидания сообщения в секундах (по умолчанию 10). + ### Проверка загрузки ```bash # Подключение к Greenplum через внешний psql клиент @@ -50,6 +97,10 @@ docker compose -f docker-compose.yml exec greenplum bash -c "su - gpadmin -c '/u 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` — шаблон переменных окружения. @@ -65,12 +116,16 @@ SELECT count(*) FROM public.orders; - `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. Образ `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. @@ -79,3 +134,22 @@ SELECT count(*) FROM public.orders; ## Ограничения и заметки - Этот стенд — учебный. Для высокой надёжности и производительности 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. + +## Что дальше (опциональные расширения) +- Добавить пример загрузки через внешние таблицы/`gpfdist` для демонстрации быстрых батчей в Greenplum. +- Показать альтернативу с GPDB 7 и `ON CONFLICT` (отдельная ветка/вариант DAG). +- Добавить пример использования Airflow Variables/Secrets Backend для передачи порогов и секретов. diff --git a/airflow-greenplum/airflow/dags/kafka_to_greenplum.py b/airflow-greenplum/airflow/dags/kafka_to_greenplum.py index 0875f19..10a079d 100644 --- a/airflow-greenplum/airflow/dags/kafka_to_greenplum.py +++ b/airflow-greenplum/airflow/dags/kafka_to_greenplum.py @@ -1,22 +1,53 @@ from __future__ import annotations -import os, json, random + +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 -import psycopg2 -from confluent_kafka import Producer, Consumer, KafkaException +from confluent_kafka import Consumer, KafkaException, Producer -# Greenplum (Postgres wire protocol) -GP_DSN = { - "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")), -} +# Пакеты для прямого подключения и батч-загрузки +import psycopg2 +from psycopg2.extras import execute_values + +# Airflow Connection ID для Greenplum (создаётся в UI или CLI). +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") KAFKA_BOOTSTRAP = os.getenv("KAFKA_BOOTSTRAP", "kafka:9092") 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")) + + +def _get_gp_conn(): + """Получаем соединение с Greenplum через Airflow Connection или через ENV-DSN. + + Если переменная `GP_USE_AIRFLOW_CONN` = false или отсутствует провайдер Postgres, + используем ENV-подключение напрямую (psycopg2). + """ + 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 _create_table(): ddl = """ @@ -29,69 +60,119 @@ def _create_table(): WITH (appendonly=true, orientation=column, compresstype=zlib) DISTRIBUTED BY (order_id); """ - with psycopg2.connect(**GP_DSN) as conn, conn.cursor() as cur: + with _get_gp_conn() as conn, conn.cursor() as cur: cur.execute(ddl) conn.commit() + def _produce(n=1000): - p = Producer({"bootstrap.servers": KAFKA_BOOTSTRAP}) - for i in range(n): + producer = Producer({"bootstrap.servers": KAFKA_BOOTSTRAP}) + for idx in range(n): payload = { - "order_id": i + 1, + "order_id": idx + 1, "order_ts": datetime.utcnow().isoformat(), "customer_id": random.randint(1, 100), "amount": round(random.uniform(10, 500), 2), } - p.produce(TOPIC, json.dumps(payload).encode("utf-8")) - p.flush() + producer.produce(TOPIC, json.dumps(payload).encode("utf-8")) + producer.flush() -def _consume_and_load(max_messages=1000, timeout_s=10): - consumer = Consumer({ - "bootstrap.servers": KAFKA_BOOTSTRAP, - "group.id": "airflow-loader-gp", - "auto.offset.reset": "earliest", - "enable.auto.commit": False, - }) + +def _flush_batch(cur, rows: List[Tuple]): + 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 psycopg2.connect(**GP_DSN) as conn, conn.cursor() as cur: - inserted = 0 - while inserted < max_messages: + with _get_gp_conn() as conn, conn.cursor() as cur: + batch: List[Tuple] = [] + consumed = 0 + while consumed < max_messages: msg = consumer.poll(timeout_s) if msg is None: + # Нет новых сообщений — сбрасываем остаток батча и выходим + if batch: + _flush_batch(cur, batch) + conn.commit() + batch.clear() break if msg.error(): raise KafkaException(msg.error()) - d = json.loads(msg.value().decode("utf-8")) - # GPDB6 не поддерживает ON CONFLICT — используем WHERE NOT EXISTS - cur.execute( - """ - INSERT INTO public.orders(order_id, order_ts, customer_id, amount) - SELECT %s, %s, %s, %s - WHERE NOT EXISTS ( - SELECT 1 FROM public.orders WHERE order_id = %s - ); - """, - (d["order_id"], d["order_ts"], d["customer_id"], d["amount"], d["order_id"]), - ) - inserted += 1 - conn.commit() - consumer.commit(); consumer.close() + data = json.loads(msg.value().decode("utf-8")) + batch.append( + ( + int(data["order_id"]), + 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_interval=None, + 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}) + 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 diff --git a/airflow-greenplum/airflow/requirements.txt b/airflow-greenplum/airflow/requirements.txt index a7da61b..f4676e7 100644 --- a/airflow-greenplum/airflow/requirements.txt +++ b/airflow-greenplum/airflow/requirements.txt @@ -1,3 +1,2 @@ confluent-kafka==2.3.0 psycopg2-binary==2.9.9 -pandas==2.1.*