From 5ca4e6b2ca6d3987833004ca7a9116ac8826b359 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Thu, 12 Mar 2026 23:07:14 +0300 Subject: [PATCH] =?UTF-8?q?feat(main):=20=D0=B7=D0=B0=D0=B3=D0=BB=D1=83?= =?UTF-8?q?=D1=88=D0=BA=D0=B8=20=D1=81=D1=82=D1=83=D0=B4=D0=B5=D0=BD=D1=87?= =?UTF-8?q?=D0=B5=D1=81=D0=BA=D0=B8=D1=85=20=D1=84=D0=B0=D0=B9=D0=BB=D0=BE?= =?UTF-8?q?=D0=B2,=20=D0=BE=D1=81=D0=BB=D0=B0=D0=B1=D0=BB=D0=B5=D0=BD?= =?UTF-8?q?=D0=B8=D0=B5=20DQ,=20=D0=B0=D0=B4=D0=B0=D0=BF=D1=82=D0=B0=D1=86?= =?UTF-8?q?=D0=B8=D1=8F=20=D1=82=D0=B5=D1=81=D1=82=D0=BE=D0=B2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - fact_flight_sales_dq: student SK (passenger/route/airplane) → NOTICE, airport_sk → порог 1% (эталон через ods.routes) - routes_dq: закомментирован RI airplane_code → ods.airplanes - ODS DAG: убрана зависимость dq_ods_airplanes → load_ods_routes - 18 студенческих файлов заменены заглушками (ODS/DDS/DM load+dq) - test_ods_sql_contract: SNAPSHOT_ENTITIES без airplanes/seats - test_dags_smoke: убран assert dq_ods_airplanes → dq_ods_routes Co-Authored-By: Claude Opus 4.6 --- airflow/dags/bookings_to_gp_ods.py | 4 +- sql/dds/dim_airplanes_dq.sql | 86 +-------------- sql/dds/dim_airplanes_load.sql | 100 +---------------- sql/dds/dim_passengers_dq.sql | 103 +---------------- sql/dds/dim_passengers_load.sql | 71 +----------- sql/dds/dim_routes_dq.sql | 160 +-------------------------- sql/dds/dim_routes_load.sql | 148 +------------------------ sql/dds/fact_flight_sales_dq.sql | 50 ++++++--- sql/dm/airport_traffic_dq.sql | 52 +-------- sql/dm/airport_traffic_load.sql | 134 +---------------------- sql/dm/monthly_overview_dq.sql | 55 +--------- sql/dm/monthly_overview_load.sql | 170 +---------------------------- sql/dm/passenger_loyalty_dq.sql | 54 +-------- sql/dm/passenger_loyalty_load.sql | 134 +---------------------- sql/dm/route_performance_dq.sql | 62 +---------- sql/dm/route_performance_load.sql | 86 +-------------- sql/ods/airplanes_dq.sql | 102 +---------------- sql/ods/airplanes_load.sql | 43 +------- sql/ods/routes_dq.sql | 34 +++--- sql/ods/seats_dq.sql | 128 +--------------------- sql/ods/seats_load.sql | 40 +------ tests/test_dags_smoke.py | 3 +- tests/test_ods_sql_contract.py | 4 +- 23 files changed, 120 insertions(+), 1703 deletions(-) diff --git a/airflow/dags/bookings_to_gp_ods.py b/airflow/dags/bookings_to_gp_ods.py index a38287d..c63f4b4 100644 --- a/airflow/dags/bookings_to_gp_ods.py +++ b/airflow/dags/bookings_to_gp_ods.py @@ -263,7 +263,9 @@ with DAG( resolve_stg_batch_id >> load_ods_airports >> dq_ods_airports resolve_stg_batch_id >> load_ods_airplanes >> dq_ods_airplanes - [dq_ods_airports, dq_ods_airplanes] >> load_ods_routes >> dq_ods_routes + # На main routes не зависит от airplanes (ods.airplanes — студенческая заглушка). + # На ветке solution: [dq_ods_airports, dq_ods_airplanes] >> load_ods_routes + dq_ods_airports >> load_ods_routes >> dq_ods_routes dq_ods_airplanes >> load_ods_seats >> dq_ods_seats dq_ods_routes >> load_ods_flights >> dq_ods_flights diff --git a/sql/dds/dim_airplanes_dq.sql b/sql/dds/dim_airplanes_dq.sql index c0d01dd..f69af41 100644 --- a/sql/dds/dim_airplanes_dq.sql +++ b/sql/dds/dim_airplanes_dq.sql @@ -1,83 +1,3 @@ --- DQ для DDS dim_airplanes. - -DO $$ -DECLARE - v_row_count BIGINT; - v_dup_sk BIGINT; - v_dup_bk BIGINT; - v_missing_bk BIGINT; - v_null_count BIGINT; -BEGIN - -- Таблица не пуста. - SELECT COUNT(*) - INTO v_row_count - FROM dds.dim_airplanes; - - IF v_row_count = 0 THEN - RAISE EXCEPTION 'DQ FAILED: dds.dim_airplanes пуста.'; - END IF; - - -- Нет дублей по SK. - SELECT COUNT(*) - COUNT(DISTINCT airplane_sk) - INTO v_dup_sk - FROM dds.dim_airplanes; - - IF v_dup_sk <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: в dds.dim_airplanes найдены дубликаты airplane_sk: %', - v_dup_sk; - END IF; - - -- Нет дублей по BK. - SELECT COUNT(*) - COUNT(DISTINCT airplane_bk) - INTO v_dup_bk - FROM dds.dim_airplanes; - - IF v_dup_bk <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: в dds.dim_airplanes найдены дубликаты airplane_bk: %', - v_dup_bk; - END IF; - - -- Покрытие ODS: все airplane_code из ODS есть в DDS. - SELECT COUNT(*) - INTO v_missing_bk - 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_bk <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: в dds.dim_airplanes отсутствуют ключи из ods.airplanes: %', - v_missing_bk; - END IF; - - -- Обязательные поля. - SELECT COUNT(*) - INTO v_null_count - FROM dds.dim_airplanes - WHERE airplane_sk IS NULL - OR airplane_bk IS NULL - OR airplane_bk = '' - OR model IS NULL - OR model = '' - OR total_seats IS NULL - OR created_at IS NULL - OR updated_at IS NULL - OR _load_id IS NULL - OR _load_id = '' - OR _load_ts IS NULL; - - IF v_null_count <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: в dds.dim_airplanes найдены NULL/пустые обязательные поля: %', - v_null_count; - END IF; - - RAISE NOTICE - 'DQ PASSED: dds.dim_airplanes ок, строк=%', - v_row_count; -END $$; +-- TODO: реализуйте проверки качества данных для dds.dim_airplanes +-- Эталонную реализацию можно найти в ветке solution. +SELECT 1; diff --git a/sql/dds/dim_airplanes_load.sql b/sql/dds/dim_airplanes_load.sql index 05999dd..60faa57 100644 --- a/sql/dds/dim_airplanes_load.sql +++ b/sql/dds/dim_airplanes_load.sql @@ -1,96 +1,4 @@ --- Загрузка DDS dim_airplanes: SCD1 UPSERT (UPDATE изменившихся + INSERT новых). - --- Statement 1: UPDATE существующих записей (если атрибуты изменились). -WITH seats_agg AS ( - SELECT - s.airplane_code, - COUNT(*)::INTEGER AS total_seats - FROM ods.seats AS s - GROUP BY s.airplane_code -), -src AS ( - SELECT - a.airplane_code, - a.model, - a.range_km, - a.speed_kmh, - COALESCE(sa.total_seats, 0) AS total_seats - FROM ods.airplanes AS a - LEFT JOIN seats_agg AS sa - ON sa.airplane_code = a.airplane_code -) -UPDATE dds.dim_airplanes AS d -SET model = s.model, - range_km = s.range_km, - speed_kmh = s.speed_kmh, - total_seats = s.total_seats, - updated_at = now(), - _load_id = '{{ run_id }}', - _load_ts = now() -FROM src AS s -WHERE d.airplane_bk = s.airplane_code - AND ( - d.model IS DISTINCT FROM s.model - OR d.range_km IS DISTINCT FROM s.range_km - OR d.speed_kmh IS DISTINCT FROM s.speed_kmh - OR d.total_seats IS DISTINCT FROM s.total_seats - ); - --- Statement 2: INSERT новых записей (MAX(sk) + ROW_NUMBER()). -WITH seats_agg AS ( - SELECT - s.airplane_code, - COUNT(*)::INTEGER AS total_seats - FROM ods.seats AS s - GROUP BY s.airplane_code -), -src AS ( - SELECT - a.airplane_code, - a.model, - a.range_km, - a.speed_kmh, - COALESCE(sa.total_seats, 0) AS total_seats - FROM ods.airplanes AS a - LEFT JOIN seats_agg AS sa - ON sa.airplane_code = a.airplane_code -), --- Учебный комментарий: Генерация SK через MAX() + ROW_NUMBER() --- Этот подход работает безопасно только потому, что Airflow запускает --- джобы загрузки для одной таблицы строго последовательно (concurrency=1). --- При параллельной загрузке возникнет состояние гонки (race condition) и возможны дубли SK. -max_sk AS ( - SELECT COALESCE(MAX(airplane_sk), 0) AS v - FROM dds.dim_airplanes -) -INSERT INTO dds.dim_airplanes ( - airplane_sk, - airplane_bk, - model, - range_km, - speed_kmh, - total_seats, - created_at, - updated_at, - _load_id, - _load_ts -) -SELECT - (SELECT v FROM max_sk) + ROW_NUMBER() OVER (ORDER BY s.airplane_code)::INTEGER, - s.airplane_code, - s.model, - s.range_km, - s.speed_kmh, - s.total_seats, - now(), - now(), - '{{ run_id }}', - now() -FROM src AS s -WHERE NOT EXISTS ( - SELECT 1 - FROM dds.dim_airplanes AS d - WHERE d.airplane_bk = s.airplane_code -); - -ANALYZE dds.dim_airplanes; +-- TODO: реализуйте загрузку dds.dim_airplanes (см. ТЗ в docs/assignment/analyst_spec.md) +-- Паттерн: SCD1 UPSERT (UPDATE изменившихся + INSERT новых). +-- Эталонную реализацию можно найти в ветке solution. +SELECT 1; diff --git a/sql/dds/dim_passengers_dq.sql b/sql/dds/dim_passengers_dq.sql index 990b365..cd96fde 100644 --- a/sql/dds/dim_passengers_dq.sql +++ b/sql/dds/dim_passengers_dq.sql @@ -1,100 +1,3 @@ --- DQ для DDS dim_passengers. - -DO $$ -DECLARE - v_row_count BIGINT; - v_src_count BIGINT; - v_dup_sk BIGINT; - v_dup_bk BIGINT; - v_missing_bk BIGINT; - v_null_count BIGINT; -BEGIN - -- Для инкрементальных периодов без новых билетов допускаем пустую dim_passengers. - SELECT COUNT(*) - INTO v_row_count - FROM dds.dim_passengers; - - SELECT COUNT(DISTINCT passenger_id) - INTO v_src_count - FROM ods.tickets - WHERE passenger_id IS NOT NULL - AND passenger_id <> ''; - - IF v_row_count = 0 AND v_src_count <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: dds.dim_passengers пуста при непустом источнике ods.tickets (passenger_id=%).', - v_src_count; - ELSIF v_row_count = 0 AND v_src_count = 0 THEN - RAISE NOTICE - 'DQ PASSED: dds.dim_passengers пуста, т.к. в ods.tickets нет passenger_id для загрузки.'; - RETURN; - END IF; - - -- Нет дублей по SK. - SELECT COUNT(*) - COUNT(DISTINCT passenger_sk) - INTO v_dup_sk - FROM dds.dim_passengers; - - IF v_dup_sk <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: в dds.dim_passengers найдены дубликаты passenger_sk: %', - v_dup_sk; - END IF; - - -- Нет дублей по BK (ID). - SELECT COUNT(*) - COUNT(DISTINCT passenger_id) - INTO v_dup_bk - FROM dds.dim_passengers; - - IF v_dup_bk <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: в dds.dim_passengers найдены дубликаты passenger_id: %', - v_dup_bk; - END IF; - - -- Покрытие ODS: все passenger_id из ods.tickets есть в DDS. - SELECT COUNT(*) - INTO v_missing_bk - FROM ( - SELECT DISTINCT passenger_id - FROM ods.tickets - WHERE passenger_id IS NOT NULL - AND passenger_id <> '' - ) AS s - WHERE NOT EXISTS ( - SELECT 1 - FROM dds.dim_passengers AS d - WHERE d.passenger_id = s.passenger_id - ); - - IF v_missing_bk <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: в dds.dim_passengers отсутствуют passenger_id из ods.tickets: %', - v_missing_bk; - END IF; - - -- Обязательные поля. - SELECT COUNT(*) - INTO v_null_count - FROM dds.dim_passengers - WHERE passenger_sk IS NULL - OR passenger_id IS NULL - OR passenger_id = '' - OR passenger_name IS NULL - OR passenger_name = '' - OR created_at IS NULL - OR updated_at IS NULL - OR _load_id IS NULL - OR _load_id = '' - OR _load_ts IS NULL; - - IF v_null_count <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: в dds.dim_passengers найдены NULL/пустые обязательные поля: %', - v_null_count; - END IF; - - RAISE NOTICE - 'DQ PASSED: dds.dim_passengers ок, строк=%', - v_row_count; -END $$; +-- TODO: реализуйте проверки качества данных для dds.dim_passengers +-- Эталонную реализацию можно найти в ветке solution. +SELECT 1; diff --git a/sql/dds/dim_passengers_load.sql b/sql/dds/dim_passengers_load.sql index 9d06bcd..145a59b 100644 --- a/sql/dds/dim_passengers_load.sql +++ b/sql/dds/dim_passengers_load.sql @@ -1,67 +1,4 @@ --- Загрузка DDS dim_passengers: SCD1 UPSERT (UPDATE изменившихся + INSERT новых). - --- Учебный комментарий: Используем TEMP TABLE для подготовки дельты. --- Это избавляет от дублирования сложной оконной функции в UPDATE и INSERT блоках. -CREATE TEMP TABLE tmp_passengers_src ON COMMIT DROP AS -SELECT - d.passenger_id, - d.passenger_name -FROM ( - SELECT - t.passenger_id, - t.passenger_name, - ROW_NUMBER() OVER ( - PARTITION BY t.passenger_id - ORDER BY t.event_ts DESC NULLS LAST, t._load_ts DESC, t.ticket_no DESC - ) AS rn - FROM ods.tickets AS t - WHERE t.passenger_id IS NOT NULL - AND t.passenger_id <> '' - AND t.passenger_name IS NOT NULL - AND t.passenger_name <> '' -) AS d -WHERE d.rn = 1; - --- Statement 1: UPDATE существующих записей (если атрибуты изменились). -UPDATE dds.dim_passengers AS d -SET passenger_name = s.passenger_name, - updated_at = now(), - _load_id = '{{ run_id }}', - _load_ts = now() -FROM tmp_passengers_src AS s -WHERE d.passenger_id = s.passenger_id - AND d.passenger_name IS DISTINCT FROM s.passenger_name; - --- Statement 2: INSERT новых записей (MAX(sk) + ROW_NUMBER()). -WITH max_sk AS ( - -- Учебный комментарий: Генерация SK через MAX() + ROW_NUMBER() - -- Этот подход работает безопасно только потому, что Airflow запускает - -- джобы загрузки для одной таблицы строго последовательно (concurrency=1). - SELECT COALESCE(MAX(passenger_sk), 0) AS v - FROM dds.dim_passengers -) -INSERT INTO dds.dim_passengers ( - passenger_sk, - passenger_id, - passenger_name, - created_at, - updated_at, - _load_id, - _load_ts -) -SELECT - (SELECT v FROM max_sk) + ROW_NUMBER() OVER (ORDER BY s.passenger_id)::INTEGER, - s.passenger_id, - s.passenger_name, - now(), - now(), - '{{ run_id }}', - now() -FROM tmp_passengers_src AS s -WHERE NOT EXISTS ( - SELECT 1 - FROM dds.dim_passengers AS d - WHERE d.passenger_id = s.passenger_id -); - -ANALYZE dds.dim_passengers; +-- TODO: реализуйте загрузку dds.dim_passengers (см. ТЗ в docs/assignment/analyst_spec.md) +-- Паттерн: SCD1 UPSERT (UPDATE изменившихся + INSERT новых). +-- Эталонную реализацию можно найти в ветке solution. +SELECT 1; diff --git a/sql/dds/dim_routes_dq.sql b/sql/dds/dim_routes_dq.sql index 6e1a6bd..6fc087e 100644 --- a/sql/dds/dim_routes_dq.sql +++ b/sql/dds/dim_routes_dq.sql @@ -1,157 +1,3 @@ --- DQ для DDS dim_routes (SCD2). - -DO $$ -DECLARE - v_row_count BIGINT; - v_dup_sk BIGINT; - v_dup_current BIGINT; - v_overlap_count BIGINT; - v_missing_count BIGINT; - v_orphan_current BIGINT; - v_null_count BIGINT; -BEGIN - -- Таблица не пуста. - SELECT COUNT(*) - INTO v_row_count - FROM dds.dim_routes; - - IF v_row_count = 0 THEN - RAISE EXCEPTION 'DQ FAILED: dds.dim_routes пуста.'; - END IF; - - -- Нет дублей по SK. - SELECT COUNT(*) - COUNT(DISTINCT route_sk) - INTO v_dup_sk - FROM dds.dim_routes; - - IF v_dup_sk <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: в dds.dim_routes найдены дубликаты route_sk: %', - v_dup_sk; - END IF; - - -- Корректность интервалов (valid_from <= valid_to для закрытых версий). - SELECT COUNT(*) - INTO v_null_count - FROM dds.dim_routes - WHERE valid_to IS NOT NULL - AND valid_from > valid_to; - - IF v_null_count <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: в dds.dim_routes найдены версии с valid_from > valid_to: %', - v_null_count; - END IF; - - -- Нет перекрытий интервалов для одного route_bk. - SELECT COUNT(*) - INTO v_overlap_count - FROM ( - SELECT 1 - FROM dds.dim_routes AS d1 - JOIN dds.dim_routes AS d2 - ON d1.route_bk = d2.route_bk - AND d1.route_sk < d2.route_sk - AND d1.valid_from < COALESCE(d2.valid_to, DATE '9999-12-31') - AND d2.valid_from < COALESCE(d1.valid_to, DATE '9999-12-31') - ) AS overlap_rows; - - IF v_overlap_count <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: в dds.dim_routes найдены перекрытия SCD2-интервалов: %', - v_overlap_count; - END IF; - - -- Не более одной текущей версии на route_bk. - SELECT COUNT(*) - INTO v_dup_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_dup_current <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: в dds.dim_routes найдены route_bk с > 1 текущей версией: %', - v_dup_current; - END IF; - - -- Покрытие ODS: все route_no имеют хотя бы одну версию в DDS. - SELECT COUNT(*) - INTO v_missing_count - FROM (SELECT DISTINCT route_no FROM ods.routes) AS o - WHERE NOT EXISTS ( - SELECT 1 - FROM dds.dim_routes AS d - WHERE d.route_bk = o.route_no - ); - - IF v_missing_count <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: в dds.dim_routes отсутствуют маршруты из ODS: %', - v_missing_count; - END IF; - - -- Current-срез DDS не содержит route_bk, которых нет в ODS. - SELECT COUNT(*) - INTO v_orphan_current - FROM ( - SELECT DISTINCT route_bk - FROM dds.dim_routes - WHERE valid_to IS NULL - ) AS d - WHERE NOT EXISTS ( - SELECT 1 - FROM ods.routes AS o - WHERE o.route_no = d.route_bk - ); - - IF v_orphan_current <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: в current-срезе dds.dim_routes есть route_bk вне ODS: %', - v_orphan_current; - END IF; - - -- Обязательные поля. - SELECT COUNT(*) - INTO v_null_count - FROM dds.dim_routes - WHERE route_sk IS NULL - OR route_bk IS NULL - OR route_bk = '' - OR departure_airport IS NULL - OR departure_airport = '' - OR arrival_airport IS NULL - OR arrival_airport = '' - OR airplane_code IS NULL - OR airplane_code = '' - OR departure_city IS NULL - OR departure_city = '' - OR arrival_city IS NULL - OR arrival_city = '' - OR airplane_model IS NULL - OR airplane_model = '' - OR total_seats IS NULL - OR hashdiff IS NULL - OR hashdiff = '' - OR valid_from IS NULL - OR created_at IS NULL - OR updated_at IS NULL - OR _load_id IS NULL - OR _load_id = '' - OR _load_ts IS NULL; - - IF v_null_count <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: в dds.dim_routes найдены NULL обязательные поля: %', - v_null_count; - END IF; - - RAISE NOTICE - 'DQ PASSED: dds.dim_routes ок, строк=% (версий)', - v_row_count; - -END $$; +-- TODO: реализуйте проверки качества данных для dds.dim_routes (SCD2) +-- Эталонную реализацию можно найти в ветке solution. +SELECT 1; diff --git a/sql/dds/dim_routes_load.sql b/sql/dds/dim_routes_load.sql index c95266b..b2c800f 100644 --- a/sql/dds/dim_routes_load.sql +++ b/sql/dds/dim_routes_load.sql @@ -1,144 +1,4 @@ --- Загрузка DDS dim_routes: SCD2 с hashdiff. - --- Учебный комментарий: Используем TEMP TABLE для подготовки дельты. --- Это избавляет от дублирования логики hashdiff в UPDATE и INSERT блоках. -CREATE TEMP TABLE tmp_routes_src ON COMMIT DROP AS -SELECT - route_no, - departure_airport, - arrival_airport, - airplane_code, - days_of_week, - departure_time, - duration, - md5( - COALESCE(departure_airport, '') || '|' || - COALESCE(arrival_airport, '') || '|' || - COALESCE(airplane_code, '') || '|' || - COALESCE(days_of_week::TEXT, '') || '|' || - COALESCE(departure_time::TEXT, '') || '|' || - COALESCE(duration::TEXT, '') - ) AS hashdiff, - ROW_NUMBER() OVER (PARTITION BY route_no ORDER BY validity DESC) AS rn -FROM ods.routes; - --- Statement 1: Закрыть устаревшие версии (valid_to = текущая дата). -UPDATE dds.dim_routes AS d -SET valid_to = CURRENT_DATE, - updated_at = now(), - _load_id = '{{ run_id }}', - _load_ts = now() -FROM tmp_routes_src AS s -WHERE s.rn = 1 - AND d.route_bk = s.route_no - AND d.valid_to IS NULL - AND d.hashdiff <> s.hashdiff; - --- Statement 1.1: Закрыть "исчезнувшие" маршруты. -UPDATE dds.dim_routes AS d -SET valid_to = CURRENT_DATE, - updated_at = now(), - _load_id = '{{ run_id }}', - _load_ts = now() -WHERE d.valid_to IS NULL - AND NOT EXISTS ( - SELECT 1 - FROM tmp_routes_src AS s - WHERE s.rn = 1 - AND s.route_no = d.route_bk - ); - --- Statement 2: Вставить новые версии (для изменённых и новых route_no). -WITH max_sk AS ( - -- Учебный комментарий: Генерация SK через MAX() + ROW_NUMBER() - -- Этот подход работает безопасно только потому, что Airflow запускает - -- джобы загрузки для одной таблицы строго последовательно (concurrency=1). - SELECT COALESCE(MAX(route_sk), 0) AS v - FROM dds.dim_routes -) -INSERT INTO dds.dim_routes ( - route_sk, - route_bk, - departure_airport, - arrival_airport, - airplane_code, - departure_city, - arrival_city, - airplane_model, - total_seats, - days_of_week, - departure_time, - duration, - hashdiff, - valid_from, - valid_to, - created_at, - updated_at, - _load_id, - _load_ts -) -SELECT - (SELECT v FROM max_sk) + ROW_NUMBER() OVER (ORDER BY s.route_no)::INTEGER, - s.route_no, - s.departure_airport, - s.arrival_airport, - s.airplane_code, - dep.city, - arr.city, - air.model, - air.total_seats, - s.days_of_week, - s.departure_time, - s.duration, - s.hashdiff, - CASE - WHEN EXISTS ( - SELECT 1 - FROM dds.dim_routes AS d2 - WHERE d2.route_bk = s.route_no - ) THEN CURRENT_DATE - ELSE '1900-01-01'::DATE - END AS valid_from, - NULL, - now(), - now(), - '{{ run_id }}', - now() -FROM tmp_routes_src AS s -LEFT JOIN dds.dim_airports AS dep ON s.departure_airport = dep.airport_bk -LEFT JOIN dds.dim_airports AS arr ON s.arrival_airport = arr.airport_bk -LEFT JOIN dds.dim_airplanes AS air ON s.airplane_code = air.airplane_bk -WHERE s.rn = 1 - AND NOT EXISTS ( - SELECT 1 - FROM dds.dim_routes AS d - WHERE d.route_bk = s.route_no - AND d.valid_to IS NULL - AND d.hashdiff = s.hashdiff - ); - --- Фаза 3: Обновление денормализованных атрибутов (refresh). --- Нужна для SCD1-изменений в dim_airports/dim_airplanes (напр. переименование города). --- Обновляем все версии (и текущие, и исторические), т.к. измерения SCD1. --- При этом мы не перезаписываем _load_id и _load_ts, чтобы не размывать lineage версий. -UPDATE dds.dim_routes AS d -SET departure_city = dep.city, - arrival_city = arr.city, - airplane_model = air.model, - total_seats = air.total_seats, - updated_at = now() -FROM dds.dim_airports AS dep, - dds.dim_airports AS arr, - dds.dim_airplanes AS air -WHERE dep.airport_bk = d.departure_airport - AND arr.airport_bk = d.arrival_airport - AND air.airplane_bk = d.airplane_code - AND ( - d.departure_city IS DISTINCT FROM dep.city - OR d.arrival_city IS DISTINCT FROM arr.city - OR d.airplane_model IS DISTINCT FROM air.model - OR d.total_seats IS DISTINCT FROM air.total_seats - ); - -ANALYZE dds.dim_routes; - +-- TODO: реализуйте загрузку dds.dim_routes (см. ТЗ в docs/assignment/analyst_spec.md) +-- Паттерн: SCD2 с hashdiff — самый нетривиальный паттерн в проекте. +-- Эталонную реализацию можно найти в ветке solution. +SELECT 1; diff --git a/sql/dds/fact_flight_sales_dq.sql b/sql/dds/fact_flight_sales_dq.sql index c4fa53c..d9f5526 100644 --- a/sql/dds/fact_flight_sales_dq.sql +++ b/sql/dds/fact_flight_sales_dq.sql @@ -7,7 +7,8 @@ DECLARE v_dup_count BIGINT; v_null_passenger BIGINT; v_null_tariff BIGINT; - v_null_route_related BIGINT; + v_null_airport BIGINT; + v_null_student BIGINT; v_null_calendar BIGINT; v_null_required BIGINT; BEGIN @@ -62,15 +63,17 @@ BEGIN v_dup_count; END IF; - -- Ссылочная целостность: passenger_sk. + -- Учебный комментарий: passenger_sk будет NULL, пока вы не реализуете dim_passengers. + -- После реализации: TRUNCATE dds.fact_flight_sales → перезагрузка → все SK заполнены. + -- Полную версию DQ (с блокировкой) см. в ветке solution. SELECT COUNT(*) INTO v_null_passenger FROM dds.fact_flight_sales WHERE passenger_sk IS NULL; IF v_null_passenger <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: в fact_flight_sales строки без passenger_sk: %', + RAISE NOTICE + 'DQ INFO: студенческий SK (passenger) NULL: %. После реализации dim_passengers: TRUNCATE fact → перезагрузка.', v_null_passenger; END IF; @@ -86,27 +89,42 @@ BEGIN v_null_tariff; END IF; - -- Route-related FK: допустимо при аномалиях, фейлим если > 1%. + -- departure_airport_sk и arrival_airport_sk заполняются через ods.routes (эталон). + -- NULL здесь — аномалия данных (пропущен маршрут в ODS), а не отсутствие студенческого кода. + -- Порог 1% — защита от единичных аномалий источника (аналогично calendar_sk). SELECT COUNT(*) - INTO v_null_route_related + INTO v_null_airport FROM dds.fact_flight_sales - WHERE route_sk IS NULL - OR departure_airport_sk IS NULL - OR arrival_airport_sk IS NULL - OR airplane_sk IS NULL; + WHERE departure_airport_sk IS NULL + OR arrival_airport_sk IS NULL; - IF v_null_route_related > 0 THEN - IF v_null_route_related * 100.0 / NULLIF(v_row_count, 0) > 1.0 THEN + IF v_null_airport > 0 THEN + IF v_null_airport * 100.0 / NULLIF(v_row_count, 0) > 1.0 THEN RAISE EXCEPTION - 'DQ FAILED: в fact_flight_sales слишком много строк с NULL в route-related FK: % (>1%%)', - v_null_route_related; + 'DQ FAILED: NULL airport_sk: % (>1%%)', + v_null_airport; ELSE RAISE NOTICE - 'DQ WARNING: в fact_flight_sales строк с NULL в route-related FK: % (<=1%%, допустимо)', - v_null_route_related; + 'DQ WARNING: NULL airport_sk: % (<=1%%, допустимо)', + v_null_airport; END IF; END IF; + -- route_sk и airplane_sk будут NULL, пока вы не реализуете dim_routes. + -- Не блокируем pipeline. + -- Полную версию DQ (с блокировкой) см. в ветке solution. + SELECT COUNT(*) + INTO v_null_student + FROM dds.fact_flight_sales + WHERE route_sk IS NULL + OR airplane_sk IS NULL; + + IF v_null_student > 0 THEN + RAISE NOTICE + 'DQ INFO: студенческие SK (route/airplane) NULL: %. После реализации dim_routes: TRUNCATE fact → перезагрузка.', + v_null_student; + END IF; + -- Calendar: допустимо если scheduled_departure IS NULL, фейлим если > 1%. SELECT COUNT(*) INTO v_null_calendar diff --git a/sql/dm/airport_traffic_dq.sql b/sql/dm/airport_traffic_dq.sql index bb8021b..1ccf863 100644 --- a/sql/dm/airport_traffic_dq.sql +++ b/sql/dm/airport_traffic_dq.sql @@ -1,49 +1,3 @@ --- DQ проверки для витрины dm.airport_traffic. --- --- Учебные цели: --- 1. Проверка арифметических инвариантов (сумма частей должна быть равна целому). --- 2. Валидация уникальности составного ключа (зерна). - -DO $$ -DECLARE - row_count INTEGER; - duplicate_count INTEGER; - invalid_total_count INTEGER; - negative_metrics_count INTEGER; -BEGIN - -- 1. Проверка на наполненность - SELECT COUNT(*) INTO row_count FROM dm.airport_traffic; - IF row_count = 0 THEN - RAISE EXCEPTION 'DQ Error: Таблица dm.airport_traffic пуста.'; - END IF; - - -- 2. Проверка на уникальность (зерно - traffic_date + airport_sk) - SELECT COUNT(*) INTO duplicate_count - FROM ( - SELECT traffic_date, airport_sk FROM dm.airport_traffic - GROUP BY traffic_date, airport_sk HAVING COUNT(*) > 1 - ) q; - IF duplicate_count > 0 THEN - RAISE EXCEPTION 'DQ Error: В dm.airport_traffic обнаружены дубликаты по (date, airport_sk) (% шт).', duplicate_count; - END IF; - - -- 3. Проверка инварианта total = departures + arrivals - SELECT COUNT(*) INTO invalid_total_count - FROM dm.airport_traffic - WHERE total_passengers != (departures_passengers + arrivals_passengers); - - IF invalid_total_count > 0 THEN - RAISE EXCEPTION 'DQ Error: В dm.airport_traffic нарушен инвариант total_passengers = dep + arr (% строк).', invalid_total_count; - END IF; - - -- 4. Проверка на отрицательные значения - SELECT COUNT(*) INTO negative_metrics_count - FROM dm.airport_traffic - WHERE departures_flights < 0 OR arrivals_flights < 0 OR departures_revenue < 0 OR arrivals_revenue < 0; - - IF negative_metrics_count > 0 THEN - RAISE EXCEPTION 'DQ Error: В dm.airport_traffic обнаружены отрицательные метрики (% строк).', negative_metrics_count; - END IF; - - RAISE NOTICE 'DQ Success: dm.airport_traffic успешно прошла все проверки (% строк).', row_count; -END $$; +-- TODO: реализуйте проверки качества данных для dm.airport_traffic +-- Эталонную реализацию можно найти в ветке solution. +SELECT 1; diff --git a/sql/dm/airport_traffic_load.sql b/sql/dm/airport_traffic_load.sql index 214113e..f4adba4 100644 --- a/sql/dm/airport_traffic_load.sql +++ b/sql/dm/airport_traffic_load.sql @@ -1,131 +1,3 @@ --- Загрузка витрины dm.airport_traffic: инкрементальный UPSERT по датам. --- --- Учебные цели: --- 1. Обработка dual-role dimensions через UNION ALL: --- Мы превращаем один билет (факт) в два "события": вылет и прилет. --- Это позволяет собрать единую статистику аэропорта в одном проходе. --- 2. Инкремент по датам (HWM): в этой витрине метрики ограничены сутками, --- поэтому пересчитываем только те дни, где появились новые факты. --- 3. Агрегация рейсов: используем COUNT(DISTINCT flight_id) для подсчета рейсов. --- 4. Денормализация (SCD1): обновление атрибутов (город) при изменениях в измерении. - --- Шаг 1: Находим даты, затронутые новыми/измененными фактами. -CREATE TEMP TABLE tmp_traffic_affected_dates ON COMMIT DROP AS -SELECT DISTINCT cal.date_actual AS traffic_date -FROM dds.fact_flight_sales f -JOIN dds.dim_calendar cal ON f.calendar_sk = cal.calendar_sk -WHERE f._load_ts > ( - SELECT COALESCE(MAX(_load_ts), '1900-01-01'::TIMESTAMP) - FROM dm.airport_traffic -); - --- Шаг 2: Unpivot (UNION ALL) и агрегация по затронутым датам. -CREATE TEMP TABLE tmp_airport_traffic_delta ON COMMIT DROP AS -WITH raw_events AS ( - -- Роль 1: Аэропорт вылета - SELECT - cal.date_actual AS traffic_date, - f.departure_airport_sk AS airport_sk, - f.flight_id, - 'departure' AS role, - (CASE WHEN f.is_boarded THEN 1 ELSE 0 END) AS is_passenger, - f.price AS revenue - FROM dds.fact_flight_sales f - JOIN dds.dim_calendar cal ON f.calendar_sk = cal.calendar_sk - WHERE cal.date_actual IN (SELECT traffic_date FROM tmp_traffic_affected_dates) - - UNION ALL - - -- Роль 2: Аэропорт прилета - SELECT - cal.date_actual AS traffic_date, - f.arrival_airport_sk AS airport_sk, - f.flight_id, - 'arrival' AS role, - (CASE WHEN f.is_boarded THEN 1 ELSE 0 END) AS is_passenger, - f.price AS revenue - FROM dds.fact_flight_sales f - JOIN dds.dim_calendar cal ON f.calendar_sk = cal.calendar_sk - WHERE cal.date_actual IN (SELECT traffic_date FROM tmp_traffic_affected_dates) -) -SELECT - e.traffic_date, - e.airport_sk, - a.airport_bk, -- Берем из dim_airports.airport_bk - a.city, - -- Агрегация вылетов - COUNT(DISTINCT CASE WHEN e.role = 'departure' THEN e.flight_id END) AS departures_flights, - SUM(CASE WHEN e.role = 'departure' THEN e.is_passenger ELSE 0 END) AS departures_passengers, - SUM(CASE WHEN e.role = 'departure' THEN e.revenue ELSE 0 END) AS departures_revenue, - -- Агрегация прилетов - COUNT(DISTINCT CASE WHEN e.role = 'arrival' THEN e.flight_id END) AS arrivals_flights, - SUM(CASE WHEN e.role = 'arrival' THEN e.is_passenger ELSE 0 END) AS arrivals_passengers, - SUM(CASE WHEN e.role = 'arrival' THEN e.revenue ELSE 0 END) AS arrivals_revenue, - -- Итого: все пассажиры (сели + вышли в этом аэропорту) - SUM(e.is_passenger) AS total_passengers -FROM raw_events e -JOIN dds.dim_airports a ON e.airport_sk = a.airport_sk -GROUP BY e.traffic_date, e.airport_sk, a.airport_bk, a.city; - --- Шаг 3: UPDATE существующих записей (включая денормализованные атрибуты). -UPDATE dm.airport_traffic AS tgt -SET - airport_bk = src.airport_bk, - city = src.city, - departures_flights = src.departures_flights, - departures_passengers = src.departures_passengers, - departures_revenue = src.departures_revenue, - arrivals_flights = src.arrivals_flights, - arrivals_passengers = src.arrivals_passengers, - arrivals_revenue = src.arrivals_revenue, - total_passengers = src.total_passengers, - updated_at = now(), - _load_id = '{{ run_id }}', - _load_ts = now() -FROM tmp_airport_traffic_delta AS src -WHERE tgt.traffic_date = src.traffic_date - AND tgt.airport_sk = src.airport_sk - AND ( - tgt.total_passengers IS DISTINCT FROM src.total_passengers - OR tgt.departures_flights IS DISTINCT FROM src.departures_flights - OR tgt.arrivals_flights IS DISTINCT FROM src.arrivals_flights - OR tgt.city IS DISTINCT FROM src.city - OR tgt.airport_bk IS DISTINCT FROM src.airport_bk - ); - --- Шаг 4: INSERT новых записей. -INSERT INTO dm.airport_traffic ( - traffic_date, - airport_sk, - airport_bk, - city, - departures_flights, - departures_passengers, - departures_revenue, - arrivals_flights, - arrivals_passengers, - arrivals_revenue, - total_passengers, - _load_id -) -SELECT - src.traffic_date, - src.airport_sk, - src.airport_bk, - src.city, - src.departures_flights, - src.departures_passengers, - src.departures_revenue, - src.arrivals_flights, - src.arrivals_passengers, - src.arrivals_revenue, - src.total_passengers, - '{{ run_id }}' AS _load_id -FROM tmp_airport_traffic_delta AS src -WHERE NOT EXISTS ( - SELECT 1 FROM dm.airport_traffic AS tgt - WHERE tgt.traffic_date = src.traffic_date - AND tgt.airport_sk = src.airport_sk -); - -ANALYZE dm.airport_traffic; +-- TODO: реализуйте загрузку витрины dm.airport_traffic (см. ТЗ в docs/assignment/analyst_spec.md) +-- Эталонную реализацию можно найти в ветке solution. +SELECT 1; diff --git a/sql/dm/monthly_overview_dq.sql b/sql/dm/monthly_overview_dq.sql index 7ca896c..dc7d5c6 100644 --- a/sql/dm/monthly_overview_dq.sql +++ b/sql/dm/monthly_overview_dq.sql @@ -1,52 +1,3 @@ --- DQ проверки для витрины dm.monthly_overview. --- --- Учебные цели: --- 1. Валидация логики двухуровневой агрегации: load factor не должен превышать 1.0 --- (в отличие от route_performance, где допускалась погрешность из-за упрощенного расчета). --- 2. Базовые санити-проверки для дат (month BETWEEN 1 AND 12). - -DO $$ -DECLARE - row_count INTEGER; - duplicate_count INTEGER; - invalid_month_count INTEGER; - invalid_load_factor_count INTEGER; -BEGIN - -- 1. Проверка на наполненность - SELECT COUNT(*) INTO row_count FROM dm.monthly_overview; - IF row_count = 0 THEN - RAISE EXCEPTION 'DQ Error: Таблица dm.monthly_overview пуста.'; - END IF; - - -- 2. Проверка на уникальность (зерно - год + месяц + самолет) - SELECT COUNT(*) INTO duplicate_count - FROM ( - SELECT year_actual, month_actual, airplane_sk FROM dm.monthly_overview - GROUP BY year_actual, month_actual, airplane_sk HAVING COUNT(*) > 1 - ) q; - IF duplicate_count > 0 THEN - RAISE EXCEPTION 'DQ Error: В dm.monthly_overview обнаружены дубликаты по (year, month, airplane_sk) (% шт).', duplicate_count; - END IF; - - -- 3. Санити-проверка месяца - SELECT COUNT(*) INTO invalid_month_count - FROM dm.monthly_overview - WHERE month_actual < 1 OR month_actual > 12; - - IF invalid_month_count > 0 THEN - RAISE EXCEPTION 'DQ Error: В dm.monthly_overview обнаружены некорректные месяцы (% строк).', invalid_month_count; - END IF; - - -- 4. Проверка точного load_factor - -- В этой витрине мы считаем его честно (по рейсам), поэтому он строго <= 1.0 - -- (исключая экзотические случаи овербукинга, но для учебного стенда ставим жесткий лимит 1.0). - SELECT COUNT(*) INTO invalid_load_factor_count - FROM dm.monthly_overview - WHERE avg_load_factor < 0 OR avg_load_factor > 1.0; - - IF invalid_load_factor_count > 0 THEN - RAISE EXCEPTION 'DQ Error: В dm.monthly_overview обнаружен avg_load_factor вне диапазона 0..1 (% строк).', invalid_load_factor_count; - END IF; - - RAISE NOTICE 'DQ Success: dm.monthly_overview успешно прошла все проверки (% строк).', row_count; -END $$; +-- TODO: реализуйте проверки качества данных для dm.monthly_overview +-- Эталонную реализацию можно найти в ветке solution. +SELECT 1; diff --git a/sql/dm/monthly_overview_load.sql b/sql/dm/monthly_overview_load.sql index 7a0aa89..c5e073d 100644 --- a/sql/dm/monthly_overview_load.sql +++ b/sql/dm/monthly_overview_load.sql @@ -1,167 +1,3 @@ --- Загрузка витрины dm.monthly_overview: инкрементальный UPSERT по месяцам. --- --- Учебные цели: --- 1. Двухуровневая агрегация: чтобы честно посчитать среднюю заполняемость (avg_load_factor), --- мы сначала считаем её для КАЖДОГО рейса, а затем берем среднее по месяцу. --- 2. Ограничения SCD1: dim_airplanes — это SCD1-измерение. Мы берем total_seats из --- его текущего состояния. Если бы самолет переоборудовали (изменили число мест) в прошлом, --- для точного исторического расчета нам потребовалось бы SCD2-измерение. - --- Шаг 1: Находим затронутые месяцы (HWM). -CREATE TEMP TABLE tmp_monthly_affected_dates ON COMMIT DROP AS -SELECT DISTINCT cal.year_actual, cal.month_actual -FROM dds.fact_flight_sales f -JOIN dds.dim_calendar cal ON f.calendar_sk = cal.calendar_sk -WHERE f._load_ts > ( - SELECT COALESCE(MAX(_load_ts), '1900-01-01'::TIMESTAMP) - FROM dm.monthly_overview -); - --- Шаг 2: Уровень 1 - Агрегация фактов до рейса (flight_id). --- Считаем точный load_factor для каждого отдельного перелета. -CREATE TEMP TABLE tmp_flight_level_metrics ON COMMIT DROP AS -SELECT - cal.year_actual, - cal.month_actual, - f.flight_id, - f.airplane_sk, - a.total_seats, -- ВНИМАНИЕ: берется текущее значение (SCD1) - COUNT(*) AS tickets_sold_per_flight, - SUM(CASE WHEN f.is_boarded THEN 1 ELSE 0 END) AS boarded_per_flight, - SUM(f.price) AS revenue_per_flight, - -- Точный load factor конкретного рейса: посаженные пассажиры / кол-во мест - (SUM(CASE WHEN f.is_boarded THEN 1 ELSE 0 END)::NUMERIC / NULLIF(a.total_seats, 0)) AS flight_load_factor -FROM dds.fact_flight_sales f -JOIN dds.dim_calendar cal ON f.calendar_sk = cal.calendar_sk -JOIN dds.dim_airplanes a ON f.airplane_sk = a.airplane_sk --- Фильтруем только затронутые месяцы -WHERE EXISTS ( - SELECT 1 FROM tmp_monthly_affected_dates tad - WHERE tad.year_actual = cal.year_actual AND tad.month_actual = cal.month_actual -) -GROUP BY cal.year_actual, cal.month_actual, f.flight_id, f.airplane_sk, a.total_seats; - --- Шаг 3: Уровень 2 - Агрегация от рейсов до месяцев. --- Плюс собираем уникальные метрики (пассажиры, маршруты) напрямую из фактов. -CREATE TEMP TABLE tmp_monthly_overview_delta ON COMMIT DROP AS -WITH flight_aggs AS ( - SELECT - year_actual, - month_actual, - airplane_sk, - COUNT(DISTINCT flight_id) AS total_flights, - SUM(tickets_sold_per_flight) AS total_tickets, - SUM(boarded_per_flight) AS total_boarded, - SUM(revenue_per_flight) AS total_revenue, - -- Честное среднее от рейсовых показателей - AVG(flight_load_factor) AS avg_load_factor - FROM tmp_flight_level_metrics - GROUP BY year_actual, month_actual, airplane_sk -), -unique_aggs AS ( - SELECT - cal.year_actual, - cal.month_actual, - f.airplane_sk, - COUNT(DISTINCT r.route_bk) AS unique_routes, -- Считаем по BK (SCD2) - COUNT(DISTINCT f.passenger_sk) AS unique_passengers - FROM dds.fact_flight_sales f - JOIN dds.dim_calendar cal ON f.calendar_sk = cal.calendar_sk - JOIN dds.dim_routes r ON f.route_sk = r.route_sk - WHERE EXISTS ( - SELECT 1 FROM tmp_monthly_affected_dates tad - WHERE tad.year_actual = cal.year_actual AND tad.month_actual = cal.month_actual - ) - GROUP BY cal.year_actual, cal.month_actual, f.airplane_sk -) -SELECT - fa.year_actual, - fa.month_actual, - fa.airplane_sk, - a.airplane_bk, - a.model AS airplane_model, - a.total_seats, - fa.total_flights, - fa.total_tickets, - fa.total_boarded, - fa.total_revenue, - (fa.total_revenue / NULLIF(fa.total_tickets, 0))::NUMERIC(10,2) AS avg_ticket_price, - fa.avg_load_factor::NUMERIC(5,4) AS avg_load_factor, - ua.unique_routes, - ua.unique_passengers -FROM flight_aggs fa -JOIN unique_aggs ua ON fa.year_actual = ua.year_actual AND fa.month_actual = ua.month_actual AND fa.airplane_sk = ua.airplane_sk -JOIN dds.dim_airplanes a ON fa.airplane_sk = a.airplane_sk; - --- Шаг 4: UPDATE существующих записей. -UPDATE dm.monthly_overview AS tgt -SET - airplane_bk = src.airplane_bk, - airplane_model = src.airplane_model, - total_seats = src.total_seats, - total_flights = src.total_flights, - total_tickets = src.total_tickets, - total_boarded = src.total_boarded, - total_revenue = src.total_revenue, - avg_ticket_price = src.avg_ticket_price, - avg_load_factor = src.avg_load_factor, - unique_routes = src.unique_routes, - unique_passengers = src.unique_passengers, - updated_at = now(), - _load_id = '{{ run_id }}', - _load_ts = now() -FROM tmp_monthly_overview_delta AS src -WHERE tgt.year_actual = src.year_actual - AND tgt.month_actual = src.month_actual - AND tgt.airplane_sk = src.airplane_sk - AND ( - tgt.total_flights IS DISTINCT FROM src.total_flights - OR tgt.total_revenue IS DISTINCT FROM src.total_revenue - OR tgt.total_boarded IS DISTINCT FROM src.total_boarded - OR tgt.unique_passengers IS DISTINCT FROM src.unique_passengers - OR tgt.airplane_model IS DISTINCT FROM src.airplane_model - ); - --- Шаг 5: INSERT новых записей. -INSERT INTO dm.monthly_overview ( - year_actual, - month_actual, - airplane_sk, - airplane_bk, - airplane_model, - total_seats, - total_flights, - total_tickets, - total_boarded, - total_revenue, - avg_ticket_price, - avg_load_factor, - unique_routes, - unique_passengers, - _load_id -) -SELECT - src.year_actual, - src.month_actual, - src.airplane_sk, - src.airplane_bk, - src.airplane_model, - src.total_seats, - src.total_flights, - src.total_tickets, - src.total_boarded, - src.total_revenue, - src.avg_ticket_price, - src.avg_load_factor, - src.unique_routes, - src.unique_passengers, - '{{ run_id }}' AS _load_id -FROM tmp_monthly_overview_delta AS src -WHERE NOT EXISTS ( - SELECT 1 FROM dm.monthly_overview AS tgt - WHERE tgt.year_actual = src.year_actual - AND tgt.month_actual = src.month_actual - AND tgt.airplane_sk = src.airplane_sk -); - -ANALYZE dm.monthly_overview; +-- TODO: реализуйте загрузку витрины dm.monthly_overview (см. ТЗ в docs/assignment/analyst_spec.md) +-- Эталонную реализацию можно найти в ветке solution. +SELECT 1; diff --git a/sql/dm/passenger_loyalty_dq.sql b/sql/dm/passenger_loyalty_dq.sql index 6e57d32..9ef7a0f 100644 --- a/sql/dm/passenger_loyalty_dq.sql +++ b/sql/dm/passenger_loyalty_dq.sql @@ -1,51 +1,3 @@ --- DQ проверки для витрины dm.passenger_loyalty. --- --- Учебные цели: --- 1. Валидация бизнес-логики дат (стаж не может быть отрицательным). --- 2. Учебная проверка ссылочной целостности (FK Integrity). - -DO $$ -DECLARE - row_count INTEGER; - duplicate_count INTEGER; - invalid_dates_count INTEGER; - orphan_keys_count INTEGER; -BEGIN - -- 1. Проверка на наполненность - SELECT COUNT(*) INTO row_count FROM dm.passenger_loyalty; - IF row_count = 0 THEN - RAISE EXCEPTION 'DQ Error: Таблица dm.passenger_loyalty пуста.'; - END IF; - - -- 2. Проверка на уникальность (зерно - passenger_sk) - SELECT COUNT(*) INTO duplicate_count - FROM ( - SELECT passenger_sk FROM dm.passenger_loyalty - GROUP BY passenger_sk HAVING COUNT(*) > 1 - ) q; - IF duplicate_count > 0 THEN - RAISE EXCEPTION 'DQ Error: В dm.passenger_loyalty обнаружены дубликаты по passenger_sk (% шт).', duplicate_count; - END IF; - - -- 3. Проверка логики дат - SELECT COUNT(*) INTO invalid_dates_count - FROM dm.passenger_loyalty - WHERE first_flight_date > last_flight_date; - - IF invalid_dates_count > 0 THEN - RAISE EXCEPTION 'DQ Error: В dm.passenger_loyalty обнаружены записи с first_date > last_date (% строк).', invalid_dates_count; - END IF; - - -- 4. Учебная проверка ссылочной целостности (FK Check) - -- В продакшене это обычно гарантируется JOIN при загрузке, но здесь мы показываем саму возможность проверки. - SELECT COUNT(*) INTO orphan_keys_count - FROM dm.passenger_loyalty tgt - LEFT JOIN dds.dim_passengers p ON tgt.passenger_sk = p.passenger_sk - WHERE p.passenger_sk IS NULL; - - IF orphan_keys_count > 0 THEN - RAISE EXCEPTION 'DQ Error: В dm.passenger_loyalty обнаружены пассажиры, отсутствующие в dim_passengers (% строк).', orphan_keys_count; - END IF; - - RAISE NOTICE 'DQ Success: dm.passenger_loyalty успешно прошла все проверки (% строк).', row_count; -END $$; +-- TODO: реализуйте проверки качества данных для dm.passenger_loyalty +-- Эталонную реализацию можно найти в ветке solution. +SELECT 1; diff --git a/sql/dm/passenger_loyalty_load.sql b/sql/dm/passenger_loyalty_load.sql index fb8e3ad..75e8427 100644 --- a/sql/dm/passenger_loyalty_load.sql +++ b/sql/dm/passenger_loyalty_load.sql @@ -1,131 +1,3 @@ --- Загрузка витрины dm.passenger_loyalty: инкрементальный UPSERT по затронутым ключам. --- --- Учебные цели: --- 1. Метод "затронутых ключей": HWM по _load_ts находит затронутые ID пассажиров, --- а затем мы ПЕРЕСЧИТЫВАЕМ всю историю именно для этого круга лиц. --- Это гарантирует точность накопительных агрегатов (total_spent, dates). --- 2. Использование DISTINCT ON (PostgreSQL-специфика): самый лаконичный способ --- найти "самое частое" (моду) в рамках группы. --- 3. Обработка NULL в фактах: фильтрация (passenger_sk IS NOT NULL). --- 4. Агрегация SCD2-измерений: при подсчете уникальных маршрутов (dim_routes) --- нужно агрегировать по BK (route_bk), т.к. один маршрут может иметь несколько SK (версий). - --- Шаг 1: Находим ID пассажиров, чьи данные изменились или добавились в фактах. -CREATE TEMP TABLE tmp_loyalty_affected_keys ON COMMIT DROP AS -SELECT DISTINCT passenger_sk -FROM dds.fact_flight_sales -WHERE _load_ts > ( - SELECT COALESCE(MAX(_load_ts), '1900-01-01'::TIMESTAMP) - FROM dm.passenger_loyalty -) -AND passenger_sk IS NOT NULL; - --- Шаг 2: Для затронутых лиц считаем ИТОГОВЫЕ агрегаты по ВСЕЙ истории фактов. -CREATE TEMP TABLE tmp_passenger_delta ON COMMIT DROP AS -WITH base_metrics AS ( - -- Агрегируем количественные метрики. - -- JOIN dds.dim_routes нужен для подсчета УНИКАЛЬНЫХ маршрутов по BK (т.к. dim_routes - SCD2). - SELECT - f.passenger_sk, - COUNT(DISTINCT f.book_ref) AS total_bookings, - COUNT(*) AS total_flights, - SUM(CASE WHEN f.is_boarded THEN 1 ELSE 0 END) AS total_boarded, - SUM(f.price) AS total_spent, - COUNT(DISTINCT r.route_bk) AS unique_routes, -- Агрегация по BK (бизнес-ключу) маршрута - MIN(cal.date_actual) AS first_flight_date, - MAX(cal.date_actual) AS last_flight_date - FROM dds.fact_flight_sales f - JOIN dds.dim_calendar cal ON f.calendar_sk = cal.calendar_sk - JOIN dds.dim_routes r ON f.route_sk = r.route_sk - WHERE f.passenger_sk IN (SELECT passenger_sk FROM tmp_loyalty_affected_keys) - GROUP BY f.passenger_sk -), -fare_modes AS ( - -- Находим самый частый класс обслуживания для каждого пассажира. - -- DISTINCT ON при ничьей выбирает произвольный вариант из топ-результатов. - SELECT DISTINCT ON (f.passenger_sk) - f.passenger_sk, - tar.fare_conditions AS favorite_fare_conditions - FROM dds.fact_flight_sales f - JOIN dds.dim_tariffs tar ON f.tariff_sk = tar.tariff_sk - WHERE f.passenger_sk IN (SELECT passenger_sk FROM tmp_loyalty_affected_keys) - GROUP BY f.passenger_sk, tar.fare_conditions - ORDER BY f.passenger_sk, COUNT(*) DESC -) -SELECT - bm.*, - fm.favorite_fare_conditions, - p.passenger_id AS passenger_bk, - p.passenger_name, - (bm.last_flight_date - bm.first_flight_date) AS days_as_customer -FROM base_metrics bm -JOIN dds.dim_passengers p ON bm.passenger_sk = p.passenger_sk -JOIN fare_modes fm ON bm.passenger_sk = fm.passenger_sk; - --- Шаг 3: UPDATE существующих записей. -UPDATE dm.passenger_loyalty AS tgt -SET - passenger_name = src.passenger_name, - total_bookings = src.total_bookings, - total_flights = src.total_flights, - total_boarded = src.total_boarded, - total_spent = src.total_spent, - avg_ticket_price = (src.total_spent / NULLIF(src.total_flights, 0))::NUMERIC(10,2), - favorite_fare_conditions = src.favorite_fare_conditions, - unique_routes = src.unique_routes, - first_flight_date = src.first_flight_date, - last_flight_date = src.last_flight_date, - days_as_customer = src.days_as_customer, - updated_at = now(), - _load_id = '{{ run_id }}', - _load_ts = now() -FROM tmp_passenger_delta AS src -WHERE tgt.passenger_sk = src.passenger_sk - AND ( - -- Обновляем только если что-то реально изменилось - tgt.total_flights IS DISTINCT FROM src.total_flights - OR tgt.total_boarded IS DISTINCT FROM src.total_boarded - OR tgt.total_spent IS DISTINCT FROM src.total_spent - OR tgt.last_flight_date IS DISTINCT FROM src.last_flight_date - OR tgt.passenger_name IS DISTINCT FROM src.passenger_name - ); - --- Шаг 4: INSERT новых пассажиров. -INSERT INTO dm.passenger_loyalty ( - passenger_sk, - passenger_bk, - passenger_name, - total_bookings, - total_flights, - total_boarded, - total_spent, - avg_ticket_price, - favorite_fare_conditions, - unique_routes, - first_flight_date, - last_flight_date, - days_as_customer, - _load_id -) -SELECT - src.passenger_sk, - src.passenger_bk, - src.passenger_name, - src.total_bookings, - src.total_flights, - src.total_boarded, - src.total_spent, - (src.total_spent / NULLIF(src.total_flights, 0))::NUMERIC(10,2), - src.favorite_fare_conditions, - src.unique_routes, - src.first_flight_date, - src.last_flight_date, - src.days_as_customer, - '{{ run_id }}' AS _load_id -FROM tmp_passenger_delta AS src -WHERE NOT EXISTS ( - SELECT 1 FROM dm.passenger_loyalty AS tgt - WHERE tgt.passenger_sk = src.passenger_sk -); - -ANALYZE dm.passenger_loyalty; +-- TODO: реализуйте загрузку витрины dm.passenger_loyalty (см. ТЗ в docs/assignment/analyst_spec.md) +-- Эталонную реализацию можно найти в ветке solution. +SELECT 1; diff --git a/sql/dm/route_performance_dq.sql b/sql/dm/route_performance_dq.sql index 6c6c509..3b52dfc 100644 --- a/sql/dm/route_performance_dq.sql +++ b/sql/dm/route_performance_dq.sql @@ -1,59 +1,3 @@ --- DQ проверки для витрины dm.route_performance. --- --- Учебные цели: --- 1. Использование PL/pgSQL блоков (DO $$ ...) для кастомных проверок. --- 2. Валидация бизнес-инвариантов после загрузки Full Rebuild. - -DO $$ -DECLARE - row_count INTEGER; - duplicate_count INTEGER; - invalid_metrics_count INTEGER; - null_attributes_count INTEGER; -BEGIN - -- 1. Проверка на наполненность - SELECT COUNT(*) INTO row_count FROM dm.route_performance; - IF row_count = 0 THEN - RAISE EXCEPTION 'DQ Error: Таблица dm.route_performance пуста после загрузки.'; - END IF; - - -- 2. Проверка на уникальность бизнес-ключа - SELECT COUNT(*) INTO duplicate_count - FROM ( - SELECT route_bk FROM dm.route_performance - GROUP BY route_bk HAVING COUNT(*) > 1 - ) AS q; - IF duplicate_count > 0 THEN - RAISE EXCEPTION 'DQ Error: В dm.route_performance обнаружены дубликаты по route_bk (% шт).', duplicate_count; - END IF; - - -- 3. Проверка бизнес-метрик (инварианты) - -- Load factor не может быть отрицательным или физически невозможным (например, более 200%) - SELECT COUNT(*) INTO invalid_metrics_count - FROM dm.route_performance - WHERE avg_load_factor < 0 OR avg_load_factor > 2.0; - - IF invalid_metrics_count > 0 THEN - RAISE EXCEPTION 'DQ Error: В dm.route_performance обнаружены некорректные показатели load factor (% строк).', invalid_metrics_count; - END IF; - - -- 4. Проверка обязательных атрибутов на NULL - SELECT COUNT(*) INTO null_attributes_count - FROM dm.route_performance - WHERE departure_city IS NULL OR arrival_city IS NULL OR airplane_model IS NULL; - - IF null_attributes_count > 0 THEN - RAISE EXCEPTION 'DQ Error: В dm.route_performance обнаружены NULL-значения в обязательных полях (% строк).', null_attributes_count; - END IF; - - -- 5. Проверка: посаженных не может быть больше, чем купивших билет - SELECT COUNT(*) INTO invalid_metrics_count - FROM dm.route_performance - WHERE total_boarded > total_tickets; - - IF invalid_metrics_count > 0 THEN - RAISE EXCEPTION 'DQ Error: В dm.route_performance total_boarded > total_tickets (% строк).', invalid_metrics_count; - END IF; - - RAISE NOTICE 'DQ Success: dm.route_performance успешно прошла все проверки (% строк).', row_count; -END $$; +-- TODO: реализуйте проверки качества данных для dm.route_performance +-- Эталонную реализацию можно найти в ветке solution. +SELECT 1; diff --git a/sql/dm/route_performance_load.sql b/sql/dm/route_performance_load.sql index c8fd3c3..49f137f 100644 --- a/sql/dm/route_performance_load.sql +++ b/sql/dm/route_performance_load.sql @@ -1,83 +1,3 @@ --- Загрузка витрины dm.route_performance: Full Rebuild. --- --- Учебные цели: --- 1. Демонстрация TRUNCATE + INSERT. Таблица маленькая (~1000 строк), --- пересоздать её с нуля дешевле, чем вычислять дельту. --- 2. Обработка SCD2: агрегируем факты по бизнес-ключу route_bk, --- чтобы собрать данные со всех исторических версий маршрута. --- 3. Денормализация атрибутов из ТЕКУЩЕЙ (актуальной) версии измерения. - --- Шаг 1: Очистка таблицы (AO Column Store не поддерживает эффективный DELETE/UPDATE). -TRUNCATE dm.route_performance; - --- Шаг 2: Агрегация всех фактов по бизнес-ключу маршрута. --- Используем временную таблицу для удобства и читаемости. -CREATE TEMP TABLE tmp_route_metrics ON COMMIT DROP AS -SELECT - r.route_bk, - COUNT(DISTINCT f.flight_id) AS total_flights, - COUNT(*) AS total_tickets, - SUM(CASE WHEN f.is_boarded THEN 1 ELSE 0 END) AS total_boarded, - SUM(f.price) AS total_revenue, - MIN(cal.date_actual) AS first_flight_date, - MAX(cal.date_actual) AS last_flight_date, - -- Доля посаженных пассажиров среди всех проданных билетов маршрута - -- (упрощение: не взвешено по рейсам) - AVG(CASE WHEN f.is_boarded THEN 1 ELSE 0 END::NUMERIC) AS avg_boarding_rate -FROM dds.fact_flight_sales f -JOIN dds.dim_routes r ON f.route_sk = r.route_sk -JOIN dds.dim_calendar cal ON f.calendar_sk = cal.calendar_sk -GROUP BY r.route_bk; - --- Шаг 3: Объединяем агрегаты с АКТУАЛЬНЫМИ атрибутами маршрута. --- Учебный комментарий: Благодаря денормализации dim_routes — один JOIN вместо четырёх. -INSERT INTO dm.route_performance ( - route_bk, - route_sk, - departure_airport_bk, - departure_city, - arrival_airport_bk, - arrival_city, - airplane_bk, - airplane_model, - total_seats, - total_flights, - total_tickets, - total_boarded, - total_revenue, - avg_ticket_price, - avg_boarding_rate, - avg_load_factor, - first_flight_date, - last_flight_date, - _load_id -) -SELECT - m.route_bk, - r_curr.route_sk, - r_curr.departure_airport AS departure_airport_bk, - r_curr.departure_city, - r_curr.arrival_airport AS arrival_airport_bk, - r_curr.arrival_city, - r_curr.airplane_code AS airplane_bk, - r_curr.airplane_model, - r_curr.total_seats, - m.total_flights, - m.total_tickets, - m.total_boarded, - m.total_revenue, - -- Средний чек - (m.total_revenue / NULLIF(m.total_tickets, 0))::NUMERIC(10,2), - m.avg_boarding_rate::NUMERIC(5,4), - -- Load Factor: общее кол-во посадок / (кол-во рейсов * мест в самолете) - ROUND( - m.total_boarded::NUMERIC / NULLIF(m.total_flights * r_curr.total_seats, 0), - 4 - ) AS avg_load_factor, - m.first_flight_date, - m.last_flight_date, - '{{ run_id }}' AS _load_id -FROM tmp_route_metrics m -JOIN dds.dim_routes r_curr ON m.route_bk = r_curr.route_bk AND r_curr.valid_to IS NULL; - -ANALYZE dm.route_performance; +-- TODO: реализуйте загрузку витрины dm.route_performance (см. ТЗ в docs/assignment/analyst_spec.md) +-- Эталонную реализацию можно найти в ветке solution. +SELECT 1; diff --git a/sql/ods/airplanes_dq.sql b/sql/ods/airplanes_dq.sql index ee9b081..c9d1889 100644 --- a/sql/ods/airplanes_dq.sql +++ b/sql/ods/airplanes_dq.sql @@ -1,99 +1,3 @@ --- DQ для ODS airplanes. - -DO $$ -DECLARE - v_batch_id TEXT := '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text; - v_stg_batch_count BIGINT; - v_dup_count BIGINT; - v_missing_keys_count BIGINT; - v_extra_keys_count BIGINT; - v_null_count BIGINT; -BEGIN - -- Для snapshot-справочников пустой батч — ошибка. - SELECT COUNT(*) - INTO v_stg_batch_count - FROM stg.airplanes - WHERE _load_id = v_batch_id; - - IF v_stg_batch_count = 0 THEN - RAISE EXCEPTION - 'DQ FAILED: _load_id=% для stg.airplanes пустой. Проверьте загрузку STG и PXF.', - v_batch_id; - END IF; - - -- В ODS не должно быть дублей по бизнес-ключу. - SELECT COUNT(*) - COUNT(DISTINCT airplane_code) - INTO v_dup_count - FROM ods.airplanes; - - IF v_dup_count <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: в ods.airplanes найдены дубликаты airplane_code: %', - v_dup_count; - END IF; - - -- Все ключи из STG текущего батча должны присутствовать в ODS. - SELECT COUNT(*) - INTO v_missing_keys_count - FROM ( - SELECT DISTINCT airplane_code - FROM stg.airplanes - WHERE _load_id = v_batch_id - ) AS s - WHERE NOT EXISTS ( - SELECT 1 - FROM ods.airplanes AS o - WHERE o.airplane_code = s.airplane_code - ); - - IF v_missing_keys_count <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: в ods.airplanes отсутствуют ключи из stg.airplanes (_load_id=%): %', - v_batch_id, - v_missing_keys_count; - END IF; - - -- В ODS не должно быть лишних ключей, которых нет в snapshot текущего батча. - SELECT COUNT(*) - INTO v_extra_keys_count - FROM ods.airplanes AS o - WHERE NOT EXISTS ( - SELECT 1 - FROM ( - SELECT DISTINCT airplane_code - FROM stg.airplanes - WHERE _load_id = v_batch_id - ) AS s - WHERE s.airplane_code = o.airplane_code - ); - - IF v_extra_keys_count <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: в ods.airplanes найдены лишние ключи вне stg _load_id=%: %', - v_batch_id, - v_extra_keys_count; - END IF; - - -- Обязательные поля в ODS. - SELECT COUNT(*) - INTO v_null_count - FROM ods.airplanes - WHERE airplane_code IS NULL - OR airplane_code = '' - OR model IS NULL - OR model = '' - OR _load_id IS NULL - OR _load_id = '' - OR _load_ts IS NULL; - - IF v_null_count <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: в ods.airplanes найдены NULL/пустые обязательные поля: %', - v_null_count; - END IF; - - RAISE NOTICE - 'DQ PASSED: ods.airplanes ок (_load_id=%): stg_batch_rows=%', - v_batch_id, - v_stg_batch_count; -END $$; +-- TODO: реализуйте проверки качества данных для ods.airplanes +-- Эталонную реализацию можно найти в ветке solution. +SELECT 1; diff --git a/sql/ods/airplanes_load.sql b/sql/ods/airplanes_load.sql index 48c1b19..4b06543 100644 --- a/sql/ods/airplanes_load.sql +++ b/sql/ods/airplanes_load.sql @@ -1,39 +1,4 @@ --- Загрузка ODS по airplanes: Полная перезагрузка (TRUNCATE + INSERT). --- Почему: для справочников-снимков в Greenplum на AO-таблицах --- эффективнее перетереть данные целиком, чем делать медленный UPDATE. - -TRUNCATE TABLE ods.airplanes; - -INSERT INTO ods.airplanes ( - airplane_code, - model, - range_km, - speed_kmh, - _load_id, - _load_ts -) -WITH src AS ( - -- Выбираем последний снимок из STG для текущего батча - SELECT - s.airplane_code, - s.model::json->>'ru' AS model, - NULLIF(s.range, '')::INTEGER AS range_km, - NULLIF(s.speed, '')::INTEGER AS speed_kmh, - ROW_NUMBER() OVER ( - PARTITION BY s.airplane_code - ORDER BY s._load_ts DESC, s.event_ts DESC NULLS LAST - ) AS rn - FROM stg.airplanes AS s - WHERE s._load_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text -) -SELECT - s.airplane_code, - s.model, - s.range_km, - s.speed_kmh, - '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, - now() -FROM src AS s -WHERE s.rn = 1; - -ANALYZE ods.airplanes; +-- TODO: реализуйте загрузку ods.airplanes (см. ТЗ в docs/assignment/analyst_spec.md) +-- Паттерн: TRUNCATE + INSERT (snapshot-справочник на AO-таблице). +-- Эталонную реализацию можно найти в ветке solution. +SELECT 1; diff --git a/sql/ods/routes_dq.sql b/sql/ods/routes_dq.sql index 4dcf325..784e5c3 100644 --- a/sql/ods/routes_dq.sql +++ b/sql/ods/routes_dq.sql @@ -140,21 +140,25 @@ BEGIN v_orphan_arrival_count; END IF; - -- Ссылочная целостность: airplane_code должен существовать в ods.airplanes. - SELECT COUNT(*) - INTO v_orphan_airplane_count - FROM ods.routes AS r - WHERE NOT EXISTS ( - SELECT 1 - FROM ods.airplanes AS a - WHERE a.airplane_code = r.airplane_code - ); - - IF v_orphan_airplane_count <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: в ods.routes найдены строки с невалидным airplane_code: %', - v_orphan_airplane_count; - END IF; + -- Проверка RI airplane_code → ods.airplanes закомментирована, + -- т.к. таблица ods.airplanes реализуется студентом. + -- После реализации — раскомментируйте этот блок. + -- Полную версию см. в ветке solution. + -- + -- SELECT COUNT(*) + -- INTO v_orphan_airplane_count + -- FROM ods.routes AS r + -- WHERE NOT EXISTS ( + -- SELECT 1 + -- FROM ods.airplanes AS a + -- WHERE a.airplane_code = r.airplane_code + -- ); + -- + -- IF v_orphan_airplane_count <> 0 THEN + -- RAISE EXCEPTION + -- 'DQ FAILED: в ods.routes найдены строки с невалидным airplane_code: %', + -- v_orphan_airplane_count; + -- END IF; RAISE NOTICE 'DQ PASSED: ods.routes ок (_load_id=%): stg_batch_rows=%', diff --git a/sql/ods/seats_dq.sql b/sql/ods/seats_dq.sql index e8c4111..95c1591 100644 --- a/sql/ods/seats_dq.sql +++ b/sql/ods/seats_dq.sql @@ -1,125 +1,3 @@ --- DQ для ODS seats. - -DO $$ -DECLARE - v_batch_id TEXT := '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text; - v_stg_batch_count BIGINT; - v_dup_count BIGINT; - v_missing_keys_count BIGINT; - v_extra_keys_count BIGINT; - v_null_count BIGINT; - v_orphan_airplane_count BIGINT; -BEGIN - -- Для snapshot-справочников пустой батч — ошибка. - SELECT COUNT(*) - INTO v_stg_batch_count - FROM stg.seats - WHERE _load_id = v_batch_id; - - IF v_stg_batch_count = 0 THEN - RAISE EXCEPTION - 'DQ FAILED: _load_id=% для stg.seats пустой. Проверьте загрузку STG и PXF.', - v_batch_id; - END IF; - - -- В ODS не должно быть дублей по составному бизнес-ключу. - SELECT COUNT(*) - INTO v_dup_count - FROM ( - SELECT airplane_code, seat_no - FROM ods.seats - GROUP BY airplane_code, seat_no - HAVING COUNT(*) > 1 - ) AS d; - - IF v_dup_count <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: в ods.seats найдены дубликаты (airplane_code, seat_no): %', - v_dup_count; - END IF; - - -- Все ключи из STG текущего батча должны присутствовать в ODS. - SELECT COUNT(*) - INTO v_missing_keys_count - FROM ( - SELECT DISTINCT airplane_code, seat_no - FROM stg.seats - WHERE _load_id = v_batch_id - ) AS s - WHERE NOT EXISTS ( - SELECT 1 - FROM ods.seats AS o - WHERE o.airplane_code = s.airplane_code - AND o.seat_no = s.seat_no - ); - - IF v_missing_keys_count <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: в ods.seats отсутствуют ключи из stg.seats (_load_id=%): %', - v_batch_id, - v_missing_keys_count; - END IF; - - -- В ODS не должно быть лишних ключей, которых нет в snapshot текущего батча. - SELECT COUNT(*) - INTO v_extra_keys_count - FROM ods.seats AS o - WHERE NOT EXISTS ( - SELECT 1 - FROM ( - SELECT DISTINCT airplane_code, seat_no - FROM stg.seats - WHERE _load_id = v_batch_id - ) AS s - WHERE s.airplane_code = o.airplane_code - AND s.seat_no = o.seat_no - ); - - IF v_extra_keys_count <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: в ods.seats найдены лишние ключи вне stg _load_id=%: %', - v_batch_id, - v_extra_keys_count; - END IF; - - -- Обязательные поля в ODS. - SELECT COUNT(*) - INTO v_null_count - FROM ods.seats - WHERE airplane_code IS NULL - OR airplane_code = '' - OR seat_no IS NULL - OR seat_no = '' - OR fare_conditions IS NULL - OR fare_conditions = '' - OR _load_id IS NULL - OR _load_id = '' - OR _load_ts IS NULL; - - IF v_null_count <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: в ods.seats найдены NULL/пустые обязательные поля: %', - v_null_count; - END IF; - - -- Ссылочная целостность: airplane_code должен существовать в ods.airplanes. - SELECT COUNT(*) - INTO v_orphan_airplane_count - FROM ods.seats AS s - WHERE NOT EXISTS ( - SELECT 1 - FROM ods.airplanes AS a - WHERE a.airplane_code = s.airplane_code - ); - - IF v_orphan_airplane_count <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: в ods.seats найдены строки с невалидным airplane_code: %', - v_orphan_airplane_count; - END IF; - - RAISE NOTICE - 'DQ PASSED: ods.seats ок (_load_id=%): stg_batch_rows=%', - v_batch_id, - v_stg_batch_count; -END $$; +-- TODO: реализуйте проверки качества данных для ods.seats +-- Эталонную реализацию можно найти в ветке solution. +SELECT 1; diff --git a/sql/ods/seats_load.sql b/sql/ods/seats_load.sql index 4d3d091..7d74fc1 100644 --- a/sql/ods/seats_load.sql +++ b/sql/ods/seats_load.sql @@ -1,36 +1,4 @@ --- Загрузка ODS по seats: Полная перезагрузка (TRUNCATE + INSERT). --- Почему: для справочников-снимков в Greenplum на AO-таблицах --- эффективнее перетереть данные целиком, чем делать медленный UPDATE. - -TRUNCATE TABLE ods.seats; - -INSERT INTO ods.seats ( - airplane_code, - seat_no, - fare_conditions, - _load_id, - _load_ts -) -WITH src AS ( - -- Выбираем последний снимок из STG для текущего батча - SELECT - s.airplane_code, - s.seat_no, - s.fare_conditions, - ROW_NUMBER() OVER ( - PARTITION BY s.airplane_code, s.seat_no - ORDER BY s._load_ts DESC, s.event_ts DESC NULLS LAST - ) AS rn - FROM stg.seats AS s - WHERE s._load_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text -) -SELECT - s.airplane_code, - s.seat_no, - s.fare_conditions, - '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, - now() -FROM src AS s -WHERE s.rn = 1; - -ANALYZE ods.seats; +-- TODO: реализуйте загрузку ods.seats (см. ТЗ в docs/assignment/analyst_spec.md) +-- Паттерн: TRUNCATE + INSERT (snapshot-справочник на AO-таблице). +-- Эталонную реализацию можно найти в ветке solution. +SELECT 1; diff --git a/tests/test_dags_smoke.py b/tests/test_dags_smoke.py index f88c68c..cfe2555 100644 --- a/tests/test_dags_smoke.py +++ b/tests/test_dags_smoke.py @@ -239,8 +239,9 @@ def test_bookings_to_gp_ods_dag_structure(): ), "airplanes не должен быть upstream для airports" # Барьеры по данным. + # На main routes не зависит от airplanes (ods.airplanes — студенческая заглушка). + # На ветке solution: _assert_reachable(dag, "dq_ods_airplanes", "dq_ods_routes") _assert_reachable(dag, "dq_ods_airports", "dq_ods_routes") - _assert_reachable(dag, "dq_ods_airplanes", "dq_ods_routes") _assert_reachable(dag, "dq_ods_airplanes", "dq_ods_seats") _assert_reachable(dag, "dq_ods_routes", "dq_ods_flights") _assert_reachable(dag, "dq_ods_flights", "dq_ods_segments") diff --git a/tests/test_ods_sql_contract.py b/tests/test_ods_sql_contract.py index 41726ff..a0b2471 100644 --- a/tests/test_ods_sql_contract.py +++ b/tests/test_ods_sql_contract.py @@ -3,7 +3,9 @@ from __future__ import annotations from pathlib import Path PROJECT_ROOT = Path(__file__).resolve().parents[1] -SNAPSHOT_ENTITIES = ("airports", "airplanes", "routes", "seats") +# На main airplanes/seats — заглушки (SELECT 1;), не содержат TRUNCATE-паттерн. +# Полный список на ветке solution: ("airports", "airplanes", "routes", "seats") +SNAPSHOT_ENTITIES = ("airports", "routes") def _read(path: str) -> str: