feat(greenplum): добавлен CSV-пайплайн и проверки качества данных

- Зачем:
  - Необходим пример загрузки данных из CSV в Greenplum через Airflow
  - Добавлены проверки качества данных для валидации загруженных данных
- Что:
  - создан DAG csv_to_greenplum для генерации и загрузки CSV в таблицу orders
  - создан DAG csv_to_greenplum_dq для проверок качества данных
  - создан DAG ddl_greenplum_base для применения DDL таблицы orders
  - добавлен модуль helpers/greenplum с функциями подключения и валидации
  - добавлен SQL-скрипт sql/base/orders_ddl.sql для создания таблицы
  - добавлены unit-тесты test_greenplum_helpers.py
- Проверка:
  - запуск DAG csv_to_greenplum в Airflow UI
  - запуск тестов pytest tests/test_greenplum_helpers.py
This commit is contained in:
2026-03-08 19:08:47 +03:00
parent 7072f88e27
commit f703d8e96d
6 changed files with 708 additions and 0 deletions
+167
View File
@@ -0,0 +1,167 @@
from __future__ import annotations
import logging
import os
import random
from datetime import UTC, datetime, timedelta
from pathlib import Path
from typing import List
import pandas as pd
from airflow.operators.python import PythonOperator
from helpers.greenplum import get_gp_conn
from airflow import DAG
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=zstd, 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.now(UTC).strftime("%Y%m%d_%H%M%S")
csv_path = csv_dir / f"orders_{timestamp}.csv"
# Генерируем данные в pandas-стиле
base_order_id = int(datetime.now(UTC).timestamp() * 1_000)
# Создаём DataFrame с использованием pandas методов
df = pd.DataFrame(
{
# Уникальные order_id начиная с базового значения
"order_id": pd.Series(
range(base_order_id, base_order_id + rows), dtype="int64"
),
# Временные метки с интервалом в 1 секунду в обратном порядке
"order_ts": pd.date_range(
end=datetime.now(UTC), periods=rows, freq="1S"
).sort_values(ascending=False),
# Случайные customer_id от 1 до 1000
"customer_id": pd.Series(
random.choices(range(1, 1001), k=rows), dtype="int64"
),
# Случайные суммы от 10 до 500 с округлением до 2 знаков
"amount": pd.Series(
[round(random.uniform(10, 500), 2) for _ in range(rows)],
dtype="float64",
),
}
)
# Сохраняем CSV без индекса
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(2017, 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
@@ -0,0 +1,97 @@
from __future__ import annotations
import logging
from datetime import datetime, timedelta
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,
)
from airflow import DAG
def _run_check(check_callable):
"""
Оборачивает проверку качества данных в контекст подключения к Greenplum.
Этот DAG предназначен для автоматической проверки качества данных
после CSV-пайплайна в таблице public.orders:
1. Проверяет существование таблицы
2. Проверяет соответствие схемы
3. Проверяет наличие данных
4. Проверяет отсутствие дубликатов
Args:
check_callable: Функция проверки, принимающая подключение к БД
"""
# Получаем имя функции для логов
check_name = check_callable.__name__.replace("assert_", "")
logging.info("🚀 Запуск проверки: %s", check_name)
with get_gp_conn() as conn:
check_callable(conn)
logging.info("✅ Проверка пройдена: %s", check_name)
def _log_dq_summary():
"""
Логирует итоговую сводку по качеству данных.
Эта задача выполняется после всех проверок и показывает общий результат.
"""
logging.info("🎉 Все проверки качества данных пройдены успешно!")
logging.info("📊 Качество данных в таблице orders соответствует требованиям.")
default_args = {"owner": "airflow", "retries": 1, "retry_delay": timedelta(seconds=30)}
with DAG(
dag_id="csv_to_greenplum_dq",
start_date=datetime(2017, 1, 1),
schedule=None,
catchup=False,
default_args=default_args,
tags=["demo", "greenplum", "quality", "csv", "dq"],
description="Проверки качества данных после CSV → public.orders в Greenplum",
) as dag:
# Задача 1: Проверка существования таблицы
check_exists = PythonOperator(
task_id="check_orders_table_exists",
python_callable=_run_check,
op_args=[assert_orders_table_exists],
)
# Задача 2: Проверка соответствия схемы таблицы
check_schema = PythonOperator(
task_id="check_orders_schema",
python_callable=_run_check,
op_args=[assert_orders_schema],
)
# Задача 3: Проверка наличия данных
check_has_rows = PythonOperator(
task_id="check_orders_has_rows",
python_callable=_run_check,
op_args=[assert_orders_have_rows],
)
# Задача 4: Проверка отсутствия дубликатов
check_no_duplicates = PythonOperator(
task_id="check_order_duplicates",
python_callable=_run_check,
op_args=[assert_orders_no_duplicates],
)
# Задача 5: Итоговая сводка
dq_summary = PythonOperator(
task_id="data_quality_summary",
python_callable=_log_dq_summary,
)
# Определяем последовательность выполнения задач
check_exists >> check_schema >> check_has_rows >> check_no_duplicates >> dq_summary
+32
View File
@@ -0,0 +1,32 @@
from __future__ import annotations
"""
Учебный DAG: применяет DDL для базовой таблицы orders в Greenplum.
Запускается вручную перед CSV‑пайплайном или после изменения схемы.
"""
from datetime import datetime, timedelta
from airflow.providers.postgres.operators.postgres import PostgresOperator
from airflow import DAG
GREENPLUM_CONN_ID = "greenplum_conn"
default_args = {"owner": "airflow", "retries": 1, "retry_delay": timedelta(seconds=30)}
with DAG(
dag_id="orders_base_ddl",
start_date=datetime(2017, 1, 1),
schedule=None,
catchup=False,
template_searchpath="/sql",
default_args=default_args,
tags=["demo", "greenplum", "ddl", "orders"],
description="Создаёт/обновляет базовую таблицу orders в схеме public",
) as dag:
apply_orders_ddl = PostgresOperator(
task_id="apply_orders_ddl",
postgres_conn_id=GREENPLUM_CONN_ID,
sql="base/orders_ddl.sql",
)
+223
View File
@@ -0,0 +1,223 @@
from __future__ import annotations
"""
LEGACY: Вспомогательные функции для прямого подключения к Greenplum через psycopg2.
Внимание: этот модуль оставлен только для поддержки базового CSV-пайплайна.
В новых DAG (ODS/DDS/DM) используйте встроенный в Airflow PostgresOperator
и штатные механизмы XCom.
"""
import logging
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",
)
# Ожидаемая схема таблицы orders для проверки качества данных
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.
Приоритет подключения:
1. Через Airflow Connection (если настроено и доступно)
2. Прямое подключение по переменным окружения (фоллбек)
Returns:
psycopg2 connection object
"""
if GP_USE_AIRFLOW_CONN:
try:
from airflow.providers.postgres.hooks.postgres import PostgresHook
hook = PostgresHook(postgres_conn_id=GP_CONN_ID)
conn = hook.get_conn()
logging.info("✅ Подключение через Airflow Connection успешно")
return conn
except Exception as e:
logging.warning("⚠️ Не удалось подключиться через Airflow Connection: %s", e)
logging.info("🔄 Переключаемся на прямое подключение по ENV переменным")
# Фоллбек на прямое подключение по переменным окружения.
# Прямое подключение по переменным окружения
conn_params = {
"dbname": os.getenv("GP_DB", "gp_dwh"),
"user": os.getenv("GP_USER", "gpadmin"),
"password": os.getenv("GP_PASSWORD", ""),
"host": os.getenv("GP_HOST", "greenplum"),
"port": int(os.getenv("GP_PORT", "5432")),
}
logging.info(
"🔗 Подключение к Greenplum: %s:%s/%s",
conn_params["host"],
conn_params["port"],
conn_params["dbname"],
)
return psycopg2.connect(**conn_params)
def assert_orders_table_exists(conn) -> None:
"""
Проверяет наличие таблицы orders в схеме public.
Args:
conn: Подключение к Greenplum
Raises:
ValueError: Если таблица не найдена
"""
logging.info("🔍 Проверяем существование таблицы public.orders...")
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 csv_to_greenplum."
)
logging.info("✅ Таблица public.orders существует")
def fetch_orders_schema(conn) -> Sequence[Tuple[str, str]]:
"""
Получает схему таблицы orders из information_schema.
Args:
conn: Подключение к Greenplum
Returns:
Список кортежей (имя_колонки, тип_данных)
"""
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 соответствует ожидаемой.
Args:
conn: Подключение к Greenplum
Raises:
ValueError: Если схема не соответствует ожидаемой
"""
logging.info("📋 Проверяем схему таблицы orders...")
schema = fetch_orders_schema(conn)
logging.info("📊 Фактическая схема: %s", list(schema))
logging.info("📊 Ожидаемая схема: %s", EXPECTED_ORDERS_SCHEMA)
if list(schema) != EXPECTED_ORDERS_SCHEMA:
raise ValueError(
f"❌ Неожиданная схема orders: {schema}. Ожидали {EXPECTED_ORDERS_SCHEMA}."
)
logging.info("✅ Схема таблицы orders соответствует ожиданиям")
def fetch_orders_count(conn) -> int:
"""
Получает количество строк в таблице orders.
Args:
conn: Подключение к Greenplum
Returns:
Количество строк в таблице
"""
with conn.cursor() as cur:
cur.execute("SELECT COUNT(*) FROM public.orders")
return cur.fetchone()[0]
def assert_orders_have_rows(conn) -> None:
"""
Проверяет, что таблица orders не пустая.
Args:
conn: Подключение к Greenplum
Raises:
ValueError: Если таблица пустая
"""
logging.info("📊 Проверяем наличие данных в таблице orders...")
row_count = fetch_orders_count(conn)
logging.info("📈 Количество строк в orders: %s", row_count)
if row_count <= 0:
raise ValueError(
"❌ Таблица public.orders пустая — запусти DAG csv_to_greenplum перед проверкой."
)
logging.info("✅ Таблица orders содержит данные (%s строк)", row_count)
def fetch_orders_duplicates(conn) -> int:
"""
Подсчитывает количество дубликатов по order_id.
Args:
conn: Подключение к Greenplum
Returns:
Количество дублирующихся order_id
"""
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.
Args:
conn: Подключение к Greenplum
Raises:
ValueError: Если обнаружены дубликаты
"""
logging.info("🔍 Проверяем отсутствие дубликатов по order_id...")
duplicates = fetch_orders_duplicates(conn)
logging.info("📊 Найдено дубликатов: %s", duplicates)
if duplicates:
raise ValueError(
f"❌ Обнаружены дубли по order_id ({duplicates} шт.) — проверь загрузку данных."
)
logging.info("✅ Дубликаты не обнаружены")