From e7e40d18f7707c0905a588f9e350c2535d8be420 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Tue, 16 Sep 2025 23:03:35 +0300 Subject: [PATCH] =?UTF-8?q?=D0=9F=D0=B5=D1=80=D0=B2=D1=8B=D0=B5=20=D0=BD?= =?UTF-8?q?=D0=B0=D0=B1=D1=80=D0=BE=D1=81=D0=BA=D0=B8=20airflow-gp?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- airflow-greenplum/.env.example | 22 +++ airflow-greenplum/Makefile | 19 +++ airflow-greenplum/README.md | 53 +++++++ .../airflow/dags/kafka_to_greenplum.py | 97 ++++++++++++ airflow-greenplum/airflow/requirements.txt | 3 + airflow-greenplum/docker-compose.yml | 140 ++++++++++++++++++ airflow-greenplum/sql/ddl_gp.sql | 10 ++ 7 files changed, 344 insertions(+) create mode 100644 airflow-greenplum/.env.example create mode 100644 airflow-greenplum/Makefile create mode 100644 airflow-greenplum/README.md create mode 100644 airflow-greenplum/airflow/dags/kafka_to_greenplum.py create mode 100644 airflow-greenplum/airflow/requirements.txt create mode 100644 airflow-greenplum/docker-compose.yml create mode 100644 airflow-greenplum/sql/ddl_gp.sql diff --git a/airflow-greenplum/.env.example b/airflow-greenplum/.env.example new file mode 100644 index 0000000..a487abf --- /dev/null +++ b/airflow-greenplum/.env.example @@ -0,0 +1,22 @@ +# ==== Compose Environment (.env) ==== +# Airflow metadata DB (Postgres) +PG_USER=airflow +PG_PASSWORD=airflow +PG_DB=airflow + +# Greenplum access from Airflow +GP_HOST=greenplum +GP_PORT=5432 +GP_DB=gpadmin +GP_USER=gpadmin +# В большинстве демо-образов TCP доступ настроен без пароля (trust). +# Если потребуется пароль, задай его здесь: +GP_PASSWORD= + +# Kafka +KAFKA_BOOTSTRAP=kafka:9092 +KAFKA_TOPIC=orders + +# Airflow admin +AIRFLOW_USER=admin +AIRFLOW_PASSWORD=admin diff --git a/airflow-greenplum/Makefile b/airflow-greenplum/Makefile new file mode 100644 index 0000000..6624e72 --- /dev/null +++ b/airflow-greenplum/Makefile @@ -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 -lc "su - gpadmin -c 'psql -p 5432 -d gpadmin'" + +ddl-gp: + docker compose -f docker-compose.yml exec greenplum bash -lc "su - gpadmin -c \"psql -d gpadmin -f /sql/ddl_gp.sql\"" diff --git a/airflow-greenplum/README.md b/airflow-greenplum/README.md new file mode 100644 index 0000000..d332ee3 --- /dev/null +++ b/airflow-greenplum/README.md @@ -0,0 +1,53 @@ +# DE Starter Kit — Greenplum + Kafka + Airflow + +Обновлено: 2025-09-15 11:12 + +Этот вариант повторяет логику Postgres-стенда, но в роли DWH — **Greenplum (single node в Docker)**. +Airflow по‑прежнему использует **Postgres** только как metadata DB (это стандартная и простая схема). + +## Что внутри +- **Greenplum** (single node) — `datagrip/greenplum:6.8` (локальный стенд для разработки). +- **Kafka (KRaft, без Zookeeper)** — генерация событий. +- **Airflow (2.9)** — оркестрация пайплайна, metadata в Postgres. +- **DAG** `kafka_to_greenplum.py` — генерирует данные → пишет в Kafka → читает и грузит в Greenplum. +- **DDL** `sql/ddl_gp.sql` — создаёт таблицу `orders` с распределением по `order_id`. + +> Примечание по версиям и надёжности: образы Greenplum в публичном Docker Hub — комьюнити/сторонние. Я выбрал `datagrip/greenplum:6.8` как популярный и поддерживаемый для локальной разработки. Для продакшен‑подобных тестов зафиксируй digest (SHA256) конкретного тега на Docker Hub и/или собери образ самостоятельно. См. раздел «Пинning и альтернативы» ниже. + +## Быстрый старт +Установим make под Вашу операционную систему +```bash +sudo apt install -y make +``` +Запустим приложение +```bash +cp .env.example .env +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`. + +### Проверка загрузки +```bash +make gp-psql # войти в psql к Greenplum +-- внутри psql: +\dt +SELECT count(*) FROM public.orders; +``` + +## Файлы +- `docker-compose.gp.yml` — сервисы Greenplum + Kafka + Airflow + Postgres (metadata). +- `.env.example` — переменные окружения. +- `airflow/dags/kafka_to_greenplum.py` — сам DAG. +- `sql/ddl_gp.sql` — DDL таблицы в Greenplum. +- `Makefile` — обёртки команд (`up`, `down`, `airflow-init`, `gp-psql`, `ddl-gp`). + +## Пинning и альтернативы +- Зафиксируй digest образа `datagrip/greenplum:6.8` (Docker Hub → Tag → «Copy digest») и замени тег на `@sha256:...` в `docker-compose.gp.yml`. +- Альтернативы: репозиторий `woblerr/docker-greenplum` позволяет **собрать свой образ** для GPDB 6/7 (подойдёт, если нужен полный контроль и повторяемость сборки). +- Особенность GPDB 6: в нём **нет** `INSERT ... ON CONFLICT`. В DAG используется безопасная для GP6 конструкция `INSERT ... WHERE NOT EXISTS` внутри транзакции. + +## Ограничения и заметки +- Этот стенд — учебный. Для высокой надёжности и производительности Greenplum обычно разворачивают кластерами на нескольких узлах, на裸‑железе/VM с отдельными дисками под сегменты. +- Для GP7 (основан на новее PostgreSQL) можно упростить загрузку, включая `ON CONFLICT`. В учебных целях мы остались на широко доступном GP6 образе. diff --git a/airflow-greenplum/airflow/dags/kafka_to_greenplum.py b/airflow-greenplum/airflow/dags/kafka_to_greenplum.py new file mode 100644 index 0000000..0875f19 --- /dev/null +++ b/airflow-greenplum/airflow/dags/kafka_to_greenplum.py @@ -0,0 +1,97 @@ +from __future__ import annotations +import os, json, random +from datetime import datetime, timedelta +from airflow import DAG +from airflow.operators.python import PythonOperator +import psycopg2 +from confluent_kafka import Producer, Consumer, KafkaException + +# 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")), +} + +KAFKA_BOOTSTRAP = os.getenv("KAFKA_BOOTSTRAP", "kafka:9092") +TOPIC = os.getenv("KAFKA_TOPIC", "orders") + +def _create_table(): + ddl = """ + CREATE TABLE IF NOT EXISTS public.orders ( + order_id BIGINT PRIMARY KEY, + 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 psycopg2.connect(**GP_DSN) 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): + payload = { + "order_id": i + 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() + +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, + }) + consumer.subscribe([TOPIC]) + + with psycopg2.connect(**GP_DSN) as conn, conn.cursor() as cur: + inserted = 0 + while inserted < max_messages: + msg = consumer.poll(timeout_s) + if msg is None: + 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() + +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, + 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 diff --git a/airflow-greenplum/airflow/requirements.txt b/airflow-greenplum/airflow/requirements.txt new file mode 100644 index 0000000..2b28044 --- /dev/null +++ b/airflow-greenplum/airflow/requirements.txt @@ -0,0 +1,3 @@ +confluent-kafka==2.3.0 +psycopg2-binary==2.9.9 +pandas==2.2.2 diff --git a/airflow-greenplum/docker-compose.yml b/airflow-greenplum/docker-compose.yml new file mode 100644 index 0000000..57bfd2c --- /dev/null +++ b/airflow-greenplum/docker-compose.yml @@ -0,0 +1,140 @@ +version: "3.9" + +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: datagrip/greenplum:6.8 + # container_name: gp_single + hostname: gpdbsne + # Порты: внешний 5432 + ports: + - "5432:5432" + volumes: + - ./sql:/sql:ro + # Простая проверка доступности: psql откликается + healthcheck: + test: ["CMD-SHELL", "pg_isready -h 127.0.0.1 -p 5432 || 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 + ports: + - 8082:8080 + environment: + - DYNAMIC_CONFIG_ENABLED=true + - KAFKA_CLUSTERS_0_NAME=KafkaCluster + - KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS=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__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 + entrypoint: ["/bin/bash", "-lc"] + command: | + set -e + until pg_isready -h pgmeta -U ${PG_USER} -d ${PG_DB}; do echo 'waiting for pgmeta'; sleep 2; done + pip install --no-cache-dir -r /opt/airflow/requirements.txt + airflow db migrate + airflow users create --username ${AIRFLOW_USER} --password ${AIRFLOW_PASSWORD} --firstname Admin --lastname User --role Admin --email admin@example.org + depends_on: + pgmeta: + condition: service_healthy + +volumes: + pgmeta: + kafka-ui-data: diff --git a/airflow-greenplum/sql/ddl_gp.sql b/airflow-greenplum/sql/ddl_gp.sql new file mode 100644 index 0000000..5ed4a37 --- /dev/null +++ b/airflow-greenplum/sql/ddl_gp.sql @@ -0,0 +1,10 @@ +-- Greenplum DDL (GPDB 6 совместимо) +-- Колонночная таблица (append-optimized) и распределение по ключу +CREATE TABLE IF NOT EXISTS public.orders ( + order_id BIGINT PRIMARY KEY, + 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);