Улучшения читаемости кода dsg
This commit is contained in:
@@ -39,21 +39,35 @@ def _generate_csv(rows: int, csv_dir: Path) -> str:
|
|||||||
timestamp = datetime.utcnow().strftime("%Y%m%d_%H%M%S")
|
timestamp = datetime.utcnow().strftime("%Y%m%d_%H%M%S")
|
||||||
csv_path = csv_dir / f"orders_{timestamp}.csv"
|
csv_path = csv_dir / f"orders_{timestamp}.csv"
|
||||||
|
|
||||||
base_order_id = int(datetime.utcnow().timestamp())
|
# Генерируем данные в pandas-стиле
|
||||||
order_ids: List[int] = list(range(base_order_id * 1_000, base_order_id * 1_000 + rows))
|
base_order_id = int(datetime.utcnow().timestamp() * 1_000)
|
||||||
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(
|
# Создаём DataFrame с использованием pandas методов
|
||||||
{
|
df = pd.DataFrame({
|
||||||
"order_id": order_ids,
|
# Уникальные order_id начиная с базового значения
|
||||||
"order_ts": order_ts,
|
"order_id": pd.Series(range(base_order_id, base_order_id + rows), dtype="int64"),
|
||||||
"customer_id": customer_ids,
|
|
||||||
"amount": amounts,
|
# Временные метки с интервалом в 1 секунду в обратном порядке
|
||||||
}
|
"order_ts": pd.date_range(
|
||||||
)
|
end=datetime.utcnow(),
|
||||||
|
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)
|
df.to_csv(csv_path, index=False)
|
||||||
logging.info("CSV сохранён: %s (строк: %s)", csv_path, len(df))
|
logging.info("CSV сохранён: %s (строк: %s)", csv_path, len(df))
|
||||||
return str(csv_path)
|
return str(csv_path)
|
||||||
|
|||||||
@@ -1,5 +1,6 @@
|
|||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import logging
|
||||||
from datetime import datetime, timedelta
|
from datetime import datetime, timedelta
|
||||||
|
|
||||||
from airflow import DAG
|
from airflow import DAG
|
||||||
@@ -15,10 +16,36 @@ from helpers.greenplum import (
|
|||||||
|
|
||||||
|
|
||||||
def _run_check(check_callable):
|
def _run_check(check_callable):
|
||||||
"""Оборачиваем проверку в контекст подключения."""
|
"""
|
||||||
|
Оборачивает проверку качества данных в контекст подключения к Greenplum.
|
||||||
|
|
||||||
|
Этот DAG предназначен для автоматической проверки качества данных в таблице 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:
|
with get_gp_conn() as conn:
|
||||||
check_callable(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)}
|
default_args = {"owner": "airflow", "retries": 1, "retry_delay": timedelta(seconds=30)}
|
||||||
|
|
||||||
@@ -29,26 +56,41 @@ with DAG(
|
|||||||
catchup=False,
|
catchup=False,
|
||||||
default_args=default_args,
|
default_args=default_args,
|
||||||
tags=["demo", "greenplum", "quality"],
|
tags=["demo", "greenplum", "quality"],
|
||||||
|
description="Автоматизированные проверки качества данных в Greenplum",
|
||||||
) as dag:
|
) as dag:
|
||||||
|
# Задача 1: Проверка существования таблицы
|
||||||
check_exists = PythonOperator(
|
check_exists = PythonOperator(
|
||||||
task_id="check_orders_table_exists",
|
task_id="check_orders_table_exists",
|
||||||
python_callable=_run_check,
|
python_callable=_run_check,
|
||||||
op_args=[assert_orders_table_exists],
|
op_args=[assert_orders_table_exists],
|
||||||
)
|
)
|
||||||
|
|
||||||
|
# Задача 2: Проверка соответствия схемы таблицы
|
||||||
check_schema = PythonOperator(
|
check_schema = PythonOperator(
|
||||||
task_id="check_orders_schema",
|
task_id="check_orders_schema",
|
||||||
python_callable=_run_check,
|
python_callable=_run_check,
|
||||||
op_args=[assert_orders_schema],
|
op_args=[assert_orders_schema],
|
||||||
)
|
)
|
||||||
|
|
||||||
|
# Задача 3: Проверка наличия данных
|
||||||
check_has_rows = PythonOperator(
|
check_has_rows = PythonOperator(
|
||||||
task_id="check_orders_has_rows",
|
task_id="check_orders_has_rows",
|
||||||
python_callable=_run_check,
|
python_callable=_run_check,
|
||||||
op_args=[assert_orders_have_rows],
|
op_args=[assert_orders_have_rows],
|
||||||
)
|
)
|
||||||
|
|
||||||
|
# Задача 4: Проверка отсутствия дубликатов
|
||||||
check_no_duplicates = PythonOperator(
|
check_no_duplicates = PythonOperator(
|
||||||
task_id="check_order_duplicates",
|
task_id="check_order_duplicates",
|
||||||
python_callable=_run_check,
|
python_callable=_run_check,
|
||||||
op_args=[assert_orders_no_duplicates],
|
op_args=[assert_orders_no_duplicates],
|
||||||
)
|
)
|
||||||
|
|
||||||
check_exists >> check_schema >> check_has_rows >> check_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
|
||||||
|
|||||||
@@ -1,5 +1,6 @@
|
|||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import logging
|
||||||
import os
|
import os
|
||||||
from typing import List, Sequence, Tuple
|
from typing import List, Sequence, Tuple
|
||||||
|
|
||||||
@@ -10,6 +11,7 @@ import psycopg2
|
|||||||
GP_CONN_ID = os.getenv("GP_CONN_ID", "greenplum_conn")
|
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")
|
GP_USE_AIRFLOW_CONN = os.getenv("GP_USE_AIRFLOW_CONN", "true").lower() in ("1", "true", "yes")
|
||||||
|
|
||||||
|
# Ожидаемая схема таблицы orders для проверки качества данных
|
||||||
EXPECTED_ORDERS_SCHEMA: List[Tuple[str, str]] = [
|
EXPECTED_ORDERS_SCHEMA: List[Tuple[str, str]] = [
|
||||||
("order_id", "bigint"),
|
("order_id", "bigint"),
|
||||||
("order_ts", "timestamp without time zone"),
|
("order_ts", "timestamp without time zone"),
|
||||||
@@ -19,28 +21,52 @@ EXPECTED_ORDERS_SCHEMA: List[Tuple[str, str]] = [
|
|||||||
|
|
||||||
|
|
||||||
def get_gp_conn():
|
def get_gp_conn():
|
||||||
"""Возвращает psycopg2 connection к Greenplum (через Airflow Connection или напрямую по ENV)."""
|
"""
|
||||||
|
Возвращает psycopg2 connection к Greenplum.
|
||||||
|
|
||||||
|
Приоритет подключения:
|
||||||
|
1. Через Airflow Connection (если настроено и доступно)
|
||||||
|
2. Прямое подключение по переменным окружения (фоллбек)
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
psycopg2 connection object
|
||||||
|
"""
|
||||||
if GP_USE_AIRFLOW_CONN:
|
if GP_USE_AIRFLOW_CONN:
|
||||||
try:
|
try:
|
||||||
from airflow.providers.postgres.hooks.postgres import PostgresHook
|
from airflow.providers.postgres.hooks.postgres import PostgresHook
|
||||||
|
|
||||||
hook = PostgresHook(postgres_conn_id=GP_CONN_ID)
|
hook = PostgresHook(postgres_conn_id=GP_CONN_ID)
|
||||||
return hook.get_conn()
|
conn = hook.get_conn()
|
||||||
except Exception:
|
logging.info("✅ Подключение через Airflow Connection успешно")
|
||||||
|
return conn
|
||||||
|
except Exception as e:
|
||||||
|
logging.warning("⚠️ Не удалось подключиться через Airflow Connection: %s", e)
|
||||||
|
logging.info("🔄 Переключаемся на прямое подключение по ENV переменным")
|
||||||
# Фоллбек на прямое подключение по переменным окружения.
|
# Фоллбек на прямое подключение по переменным окружения.
|
||||||
pass
|
|
||||||
|
|
||||||
return psycopg2.connect(
|
# Прямое подключение по переменным окружения
|
||||||
dbname=os.getenv("GP_DB", "gpadmin"),
|
conn_params = {
|
||||||
user=os.getenv("GP_USER", "gpadmin"),
|
"dbname": os.getenv("GP_DB", "gpadmin"),
|
||||||
password=os.getenv("GP_PASSWORD", ""),
|
"user": os.getenv("GP_USER", "gpadmin"),
|
||||||
host=os.getenv("GP_HOST", "greenplum"),
|
"password": os.getenv("GP_PASSWORD", ""),
|
||||||
port=int(os.getenv("GP_PORT", "5432")),
|
"host": os.getenv("GP_HOST", "greenplum"),
|
||||||
)
|
"port": int(os.getenv("GP_PORT", "5432")),
|
||||||
|
}
|
||||||
|
logging.info("🔗 Подключение к Greenplum: %s:%s", conn_params["host"], conn_params["port"])
|
||||||
|
return psycopg2.connect(**conn_params)
|
||||||
|
|
||||||
|
|
||||||
def assert_orders_table_exists(conn) -> None:
|
def assert_orders_table_exists(conn) -> None:
|
||||||
"""Проверяет наличие таблицы orders в схеме public."""
|
"""
|
||||||
|
Проверяет наличие таблицы orders в схеме public.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
conn: Подключение к Greenplum
|
||||||
|
|
||||||
|
Raises:
|
||||||
|
ValueError: Если таблица не найдена
|
||||||
|
"""
|
||||||
|
logging.info("🔍 Проверяем существование таблицы public.orders...")
|
||||||
with conn.cursor() as cur:
|
with conn.cursor() as cur:
|
||||||
cur.execute(
|
cur.execute(
|
||||||
"""
|
"""
|
||||||
@@ -50,10 +76,20 @@ def assert_orders_table_exists(conn) -> None:
|
|||||||
"""
|
"""
|
||||||
)
|
)
|
||||||
if cur.fetchone() is None:
|
if cur.fetchone() is None:
|
||||||
raise ValueError("Таблица public.orders не найдена; запусти DAG csv_to_greenplum.")
|
raise ValueError("❌ Таблица public.orders не найдена; запусти DAG csv_to_greenplum.")
|
||||||
|
logging.info("✅ Таблица public.orders существует")
|
||||||
|
|
||||||
|
|
||||||
def fetch_orders_schema(conn) -> Sequence[Tuple[str, str]]:
|
def fetch_orders_schema(conn) -> Sequence[Tuple[str, str]]:
|
||||||
|
"""
|
||||||
|
Получает схему таблицы orders из information_schema.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
conn: Подключение к Greenplum
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
Список кортежей (имя_колонки, тип_данных)
|
||||||
|
"""
|
||||||
with conn.cursor() as cur:
|
with conn.cursor() as cur:
|
||||||
cur.execute(
|
cur.execute(
|
||||||
"""
|
"""
|
||||||
@@ -67,25 +103,69 @@ def fetch_orders_schema(conn) -> Sequence[Tuple[str, str]]:
|
|||||||
|
|
||||||
|
|
||||||
def assert_orders_schema(conn) -> None:
|
def assert_orders_schema(conn) -> None:
|
||||||
"""Проверяет, что схема таблицы orders соответствует ожидаемой."""
|
"""
|
||||||
|
Проверяет, что схема таблицы orders соответствует ожидаемой.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
conn: Подключение к Greenplum
|
||||||
|
|
||||||
|
Raises:
|
||||||
|
ValueError: Если схема не соответствует ожидаемой
|
||||||
|
"""
|
||||||
|
logging.info("📋 Проверяем схему таблицы orders...")
|
||||||
schema = fetch_orders_schema(conn)
|
schema = fetch_orders_schema(conn)
|
||||||
|
logging.info("📊 Фактическая схема: %s", list(schema))
|
||||||
|
logging.info("📊 Ожидаемая схема: %s", EXPECTED_ORDERS_SCHEMA)
|
||||||
|
|
||||||
if list(schema) != EXPECTED_ORDERS_SCHEMA:
|
if list(schema) != EXPECTED_ORDERS_SCHEMA:
|
||||||
raise ValueError(f"Неожиданная схема orders: {schema}. Ожидали {EXPECTED_ORDERS_SCHEMA}.")
|
raise ValueError(f"❌ Неожиданная схема orders: {schema}. Ожидали {EXPECTED_ORDERS_SCHEMA}.")
|
||||||
|
logging.info("✅ Схема таблицы orders соответствует ожиданиям")
|
||||||
|
|
||||||
|
|
||||||
def fetch_orders_count(conn) -> int:
|
def fetch_orders_count(conn) -> int:
|
||||||
|
"""
|
||||||
|
Получает количество строк в таблице orders.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
conn: Подключение к Greenplum
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
Количество строк в таблице
|
||||||
|
"""
|
||||||
with conn.cursor() as cur:
|
with conn.cursor() as cur:
|
||||||
cur.execute("SELECT COUNT(*) FROM public.orders")
|
cur.execute("SELECT COUNT(*) FROM public.orders")
|
||||||
return cur.fetchone()[0]
|
return cur.fetchone()[0]
|
||||||
|
|
||||||
|
|
||||||
def assert_orders_have_rows(conn) -> None:
|
def assert_orders_have_rows(conn) -> None:
|
||||||
"""Проверяет, что таблица orders не пустая."""
|
"""
|
||||||
if fetch_orders_count(conn) <= 0:
|
Проверяет, что таблица orders не пустая.
|
||||||
raise ValueError("Таблица public.orders пустая — запусти DAG csv_to_greenplum перед проверкой.")
|
|
||||||
|
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:
|
def fetch_orders_duplicates(conn) -> int:
|
||||||
|
"""
|
||||||
|
Подсчитывает количество дубликатов по order_id.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
conn: Подключение к Greenplum
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
Количество дублирующихся order_id
|
||||||
|
"""
|
||||||
with conn.cursor() as cur:
|
with conn.cursor() as cur:
|
||||||
cur.execute(
|
cur.execute(
|
||||||
"""
|
"""
|
||||||
@@ -101,7 +181,19 @@ def fetch_orders_duplicates(conn) -> int:
|
|||||||
|
|
||||||
|
|
||||||
def assert_orders_no_duplicates(conn) -> None:
|
def assert_orders_no_duplicates(conn) -> None:
|
||||||
"""Проверяет, что в таблице нет дублей по order_id."""
|
"""
|
||||||
|
Проверяет, что в таблице нет дублей по order_id.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
conn: Подключение к Greenplum
|
||||||
|
|
||||||
|
Raises:
|
||||||
|
ValueError: Если обнаружены дубликаты
|
||||||
|
"""
|
||||||
|
logging.info("🔍 Проверяем отсутствие дубликатов по order_id...")
|
||||||
duplicates = fetch_orders_duplicates(conn)
|
duplicates = fetch_orders_duplicates(conn)
|
||||||
|
logging.info("📊 Найдено дубликатов: %s", duplicates)
|
||||||
|
|
||||||
if duplicates:
|
if duplicates:
|
||||||
raise ValueError(f"Обнаружены дубли по order_id ({duplicates} шт.) — проверь загрузку данных.")
|
raise ValueError(f"❌ Обнаружены дубли по order_id ({duplicates} шт.) — проверь загрузку данных.")
|
||||||
|
logging.info("✅ Дубликаты не обнаружены")
|
||||||
|
|||||||
Reference in New Issue
Block a user