diff --git a/airflow/dags/bookings_validate.py b/airflow/dags/bookings_validate.py new file mode 100644 index 0000000..7bca8d0 --- /dev/null +++ b/airflow/dags/bookings_validate.py @@ -0,0 +1,150 @@ +""" +Валидационный DAG: самопроверка студенческих заданий. + +Запускается вручную в Airflow UI после реализации заданий. +Таски сгруппированы по слоям — студент видит, где именно проблема. + +Структура: + validate_ods — проверки ODS: rowcount, дубли BK, NULL в PK + validate_dds — проверки DDS: SCD1-измерения, SCD2 dim_routes (активный тест!) + validate_dm — проверки DM-витрин: не пусты, нет NULL в ключах + +Все три группы запускаются параллельно — инкрементальная обратная связь: +студент может проверить только ODS, пока DDS/DM ещё не реализованы. +""" + +from datetime import timedelta +from logging import getLogger + +import pendulum +from airflow.providers.postgres.operators.postgres import PostgresOperator +from airflow.utils.task_group import TaskGroup + +from airflow import DAG + +GREENPLUM_CONN_ID = "greenplum_conn" +log = getLogger(__name__) + +default_args = { + "owner": "airflow", + "retries": 0, # Без ретраев — студент должен увидеть ошибку сразу + "retry_delay": timedelta(seconds=10), +} + +with DAG( + dag_id="bookings_validate", + start_date=pendulum.datetime(2017, 1, 1, tz="UTC"), + schedule=None, # Только ручной запуск + catchup=False, + max_active_runs=1, + template_searchpath="/sql", + default_args=default_args, + tags=["demo", "bookings", "greenplum", "validate"], + description="Валидация студенческих заданий: ODS, DDS, DM", +) as dag: + + with TaskGroup("validate_ods") as validate_ods: + check_ods_airplanes = PostgresOperator( + task_id="check_ods_airplanes_rowcount", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="validate/ods_airplanes_rowcount.sql", + ) + check_ods_seats = PostgresOperator( + task_id="check_ods_seats_rowcount", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="validate/ods_seats_rowcount.sql", + ) + check_ods_dup_bk = PostgresOperator( + task_id="check_ods_no_dup_bk", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="validate/ods_no_dup_bk.sql", + ) + check_ods_pks = PostgresOperator( + task_id="check_ods_no_null_pks", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="validate/ods_no_null_pks.sql", + ) + + with TaskGroup("validate_dds") as validate_dds: + check_dim_airplanes = PostgresOperator( + task_id="check_dim_airplanes_exists", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="validate/dim_airplanes_exists.sql", + ) + check_dim_passengers = PostgresOperator( + task_id="check_dim_passengers_exists", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="validate/dim_passengers_exists.sql", + ) + check_dim_passengers_dup = PostgresOperator( + task_id="check_dim_passengers_no_dup_bk", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="validate/dim_passengers_no_dup_bk.sql", + ) + check_dim_routes = PostgresOperator( + task_id="check_dim_routes_exists", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="validate/dim_routes_exists.sql", + ) + + # SCD2 активный тест: цепочка backup → mutate → load → check → restore + scd2_backup = PostgresOperator( + task_id="scd2_backup", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="validate/dim_routes_scd2_backup.sql", + ) + scd2_mutate = PostgresOperator( + task_id="scd2_mutate", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="validate/dim_routes_scd2_mutate.sql", + ) + scd2_run_load = PostgresOperator( + task_id="scd2_run_student_load", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="dds/dim_routes_load.sql", # студенческий load-скрипт! + ) + scd2_check = PostgresOperator( + task_id="scd2_check", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="validate/dim_routes_scd2_check.sql", + ) + scd2_restore = PostgresOperator( + task_id="scd2_restore", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="validate/dim_routes_scd2_restore.sql", + trigger_rule="all_done", # Откат ВСЕГДА, даже если check упал + ) + + check_dim_routes_gaps = PostgresOperator( + task_id="check_dim_routes_no_gaps", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="validate/dim_routes_no_gaps.sql", + ) + + # SCD2 цепочка + scd2_backup >> scd2_mutate >> scd2_run_load >> scd2_check >> scd2_restore + + # no_gaps запускается после restore (на чистых данных) + scd2_restore >> check_dim_routes_gaps + + with TaskGroup("validate_dm") as validate_dm: + check_airport_traffic = PostgresOperator( + task_id="check_airport_traffic_exists", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="validate/airport_traffic_exists.sql", + ) + check_route_performance = PostgresOperator( + task_id="check_route_performance_exists", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="validate/route_performance_exists.sql", + ) + check_monthly_overview = PostgresOperator( + task_id="check_monthly_overview_exists", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="validate/monthly_overview_exists.sql", + ) + check_passenger_loyalty = PostgresOperator( + task_id="check_passenger_loyalty_exists", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="validate/passenger_loyalty_exists.sql", + ) diff --git a/sql/validate/airport_traffic_exists.sql b/sql/validate/airport_traffic_exists.sql new file mode 100644 index 0000000..323ab51 --- /dev/null +++ b/sql/validate/airport_traffic_exists.sql @@ -0,0 +1,23 @@ +-- Проверка dm.airport_traffic: +-- 1. Таблица не пуста +-- 2. Нет NULL в ключевых полях (traffic_date, airport_sk) +DO $$ +DECLARE + v_count BIGINT; + v_null_pk BIGINT; +BEGIN + SELECT COUNT(*) INTO v_count FROM dm.airport_traffic; + IF v_count = 0 THEN + RAISE EXCEPTION 'FAILED: dm.airport_traffic пуста. Реализуйте загрузку: sql/dm/airport_traffic_load.sql'; + END IF; + + SELECT COUNT(*) INTO v_null_pk + FROM dm.airport_traffic + WHERE traffic_date IS NULL OR airport_sk IS NULL; + + IF v_null_pk > 0 THEN + RAISE EXCEPTION 'FAILED: dm.airport_traffic содержит % строк с NULL в ключе (traffic_date, airport_sk).', v_null_pk; + END IF; + + RAISE NOTICE 'PASSED: dm.airport_traffic содержит % строк', v_count; +END $$; diff --git a/sql/validate/dim_airplanes_exists.sql b/sql/validate/dim_airplanes_exists.sql new file mode 100644 index 0000000..912bc10 --- /dev/null +++ b/sql/validate/dim_airplanes_exists.sql @@ -0,0 +1,41 @@ +-- Проверка dds.dim_airplanes: +-- 1. Таблица не пуста +-- 2. Нет дублей по airplane_bk (SCD1 — UPSERT должен это гарантировать) +-- 3. Покрытие ODS: все airplane_code из ods.airplanes есть в измерении +-- 4. total_seats заполнен (вычисляется агрегацией из ods.seats; NULL = потерян JOIN) +DO $$ +DECLARE + v_count BIGINT; + v_dup BIGINT; + v_missing BIGINT; + v_null_seats BIGINT; +BEGIN + SELECT COUNT(*) INTO v_count FROM dds.dim_airplanes; + IF v_count = 0 THEN + RAISE EXCEPTION 'FAILED: dds.dim_airplanes пуста. Реализуйте загрузку: sql/dds/dim_airplanes_load.sql'; + END IF; + + SELECT COUNT(*) - COUNT(DISTINCT airplane_bk) INTO v_dup FROM dds.dim_airplanes; + IF v_dup > 0 THEN + RAISE EXCEPTION 'FAILED: dds.dim_airplanes содержит % дублей по airplane_bk. Проверьте SCD1-логику (UPSERT).', v_dup; + END IF; + + SELECT COUNT(*) INTO v_missing + FROM (SELECT DISTINCT airplane_code FROM ods.airplanes) AS s + WHERE NOT EXISTS ( + SELECT 1 FROM dds.dim_airplanes AS d WHERE d.airplane_bk = s.airplane_code + ); + IF v_missing > 0 THEN + RAISE EXCEPTION 'FAILED: % самолётов из ods.airplanes отсутствуют в dds.dim_airplanes.', v_missing; + END IF; + + -- total_seats вычисляется через JOIN с ods.seats; если JOIN забыт — будет NULL + SELECT COUNT(*) INTO v_null_seats + FROM dds.dim_airplanes WHERE total_seats IS NULL; + IF v_null_seats > 0 THEN + RAISE EXCEPTION E'FAILED: % строк в dds.dim_airplanes имеют NULL в total_seats.\n' + 'Подсказка: total_seats считается агрегацией из ods.seats — проверьте JOIN в dim_airplanes_load.sql.', v_null_seats; + END IF; + + RAISE NOTICE 'PASSED: dds.dim_airplanes содержит % строк, покрывает все BK из ODS, total_seats заполнен', v_count; +END $$; diff --git a/sql/validate/dim_passengers_exists.sql b/sql/validate/dim_passengers_exists.sql new file mode 100644 index 0000000..2552311 --- /dev/null +++ b/sql/validate/dim_passengers_exists.sql @@ -0,0 +1,33 @@ +-- Проверка dds.dim_passengers: +-- 1. Таблица не пуста +-- 2. Нет дублей по passenger_id (SCD1 — типичная ошибка: INSERT без EXISTS) +-- 3. Покрытие ODS: все уникальные passenger_id из ods.tickets есть в измерении +DO $$ +DECLARE + v_count BIGINT; + v_dup BIGINT; + v_missing BIGINT; +BEGIN + SELECT COUNT(*) INTO v_count FROM dds.dim_passengers; + IF v_count = 0 THEN + RAISE EXCEPTION 'FAILED: dds.dim_passengers пуста. Реализуйте загрузку: sql/dds/dim_passengers_load.sql'; + END IF; + + SELECT COUNT(*) - COUNT(DISTINCT passenger_id) INTO v_dup FROM dds.dim_passengers; + IF v_dup > 0 THEN + RAISE EXCEPTION E'FAILED: dds.dim_passengers содержит % дублей по passenger_id.\n' + 'Подсказка: INSERT без проверки EXISTS создаёт дубли при повторном запуске.\n' + 'Используйте INSERT ... WHERE NOT EXISTS или ON CONFLICT DO UPDATE.', v_dup; + END IF; + + SELECT COUNT(*) INTO v_missing + FROM (SELECT DISTINCT passenger_id FROM ods.tickets) AS s + WHERE NOT EXISTS ( + SELECT 1 FROM dds.dim_passengers AS d WHERE d.passenger_id = s.passenger_id + ); + IF v_missing > 0 THEN + RAISE EXCEPTION 'FAILED: % пассажиров из ods.tickets отсутствуют в dds.dim_passengers.', v_missing; + END IF; + + RAISE NOTICE 'PASSED: dds.dim_passengers содержит % строк, нет дублей, покрывает все BK из ODS', v_count; +END $$; diff --git a/sql/validate/dim_passengers_no_dup_bk.sql b/sql/validate/dim_passengers_no_dup_bk.sql new file mode 100644 index 0000000..f0dc687 --- /dev/null +++ b/sql/validate/dim_passengers_no_dup_bk.sql @@ -0,0 +1,22 @@ +-- Проверка: нет дублей по passenger_id в dds.dim_passengers. +-- +-- Отдельный таск для точной диагностики — студент сразу видит причину проблемы. +-- Типичная ошибка: INSERT без проверки EXISTS при SCD1-загрузке. +DO $$ +DECLARE + v_dup BIGINT; +BEGIN + SELECT COUNT(*) INTO v_dup + FROM ( + SELECT passenger_id FROM dds.dim_passengers + GROUP BY passenger_id HAVING COUNT(*) > 1 + ) AS d; + + IF v_dup > 0 THEN + RAISE EXCEPTION E'FAILED: dds.dim_passengers содержит % дублирующихся passenger_id.\n' + 'Подсказка: при SCD1 нужно INSERT ... WHERE NOT EXISTS или ON CONFLICT DO NOTHING/UPDATE.\n' + 'Если таск check_dim_passengers_exists тоже упал — начните с него.', v_dup; + END IF; + + RAISE NOTICE 'PASSED: Нет дублей по passenger_id в dds.dim_passengers'; +END $$; diff --git a/sql/validate/dim_routes_exists.sql b/sql/validate/dim_routes_exists.sql new file mode 100644 index 0000000..ad2f201 --- /dev/null +++ b/sql/validate/dim_routes_exists.sql @@ -0,0 +1,48 @@ +-- Проверка dds.dim_routes: +-- 1. Таблица не пуста +-- 2. Покрытие ODS: у каждого route_no из ods.routes есть открытая текущая версия (valid_to IS NULL) +-- 3. SCD2-инвариант: не более одной текущей версии на route_bk (valid_to IS NULL) +-- +-- Почему покрытие проверяем через valid_to IS NULL, а не просто EXISTS: +-- закрытая версия (valid_to IS NOT NULL) без открытой означает «маршрут есть в ODS, +-- но в DDS только архив» — баг вида "старую версию закрыли, новую не вставили". +-- Такой маршрут теряется в point-in-time JOIN'ах и витринах. +DO $$ +DECLARE + v_count BIGINT; + v_missing BIGINT; + v_multi_current BIGINT; +BEGIN + SELECT COUNT(*) INTO v_count FROM dds.dim_routes; + IF v_count = 0 THEN + RAISE EXCEPTION 'FAILED: dds.dim_routes пуста. Реализуйте загрузку: sql/dds/dim_routes_load.sql'; + END IF; + + -- Покрытие ODS: у каждого маршрута должна быть открытая (текущая) версия + SELECT COUNT(*) INTO v_missing + FROM (SELECT DISTINCT route_no FROM ods.routes) AS s + WHERE NOT EXISTS ( + SELECT 1 FROM dds.dim_routes AS d + WHERE d.route_bk = s.route_no AND d.valid_to IS NULL + ); + IF v_missing > 0 THEN + RAISE EXCEPTION E'FAILED: % маршрутов из ods.routes не имеют текущей версии в dds.dim_routes (valid_to IS NULL).\n' + 'Подсказка: проверьте, что после UPDATE (закрытие старой версии) выполняется INSERT новой.', v_missing; + END IF; + + -- SCD2-инвариант: не более одной текущей версии на route_bk + SELECT COUNT(*) INTO v_multi_current + FROM ( + SELECT route_bk FROM dds.dim_routes + WHERE valid_to IS NULL + GROUP BY route_bk HAVING COUNT(*) > 1 + ) AS d; + + IF v_multi_current > 0 THEN + RAISE EXCEPTION E'FAILED: % маршрутов имеют более одной текущей версии (valid_to IS NULL).\n' + 'Подсказка: при INSERT новой версии проверяйте NOT EXISTS ... AND hashdiff = ...\n' + 'чтобы не создавать дубликат, если hashdiff не изменился.', v_multi_current; + END IF; + + RAISE NOTICE 'PASSED: dds.dim_routes содержит % строк, покрывает ODS, по одной текущей версии на маршрут', v_count; +END $$; diff --git a/sql/validate/dim_routes_no_gaps.sql b/sql/validate/dim_routes_no_gaps.sql new file mode 100644 index 0000000..0a6cbc1 --- /dev/null +++ b/sql/validate/dim_routes_no_gaps.sql @@ -0,0 +1,28 @@ +-- Проверка: нет «дыр» в SCD2-интервалах dim_routes. +-- +-- Для каждого route_bk с несколькими версиями проверяем, что +-- valid_to предыдущей версии = valid_from следующей. +-- +-- Важно: если маршрут «исчез» из ODS, его текущая версия закрывается +-- (valid_to = CURRENT_DATE), но новая не вставляется — это корректно. +-- Проверяем «дыру» только если следующая версия существует. +DO $$ +DECLARE + v_gaps BIGINT; +BEGIN + SELECT COUNT(*) INTO v_gaps + FROM ( + SELECT route_bk, valid_to, + LEAD(valid_from) OVER (PARTITION BY route_bk ORDER BY valid_from) AS next_valid_from + FROM dds.dim_routes + ) AS t + WHERE valid_to IS NOT NULL + AND next_valid_from IS NOT NULL -- следующая версия существует (не «исчезнувший» маршрут) + AND valid_to <> next_valid_from; + + IF v_gaps > 0 THEN + RAISE EXCEPTION 'FAILED: В dds.dim_routes найдено % «дыр» между версиями SCD2. valid_to старой версии должен совпадать с valid_from новой.', v_gaps; + END IF; + + RAISE NOTICE 'PASSED: Нет «дыр» в SCD2-версиях dim_routes'; +END $$; diff --git a/sql/validate/dim_routes_scd2_backup.sql b/sql/validate/dim_routes_scd2_backup.sql new file mode 100644 index 0000000..3daf917 --- /dev/null +++ b/sql/validate/dim_routes_scd2_backup.sql @@ -0,0 +1,12 @@ +-- Бэкап текущего состояния перед активным тестом SCD2. +-- Используем обычные таблицы (не TEMP) — между тасками Airflow +-- TEMP-таблицы не сохраняются (каждый таск = отдельная транзакция). + +DROP TABLE IF EXISTS _validate_bk_ods_routes; +CREATE TABLE _validate_bk_ods_routes AS SELECT * FROM ods.routes; + +DROP TABLE IF EXISTS _validate_bk_dim_routes; +CREATE TABLE _validate_bk_dim_routes AS SELECT * FROM dds.dim_routes; + +-- Cleanup служебной таблицы от предыдущего запуска (на случай если restore не доехал) +DROP TABLE IF EXISTS _validate_scd2_target; diff --git a/sql/validate/dim_routes_scd2_check.sql b/sql/validate/dim_routes_scd2_check.sql new file mode 100644 index 0000000..3b7d02d --- /dev/null +++ b/sql/validate/dim_routes_scd2_check.sql @@ -0,0 +1,95 @@ +-- Проверяем, что SCD2-логика студента сработала корректно после мутации. +-- Читаем route_no тестового маршрута из служебной таблицы _validate_scd2_target +-- (создана на шаге mutate), а не угадываем по побочным эффектам. +DO $$ +DECLARE + v_route TEXT; + v_version_count BIGINT; + v_closed_count BIGINT; + v_open_count BIGINT; + v_old_hash TEXT; + v_new_hash TEXT; + v_gap_count BIGINT; +BEGIN + -- Читаем тестовый маршрут из служебной таблицы + SELECT route_no INTO v_route FROM _validate_scd2_target LIMIT 1; + + IF v_route IS NULL THEN + RAISE EXCEPTION 'FAILED: служебная таблица _validate_scd2_target пуста. Шаг mutate не выполнился?'; + END IF; + + -- Если в dim_routes нет маршрута — load не запустился или упал + IF NOT EXISTS (SELECT 1 FROM dds.dim_routes WHERE route_bk = v_route) THEN + RAISE EXCEPTION E'FAILED: dds.dim_routes не содержит маршрут % после запуска load.\n' + 'Проверьте sql/dds/dim_routes_load.sql.', v_route; + END IF; + + -- Проверка a: Ровно 2 версии тестового маршрута (было 1, стало 2 после мутации) + SELECT COUNT(*) INTO v_version_count + FROM dds.dim_routes WHERE route_bk = v_route; + + IF v_version_count < 2 THEN + RAISE EXCEPTION E'FAILED: Маршрут % — найдена % версия (ожидается 2: старая закрытая + новая открытая).\n' + 'SCD2 должен был создать новую версию после изменения departure_time.', v_route, v_version_count; + END IF; + + IF v_version_count > 2 THEN + RAISE EXCEPTION E'FAILED: Маршрут % — найдено % версий (ожидается 2).\n' + 'Возможно, load создаёт лишние дубликаты. Проверьте условие NOT EXISTS при INSERT.', v_route, v_version_count; + END IF; + + -- Проверка b: Старая версия закрыта (valid_to IS NOT NULL) + SELECT COUNT(*) INTO v_closed_count + FROM dds.dim_routes WHERE route_bk = v_route AND valid_to IS NOT NULL; + + IF v_closed_count = 0 THEN + RAISE EXCEPTION E'FAILED: Маршрут % — SCD2 не закрыл старую версию (valid_to IS NULL у всех версий).\n' + 'Подсказка: hashdiff изменился (departure_time сдвинут на 1 час),\n' + 'но ваш load-скрипт не обнаружил это изменение.\n' + 'Проверьте:\n' + ' 1. Формулу hashdiff — включает ли она departure_time?\n' + ' 2. Логику сравнения hashdiff (UPDATE ... SET valid_to = CURRENT_DATE WHERE hashdiff <> новый_hashdiff)', v_route; + END IF; + + -- Проверка c: Новая версия открыта (valid_to IS NULL) + SELECT COUNT(*) INTO v_open_count + FROM dds.dim_routes WHERE route_bk = v_route AND valid_to IS NULL; + + IF v_open_count <> 1 THEN + RAISE EXCEPTION E'FAILED: Маршрут % — ожидается ровно 1 открытая версия (valid_to IS NULL), найдено %.\n' + 'Подсказка: SCD2 должен вставить новую строку с valid_to = NULL.', v_route, v_open_count; + END IF; + + -- Проверка d: hashdiff старой ≠ hashdiff новой (мутация действительно отразилась в хеше) + SELECT hashdiff INTO v_old_hash + FROM dds.dim_routes WHERE route_bk = v_route AND valid_to IS NOT NULL + ORDER BY valid_from DESC LIMIT 1; + + SELECT hashdiff INTO v_new_hash + FROM dds.dim_routes WHERE route_bk = v_route AND valid_to IS NULL; + + IF v_old_hash = v_new_hash THEN + RAISE EXCEPTION E'FAILED: Маршрут % — hashdiff старой и новой версий совпадают.\n' + 'Мутация сдвинула departure_time на 1 час, но hashdiff не изменился.\n' + 'Проверьте, что departure_time входит в формулу hashdiff.', v_route; + END IF; + + -- Проверка e: Нет «дыры» между valid_to старой и valid_from новой + SELECT COUNT(*) INTO v_gap_count + FROM dds.dim_routes AS old_v + JOIN dds.dim_routes AS new_v + ON old_v.route_bk = new_v.route_bk + WHERE old_v.route_bk = v_route + AND old_v.valid_to IS NOT NULL + AND new_v.valid_to IS NULL + AND old_v.valid_to <> new_v.valid_from; + + IF v_gap_count > 0 THEN + RAISE EXCEPTION E'FAILED: Маршрут % — «дыра» между версиями:\n' + 'valid_to старой ≠ valid_from новой.\n' + 'Подсказка: полуоткрытый интервал [valid_from, valid_to).\n' + 'valid_from новой версии должен = valid_to старой (обычно CURRENT_DATE).', v_route; + END IF; + + RAISE NOTICE 'PASSED: SCD2 корректен для маршрута %. 2 версии, hashdiff различаются, «дыр» нет.', v_route; +END $$; diff --git a/sql/validate/dim_routes_scd2_mutate.sql b/sql/validate/dim_routes_scd2_mutate.sql new file mode 100644 index 0000000..ab93d27 --- /dev/null +++ b/sql/validate/dim_routes_scd2_mutate.sql @@ -0,0 +1,35 @@ +-- Мутация: сдвигаем departure_time у одного маршрута на 1 час. +-- Это должно изменить hashdiff → SCD2 должен закрыть старую версию и создать новую. +-- +-- Сохраняем route_no тестового маршрута в служебную таблицу _validate_scd2_target, +-- чтобы check-скрипт точно знал, какой маршрут проверять (а не угадывал по побочным эффектам). +DO $$ +DECLARE + v_route TEXT; + v_old_time TIME; +BEGIN + -- Берём первый маршрут, у которого departure_time заполнен + SELECT route_no, departure_time + INTO v_route, v_old_time + FROM ods.routes + WHERE departure_time IS NOT NULL + ORDER BY route_no + LIMIT 1; + + IF v_route IS NULL THEN + RAISE EXCEPTION 'FAILED: ods.routes пуста или нет маршрутов с departure_time. Загрузите STG→ODS перед проверкой.'; + END IF; + + -- Запоминаем тестовый маршрут в служебную таблицу + DROP TABLE IF EXISTS _validate_scd2_target; + CREATE TABLE _validate_scd2_target AS + SELECT v_route AS route_no; + + -- Сдвигаем время на 1 час у всех записей этого маршрута + UPDATE ods.routes + SET departure_time = departure_time + INTERVAL '1 hour' + WHERE route_no = v_route; + + RAISE NOTICE 'SCD2 TEST: маршрут % — departure_time сдвинут с % на %', + v_route, v_old_time, v_old_time + INTERVAL '1 hour'; +END $$; diff --git a/sql/validate/dim_routes_scd2_restore.sql b/sql/validate/dim_routes_scd2_restore.sql new file mode 100644 index 0000000..8ff75f6 --- /dev/null +++ b/sql/validate/dim_routes_scd2_restore.sql @@ -0,0 +1,42 @@ +-- Откат данных после активного теста SCD2. +-- Выполняется ВСЕГДА (trigger_rule="all_done"), даже если check упал. +-- +-- Безопасность: если backup-шаг не создал таблицы (сбой на backup), +-- откат НЕ трогает live-данные — просто чистит служебные таблицы. +-- Это гарантирует, что restore никогда не сломает ods.routes / dds.dim_routes. +-- +-- Паттерн: setup → act → assert → teardown (стандарт интеграционных тестов). + +DO $$ +DECLARE + v_has_ods_backup BOOLEAN; + v_has_dim_backup BOOLEAN; +BEGIN + -- Проверяем существование backup-таблиц через to_regclass + -- (ищет по search_path — совпадает с тем, как CREATE TABLE их создал) + v_has_ods_backup := to_regclass('_validate_bk_ods_routes') IS NOT NULL; + v_has_dim_backup := to_regclass('_validate_bk_dim_routes') IS NOT NULL; + + -- Восстанавливаем ods.routes только если бэкап существует + IF v_has_ods_backup THEN + TRUNCATE ods.routes; + INSERT INTO ods.routes SELECT * FROM _validate_bk_ods_routes; + RAISE NOTICE 'RESTORE: ods.routes восстановлена из бэкапа'; + ELSE + RAISE NOTICE 'RESTORE: бэкап ods.routes не найден — пропускаем (backup-шаг не завершился?)'; + END IF; + + -- Восстанавливаем dds.dim_routes только если бэкап существует + IF v_has_dim_backup THEN + TRUNCATE dds.dim_routes; + INSERT INTO dds.dim_routes SELECT * FROM _validate_bk_dim_routes; + RAISE NOTICE 'RESTORE: dds.dim_routes восстановлена из бэкапа'; + ELSE + RAISE NOTICE 'RESTORE: бэкап dds.dim_routes не найден — пропускаем'; + END IF; +END $$; + +-- Cleanup служебных таблиц (безусловно, IF EXISTS) +DROP TABLE IF EXISTS _validate_bk_ods_routes; +DROP TABLE IF EXISTS _validate_bk_dim_routes; +DROP TABLE IF EXISTS _validate_scd2_target; diff --git a/sql/validate/monthly_overview_exists.sql b/sql/validate/monthly_overview_exists.sql new file mode 100644 index 0000000..be4d52e --- /dev/null +++ b/sql/validate/monthly_overview_exists.sql @@ -0,0 +1,23 @@ +-- Проверка dm.monthly_overview: +-- 1. Таблица не пуста +-- 2. Нет NULL в ключевых полях (year_actual, month_actual, airplane_sk) +DO $$ +DECLARE + v_count BIGINT; + v_null_pk BIGINT; +BEGIN + SELECT COUNT(*) INTO v_count FROM dm.monthly_overview; + IF v_count = 0 THEN + RAISE EXCEPTION 'FAILED: dm.monthly_overview пуста. Реализуйте загрузку: sql/dm/monthly_overview_load.sql'; + END IF; + + SELECT COUNT(*) INTO v_null_pk + FROM dm.monthly_overview + WHERE year_actual IS NULL OR month_actual IS NULL OR airplane_sk IS NULL; + + IF v_null_pk > 0 THEN + RAISE EXCEPTION 'FAILED: dm.monthly_overview содержит % строк с NULL в ключе (year_actual, month_actual, airplane_sk).', v_null_pk; + END IF; + + RAISE NOTICE 'PASSED: dm.monthly_overview содержит % строк', v_count; +END $$; diff --git a/sql/validate/ods_airplanes_rowcount.sql b/sql/validate/ods_airplanes_rowcount.sql new file mode 100644 index 0000000..18af086 --- /dev/null +++ b/sql/validate/ods_airplanes_rowcount.sql @@ -0,0 +1,56 @@ +-- Проверка: ODS airplanes содержит все BK из STG-батча. +-- +-- Логика: ODS — TRUNCATE+INSERT snapshot. STG — append-only история всех батчей. +-- Сравниваем не счётчики строк, а точное множество BK для того батча, +-- который ODS фактически загрузил (определяем по _load_id из ods.airplanes). +DO $$ +DECLARE + v_batch_count BIGINT; + v_batch TEXT; + v_ods_count BIGINT; + v_missing_in_ods BIGINT; + v_extra_in_ods BIGINT; +BEGIN + -- ODS пуста? + SELECT COUNT(*) INTO v_ods_count FROM ods.airplanes; + IF v_ods_count = 0 THEN + RAISE EXCEPTION 'FAILED: ods.airplanes пуста. Реализуйте загрузку: sql/ods/airplanes_load.sql'; + END IF; + + -- Инвариант: ODS после TRUNCATE+INSERT содержит ровно один _load_id + SELECT COUNT(DISTINCT _load_id) INTO v_batch_count FROM ods.airplanes; + IF v_batch_count <> 1 THEN + RAISE EXCEPTION 'FAILED: ods.airplanes содержит % разных _load_id (ожидается 1 после TRUNCATE+INSERT). Проверьте, что load начинается с TRUNCATE.', v_batch_count; + END IF; + + -- Определяем батч, из которого загружена ODS + SELECT DISTINCT _load_id INTO v_batch FROM ods.airplanes; + + -- BK есть в STG-батче, но нет в ODS (потеряны при загрузке) + SELECT COUNT(*) INTO v_missing_in_ods + FROM ( + SELECT DISTINCT airplane_code FROM stg.airplanes WHERE _load_id = v_batch + ) AS stg_bk + WHERE NOT EXISTS ( + SELECT 1 FROM ods.airplanes AS o WHERE o.airplane_code = stg_bk.airplane_code + ); + + IF v_missing_in_ods > 0 THEN + RAISE EXCEPTION 'FAILED: % самолётов из STG-батча (%) отсутствуют в ods.airplanes. Проверьте логику TRUNCATE+INSERT.', v_missing_in_ods, v_batch; + END IF; + + -- BK есть в ODS, но нет в STG-батче (откуда взялись?) + SELECT COUNT(*) INTO v_extra_in_ods + FROM ( + SELECT DISTINCT airplane_code FROM ods.airplanes + ) AS ods_bk + WHERE NOT EXISTS ( + SELECT 1 FROM stg.airplanes AS s WHERE s._load_id = v_batch AND s.airplane_code = ods_bk.airplane_code + ); + + IF v_extra_in_ods > 0 THEN + RAISE EXCEPTION 'FAILED: % самолётов в ods.airplanes отсутствуют в STG-батче (%). Возможно, TRUNCATE не выполнился перед INSERT.', v_extra_in_ods, v_batch; + END IF; + + RAISE NOTICE 'PASSED: ods.airplanes содержит % самолётов, множество BK = STG-батч %', v_ods_count, v_batch; +END $$; diff --git a/sql/validate/ods_no_dup_bk.sql b/sql/validate/ods_no_dup_bk.sql new file mode 100644 index 0000000..27e23fd --- /dev/null +++ b/sql/validate/ods_no_dup_bk.sql @@ -0,0 +1,55 @@ +-- Проверка ODS-инвариантов: +-- 1. Нет дублей по BK в ods.airplanes и ods.seats. +-- 2. Обе таблицы загружены из одного согласованного батча (_load_id совпадают). +-- +-- Почему важна согласованность батча: +-- ods.airplanes и ods.seats — части одного snapshot'а (TRUNCATE+INSERT из одного STG-батча). +-- Если они собраны из разных батчей (например, airplanes перезагрузили, а seats — нет), +-- JOIN между ними даст неконсистентный срез и невалидный total_seats в dim_airplanes. +DO $$ +DECLARE + v_dup_airplanes BIGINT; + v_dup_seats BIGINT; + v_load_id_airplanes TEXT; + v_load_id_seats TEXT; +BEGIN + -- airplanes: BK = airplane_code + SELECT COUNT(*) INTO v_dup_airplanes + FROM ( + SELECT airplane_code FROM ods.airplanes + GROUP BY airplane_code HAVING COUNT(*) > 1 + ) AS d; + + IF v_dup_airplanes > 0 THEN + RAISE EXCEPTION 'FAILED: ods.airplanes содержит % дублирующихся airplane_code. Проверьте, что load начинается с TRUNCATE.', v_dup_airplanes; + END IF; + + -- seats: BK = (airplane_code, seat_no) + SELECT COUNT(*) INTO v_dup_seats + FROM ( + SELECT airplane_code, seat_no FROM ods.seats + GROUP BY airplane_code, seat_no HAVING COUNT(*) > 1 + ) AS d; + + IF v_dup_seats > 0 THEN + RAISE EXCEPTION 'FAILED: ods.seats содержит % дублирующихся (airplane_code, seat_no). Проверьте, что load начинается с TRUNCATE.', v_dup_seats; + END IF; + + -- Согласованность батча: обе таблицы должны быть загружены из одного _load_id + SELECT DISTINCT _load_id INTO v_load_id_airplanes FROM ods.airplanes; + SELECT DISTINCT _load_id INTO v_load_id_seats FROM ods.seats; + + IF v_load_id_airplanes IS NOT NULL + AND v_load_id_seats IS NOT NULL + AND v_load_id_airplanes <> v_load_id_seats + THEN + RAISE EXCEPTION E'FAILED: ods.airplanes и ods.seats загружены из разных батчей.\n' + ' airplanes._load_id = %\n' + ' seats._load_id = %\n' + 'ODS — единый snapshot: обе таблицы должны содержать один и тот же _load_id.\n' + 'Запустите оба load-скрипта (airplanes_load.sql и seats_load.sql) в одном прогоне.', + v_load_id_airplanes, v_load_id_seats; + END IF; + + RAISE NOTICE 'PASSED: Нет дублей по BK в ods.airplanes и ods.seats; батч согласован (%)', v_load_id_airplanes; +END $$; diff --git a/sql/validate/ods_no_null_pks.sql b/sql/validate/ods_no_null_pks.sql new file mode 100644 index 0000000..afd18ca --- /dev/null +++ b/sql/validate/ods_no_null_pks.sql @@ -0,0 +1,26 @@ +-- Проверка: PK-поля не содержат NULL в ods.airplanes и ods.seats. +-- +-- NULL в PK ломает JOIN'ы и агрегаты в DDS/DM — такие строки «теряются» тихо. +DO $$ +DECLARE + v_null_airplanes BIGINT; + v_null_seats BIGINT; +BEGIN + -- airplanes: PK = airplane_code + SELECT COUNT(*) INTO v_null_airplanes + FROM ods.airplanes WHERE airplane_code IS NULL; + + IF v_null_airplanes > 0 THEN + RAISE EXCEPTION 'FAILED: ods.airplanes содержит % строк с NULL в airplane_code.', v_null_airplanes; + END IF; + + -- seats: PK = (airplane_code, seat_no) + SELECT COUNT(*) INTO v_null_seats + FROM ods.seats WHERE airplane_code IS NULL OR seat_no IS NULL; + + IF v_null_seats > 0 THEN + RAISE EXCEPTION 'FAILED: ods.seats содержит % строк с NULL в PK (airplane_code, seat_no).', v_null_seats; + END IF; + + RAISE NOTICE 'PASSED: PK не содержат NULL в ods.airplanes и ods.seats'; +END $$; diff --git a/sql/validate/ods_seats_rowcount.sql b/sql/validate/ods_seats_rowcount.sql new file mode 100644 index 0000000..5deb0e7 --- /dev/null +++ b/sql/validate/ods_seats_rowcount.sql @@ -0,0 +1,55 @@ +-- Проверка: ODS seats содержит все BK из STG-батча. +-- +-- Аналог ods_airplanes_rowcount.sql, но для seats с составным BK (airplane_code, seat_no). +DO $$ +DECLARE + v_batch_count BIGINT; + v_batch TEXT; + v_ods_count BIGINT; + v_missing_in_ods BIGINT; + v_extra_in_ods BIGINT; +BEGIN + -- ODS пуста? + SELECT COUNT(*) INTO v_ods_count FROM ods.seats; + IF v_ods_count = 0 THEN + RAISE EXCEPTION 'FAILED: ods.seats пуста. Реализуйте загрузку: sql/ods/seats_load.sql'; + END IF; + + -- Инвариант: ODS после TRUNCATE+INSERT содержит ровно один _load_id + SELECT COUNT(DISTINCT _load_id) INTO v_batch_count FROM ods.seats; + IF v_batch_count <> 1 THEN + RAISE EXCEPTION 'FAILED: ods.seats содержит % разных _load_id (ожидается 1 после TRUNCATE+INSERT). Проверьте, что load начинается с TRUNCATE.', v_batch_count; + END IF; + + SELECT DISTINCT _load_id INTO v_batch FROM ods.seats; + + -- BK есть в STG-батче, но нет в ODS (потеряны при загрузке) + SELECT COUNT(*) INTO v_missing_in_ods + FROM ( + SELECT DISTINCT airplane_code, seat_no FROM stg.seats WHERE _load_id = v_batch + ) AS stg_bk + WHERE NOT EXISTS ( + SELECT 1 FROM ods.seats AS o + WHERE o.airplane_code = stg_bk.airplane_code AND o.seat_no = stg_bk.seat_no + ); + + IF v_missing_in_ods > 0 THEN + RAISE EXCEPTION 'FAILED: % мест из STG-батча (%) отсутствуют в ods.seats. Проверьте логику TRUNCATE+INSERT.', v_missing_in_ods, v_batch; + END IF; + + -- BK есть в ODS, но нет в STG-батче + SELECT COUNT(*) INTO v_extra_in_ods + FROM ( + SELECT DISTINCT airplane_code, seat_no FROM ods.seats + ) AS ods_bk + WHERE NOT EXISTS ( + SELECT 1 FROM stg.seats AS s + WHERE s._load_id = v_batch AND s.airplane_code = ods_bk.airplane_code AND s.seat_no = ods_bk.seat_no + ); + + IF v_extra_in_ods > 0 THEN + RAISE EXCEPTION 'FAILED: % мест в ods.seats отсутствуют в STG-батче (%). Возможно, TRUNCATE не выполнился перед INSERT.', v_extra_in_ods, v_batch; + END IF; + + RAISE NOTICE 'PASSED: ods.seats содержит % мест, множество BK = STG-батч %', v_ods_count, v_batch; +END $$; diff --git a/sql/validate/passenger_loyalty_exists.sql b/sql/validate/passenger_loyalty_exists.sql new file mode 100644 index 0000000..4cda129 --- /dev/null +++ b/sql/validate/passenger_loyalty_exists.sql @@ -0,0 +1,23 @@ +-- Проверка dm.passenger_loyalty: +-- 1. Таблица не пуста +-- 2. Нет NULL в ключевом поле (passenger_sk) +DO $$ +DECLARE + v_count BIGINT; + v_null_pk BIGINT; +BEGIN + SELECT COUNT(*) INTO v_count FROM dm.passenger_loyalty; + IF v_count = 0 THEN + RAISE EXCEPTION 'FAILED: dm.passenger_loyalty пуста. Реализуйте загрузку: sql/dm/passenger_loyalty_load.sql'; + END IF; + + SELECT COUNT(*) INTO v_null_pk + FROM dm.passenger_loyalty + WHERE passenger_sk IS NULL; + + IF v_null_pk > 0 THEN + RAISE EXCEPTION 'FAILED: dm.passenger_loyalty содержит % строк с NULL в ключе (passenger_sk).', v_null_pk; + END IF; + + RAISE NOTICE 'PASSED: dm.passenger_loyalty содержит % строк', v_count; +END $$; diff --git a/sql/validate/route_performance_exists.sql b/sql/validate/route_performance_exists.sql new file mode 100644 index 0000000..4e4593e --- /dev/null +++ b/sql/validate/route_performance_exists.sql @@ -0,0 +1,23 @@ +-- Проверка dm.route_performance: +-- 1. Таблица не пуста +-- 2. Нет NULL в ключевом поле (route_bk) +DO $$ +DECLARE + v_count BIGINT; + v_null_pk BIGINT; +BEGIN + SELECT COUNT(*) INTO v_count FROM dm.route_performance; + IF v_count = 0 THEN + RAISE EXCEPTION 'FAILED: dm.route_performance пуста. Реализуйте загрузку: sql/dm/route_performance_load.sql'; + END IF; + + SELECT COUNT(*) INTO v_null_pk + FROM dm.route_performance + WHERE route_bk IS NULL; + + IF v_null_pk > 0 THEN + RAISE EXCEPTION 'FAILED: dm.route_performance содержит % строк с NULL в ключе (route_bk).', v_null_pk; + END IF; + + RAISE NOTICE 'PASSED: dm.route_performance содержит % строк', v_count; +END $$; diff --git a/tests/test_dags_smoke.py b/tests/test_dags_smoke.py index 8d645e5..f88c68c 100644 --- a/tests/test_dags_smoke.py +++ b/tests/test_dags_smoke.py @@ -369,6 +369,66 @@ def test_bookings_dm_ddl_dag_structure(): ) +class TestBookingsValidate: + """Smoke-тесты DAG bookings_validate.""" + + def test_dag_loads(self): + dag = _load_dag("airflow.dags.bookings_validate") + assert dag is not None + + def test_expected_tasks(self): + dag = _load_dag("airflow.dags.bookings_validate") + expected = { + # ODS + "validate_ods.check_ods_airplanes_rowcount", + "validate_ods.check_ods_seats_rowcount", + "validate_ods.check_ods_no_dup_bk", + "validate_ods.check_ods_no_null_pks", + # DDS + "validate_dds.check_dim_airplanes_exists", + "validate_dds.check_dim_passengers_exists", + "validate_dds.check_dim_passengers_no_dup_bk", + "validate_dds.check_dim_routes_exists", + "validate_dds.scd2_backup", + "validate_dds.scd2_mutate", + "validate_dds.scd2_run_student_load", + "validate_dds.scd2_check", + "validate_dds.scd2_restore", + "validate_dds.check_dim_routes_no_gaps", + # DM + "validate_dm.check_airport_traffic_exists", + "validate_dm.check_route_performance_exists", + "validate_dm.check_monthly_overview_exists", + "validate_dm.check_passenger_loyalty_exists", + } + assert expected.issubset(dag.task_dict.keys()) + + def test_scd2_chain(self): + """SCD2 цепочка backup → mutate → load → check → restore.""" + dag = _load_dag("airflow.dags.bookings_validate") + _assert_direct_edge(dag, "validate_dds.scd2_backup", "validate_dds.scd2_mutate") + _assert_direct_edge( + dag, "validate_dds.scd2_mutate", "validate_dds.scd2_run_student_load" + ) + _assert_direct_edge( + dag, "validate_dds.scd2_run_student_load", "validate_dds.scd2_check" + ) + _assert_direct_edge(dag, "validate_dds.scd2_check", "validate_dds.scd2_restore") + + def test_no_gaps_after_restore(self): + """no_gaps должен выполняться после restore (на чистых данных).""" + dag = _load_dag("airflow.dags.bookings_validate") + _assert_direct_edge( + dag, "validate_dds.scd2_restore", "validate_dds.check_dim_routes_no_gaps" + ) + + def test_restore_trigger_rule(self): + """restore должен выполняться всегда (all_done), даже если check упал.""" + dag = _load_dag("airflow.dags.bookings_validate") + restore_task = dag.task_dict["validate_dds.scd2_restore"] + assert restore_task.trigger_rule == "all_done" + + def test_bookings_to_gp_dm_dag_structure(): """Проверка структуры DAG bookings_to_gp_dm.""" dag = _load_dag("airflow.dags.bookings_to_gp_dm")