From dfc1f2ca8452aa4c593e72258c77fb8ac7998d2f Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sun, 12 Oct 2025 17:24:50 +0300 Subject: [PATCH] =?UTF-8?q?=D0=94=D0=BE=D0=B1=D0=B0=D0=B2=D0=BB=D0=B5?= =?UTF-8?q?=D0=BD=20=D0=BA=D0=BE=D0=BD=D1=82=D1=80=D0=BE=D0=BB=D1=8C=20?= =?UTF-8?q?=D0=BA=D0=B0=D1=87=D0=B5=D1=81=D1=82=D0=B2=D0=B0=20=D0=B4=D0=B0?= =?UTF-8?q?=D0=BD=D0=BD=D1=8B=D1=85?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- airflow-greenplum/.gitignore | 1 + airflow-greenplum/README.md | 10 +- .../airflow/dags/data_quality_greenplum.py | 54 +++++++++ .../airflow/dags/helpers/greenplum.py | 107 ++++++++++++++++++ .../airflow/dags/kafka_to_greenplum.py | 41 +------ airflow-greenplum/sql/ddl_gp.sql | 8 +- 6 files changed, 177 insertions(+), 44 deletions(-) create mode 100644 airflow-greenplum/airflow/dags/data_quality_greenplum.py create mode 100644 airflow-greenplum/airflow/dags/helpers/greenplum.py diff --git a/airflow-greenplum/.gitignore b/airflow-greenplum/.gitignore index 1c94086..3f72d2c 100644 --- a/airflow-greenplum/.gitignore +++ b/airflow-greenplum/.gitignore @@ -3,3 +3,4 @@ # Do not commit secrets .env .env.* +/airflow/dags/__pycache__ \ No newline at end of file diff --git a/airflow-greenplum/README.md b/airflow-greenplum/README.md index c2f25af..8311054 100644 --- a/airflow-greenplum/README.md +++ b/airflow-greenplum/README.md @@ -15,7 +15,8 @@ Airflow по‑прежнему использует **Postgres** только - **Kafka (KRaft, без Zookeeper)** — генерация событий. - **Airflow (2.9)** — оркестрация пайплайна, metadata в Postgres. - **DAG** `kafka_to_greenplum.py` — генерирует данные → пишет в Kafka → читает и грузит в Greenplum. -- **DDL** `sql/ddl_gp.sql` — создаёт таблицу `orders` с распределением по `order_id`. +- **DAG** `greenplum_data_quality.py` — выполняет проверки качества данных (наличие таблицы, схема, заполненность, дубли). +- **DDL** `sql/ddl_gp.sql` — создаёт колонночную таблицу `orders` без PRIMARY KEY (AO-таблицы GP6 не поддерживают его), распределённую по `order_id`; контроль дублей реализован в DAG. > Примечание по версиям и надёжности: используется образ `woblerr/greenplum:6.27.1` с поддержкой переменных окружения и fallback значениями. Для продакшен‑подобных тестов зафиксируй digest (SHA256) конкретного тега на Docker Hub. @@ -149,7 +150,6 @@ SELECT count(*) FROM public.orders; - Нет топика `orders`: создай его через Kafka UI (или перезапусти DAG после включения авто‑создания топиков). - `make` отсутствует на Windows: используй команды `docker compose` из раздела «Альтернатива без Make» или установи Git Bash/WSL. -## Что дальше (опциональные расширения) -- Добавить пример загрузки через внешние таблицы/`gpfdist` для демонстрации быстрых батчей в Greenplum. -- Показать альтернативу с GPDB 7 и `ON CONFLICT` (отдельная ветка/вариант DAG). -- Добавить пример использования Airflow Variables/Secrets Backend для передачи порогов и секретов. +## Проверка данных +- Подними стенд (`make up && make airflow-init`) и запусти DAG `kafka_to_greenplum`, чтобы заполнить таблицу `orders`. +- Активируй и запусти DAG `greenplum_data_quality` — он последовательно проверит наличие таблицы, схему, объём данных и отсутствие дублей. Все проверки выполняются внутри Airflow и используют те же настройки подключений. diff --git a/airflow-greenplum/airflow/dags/data_quality_greenplum.py b/airflow-greenplum/airflow/dags/data_quality_greenplum.py new file mode 100644 index 0000000..409f3fe --- /dev/null +++ b/airflow-greenplum/airflow/dags/data_quality_greenplum.py @@ -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 diff --git a/airflow-greenplum/airflow/dags/helpers/greenplum.py b/airflow-greenplum/airflow/dags/helpers/greenplum.py new file mode 100644 index 0000000..9a3b0ad --- /dev/null +++ b/airflow-greenplum/airflow/dags/helpers/greenplum.py @@ -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} шт.) — проверь загрузку данных.") diff --git a/airflow-greenplum/airflow/dags/kafka_to_greenplum.py b/airflow-greenplum/airflow/dags/kafka_to_greenplum.py index 10a079d..a3e92f0 100644 --- a/airflow-greenplum/airflow/dags/kafka_to_greenplum.py +++ b/airflow-greenplum/airflow/dags/kafka_to_greenplum.py @@ -9,50 +9,18 @@ from typing import List, Tuple, Optional from airflow import DAG from airflow.operators.python import PythonOperator from confluent_kafka import Consumer, KafkaException, Producer - -# Пакеты для прямого подключения и батч-загрузки -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") +from helpers.greenplum import get_gp_conn 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 = """ CREATE TABLE IF NOT EXISTS public.orders ( - order_id BIGINT PRIMARY KEY, + order_id BIGINT, order_ts TIMESTAMP NOT NULL, customer_id BIGINT NOT NULL, amount NUMERIC(12,2) NOT NULL @@ -60,7 +28,7 @@ def _create_table(): WITH (appendonly=true, orientation=column, compresstype=zlib) DISTRIBUTED BY (order_id); """ - with _get_gp_conn() as conn, conn.cursor() as cur: + with get_gp_conn() as conn, conn.cursor() as cur: cur.execute(ddl) conn.commit() @@ -79,6 +47,7 @@ def _produce(n=1000): 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) @@ -114,7 +83,7 @@ def _consume_and_load(max_messages=1000, timeout_s: Optional[int] = None): ) consumer.subscribe([TOPIC]) - with _get_gp_conn() as conn, conn.cursor() as cur: + with get_gp_conn() as conn, conn.cursor() as cur: batch: List[Tuple] = [] consumed = 0 while consumed < max_messages: diff --git a/airflow-greenplum/sql/ddl_gp.sql b/airflow-greenplum/sql/ddl_gp.sql index 5ed4a37..1047666 100644 --- a/airflow-greenplum/sql/ddl_gp.sql +++ b/airflow-greenplum/sql/ddl_gp.sql @@ -1,10 +1,12 @@ -- Greenplum DDL (GPDB 6 совместимо) --- Колонночная таблица (append-optimized) и распределение по ключу +-- Колонночная таблица (append-optimized) и распределение по ключу. +-- Внимание: append-optimized таблицы не поддерживают UNIQUE/PRIMARY KEY, +-- поэтому контроль дублей выполняем в DAG при загрузке. CREATE TABLE IF NOT EXISTS public.orders ( - order_id BIGINT PRIMARY KEY, + 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) +WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) DISTRIBUTED BY (order_id);