refactor(csv): удален легаси CSV-пайплайн и связанные с ним файлы

- Зачем:
  - Пример базовой загрузки CSV перенесен в отдельный репозиторий `airflow-manual` для разделения учебных треков.
- Что:
  - удалены DAG-файлы `csv_to_greenplum` и вспомогательные скрипты `helpers/greenplum.py`, `orders_ddl.sql`.
  - из `docker-compose.yml` и `.env.example` удалены переменные и тома (`airflow_data`), необходимые для CSV.
  - очищена документация (`README.md`, `TESTING.md`, `educational-tasks.md`) и тесты (`test_dags_smoke.py`, `conftest.py`).
  - отмечен выполненным 'Этап 1' в `TODO.md`.
- Проверка:
  - `make test` проходит успешно (smoke-тесты оставшихся DAG-ов не затронуты).
This commit is contained in:
2026-03-08 19:13:19 +03:00
parent 45c1d37da5
commit a6a3559ab9
19 changed files with 21 additions and 866 deletions
-4
View File
@@ -43,7 +43,3 @@ GP_USE_AIRFLOW_CONN=true
PXF_SEED_OVERWRITE=0
# PXF_SYNC_ON_START=1 — выполнять `pxf cluster sync` при старте контейнера (дольше, но гарантирует актуальные конфиги)
PXF_SYNC_ON_START=0
# CSV Pipeline
CSV_DIR=/opt/airflow/data
CSV_ROWS=1000
+3 -3
View File
@@ -9,7 +9,7 @@
- **Фокус на «Почему»:** При использовании специфичных паттернов DWH (например, `delete + insert` для инкремента в Greenplum вместо `merge`) — добавляйте краткий комментарий, объясняющий этот выбор студентам.
## 2. Карта проекта (Навигация для Агента)
- `airflow/dags/` — DAG-файлы (напр. `csv_to_greenplum.py`).
- `airflow/dags/` — DAG-файлы (напр. `bookings_to_gp_stage.py`).
- `sql/` — DDL и SQL-скрипты. Разделены на слои: src/ (исходные системы), stg/ (стейджинг), ods/ (операционное хранилище), dds/ (детальное хранилище), dm/ (слой витрин).
- *Правило ИИ:* DDL таблиц хранится строго рядом с объектом (напр. `sql/stg/bookings_ddl.sql`).
- `docs/internal/naming_conventions.md` — Единый источник истины для нейминга служебных и SCD-полей. *Правило ИИ: Всегда сверяться с этим файлом при генерации новых DDL/SQL.*
@@ -54,7 +54,7 @@
## Тестирование
- Тесты лежат в `tests/` (pytest). Запуск: `make test`.
- Есть юнит‑тесты для `helpers/greenplum.py` и smoke‑тесты DAG‑структуры (`tests/test_dags_smoke.py`).
- Есть smoke‑тесты DAG‑структуры (`tests/test_dags_smoke.py`).
- Smoke‑тесты DAG автоматически пропускаются, если Airflow не установлен в venv.
- Для ручного прогона стенда см. `TESTING.md` (пошаговый чек‑лист для студентов).
- Для программной проверки DAG (без браузера) — см. `docs/agent-dag-testing.md`: CLI, REST API, проверка параллельности, запросы в Greenplum.
@@ -66,7 +66,7 @@
- При изменении схемы/поведения — обновляйте `README.md` и `sql/ddl_gp.sql`.
## Безопасность и конфигурация
- Все настройки — через `.env`; креды в коде не хардкодим. Частые переменные: `GP_*`, `PG_*`, `AIRFLOW_*`, `CSV_*`.
- Все настройки — через `.env`; креды в коде не хардкодим. Частые переменные: `GP_*`, `PG_*`, `AIRFLOW_*`.
- `make clean` удаляет тома — предупреждайте студентов, что данные пропадут.
## Для агента (особенности аудитории)
-7
View File
@@ -22,7 +22,6 @@
- **Greenplum** (singlenode для обучения; внешний порт по умолчанию `5435`)
- **bookings-db** (Postgres с демо‑БД `demo`; внешний порт по умолчанию `5434`)
- **PXF** как “транспорт” между Postgres и Greenplum (уже настроен в образе)
- Побочный пример: загрузка данных через **pandas/CSV** (`csv_to_greenplum`)
## Требования
@@ -104,12 +103,6 @@ SELECT COUNT(*) FROM dds.fact_flight_sales;
- `bookings_dds_ddl` — создаёт/обновляет DDS-таблицы (`dim_*`, `fact_flight_sales`) по домену bookings.
- `bookings_to_gp_dds` — загружает данные из ODS в DDS (SCD1/SCD2 + факт) и выполняет DQ‑проверки.
Вспомогательные (побочный трек с CSV):
- `orders_base_ddl` — создаёт таблицу `public.orders` для CSV‑пайплайна;
- `csv_to_greenplum` — pandas → CSV → Greenplum (пример загрузки без источника‑БД);
- `csv_to_greenplum_dq` — проверки качества данных для `public.orders`.
## Полезные команды
```bash
-13
View File
@@ -28,14 +28,6 @@
1. Открыть http://localhost:8080 (admin/admin).
2. (опционально) Зайти в Admin → Connections и убедиться, что DAG’и видят подключения:
- `greenplum_conn` и `bookings_db` задаются через переменные `AIRFLOW_CONN_...` в docker-compose и могут не отображаться в списке, но `airflow connections get greenplum_conn` / `bookings_db` внутри контейнера должны отрабатывать без ошибок.
3. DAG `csv_to_greenplum`:
- Включить переключатель.
- Нажать «Trigger DAG».
- Контроль: все таски Success, в `data/` появился CSV, в логах `load_csv_to_greenplum` видно `INSERT`.
- В Greenplum (см. п.5) убедиться в наличии строк `(SELECT COUNT(*) ...)`.
4. DAG `csv_to_greenplum_dq`:
- Запустить вручную после первого DAG.
- Проверить, что все 5 задач Success и логи содержат `Проверка пройдена`.
- DAG `bookings_to_gp_stage` (полная проверка цепочки bookings → Greenplum STG):
- предварительно выполнить один раз: `make bookings-init` (установка демобазы `demo` в контейнере `bookings-db`) и `make ddl-gp` (создаёт STG/ODS/DDS слои в Greenplum, включая внешние `*_ext` через PXF);
@@ -59,18 +51,13 @@
- `docker compose exec greenplum bash -lc "su - gpadmin -c '/usr/local/pxf/bin/pxf cluster status'"`
- Команды внутри psql:
- `\dt public.*` — таблицы схему public.
- `SELECT COUNT(*) FROM public.orders;` — оценка объёма.
- `SELECT * FROM public.orders LIMIT 5;` — визуальная проверка.
- `SELECT order_id FROM public.orders GROUP BY 1 HAVING COUNT(*) > 1;` — поиск дублей.
- (после настройки PXF) `SELECT COUNT(*) FROM public.ext_bookings_bookings;` — проверка чтения из демо-БД bookings через PXF.
- (после настройки PXF) `SELECT * FROM public.ext_bookings_bookings LIMIT 5;` — визуальное сравнение с таблицей `bookings.bookings` в исходной БД.
- Завершить `\q`.
## 6. Негативные сценарии и fallback
- **Пустая таблица**: запустить `csv_to_greenplum_dq` до `csv_to_greenplum`. Ожидается ошибка на таске `check_orders_has_rows`.
- **Проблемы с подключением**: временно изменить `GP_HOST` или `GP_PORT` на несуществующий, перезапустить `make up`, убедиться, что DAG падает с понятной ошибкой (`psycopg2.OperationalError`).
- **Fallback без Airflow Connection**: установить `GP_USE_AIRFLOW_CONN=false`, перезапустить стек (`make down && make up`), удостовериться, что загрузка и DQ работают через ENV.
- **Дубликаты**: дважды вызвать `csv_to_greenplum` — ожидаем, что количество строк в `public.orders` не увеличится на размер CSV, а DAG `csv_to_greenplum_dq` не найдёт дублей.
- **PXF и демобаза bookings** (после настройки PXF и выполнения `make ddl-gp`): временно остановить `bookings-db` (`docker compose stop bookings-db`) и попробовать выполнить `SELECT COUNT(*) FROM public.ext_bookings_bookings;` в `make gp-psql` — ожидается ошибка подключения. Затем запустить `bookings-db` (`docker compose start bookings-db`) и убедиться, что запрос снова работает.
## 7. Быстрый reset (если «что-то сломалось»)
+3 -3
View File
@@ -16,11 +16,11 @@
**Инструмент:** Sonnet / Gemini / ChatGPT — механическая работа, перенос файлов.
- [ ] Перенести в [airflow-manual](https://github.com/dementev-dev/airflow-manual):
- [x] Перенести в [airflow-manual](https://github.com/dementev-dev/airflow-manual):
`csv_to_greenplum.py`, `csv_to_greenplum_dq.py`, `ddl_greenplum_base.py`,
`helpers/greenplum.py`, `sql/base/orders_ddl.sql`, связанные тесты
- [ ] Убрать CSV-зависимости из docker-compose / .env (`CSV_DIR`, `CSV_ROWS`)
- [ ] Обновить README (убрать упоминания CSV-пайплайна)
- [x] Убрать CSV-зависимости из docker-compose / .env (`CSV_DIR`, `CSV_ROWS`)
- [x] Обновить README (убрать упоминания CSV-пайплайна)
### Этап 1.5. Полировка эталона
-167
View File
@@ -1,167 +0,0 @@
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
-97
View File
@@ -1,97 +0,0 @@
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
@@ -1,32 +0,0 @@
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
@@ -1,223 +0,0 @@
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("✅ Дубликаты не обнаружены")
+1 -6
View File
@@ -11,7 +11,6 @@ x-airflow-common-env: &airflow-env
x-airflow-common-volumes: &airflow-volumes
- ./airflow/dags:/opt/airflow/dags
- ./sql:/sql:ro
- airflow_data:/opt/airflow/data
- airflow_logs:/opt/airflow/logs
x-airflow-common-depends: &airflow-depends
@@ -110,7 +109,6 @@ services:
volumes:
- ./airflow/dags:/opt/airflow/dags
- ./sql:/sql:ro
- airflow_data:/opt/airflow/data
- airflow_logs:/opt/airflow/logs
- ./airflow/requirements.txt:/opt/airflow/requirements.txt
healthcheck:
@@ -135,7 +133,6 @@ services:
volumes:
- ./airflow/dags:/opt/airflow/dags
- ./sql:/sql:ro
- airflow_data:/opt/airflow/data
- airflow_logs:/opt/airflow/logs
- ./airflow/requirements.txt:/opt/airflow/requirements.txt
depends_on:
@@ -154,12 +151,11 @@ services:
volumes:
- ./airflow/dags:/opt/airflow/dags
- ./sql:/sql:ro
- airflow_data:/opt/airflow/data
- airflow_logs:/opt/airflow/logs
command: >
bash -lc "
set -e;
mkdir -p /opt/airflow/data /opt/airflow/logs && chown -R airflow:root /opt/airflow/data /opt/airflow/logs;
mkdir -p /opt/airflow/logs && chown -R airflow:root /opt/airflow/logs;
# Дожидаемся готовности БД ретрая миграции
for i in {1..30}; do
su -s /bin/bash airflow -c \"PATH='/home/airflow/.local/bin:$${PATH}' airflow db migrate\" && break || echo 'waiting for pgmeta' && sleep 3;
@@ -175,5 +171,4 @@ volumes:
pgmeta:
bookings_data:
greenplum_data:
airflow_data:
airflow_logs:
+1 -3
View File
@@ -91,9 +91,7 @@ GP-специфичная best practice, которую забывают даж
- `_resolve_stg_batch_id` с INTERSECT по 4 таблицам — нет комментария **зачем** нужна согласованность
- **Нужно**: комментарий в `airflow/dags/bookings_to_gp_ods.py` перед SQL-запросом
- [x] **`helpers/greenplum.py` без пометки «legacy»**
- Использует прямой psycopg2 + ENV — противоречит PostgresOperator-подходу
- **Нужно**: docstring «LEGACY: только для CSV-пайплайна» в `airflow/dags/helpers/greenplum.py`
- [x] **`helpers/greenplum.py`** — удалён вместе с CSV-пайплайном (перенесён в airflow-manual)
### P2: Средние усилия, заметное улучшение качества
+1 -1
View File
@@ -54,7 +54,7 @@ make gp-psql
PXF (Platform Extension Framework) — компонент Greenplum для работы с внешними источниками.
В этом стенде PXF используется для чтения таблицы `bookings.bookings` из Postgres прямо из Greenplum
через внешнюю таблицу `stg.bookings_ext`. Поэтому загрузка в `stg.bookings` выглядит как обычный
`INSERT ... SELECT` без промежуточных CSV.
`INSERT ... SELECT` без промежуточных файлов.
## Greenplum + PXF: свой образ
+10 -55
View File
@@ -1,63 +1,18 @@
# Учебные задания по стенду
Этот документ собирает в одном месте задания для менти.
Он разбит на блоки: от базовой работы с CSV‑pipeline до более продвинутого сценария с демо‑БД bookings и слоем STG в Greenplum.
Он разбит на блоки: от архитектуры Greenplum и демо‑БД bookings до реализации аналитических слоев DWH.
Если вы только начинаете, выполняйте задания по порядку. К разделу про bookings можно вернуться позже.
Если вы только начинаете, выполняйте задания по порядку.
---
## 1. Базовый CSVpipeline (csv_to_greenplum)
Основная цель этого блока — понять, как устроен простой ETL: генерация данных через pandas, сохранение в CSV и загрузка в Greenplum.
### 1.1. Разбор готового pipeline
1. Найдите DAG `csv_to_greenplum` в `airflow/dags/csv_to_greenplum.py`.
2. Ответьте себе на вопросы (можно коротко в отдельном файле/блокноте):
- какие задачи (tasks) входят в DAG и что делает каждая из них;
- какие таблицы создаются в Greenplum;
- где физически лежат CSV‑файлы;
- какие параметры управляют размером датасета.
3. Поднимите стенд и запустите DAG:
- `make up` (Airflow инициализируется автоматически при первом старте)
- включите и запустите DAG `csv_to_greenplum` в Airflow UI.
4. Проверьте результат в Greenplum:
- `make gp-psql`
- `SELECT COUNT(*) FROM public.orders;`
- `SELECT * FROM public.orders LIMIT 5;`
### 1.2. Изменение параметров генерации
1. Найдите, где задаётся количество строк для генерации (`CSV_ROWS` в `.env` и параметр в DAG).
2. Поставьте другое значение и перезапустите DAG:
- оцените, как изменилось количество строк в `public.orders`;
- убедитесь, что пайплайн по‑прежнему работает без ошибок.
3. Попробуйте изменить схему данных (добавить колонку в CSV и таблицу в Greenplum):
- добавьте новую колонку в генерацию pandas;
- обновите DDL/SQL, чтобы колонка появилась в таблице `public.orders`;
- перезапустите DAG и убедитесь, что новая колонка заполняется.
### 1.3. Собственные проверки качества данных
1. Найдите DAG `csv_to_greenplum_dq` в `airflow/dags/csv_to_greenplum_dq.py`.
2. Посмотрите, какие проверки уже реализованы (наличие таблицы, схема, дубликаты).
3. Добавьте ещё одну простую проверку, например:
- проверка, что в таблице `public.orders` не больше N строк;
- проверка, что поле (например, `order_price`) не содержит отрицательных значений;
- проверка, что нет строк с `NULL` в ключевых колонках.
4. Запустите DAG `csv_to_greenplum_dq` и убедитесь, что:
- новая проверка проходит на «хороших» данных;
- при нарушении условия DAG падает с понятной ошибкой.
---
## 2. Greenplum и модель данных (введение)
## 1. Greenplum и модель данных (введение)
В следующих заданиях мы будем опираться на демо‑БД bookings (Postgres) и слой STG в Greenplum.
На этом этапе достаточно бегло посмотреть на структуру и понять общую идею, детальная проработка пойдёт позже.
### 2.1. Знакомство с демо‑БД bookings
### 1.1. Знакомство с демо‑БД bookings
1. Прочитайте `bookings/README.md` — какие сервисы и команды относятся к демобазе.
2. Поднимите стенд и выполните:
@@ -70,7 +25,7 @@
- какие типы колонок используются;
- какие поля выглядят как ключи, даты, суммы.
### 2.2. Знакомство с STG в Greenplum
### 1.2. Знакомство с STG в Greenplum
1. Прочитайте `sql/stg/bookings_ddl.sql` и краткое описание потока `docs/bookings_to_gp_stage.md` (если интересно — `docs/internal/bookings_stg_design.md`).
2. Ответьте себе на вопросы:
@@ -81,7 +36,7 @@
- `\dn` и `\dt stg.*`
- `SELECT * FROM stg.bookings LIMIT 5;` (после запуска соответствующего DAG).
### 2.3. Как генерируются учебные данные bookings
### 1.3. Как генерируются учебные данные bookings
1. Откройте файл `bookings/generate_next_day.sql` и ответьте себе на вопросы:
- с какой даты начинается генерация данных (посмотрите на GUC `bookings.start_date` и переменную `v_start_cfg`);
@@ -99,12 +54,12 @@
---
## 3. DAG bookings_to_gp_stage (заготовка заданий)
## 2. DAG bookings_to_gp_stage (заготовка заданий)
Этот DAG показывает путь данных от демо‑БД bookings в Postgres до сырого слоя STG в Greenplum.
Сейчас он уже реализован как учебный пример, а в будущем вокруг него появятся отдельные задания по моделированию DWH.
### 3.1. Что есть сейчас
### 2.1. Что есть сейчас
1. Откройте `airflow/dags/bookings_to_gp_stage.py`.
2. Найдите в коде ссылки на SQL‑файлы:
@@ -121,9 +76,9 @@
На этом этапе достаточно понять общую цепочку. Детальные задания по переработке модели данных и построению ODS/DDS/DM слоёв будут добавлены позже.
### 3.2. Идеи для будущих заданий (черновик)
### 2.2. Идеи для будущих заданий (черновик)
> Ниже — набросок задач, к которым мы вернёмся, когда базовые темы по Airflow и CSVpipeline будут освоены.
> Ниже — набросок задач, к которым мы вернёмся, когда базовые темы по Airflow будут освоены.
Планируемые направления:
+1 -1
View File
@@ -284,7 +284,7 @@ airflow-init:
- Проверить наличие DAG'ов в списке
2. **Проверить запуск DAG'ов**
- Выбрать любой DAG (например, `csv_to_greenplum`)
- Выбрать любой DAG (например, `bookings_to_gp_stage`)
- Запустить вручную через UI
- Проверить успешность выполнения задач
-11
View File
@@ -1,11 +0,0 @@
-- DDL для базовой таблицы orders, которую использует CSV‑pipeline.
-- Выполняется идемпотентно: таблица создаётся, если ещё не существует.
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);
+1 -4
View File
@@ -1,6 +1,6 @@
-- Главный входной DDL-скрипт для Greenplum в учебном стенде.
-- Выполняется из контейнера командой `make ddl-gp` и создаёт/обновляет
-- все объекты, которые нужны базовым DAG (csv_to_greenplum, bookings_to_gp_stage);
-- все объекты, которые нужны DAG'ам bookings ETL-пайплайна;
-- подключает файловые DDL через \i, чтобы сохранять единый входной скрипт.
--
-- Чтобы не ломать задания, новые объекты лучше добавлять в отдельные файлы
@@ -8,9 +8,6 @@
--
-- Подробнее про STG/bookings: см. docs/bookings_to_gp_stage.md.
-- Таблица для CSV‑пайплайна (csv_to_greenplum).
\i base/orders_ddl.sql
-- Внешняя таблица для чтения данных из демо-БД bookings через PXF (JDBC).
-- Источник: таблица bookings.bookings в базе demo (Postgres, сервис bookings-db).
DROP EXTERNAL TABLE IF EXISTS public.ext_bookings_bookings;
-11
View File
@@ -55,14 +55,3 @@ def _ensure_stub_module(full_name: str) -> ModuleType:
module = sys.modules[path]
assert isinstance(module, ModuleType)
return module
def patch_postgres_hook(monkeypatch, hook_cls: Type) -> None:
"""
Patch PostgresHook so that helpers.greenplum can be exercised without real Airflow.
"""
try:
module = importlib.import_module("airflow.providers.postgres.hooks.postgres")
except ModuleNotFoundError:
module = _ensure_stub_module("airflow.providers.postgres.hooks.postgres")
monkeypatch.setattr(module, "PostgresHook", hook_cls, raising=False)
-47
View File
@@ -41,53 +41,6 @@ def _assert_reachable(dag, upstream_task_id: str, downstream_task_id: str) -> No
), f"Expected {downstream_task_id} to be downstream of {upstream_task_id}"
def test_csv_to_greenplum_dag_structure():
dag = _load_dag("airflow.dags.csv_to_greenplum")
# tasks
expected_tasks = {
"create_orders_table",
"generate_csv",
"preview_csv",
"load_csv_to_greenplum",
}
assert expected_tasks.issubset(dag.task_dict.keys())
# linear dependencies
t1 = dag.get_task("create_orders_table")
t2 = dag.get_task("generate_csv")
t3 = dag.get_task("preview_csv")
t4 = dag.get_task("load_csv_to_greenplum")
assert t2 in t1.get_direct_relatives(upstream=False)
assert t3 in t2.get_direct_relatives(upstream=False)
assert t4 in t3.get_direct_relatives(upstream=False)
def test_csv_to_greenplum_dq_dag_structure():
dag = _load_dag("airflow.dags.csv_to_greenplum_dq")
expected_tasks = {
"check_orders_table_exists",
"check_orders_schema",
"check_orders_has_rows",
"check_order_duplicates",
"data_quality_summary",
}
assert expected_tasks.issubset(dag.task_dict.keys())
e = dag.get_task("check_orders_table_exists")
s = dag.get_task("check_orders_schema")
h = dag.get_task("check_orders_has_rows")
d = dag.get_task("check_order_duplicates")
q = dag.get_task("data_quality_summary")
assert s in e.get_direct_relatives(upstream=False)
assert h in s.get_direct_relatives(upstream=False)
assert d in h.get_direct_relatives(upstream=False)
assert q in d.get_direct_relatives(upstream=False)
def test_bookings_stg_ddl_dag_structure():
"""Проверка структуры DAG bookings_stg_ddl."""
dag = _load_dag("airflow.dags.bookings_stg_ddl")
-178
View File
@@ -1,178 +0,0 @@
from __future__ import annotations
from dataclasses import dataclass
from typing import Any, List, Sequence
import pytest
import airflow.dags.helpers.greenplum as greenplum
from tests.conftest import patch_postgres_hook
@dataclass
class FakeCursor:
fetchone_value: Any = None
fetchall_value: Sequence[Any] | None = None
rowcount: int | None = None
def __post_init__(self) -> None:
self.queries: List[Any] = []
def execute(self, query: str, params: Any | None = None) -> None:
self.queries.append((query, params))
def fetchone(self) -> Any:
return self.fetchone_value
def fetchall(self) -> Sequence[Any] | None:
return self.fetchall_value
def __enter__(self) -> FakeCursor:
return self
def __exit__(self, exc_type, exc, tb) -> None:
return None
class FakeConn:
def __init__(self, cursors: Sequence[FakeCursor]) -> None:
self._cursors = list(cursors)
self._index = 0
self.commits = 0
def cursor(self) -> FakeCursor:
cursor = self._cursors[self._index]
self._index += 1
return cursor
def commit(self) -> None:
self.commits += 1
def test_get_gp_conn_uses_airflow_hook(monkeypatch) -> None:
class FakeHook:
def __init__(self, postgres_conn_id: str) -> None:
self.postgres_conn_id = postgres_conn_id
def get_conn(self) -> str:
return "hook_connection"
patch_postgres_hook(monkeypatch, FakeHook)
monkeypatch.setattr(greenplum, "GP_CONN_ID", "demo_conn", raising=False)
monkeypatch.setattr(greenplum, "GP_USE_AIRFLOW_CONN", True, raising=False)
conn = greenplum.get_gp_conn()
assert conn == "hook_connection"
def test_get_gp_conn_fallback_to_psycopg(monkeypatch) -> None:
class BrokenHook:
def __init__(self, postgres_conn_id: str) -> None:
self.postgres_conn_id = postgres_conn_id
def get_conn(self):
raise RuntimeError("boom")
patch_postgres_hook(monkeypatch, BrokenHook)
monkeypatch.setattr(greenplum, "GP_USE_AIRFLOW_CONN", True, raising=False)
monkeypatch.setattr(greenplum, "GP_CONN_ID", "demo_conn", raising=False)
monkeypatch.setenv("GP_DB", "demo_db")
monkeypatch.setenv("GP_USER", "demo_user")
monkeypatch.setenv("GP_PASSWORD", "secret")
monkeypatch.setenv("GP_HOST", "greenplum-host")
monkeypatch.setenv("GP_PORT", "5434")
captured_kwargs = {}
def fake_connect(**kwargs):
captured_kwargs.update(kwargs)
return "psycopg_connection"
monkeypatch.setattr(greenplum.psycopg2, "connect", fake_connect)
conn = greenplum.get_gp_conn()
assert conn == "psycopg_connection"
assert captured_kwargs == {
"dbname": "demo_db",
"user": "demo_user",
"password": "secret",
"host": "greenplum-host",
"port": 5434,
}
def test_get_gp_conn_without_airflow(monkeypatch) -> None:
monkeypatch.setattr(greenplum, "GP_USE_AIRFLOW_CONN", False, raising=False)
monkeypatch.setenv("GP_DB", "demo_db")
monkeypatch.setenv("GP_USER", "demo_user")
monkeypatch.setenv("GP_PASSWORD", "secret")
monkeypatch.setenv("GP_HOST", "greenplum-host")
monkeypatch.setenv("GP_PORT", "5435")
captured_kwargs = {}
def fake_connect(**kwargs):
captured_kwargs.update(kwargs)
return "direct_psycopg"
monkeypatch.setattr(greenplum.psycopg2, "connect", fake_connect)
conn = greenplum.get_gp_conn()
assert conn == "direct_psycopg"
assert captured_kwargs["port"] == 5435
def test_assert_orders_table_exists_ok() -> None:
conn = FakeConn([FakeCursor(fetchone_value=(1,))])
greenplum.assert_orders_table_exists(conn)
def test_assert_orders_table_exists_missing() -> None:
conn = FakeConn([FakeCursor(fetchone_value=None)])
with pytest.raises(ValueError):
greenplum.assert_orders_table_exists(conn)
def test_assert_orders_schema_ok() -> None:
expected = list(greenplum.EXPECTED_ORDERS_SCHEMA)
conn = FakeConn([FakeCursor(fetchall_value=expected)])
greenplum.assert_orders_schema(conn)
def test_assert_orders_schema_mismatch() -> None:
conn = FakeConn([FakeCursor(fetchall_value=[("order_id", "bigint")])])
with pytest.raises(ValueError):
greenplum.assert_orders_schema(conn)
def test_assert_orders_have_rows_ok() -> None:
conn = FakeConn([FakeCursor(fetchone_value=(5,))])
greenplum.assert_orders_have_rows(conn)
def test_assert_orders_have_rows_empty() -> None:
conn = FakeConn([FakeCursor(fetchone_value=(0,))])
with pytest.raises(ValueError):
greenplum.assert_orders_have_rows(conn)
def test_assert_orders_no_duplicates_ok() -> None:
conn = FakeConn([FakeCursor(fetchone_value=(0,))])
greenplum.assert_orders_no_duplicates(conn)
def test_assert_orders_no_duplicates_detected() -> None:
conn = FakeConn([FakeCursor(fetchone_value=(3,))])
with pytest.raises(ValueError):
greenplum.assert_orders_no_duplicates(conn)