Результат прогона линтера

This commit is contained in:
2025-10-18 23:56:37 +03:00
parent b2f40e4dd2
commit f1e01274de
5 changed files with 91 additions and 72 deletions
+36 -29
View File
@@ -8,10 +8,11 @@ from pathlib import Path
from typing import List from typing import List
import pandas as pd import pandas as pd
from airflow import DAG
from airflow.operators.python import PythonOperator from airflow.operators.python import PythonOperator
from helpers.greenplum import get_gp_conn from helpers.greenplum import get_gp_conn
from airflow import DAG
CSV_DIR = Path(os.getenv("CSV_DIR", "/opt/airflow/data")) CSV_DIR = Path(os.getenv("CSV_DIR", "/opt/airflow/data"))
CSV_ROWS = int(os.getenv("CSV_ROWS", "1000")) CSV_ROWS = int(os.getenv("CSV_ROWS", "1000"))
@@ -41,32 +42,30 @@ def _generate_csv(rows: int, csv_dir: Path) -> str:
# Генерируем данные в pandas-стиле # Генерируем данные в pandas-стиле
base_order_id = int(datetime.utcnow().timestamp() * 1_000) base_order_id = int(datetime.utcnow().timestamp() * 1_000)
# Создаём DataFrame с использованием pandas методов # Создаём DataFrame с использованием pandas методов
df = pd.DataFrame({ df = pd.DataFrame(
# Уникальные order_id начиная с базового значения {
"order_id": pd.Series(range(base_order_id, base_order_id + rows), dtype="int64"), # Уникальные order_id начиная с базового значения
"order_id": pd.Series(
# Временные метки с интервалом в 1 секунду в обратном порядке range(base_order_id, base_order_id + rows), dtype="int64"
"order_ts": pd.date_range( ),
end=datetime.utcnow(), # Временные метки с интервалом в 1 секунду в обратном порядке
periods=rows, "order_ts": pd.date_range(
freq="1S" end=datetime.utcnow(), periods=rows, freq="1S"
).sort_values(ascending=False), ).sort_values(ascending=False),
# Случайные customer_id от 1 до 1000
# Случайные customer_id от 1 до 1000 "customer_id": pd.Series(
"customer_id": pd.Series( random.choices(range(1, 1001), k=rows), dtype="int64"
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)],
# Случайные суммы от 10 до 500 с округлением до 2 знаков dtype="float64",
"amount": pd.Series( ),
[round(random.uniform(10, 500), 2) for _ in range(rows)], }
dtype="float64" )
)
})
# Сохраняем CSV без индекса # Сохраняем 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))
@@ -77,7 +76,9 @@ def _preview_csv(csv_path: str, sample_rows: int = 5) -> None:
"""Отображает предпросмотр CSV через pandas (head и describe).""" """Отображает предпросмотр CSV через pandas (head и describe)."""
df = pd.read_csv(csv_path) df = pd.read_csv(csv_path)
df["order_ts"] = pd.to_datetime(df["order_ts"], errors="coerce") 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)) logging.info(
"Первые %s строк:\n%s", sample_rows, df.head(sample_rows).to_string(index=False)
)
numeric_summary = df.describe(include="number") numeric_summary = df.describe(include="number")
logging.info("Числовая статистика:\n%s", numeric_summary.to_string()) logging.info("Числовая статистика:\n%s", numeric_summary.to_string())
if df["order_ts"].notna().any(): if df["order_ts"].notna().any():
@@ -94,8 +95,14 @@ def _load_csv(csv_path: str) -> None:
if not csv_file.exists(): if not csv_file.exists():
raise FileNotFoundError(f"CSV не найден: {csv_file}") 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: with (
cur.execute("CREATE TEMP TABLE tmp_orders (LIKE public.orders INCLUDING DEFAULTS) ON COMMIT DROP;") 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( cur.copy_expert(
"COPY tmp_orders (order_id, order_ts, customer_id, amount) FROM STDIN WITH CSV HEADER", "COPY tmp_orders (order_id, order_ts, customer_id, amount) FROM STDIN WITH CSV HEADER",
f, f,
+13 -16
View File
@@ -3,38 +3,35 @@ from __future__ import annotations
import logging import logging
from datetime import datetime, timedelta from datetime import datetime, timedelta
from airflow import DAG
from airflow.operators.python import PythonOperator 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 helpers.greenplum import ( from airflow import DAG
assert_orders_have_rows,
assert_orders_no_duplicates,
assert_orders_schema,
assert_orders_table_exists,
get_gp_conn,
)
def _run_check(check_callable): def _run_check(check_callable):
""" """
Оборачивает проверку качества данных в контекст подключения к Greenplum. Оборачивает проверку качества данных в контекст подключения к Greenplum.
Этот DAG предназначен для автоматической проверки качества данных в таблице orders: Этот DAG предназначен для автоматической проверки качества данных в таблице orders:
1. Проверяет существование таблицы 1. Проверяет существование таблицы
2. Проверяет соответствие схемы 2. Проверяет соответствие схемы
3. Проверяет наличие данных 3. Проверяет наличие данных
4. Проверяет отсутствие дубликатов 4. Проверяет отсутствие дубликатов
Args: Args:
check_callable: Функция проверки, принимающая подключение к БД check_callable: Функция проверки, принимающая подключение к БД
""" """
# Получаем имя функции для логов # Получаем имя функции для логов
check_name = check_callable.__name__.replace("assert_", "") check_name = check_callable.__name__.replace("assert_", "")
logging.info("🚀 Запуск проверки: %s", check_name) 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) logging.info("✅ Проверка пройдена: %s", check_name)
@@ -64,28 +61,28 @@ with DAG(
python_callable=_run_check, python_callable=_run_check,
op_args=[assert_orders_table_exists], op_args=[assert_orders_table_exists],
) )
# Задача 2: Проверка соответствия схемы таблицы # Задача 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: Проверка наличия данных # Задача 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: Проверка отсутствия дубликатов # Задача 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],
) )
# Задача 5: Итоговая сводка # Задача 5: Итоговая сводка
dq_summary = PythonOperator( dq_summary = PythonOperator(
task_id="data_quality_summary", task_id="data_quality_summary",
+39 -25
View File
@@ -9,7 +9,11 @@ import psycopg2
# Настройки для подключения к Greenplum. По умолчанию используем Airflow Connection, # Настройки для подключения к Greenplum. По умолчанию используем Airflow Connection,
# но при проблемах можно переключиться на ENV-подключение, установив GP_USE_AIRFLOW_CONN=false. # но при проблемах можно переключиться на ENV-подключение, установив GP_USE_AIRFLOW_CONN=false.
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 для проверки качества данных # Ожидаемая схема таблицы orders для проверки качества данных
EXPECTED_ORDERS_SCHEMA: List[Tuple[str, str]] = [ EXPECTED_ORDERS_SCHEMA: List[Tuple[str, str]] = [
@@ -23,11 +27,11 @@ EXPECTED_ORDERS_SCHEMA: List[Tuple[str, str]] = [
def get_gp_conn(): def get_gp_conn():
""" """
Возвращает psycopg2 connection к Greenplum. Возвращает psycopg2 connection к Greenplum.
Приоритет подключения: Приоритет подключения:
1. Через Airflow Connection (если настроено и доступно) 1. Через Airflow Connection (если настроено и доступно)
2. Прямое подключение по переменным окружения (фоллбек) 2. Прямое подключение по переменным окружения (фоллбек)
Returns: Returns:
psycopg2 connection object psycopg2 connection object
""" """
@@ -52,17 +56,19 @@ def get_gp_conn():
"host": os.getenv("GP_HOST", "greenplum"), "host": os.getenv("GP_HOST", "greenplum"),
"port": int(os.getenv("GP_PORT", "5432")), "port": int(os.getenv("GP_PORT", "5432")),
} }
logging.info("🔗 Подключение к Greenplum: %s:%s", conn_params["host"], conn_params["port"]) logging.info(
"🔗 Подключение к Greenplum: %s:%s", conn_params["host"], conn_params["port"]
)
return psycopg2.connect(**conn_params) return psycopg2.connect(**conn_params)
def assert_orders_table_exists(conn) -> None: def assert_orders_table_exists(conn) -> None:
""" """
Проверяет наличие таблицы orders в схеме public. Проверяет наличие таблицы orders в схеме public.
Args: Args:
conn: Подключение к Greenplum conn: Подключение к Greenplum
Raises: Raises:
ValueError: Если таблица не найдена ValueError: Если таблица не найдена
""" """
@@ -76,17 +82,19 @@ 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 существует") 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. Получает схему таблицы orders из information_schema.
Args: Args:
conn: Подключение к Greenplum conn: Подключение к Greenplum
Returns: Returns:
Список кортежей (имя_колонки, тип_данных) Список кортежей (имя_колонки, тип_данных)
""" """
@@ -105,10 +113,10 @@ def fetch_orders_schema(conn) -> Sequence[Tuple[str, str]]:
def assert_orders_schema(conn) -> None: def assert_orders_schema(conn) -> None:
""" """
Проверяет, что схема таблицы orders соответствует ожидаемой. Проверяет, что схема таблицы orders соответствует ожидаемой.
Args: Args:
conn: Подключение к Greenplum conn: Подключение к Greenplum
Raises: Raises:
ValueError: Если схема не соответствует ожидаемой ValueError: Если схема не соответствует ожидаемой
""" """
@@ -116,19 +124,21 @@ def assert_orders_schema(conn) -> None:
schema = fetch_orders_schema(conn) schema = fetch_orders_schema(conn)
logging.info("📊 Фактическая схема: %s", list(schema)) logging.info("📊 Фактическая схема: %s", list(schema))
logging.info("📊 Ожидаемая схема: %s", EXPECTED_ORDERS_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 соответствует ожиданиям") logging.info("✅ Схема таблицы orders соответствует ожиданиям")
def fetch_orders_count(conn) -> int: def fetch_orders_count(conn) -> int:
""" """
Получает количество строк в таблице orders. Получает количество строк в таблице orders.
Args: Args:
conn: Подключение к Greenplum conn: Подключение к Greenplum
Returns: Returns:
Количество строк в таблице Количество строк в таблице
""" """
@@ -140,29 +150,31 @@ def fetch_orders_count(conn) -> int:
def assert_orders_have_rows(conn) -> None: def assert_orders_have_rows(conn) -> None:
""" """
Проверяет, что таблица orders не пустая. Проверяет, что таблица orders не пустая.
Args: Args:
conn: Подключение к Greenplum conn: Подключение к Greenplum
Raises: Raises:
ValueError: Если таблица пустая ValueError: Если таблица пустая
""" """
logging.info("📊 Проверяем наличие данных в таблице orders...") logging.info("📊 Проверяем наличие данных в таблице orders...")
row_count = fetch_orders_count(conn) row_count = fetch_orders_count(conn)
logging.info("📈 Количество строк в orders: %s", row_count) logging.info("📈 Количество строк в orders: %s", row_count)
if row_count <= 0: if row_count <= 0:
raise ValueError("❌ Таблица public.orders пустая — запусти DAG csv_to_greenplum перед проверкой.") raise ValueError(
"❌ Таблица public.orders пустая — запусти DAG csv_to_greenplum перед проверкой."
)
logging.info("✅ Таблица orders содержит данные (%s строк)", row_count) logging.info("✅ Таблица orders содержит данные (%s строк)", row_count)
def fetch_orders_duplicates(conn) -> int: def fetch_orders_duplicates(conn) -> int:
""" """
Подсчитывает количество дубликатов по order_id. Подсчитывает количество дубликатов по order_id.
Args: Args:
conn: Подключение к Greenplum conn: Подключение к Greenplum
Returns: Returns:
Количество дублирующихся order_id Количество дублирующихся order_id
""" """
@@ -183,17 +195,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: Args:
conn: Подключение к Greenplum conn: Подключение к Greenplum
Raises: Raises:
ValueError: Если обнаружены дубликаты ValueError: Если обнаружены дубликаты
""" """
logging.info("🔍 Проверяем отсутствие дубликатов по order_id...") logging.info("🔍 Проверяем отсутствие дубликатов по order_id...")
duplicates = fetch_orders_duplicates(conn) duplicates = fetch_orders_duplicates(conn)
logging.info("📊 Найдено дубликатов: %s", duplicates) logging.info("📊 Найдено дубликатов: %s", duplicates)
if duplicates: if duplicates:
raise ValueError(f"❌ Обнаружены дубли по order_id ({duplicates} шт.) — проверь загрузку данных.") raise ValueError(
f"❌ Обнаружены дубли по order_id ({duplicates} шт.) — проверь загрузку данных."
)
logging.info("✅ Дубликаты не обнаружены") logging.info("✅ Дубликаты не обнаружены")
-1
View File
@@ -6,7 +6,6 @@ from pathlib import Path
from types import ModuleType from types import ModuleType
from typing import Type from typing import Type
PROJECT_ROOT = Path(__file__).resolve().parents[1] PROJECT_ROOT = Path(__file__).resolve().parents[1]
if str(PROJECT_ROOT) not in sys.path: if str(PROJECT_ROOT) not in sys.path:
sys.path.append(str(PROJECT_ROOT)) sys.path.append(str(PROJECT_ROOT))
+3 -1
View File
@@ -14,7 +14,9 @@ def _airflow_available() -> bool:
return hasattr(af, "DAG") return hasattr(af, "DAG")
pytestmark = pytest.mark.skipif(not _airflow_available(), reason="Airflow is not installed for DAG smoke tests") pytestmark = pytest.mark.skipif(
not _airflow_available(), reason="Airflow is not installed for DAG smoke tests"
)
def _load_dag(module_name: str): def _load_dag(module_name: str):