diff --git a/airflow-greenplum/.env.example b/airflow-greenplum/.env.example index a220e1a..22959ab 100644 --- a/airflow-greenplum/.env.example +++ b/airflow-greenplum/.env.example @@ -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 diff --git a/airflow-greenplum/.gitignore b/airflow-greenplum/.gitignore index 0c90d68..64e6841 100644 --- a/airflow-greenplum/.gitignore +++ b/airflow-greenplum/.gitignore @@ -6,3 +6,4 @@ __pycache__/ */__pycache__/ *.pyc +data/ diff --git a/airflow-greenplum/README.md b/airflow-greenplum/README.md index e8cc850..9f22a60 100644 --- a/airflow-greenplum/README.md +++ b/airflow-greenplum/README.md @@ -1,12 +1,13 @@ -# DE Starter Kit — Greenplum + Kafka + Airflow +# DE Starter Kit — Airflow + Greenplum + CSV -Добро пожаловать в учебный стенд для изучения основ Data Engineering! Этот проект поможет вам освоить ключевые инструменты современных data pipeline: **Airflow** для оркестрации, **Kafka** для потоковой передачи данных и **Greenplum** как аналитическую базу данных. +Добро пожаловать в учебный стенд для изучения основ Data Engineering! Этот проект поможет вам освоить ключевые инструменты современных data pipeline: **Airflow** для оркестрации, **pandas/CSV** для подготовки данных и **Greenplum** как аналитическую базу данных. ## 🎯 Что вы узнаете - Как настроить локальный стек данных с помощью Docker - Как Airflow управляет workflow и координирует задачи -- Как данные перемещаются из Kafka в аналитическую базу +- Как генерировать датасеты через pandas и сохранять их в CSV +- Как загружать данные в Greenplum пакетами и избегать дублей - Как проверять качество данных в автоматизированных pipeline - Основы проектирования ETL/ELT процессов @@ -38,22 +39,18 @@ docker compose run --rm airflow-init ### Шаг 3: Первый запуск pipeline 1. Откройте Airflow UI: **http://localhost:8080** (логин/пароль: admin/admin) -2. Найдите DAG с названием **kafka_to_greenplum** +2. Найдите DAG с названием **csv_to_greenplum** 3. Нажмите на переключатель слева от названия DAG, чтобы включить его 4. Нажмите кнопку **Trigger** (значок воспроизведения ▶️) 🎉 **Поздравляем!** Вы только что запустили свой первый data pipeline: -- Система сгенерировала 1000 тестовых заказов -- Данные отправились в Kafka (систему потоковой передачи сообщений) -- Airflow прочитал данные из Kafka и загрузил их в Greenplum +- Система сгенерировала 1000 тестовых заказов при помощи pandas +- Датасет сохранился в CSV-файл в каталоге `./data` +- Airflow загрузил данные из CSV в Greenplum без дублей по `order_id` ### Шаг 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'" @@ -61,8 +58,13 @@ 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. + --- ## 🛠️ Подробная настройка (для уверенных пользователей) @@ -105,12 +107,12 @@ make gp-psql # Подключение к Greenplum ### Основные компоненты - **Greenplum** — аналитическая база данных для хранения и анализа данных -- **Kafka** — система потоковой передачи данных (без Zookeeper) - **Airflow** — оркестратор workflow и задач - **Postgres** — база метаданных для Airflow +- **pandas** — библиотека для генерации и анализа данных в формате CSV ### Готовые DAG (workflow) -- **kafka_to_greenplum** — базовый pipeline: генерация → Kafka → Greenplum +- **csv_to_greenplum** — базовый pipeline: pandas → CSV → Greenplum - **greenplum_data_quality** — проверки качества данных (наличие таблицы, схема, дубликаты) ### Полезные команды @@ -138,9 +140,9 @@ make logs # Следить за логами Airflow - `GP_DB` — база данных (по умолчанию: gpadmin) - `GP_PORT` — порт (по умолчанию: 5432) -### Kafka -- `KAFKA_TOPIC` — имя топика (по умолчанию: orders) -- `KAFKA_BATCH_SIZE` — размер пакета при загрузке (по умолчанию: 500) +### CSV pipeline +- `CSV_DIR` — путь к каталогу с CSV внутри контейнеров Airflow (по умолчанию: `/opt/airflow/data`) +- `CSV_ROWS` — количество строк, генерируемых DAG (по умолчанию: 1000) ### Airflow - `GP_CONN_ID` — ID подключения (по умолчанию: greenplum_conn) @@ -151,10 +153,11 @@ make logs # Следить за логами Airflow ### Архитектура pipeline -**Поток данных в DAG `kafka_to_greenplum`:** -1. `create_table` — создает таблицу `public.orders` в Greenplum -2. `produce_messages` — генерирует и отправляет сообщения в Kafka -3. `consume_and_load` — читает из Kafka и загружает в Greenplum батчами +**Поток данных в 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` > 💡 **Безопасность повторного запуска:** Pipeline защищен от дубликатов, поэтому его можно запускать многократно. @@ -180,7 +183,7 @@ make logs # Следить за логами Airflow |----------|---------| | Airflow UI не открывается | Дождитесь сообщения `Listening at: http://0.0.0.0:8080` в логах (`make logs`) | | Ошибка подключения к Greenplum | Убедитесь, что контейнер `greenplum` стал статусом `healthy` (проверьте `docker compose ps`) | -| Нет топика `orders` в Kafka | Создайте топик через Kafka UI или дождитесь авто-создания при первом запуске DAG | +| Нет файла в `./data` после запуска DAG | Проверьте логи задачи `generate_csv`, убедитесь, что `CSV_DIR` смонтирован в docker-compose | | Команда `make` не найдена | Используйте полные команды `docker compose` или установите make | --- @@ -193,7 +196,7 @@ make logs # Следить за логами Airflow ├── Makefile # Удобные команды для работы ├── airflow/ │ └── dags/ # Файлы workflow (DAG) -│ ├── kafka_to_greenplum.py +│ ├── csv_to_greenplum.py │ └── data_quality_greenplum.py └── sql/ └── ddl_gp.sql # Создание таблицы в Greenplum diff --git a/airflow-greenplum/airflow/dags/csv_to_greenplum.py b/airflow-greenplum/airflow/dags/csv_to_greenplum.py new file mode 100644 index 0000000..fe51ac0 --- /dev/null +++ b/airflow-greenplum/airflow/dags/csv_to_greenplum.py @@ -0,0 +1,146 @@ +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" + + base_order_id = int(datetime.utcnow().timestamp()) + order_ids: List[int] = list(range(base_order_id * 1_000, base_order_id * 1_000 + rows)) + now = datetime.utcnow() + order_ts = [(now - timedelta(seconds=i)).isoformat() for i in range(rows)] + customer_ids = [random.randint(1, 1_000) for _ in range(rows)] + amounts = [round(random.uniform(10, 500), 2) for _ in range(rows)] + + df = pd.DataFrame( + { + "order_id": order_ids, + "order_ts": order_ts, + "customer_id": customer_ids, + "amount": amounts, + } + ) + 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 diff --git a/airflow-greenplum/airflow/dags/helpers/greenplum.py b/airflow-greenplum/airflow/dags/helpers/greenplum.py index 9a3b0ad..0709bd8 100644 --- a/airflow-greenplum/airflow/dags/helpers/greenplum.py +++ b/airflow-greenplum/airflow/dags/helpers/greenplum.py @@ -50,7 +50,7 @@ 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.") def fetch_orders_schema(conn) -> Sequence[Tuple[str, str]]: @@ -82,7 +82,7 @@ def fetch_orders_count(conn) -> int: def assert_orders_have_rows(conn) -> None: """Проверяет, что таблица orders не пустая.""" if fetch_orders_count(conn) <= 0: - raise ValueError("Таблица public.orders пустая — запусти DAG kafka_to_greenplum перед проверкой.") + raise ValueError("Таблица public.orders пустая — запусти DAG csv_to_greenplum перед проверкой.") def fetch_orders_duplicates(conn) -> int: diff --git a/airflow-greenplum/airflow/dags/kafka_to_greenplum.py b/airflow-greenplum/airflow/dags/kafka_to_greenplum.py deleted file mode 100644 index 969f2b9..0000000 --- a/airflow-greenplum/airflow/dags/kafka_to_greenplum.py +++ /dev/null @@ -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 diff --git a/airflow-greenplum/airflow/requirements.txt b/airflow-greenplum/airflow/requirements.txt index f4676e7..f519dd1 100644 --- a/airflow-greenplum/airflow/requirements.txt +++ b/airflow-greenplum/airflow/requirements.txt @@ -1,2 +1,2 @@ -confluent-kafka==2.3.0 psycopg2-binary==2.9.9 +pandas==2.1.4 diff --git a/airflow-greenplum/docker-compose.yml b/airflow-greenplum/docker-compose.yml index 774f591..73e9613 100644 --- a/airflow-greenplum/docker-compose.yml +++ b/airflow-greenplum/docker-compose.yml @@ -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: