From 7abeb67fba4fffb8a161ee02ca9a7ce93e8bb720 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sun, 18 Jan 2026 11:04:46 +0300 Subject: [PATCH] =?UTF-8?q?=D0=93=D0=B5=D0=BD=D0=B5=D1=80=D0=B0=D1=86?= =?UTF-8?q?=D0=B8=D1=8F=20dds=20=D1=81=D0=BB=D0=BE=D1=8F=20=D0=BF=D0=BE=20?= =?UTF-8?q?=D0=A2=D0=97=20-=20=D0=B1=D0=B5=D0=B7=20=D1=82=D0=B5=D1=81?= =?UTF-8?q?=D1=82=D0=BE=D0=B2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- airflow/dags/bookings_stg_ddl.py | 61 ++++++++++++- airflow/dags/bookings_to_gp_stage.py | 107 +++++++++++++++++++++- docs/internal/db_schema.md | 127 ++++++++++++++++++++++++++- sql/ddl_gp.sql | 11 +++ sql/stg/airplanes_ddl.sql | 37 ++++++++ sql/stg/airplanes_dq.sql | 67 ++++++++++++++ sql/stg/airplanes_load.sql | 32 +++++++ sql/stg/airports_ddl.sql | 40 +++++++++ sql/stg/airports_dq.sql | 69 +++++++++++++++ sql/stg/airports_load.sql | 36 ++++++++ sql/stg/boarding_passes_ddl.sql | 38 ++++++++ sql/stg/boarding_passes_dq.sql | 99 +++++++++++++++++++++ sql/stg/boarding_passes_load.sql | 36 ++++++++ sql/stg/flights_ddl.sql | 42 +++++++++ sql/stg/flights_dq.sql | 95 ++++++++++++++++++++ sql/stg/flights_load.sql | 48 ++++++++++ sql/stg/routes_ddl.sql | 44 ++++++++++ sql/stg/routes_dq.sql | 116 ++++++++++++++++++++++++ sql/stg/routes_load.sql | 41 +++++++++ sql/stg/seats_ddl.sql | 34 +++++++ sql/stg/seats_dq.sql | 84 ++++++++++++++++++ sql/stg/seats_load.sql | 31 +++++++ sql/stg/segments_ddl.sql | 36 ++++++++ sql/stg/segments_dq.sql | 113 ++++++++++++++++++++++++ sql/stg/segments_load.sql | 45 ++++++++++ tests/test_dags_smoke.py | 82 +++++++++++++++++ 26 files changed, 1566 insertions(+), 5 deletions(-) create mode 100644 sql/stg/airplanes_ddl.sql create mode 100644 sql/stg/airplanes_dq.sql create mode 100644 sql/stg/airplanes_load.sql create mode 100644 sql/stg/airports_ddl.sql create mode 100644 sql/stg/airports_dq.sql create mode 100644 sql/stg/airports_load.sql create mode 100644 sql/stg/boarding_passes_ddl.sql create mode 100644 sql/stg/boarding_passes_dq.sql create mode 100644 sql/stg/boarding_passes_load.sql create mode 100644 sql/stg/flights_ddl.sql create mode 100644 sql/stg/flights_dq.sql create mode 100644 sql/stg/flights_load.sql create mode 100644 sql/stg/routes_ddl.sql create mode 100644 sql/stg/routes_dq.sql create mode 100644 sql/stg/routes_load.sql create mode 100644 sql/stg/seats_ddl.sql create mode 100644 sql/stg/seats_dq.sql create mode 100644 sql/stg/seats_load.sql create mode 100644 sql/stg/segments_ddl.sql create mode 100644 sql/stg/segments_dq.sql create mode 100644 sql/stg/segments_load.sql diff --git a/airflow/dags/bookings_stg_ddl.py b/airflow/dags/bookings_stg_ddl.py index 0fc3a95..f205e79 100644 --- a/airflow/dags/bookings_stg_ddl.py +++ b/airflow/dags/bookings_stg_ddl.py @@ -38,4 +38,63 @@ with DAG( sql="stg/tickets_ddl.sql", ) - apply_stg_bookings_ddl >> apply_stg_tickets_ddl + # DDL для справочников + apply_stg_airports_ddl = PostgresOperator( + task_id="apply_stg_airports_ddl", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="stg/airports_ddl.sql", + ) + + apply_stg_airplanes_ddl = PostgresOperator( + task_id="apply_stg_airplanes_ddl", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="stg/airplanes_ddl.sql", + ) + + apply_stg_routes_ddl = PostgresOperator( + task_id="apply_stg_routes_ddl", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="stg/routes_ddl.sql", + ) + + apply_stg_seats_ddl = PostgresOperator( + task_id="apply_stg_seats_ddl", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="stg/seats_ddl.sql", + ) + + # DDL для транзакционных таблиц + apply_stg_flights_ddl = PostgresOperator( + task_id="apply_stg_flights_ddl", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="stg/flights_ddl.sql", + ) + + apply_stg_segments_ddl = PostgresOperator( + task_id="apply_stg_segments_ddl", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="stg/segments_ddl.sql", + ) + + apply_stg_boarding_passes_ddl = PostgresOperator( + task_id="apply_stg_boarding_passes_ddl", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="stg/boarding_passes_ddl.sql", + ) + + # Сначала создаются справочники, затем транзакционные таблицы + ( + apply_stg_bookings_ddl + >> apply_stg_tickets_ddl + >> [ + apply_stg_airports_ddl, + apply_stg_airplanes_ddl, + apply_stg_routes_ddl, + apply_stg_seats_ddl, + ] + >> [ + apply_stg_flights_ddl, + apply_stg_segments_ddl, + apply_stg_boarding_passes_ddl, + ] + ) diff --git a/airflow/dags/bookings_to_gp_stage.py b/airflow/dags/bookings_to_gp_stage.py index 8a5e21a..d978de8 100644 --- a/airflow/dags/bookings_to_gp_stage.py +++ b/airflow/dags/bookings_to_gp_stage.py @@ -94,11 +94,116 @@ with DAG( sql="stg/tickets_dq.sql", ) + # Загрузка справочников (full load) + load_airports_to_stg = PostgresOperator( + task_id="load_airports_to_stg", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="stg/airports_load.sql", + ) + + check_airports_dq = PostgresOperator( + task_id="check_airports_dq", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="stg/airports_dq.sql", + ) + + load_airplanes_to_stg = PostgresOperator( + task_id="load_airplanes_to_stg", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="stg/airplanes_load.sql", + ) + + check_airplanes_dq = PostgresOperator( + task_id="check_airplanes_dq", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="stg/airplanes_dq.sql", + ) + + load_routes_to_stg = PostgresOperator( + task_id="load_routes_to_stg", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="stg/routes_load.sql", + ) + + check_routes_dq = PostgresOperator( + task_id="check_routes_dq", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="stg/routes_dq.sql", + ) + + load_seats_to_stg = PostgresOperator( + task_id="load_seats_to_stg", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="stg/seats_load.sql", + ) + + check_seats_dq = PostgresOperator( + task_id="check_seats_dq", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="stg/seats_dq.sql", + ) + + # Загрузка транзакций (инкремент/full snapshot) + load_flights_to_stg = PostgresOperator( + task_id="load_flights_to_stg", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="stg/flights_load.sql", + ) + + check_flights_dq = PostgresOperator( + task_id="check_flights_dq", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="stg/flights_dq.sql", + ) + + load_segments_to_stg = PostgresOperator( + task_id="load_segments_to_stg", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="stg/segments_load.sql", + ) + + check_segments_dq = PostgresOperator( + task_id="check_segments_dq", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="stg/segments_dq.sql", + ) + + load_boarding_passes_to_stg = PostgresOperator( + task_id="load_boarding_passes_to_stg", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="stg/boarding_passes_load.sql", + ) + + check_boarding_passes_dq = PostgresOperator( + task_id="check_boarding_passes_dq", + postgres_conn_id=GREENPLUM_CONN_ID, + sql="stg/boarding_passes_dq.sql", + ) + # 6. Финальный лог/сводка finish_summary = PythonOperator( task_id="finish_summary", python_callable=_finish_summary, ) + # Сначала загружаются и проверяются bookings и tickets generate_bookings_day >> load_bookings_to_stg >> check_row_counts - check_row_counts >> load_tickets_to_stg >> check_tickets_dq >> finish_summary + check_row_counts >> load_tickets_to_stg >> check_tickets_dq + + # Затем загружаются справочники + check_tickets_dq >> [ + (load_airports_to_stg >> check_airports_dq), + (load_airplanes_to_stg >> check_airplanes_dq), + (load_routes_to_stg >> check_routes_dq), + (load_seats_to_stg >> check_seats_dq), + ] + + # Затем загружаются транзакции + [check_airports_dq, check_airplanes_dq, check_routes_dq, check_seats_dq] >> [ + (load_flights_to_stg >> check_flights_dq), + (load_segments_to_stg >> check_segments_dq), + (load_boarding_passes_to_stg >> check_boarding_passes_dq), + ] + + # В конце финальный лог + [check_flights_dq, check_segments_dq, check_boarding_passes_dq] >> finish_summary diff --git a/docs/internal/db_schema.md b/docs/internal/db_schema.md index cfcc451..9d3693b 100644 --- a/docs/internal/db_schema.md +++ b/docs/internal/db_schema.md @@ -1,6 +1,6 @@ # Схема БД DWH (Bookings → Greenplum) -> **Статус:** Проект в разработке. Реализован только STG слой (частично: bookings, tickets). +> **Статус:** Проект в разработке. Реализован STG слой полностью (все 9 таблиц). ## Обзор @@ -28,7 +28,7 @@ | Слой | Статус | Реализовано | |------|--------|-------------| | **Source** | ✅ Готово | Демо-БД bookings (Postgres) | -| **STG** | ⚠️ В процессе | 2 из 9 таблиц (bookings, tickets) | +| **STG** | ✅ Готово | 9 из 9 таблиц (bookings, tickets, airports, airplanes, routes, seats, flights, segments, boarding_passes) | | **ODS** | ❌ Не реализован | Планируется | | **DDS** | ❌ Не реализован | Планируется | @@ -90,6 +90,127 @@ --- +### Детальное описание таблиц STG слоя + +#### stg.bookings (транзакции, инкремент) +- **Источник:** `bookings.bookings` (через PXF) +- **Ключ распределения:** `book_ref` +- **Бизнес-колонки:** + - `book_ref TEXT` - номер бронирования + - `book_date TEXT` - дата бронирования + - `total_amount TEXT` - общая сумма +- **Технические колонки:** `src_created_at_ts` (=book_date), `load_dttm`, `batch_id` +- **Стратегия загрузки:** Инкремент по `book_date` +- **DQ проверки:** count (окно инкремента), дубликаты book_ref, NULL обязательных полей + +#### stg.tickets (транзакции, инкремент) +- **Источник:** `bookings.tickets` (через PXF) +- **Ключ распределения:** `ticket_no` +- **Бизнес-колонки:** + - `ticket_no TEXT` - номер билета + - `book_ref TEXT` - номер бронирования + - `passenger_id TEXT` - идентификатор пассажира + - `passenger_name TEXT` - имя пассажира + - `contact_data TEXT` - контактные данные (JSONB) +- **Технические колонки:** `src_created_at_ts` (из book_date через bookings), `load_dttm`, `batch_id` +- **Стратегия загрузки:** Инкремент по `book_date` (через bookings) +- **DQ проверки:** count (окно инкремента), дубликаты ticket_no, NULL обязательных полей, ссылочная целостность + +#### stg.airports (справочник, full load) +- **Источник:** `bookings.airports_data` (через PXF) +- **Ключ распределения:** `airport_code` +- **Бизнес-колонки:** + - `airport_code TEXT` - код аэропорта + - `airport_name TEXT` - название аэропорта (из JSONB) + - `city TEXT` - город (из JSONB) + - `country TEXT` - страна (из JSONB) + - `coordinates TEXT` - координаты + - `timezone TEXT` - часовой пояс +- **Технические колонки:** `src_created_at_ts`, `load_dttm`, `batch_id` +- **Стратегия загрузки:** Full load (все строки при каждом запуске) +- **DQ проверки:** count, дубликаты airport_code, NULL обязательных полей + +#### stg.airplanes (справочник, full load) +- **Источник:** `bookings.airplanes_data` (через PXF) +- **Ключ распределения:** `airplane_code` +- **Бизнес-колонки:** + - `airplane_code TEXT` - код самолёта + - `model TEXT` - модель (из JSONB) + - `range TEXT` - дальность полёта + - `speed TEXT` - скорость +- **Технические колонки:** `src_created_at_ts`, `load_dttm`, `batch_id` +- **Стратегия загрузки:** Full load +- **DQ проверки:** count, дубликаты airplane_code, NULL обязательных полей + +#### stg.routes (справочник, full load) +- **Источник:** `bookings.routes` (через PXF) +- **Ключ распределения:** `route_no` +- **Бизнес-колонки:** + - `route_no TEXT` - номер маршрута + - `validity TEXT` - период действия (из tstzrange) + - `departure_airport TEXT` - аэропорт вылета + - `arrival_airport TEXT` - аэропорт прилёта + - `airplane_code TEXT` - код самолёта + - `days_of_week TEXT` - дни недели (из int[]) + - `scheduled_time TEXT` - плановое время + - `duration TEXT` - длительность +- **Технические колонки:** `src_created_at_ts`, `load_dttm`, `batch_id` +- **Стратегия загрузки:** Full load +- **DQ проверки:** count, дубликаты (route_no, validity), NULL обязательных полей, ссылочная целостность + +#### stg.seats (справочник, full load) +- **Источник:** `bookings.seats` (через PXF) +- **Ключ распределения:** `airplane_code` (co-location с airplanes) +- **Бизнес-колонки:** + - `airplane_code TEXT` - код самолёта + - `seat_no TEXT` - номер места + - `fare_conditions TEXT` - класс обслуживания +- **Технические колонки:** `src_created_at_ts`, `load_dttm`, `batch_id` +- **Стратегия загрузки:** Full load +- **DQ проверки:** count, дубликаты (airplane_code, seat_no), NULL обязательных полей, ссылочная целостность + +#### stg.flights (транзакции, инкремент) +- **Источник:** `bookings.flights` (через PXF) +- **Ключ распределения:** `flight_id` +- **Бизнес-колонки:** + - `flight_id TEXT` - идентификатор рейса + - `route_no TEXT` - номер маршрута + - `status TEXT` - статус + - `scheduled_departure TEXT` - плановое время вылета + - `scheduled_arrival TEXT` - плановое время прилёта + - `actual_departure TEXT` - фактическое время вылета + - `actual_arrival TEXT` - фактическое время прилёта +- **Технические колонки:** `src_created_at_ts` (=scheduled_departure), `load_dttm`, `batch_id` +- **Стратегия загрузки:** Инкремент по `scheduled_departure` +- **DQ проверки:** count (окно инкремента), дубликаты flight_id, NULL обязательных полей, ссылочная целостность + +#### stg.segments (транзакции, инкремент) +- **Источник:** `bookings.segments` (через PXF) +- **Ключ распределения:** `ticket_no` (co-location с tickets) +- **Бизнес-колонки:** + - `ticket_no TEXT` - номер билета + - `flight_id TEXT` - идентификатор рейса + - `fare_conditions TEXT` - класс обслуживания + - `price TEXT` - цена +- **Технические колонки:** `src_created_at_ts` (из book_date через tickets), `load_dttm`, `batch_id` +- **Стратегия загрузки:** Инкремент по `book_date` (через tickets) +- **DQ проверки:** count (окно инкремента), дубликаты (ticket_no, flight_id), NULL обязательных полей, ссылочная целостность + +#### stg.boarding_passes (транзакции, full snapshot) +- **Источник:** `bookings.boarding_passes` (через PXF) +- **Ключ распределения:** `ticket_no` (co-location с tickets/segments) +- **Бизнес-колонки:** + - `ticket_no TEXT` - номер билета + - `flight_id TEXT` - идентификатор рейса + - `seat_no TEXT` - номер места + - `boarding_no TEXT` - номер посадки + - `boarding_time TEXT` - время посадки +- **Технические колонки:** `src_created_at_ts` (=now()), `load_dttm`, `batch_id` +- **Стратегия загрузки:** Full snapshot (все строки при каждом запуске) +- **DQ проверки:** count, дубликаты (ticket_no, flight_id), NULL обязательных полей, ссылочная целостность + +--- + ## Полная схема потоков данных (Data Lineage) ```mermaid @@ -299,7 +420,7 @@ graph LR ## TODO -- [ ] Реализовать STG слой полностью (все 9 таблиц) +- [x] Реализовать STG слой полностью (все 9 таблиц) - [ ] Реализовать ODS слой - [ ] Реализовать DDS слой (измерения и факт) - [ ] Создать DAG для загрузки ODS diff --git a/sql/ddl_gp.sql b/sql/ddl_gp.sql index e95f0e7..40da916 100644 --- a/sql/ddl_gp.sql +++ b/sql/ddl_gp.sql @@ -26,3 +26,14 @@ FORMAT 'CUSTOM' (formatter='pxfwritable_import'); -- Здесь подключаем их через psql \i, чтобы сохранить единый входной скрипт. \i stg/bookings_ddl.sql \i stg/tickets_ddl.sql + +-- DDL для новых таблиц STG слоя (справочники) +\i stg/airports_ddl.sql +\i stg/airplanes_ddl.sql +\i stg/routes_ddl.sql +\i stg/seats_ddl.sql + +-- DDL для новых таблиц STG слоя (транзакции) +\i stg/flights_ddl.sql +\i stg/segments_ddl.sql +\i stg/boarding_passes_ddl.sql diff --git a/sql/stg/airplanes_ddl.sql b/sql/stg/airplanes_ddl.sql new file mode 100644 index 0000000..d29196d --- /dev/null +++ b/sql/stg/airplanes_ddl.sql @@ -0,0 +1,37 @@ +-- DDL для слоя STG по таблице airplanes (справочник). +-- Используется как из общего скрипта ddl_gp.sql (через \i), +-- так и может выполняться отдельно при изменении схемы. + +-- Схема stg для сырого слоя DWH. +CREATE SCHEMA IF NOT EXISTS stg; + +-- Внешняя таблица в схеме stg для чтения данных из bookings.airplanes_data через PXF. +DROP EXTERNAL TABLE IF EXISTS stg.airplanes_ext; +CREATE EXTERNAL TABLE stg.airplanes_ext ( + airplane_code TEXT, + model JSONB, + range INTEGER, + speed INTEGER +) +LOCATION ('pxf://bookings.airplanes_data?PROFILE=JDBC&SERVER=bookings-db') +FORMAT 'CUSTOM' (formatter='pxfwritable_import'); + +-- Внутренняя таблица stg.airplanes — сырой слой, все бизнес-колонки как TEXT. +CREATE TABLE IF NOT EXISTS stg.airplanes ( + airplane_code TEXT, + model TEXT, + range TEXT, + speed TEXT, + src_created_at_ts TIMESTAMP, + load_dttm TIMESTAMP NOT NULL DEFAULT now(), + batch_id TEXT +) +WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) +-- Ключ распределения: airplane_code +-- Обоснование: airplane_code — это уникальный идентификатор самолёта. +-- Использование airplane_code обеспечивает: +-- 1. Равномерное распределение данных по сегментам (airplane_code имеет высокую кардинальность) +-- 2. Co-location данных airplanes и routes при JOIN по airplane_code +-- 3. Co-location данных airplanes и seats при JOIN по airplane_code +-- 4. Оптимизацию запросов, которые фильтруют или группируют по airplane_code +DISTRIBUTED BY (airplane_code); diff --git a/sql/stg/airplanes_dq.sql b/sql/stg/airplanes_dq.sql new file mode 100644 index 0000000..61d84f9 --- /dev/null +++ b/sql/stg/airplanes_dq.sql @@ -0,0 +1,67 @@ +-- Проверки качества данных для airplanes (справочник) + +DO $$ +DECLARE + v_batch_id TEXT := '{{ run_id }}'::text; + v_src_count BIGINT; + v_stg_count BIGINT; + v_dup_count BIGINT; + v_null_count BIGINT; +BEGIN + -- Источник: считаем все строки во внешней таблице + SELECT COUNT(*) + INTO v_src_count + FROM stg.airplanes_ext; + + IF v_src_count = 0 THEN + RAISE EXCEPTION + 'В источнике airplanes_ext нет строк.'; + END IF; + + -- Считаем строки, реально вставленные в stg.airplanes в этом батче + SELECT COUNT(*) + INTO v_stg_count + FROM stg.airplanes + WHERE batch_id = v_batch_id; + + IF v_src_count <> v_stg_count THEN + RAISE EXCEPTION + 'DQ FAILED: несовпадение количества строк. Источник: %, STG: %', + v_src_count, + v_stg_count; + END IF; + + -- Проверка на дубликаты airplane_code + SELECT COUNT(*) - COUNT(DISTINCT airplane_code) + INTO v_dup_count + FROM stg.airplanes AS a + WHERE a.batch_id = v_batch_id; + + IF v_dup_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: найдены дубликаты airplane_code (batch_id=%): %', + v_batch_id, + v_dup_count; + END IF; + + -- Проверка обязательных полей: airplane_code, model + SELECT COUNT(*) + INTO v_null_count + FROM stg.airplanes AS a + WHERE a.batch_id = v_batch_id + AND (a.airplane_code IS NULL OR a.airplane_code = '' + OR a.model IS NULL OR a.model = ''); + + IF v_null_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: найдены строки с NULL в обязательных полях (airplane_code, model) (batch_id=%): %', + v_batch_id, + v_null_count; + END IF; + + RAISE NOTICE + 'DQ PASSED: airplanes ок (batch_id=%): source=% stg=%', + v_batch_id, + v_src_count, + v_stg_count; +END $$; diff --git a/sql/stg/airplanes_load.sql b/sql/stg/airplanes_load.sql new file mode 100644 index 0000000..b140add --- /dev/null +++ b/sql/stg/airplanes_load.sql @@ -0,0 +1,32 @@ +-- Загрузка всех строк из stg.airplanes_ext в stg.airplanes (full load). +-- Используем batch_id для отслеживания загрузки. + +INSERT INTO stg.airplanes ( + airplane_code, + model, + range, + speed, + src_created_at_ts, + load_dttm, + batch_id +) +SELECT + ext.airplane_code::text, + ext.model::text, + ext.range::text, + ext.speed::text, + now()::timestamp, + now()::timestamp, + '{{ run_id }}'::text +FROM stg.airplanes_ext AS ext +WHERE NOT EXISTS ( + -- Защита от дублей в рамках одного batch_id + SELECT 1 + FROM stg.airplanes AS a + WHERE a.batch_id = '{{ run_id }}'::text + AND a.airplane_code = ext.airplane_code::text +); + +-- Обновляем статистику для оптимизатора Greenplum +-- Это критично для корректной работы оптимизатора и выбора оптимального плана выполнения +ANALYZE stg.airplanes; diff --git a/sql/stg/airports_ddl.sql b/sql/stg/airports_ddl.sql new file mode 100644 index 0000000..3f27dd2 --- /dev/null +++ b/sql/stg/airports_ddl.sql @@ -0,0 +1,40 @@ +-- DDL для слоя STG по таблице airports (справочник). +-- Используется как из общего скрипта ddl_gp.sql (через \i), +-- так и может выполняться отдельно при изменении схемы. + +-- Схема stg для сырого слоя DWH. +CREATE SCHEMA IF NOT EXISTS stg; + +-- Внешняя таблица в схеме stg для чтения данных из bookings.airports_data через PXF. +DROP EXTERNAL TABLE IF EXISTS stg.airports_ext; +CREATE EXTERNAL TABLE stg.airports_ext ( + airport_code TEXT, + airport_name JSONB, + city JSONB, + country JSONB, + coordinates POINT, + timezone TEXT +) +LOCATION ('pxf://bookings.airports_data?PROFILE=JDBC&SERVER=bookings-db') +FORMAT 'CUSTOM' (formatter='pxfwritable_import'); + +-- Внутренняя таблица stg.airports — сырой слой, все бизнес-колонки как TEXT. +CREATE TABLE IF NOT EXISTS stg.airports ( + airport_code TEXT, + airport_name TEXT, + city TEXT, + country TEXT, + coordinates TEXT, + timezone TEXT, + src_created_at_ts TIMESTAMP, + load_dttm TIMESTAMP NOT NULL DEFAULT now(), + batch_id TEXT +) +WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) +-- Ключ распределения: airport_code +-- Обоснование: airport_code — это уникальный идентификатор аэропорта. +-- Использование airport_code обеспечивает: +-- 1. Равномерное распределение данных по сегментам (airport_code имеет высокую кардинальность) +-- 2. Co-location данных airports и routes при JOIN по departure_airport/arrival_airport +-- 3. Оптимизацию запросов, которые фильтруют или группируют по airport_code +DISTRIBUTED BY (airport_code); diff --git a/sql/stg/airports_dq.sql b/sql/stg/airports_dq.sql new file mode 100644 index 0000000..7bea05f --- /dev/null +++ b/sql/stg/airports_dq.sql @@ -0,0 +1,69 @@ +-- Проверки качества данных для airports (справочник) + +DO $$ +DECLARE + v_batch_id TEXT := '{{ run_id }}'::text; + v_src_count BIGINT; + v_stg_count BIGINT; + v_dup_count BIGINT; + v_null_count BIGINT; +BEGIN + -- Источник: считаем все строки во внешней таблице + SELECT COUNT(*) + INTO v_src_count + FROM stg.airports_ext; + + IF v_src_count = 0 THEN + RAISE EXCEPTION + 'В источнике airports_ext нет строк.'; + END IF; + + -- Считаем строки, реально вставленные в stg.airports в этом батче + SELECT COUNT(*) + INTO v_stg_count + FROM stg.airports + WHERE batch_id = v_batch_id; + + IF v_src_count <> v_stg_count THEN + RAISE EXCEPTION + 'DQ FAILED: несовпадение количества строк. Источник: %, STG: %', + v_src_count, + v_stg_count; + END IF; + + -- Проверка на дубликаты airport_code + SELECT COUNT(*) - COUNT(DISTINCT airport_code) + INTO v_dup_count + FROM stg.airports AS a + WHERE a.batch_id = v_batch_id; + + IF v_dup_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: найдены дубликаты airport_code (batch_id=%): %', + v_batch_id, + v_dup_count; + END IF; + + -- Проверка обязательных полей: airport_code, airport_name, city, timezone + SELECT COUNT(*) + INTO v_null_count + FROM stg.airports AS a + WHERE a.batch_id = v_batch_id + AND (a.airport_code IS NULL OR a.airport_code = '' + OR a.airport_name IS NULL OR a.airport_name = '' + OR a.city IS NULL OR a.city = '' + OR a.timezone IS NULL OR a.timezone = ''); + + IF v_null_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: найдены строки с NULL в обязательных полях (airport_code, airport_name, city, timezone) (batch_id=%): %', + v_batch_id, + v_null_count; + END IF; + + RAISE NOTICE + 'DQ PASSED: airports ок (batch_id=%): source=% stg=%', + v_batch_id, + v_src_count, + v_stg_count; +END $$; diff --git a/sql/stg/airports_load.sql b/sql/stg/airports_load.sql new file mode 100644 index 0000000..3197c2d --- /dev/null +++ b/sql/stg/airports_load.sql @@ -0,0 +1,36 @@ +-- Загрузка всех строк из stg.airports_ext в stg.airports (full load). +-- Используем batch_id для отслеживания загрузки. + +INSERT INTO stg.airports ( + airport_code, + airport_name, + city, + country, + coordinates, + timezone, + src_created_at_ts, + load_dttm, + batch_id +) +SELECT + ext.airport_code::text, + ext.airport_name::text, + ext.city::text, + ext.country::text, + ext.coordinates::text, + ext.timezone::text, + now()::timestamp, + now()::timestamp, + '{{ run_id }}'::text +FROM stg.airports_ext AS ext +WHERE NOT EXISTS ( + -- Защита от дублей в рамках одного batch_id + SELECT 1 + FROM stg.airports AS a + WHERE a.batch_id = '{{ run_id }}'::text + AND a.airport_code = ext.airport_code::text +); + +-- Обновляем статистику для оптимизатора Greenplum +-- Это критично для корректной работы оптимизатора и выбора оптимального плана выполнения +ANALYZE stg.airports; diff --git a/sql/stg/boarding_passes_ddl.sql b/sql/stg/boarding_passes_ddl.sql new file mode 100644 index 0000000..659ed71 --- /dev/null +++ b/sql/stg/boarding_passes_ddl.sql @@ -0,0 +1,38 @@ +-- DDL для слоя STG по таблице boarding_passes. +-- Используется как из общего скрипта ddl_gp.sql (через \i), +-- так и может выполняться отдельно при изменении схемы. + +-- Схема stg для сырого слоя DWH. +CREATE SCHEMA IF NOT EXISTS stg; + +-- Внешняя таблица в схеме stg для чтения данных из bookings.boarding_passes через PXF. +DROP EXTERNAL TABLE IF EXISTS stg.boarding_passes_ext; +CREATE EXTERNAL TABLE stg.boarding_passes_ext ( + ticket_no TEXT, + flight_id TEXT, + seat_no TEXT, + boarding_no INTEGER, + boarding_time TIMESTAMP +) +LOCATION ('pxf://bookings.boarding_passes?PROFILE=JDBC&SERVER=bookings-db') +FORMAT 'CUSTOM' (formatter='pxfwritable_import'); + +-- Внутренняя таблица stg.boarding_passes — сырой слой, все бизнес-колонки как TEXT. +CREATE TABLE IF NOT EXISTS stg.boarding_passes ( + ticket_no TEXT, + flight_id TEXT, + seat_no TEXT, + boarding_no TEXT, + boarding_time TEXT, + src_created_at_ts TIMESTAMP, + load_dttm TIMESTAMP NOT NULL DEFAULT now(), + batch_id TEXT +) +WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) +-- Ключ распределения: ticket_no +-- Обоснование: ticket_no — это основной бизнес-ключ для билетов. +-- Использование ticket_no обеспечивает: +-- 1. Co-location данных boarding_passes и tickets при JOIN по ticket_no +-- 2. Co-location данных boarding_passes и segments при JOIN по ticket_no +-- 3. Равномерное распределение данных по сегментам (ticket_no имеет высокую кардинальность) +DISTRIBUTED BY (ticket_no); diff --git a/sql/stg/boarding_passes_dq.sql b/sql/stg/boarding_passes_dq.sql new file mode 100644 index 0000000..a19e710 --- /dev/null +++ b/sql/stg/boarding_passes_dq.sql @@ -0,0 +1,99 @@ +-- Проверки качества данных для boarding_passes + +DO $$ +DECLARE + v_batch_id TEXT := '{{ run_id }}'::text; + v_src_count BIGINT; + v_stg_count BIGINT; + v_dup_count BIGINT; + v_null_count BIGINT; + v_orphan_ticket_count BIGINT; + v_orphan_segment_count BIGINT; +BEGIN + -- Источник: считаем все строки во внешней таблице + SELECT COUNT(*) + INTO v_src_count + FROM stg.boarding_passes_ext; + + IF v_src_count = 0 THEN + RAISE EXCEPTION + 'В источнике boarding_passes_ext нет строк.'; + END IF; + + -- Считаем строки, реально вставленные в stg.boarding_passes в этом батче + SELECT COUNT(*) + INTO v_stg_count + FROM stg.boarding_passes + WHERE batch_id = v_batch_id; + + IF v_src_count <> v_stg_count THEN + RAISE EXCEPTION + 'DQ FAILED: несовпадение количества строк. Источник: %, STG: %', + v_src_count, + v_stg_count; + END IF; + + -- Проверка на дубликаты (ticket_no, flight_id) + SELECT COUNT(*) - COUNT(DISTINCT ticket_no || '|' || flight_id) + INTO v_dup_count + FROM stg.boarding_passes AS bp + WHERE bp.batch_id = v_batch_id; + + IF v_dup_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: найдены дубликаты (ticket_no, flight_id) (batch_id=%): %', + v_batch_id, + v_dup_count; + END IF; + + -- Проверка обязательных полей (ticket_no, flight_id) + SELECT COUNT(*) + INTO v_null_count + FROM stg.boarding_passes AS bp + WHERE bp.batch_id = v_batch_id + AND (bp.ticket_no IS NULL OR bp.ticket_no = '' + OR bp.flight_id IS NULL OR bp.flight_id = ''); + + IF v_null_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: найдены строки с NULL в обязательных полях (batch_id=%): %', + v_batch_id, + v_null_count; + END IF; + + -- Проверка ссылочной целостности: все boarding_passes должны иметь соответствующие tickets + SELECT COUNT(*) + INTO v_orphan_ticket_count + FROM stg.boarding_passes AS bp + LEFT JOIN stg.tickets AS t ON bp.ticket_no = t.ticket_no + WHERE bp.batch_id = v_batch_id + AND t.ticket_no IS NULL; + + IF v_orphan_ticket_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: найдены boarding_passes без соответствующих tickets (batch_id=%): %', + v_batch_id, + v_orphan_ticket_count; + END IF; + + -- Проверка ссылочной целостности: все boarding_passes должны иметь соответствующие segments + SELECT COUNT(*) + INTO v_orphan_segment_count + FROM stg.boarding_passes AS bp + LEFT JOIN stg.segments AS s ON bp.ticket_no = s.ticket_no AND bp.flight_id = s.flight_id + WHERE bp.batch_id = v_batch_id + AND s.ticket_no IS NULL; + + IF v_orphan_segment_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: найдены boarding_passes без соответствующих segments (batch_id=%): %', + v_batch_id, + v_orphan_segment_count; + END IF; + + RAISE NOTICE + 'DQ PASSED: boarding_passes ок (batch_id=%): source=% stg=%', + v_batch_id, + v_src_count, + v_stg_count; +END $$; diff --git a/sql/stg/boarding_passes_load.sql b/sql/stg/boarding_passes_load.sql new file mode 100644 index 0000000..1d94b7d --- /dev/null +++ b/sql/stg/boarding_passes_load.sql @@ -0,0 +1,36 @@ +-- Загрузка всех строк из stg.boarding_passes_ext в stg.boarding_passes. +-- Используем full snapshot: все строки при каждом запуске. +-- Используем batch_id для отслеживания загрузки. + +INSERT INTO stg.boarding_passes ( + ticket_no, + flight_id, + seat_no, + boarding_no, + boarding_time, + src_created_at_ts, + load_dttm, + batch_id +) +SELECT + ext.ticket_no, + ext.flight_id, + ext.seat_no, + ext.boarding_no::text, + ext.boarding_time::text, + now()::timestamp, + now(), + '{{ run_id }}'::text +FROM stg.boarding_passes_ext AS ext +WHERE NOT EXISTS ( + -- Защита от дублей в рамках одного batch_id + SELECT 1 + FROM stg.boarding_passes AS bp + WHERE bp.batch_id = '{{ run_id }}'::text + AND bp.ticket_no = ext.ticket_no + AND bp.flight_id = ext.flight_id +); + +-- Обновляем статистику для оптимизатора Greenplum +-- Это критично для корректной работы оптимизатора и выбора оптимального плана выполнения +ANALYZE stg.boarding_passes; diff --git a/sql/stg/flights_ddl.sql b/sql/stg/flights_ddl.sql new file mode 100644 index 0000000..546de1a --- /dev/null +++ b/sql/stg/flights_ddl.sql @@ -0,0 +1,42 @@ +-- DDL для слоя STG по таблице flights. +-- Используется как из общего скрипта ddl_gp.sql (через \i), +-- так и может выполняться отдельно при изменении схемы. + +-- Схема stg для сырого слоя DWH. +CREATE SCHEMA IF NOT EXISTS stg; + +-- Внешняя таблица в схеме stg для чтения данных из bookings.flights через PXF. +DROP EXTERNAL TABLE IF EXISTS stg.flights_ext; +CREATE EXTERNAL TABLE stg.flights_ext ( + flight_id TEXT, + route_no TEXT, + status TEXT, + scheduled_departure TIMESTAMP, + scheduled_arrival TIMESTAMP, + actual_departure TIMESTAMP, + actual_arrival TIMESTAMP +) +LOCATION ('pxf://bookings.flights?PROFILE=JDBC&SERVER=bookings-db') +FORMAT 'CUSTOM' (formatter='pxfwritable_import'); + +-- Внутренняя таблица stg.flights — сырой слой, все бизнес-колонки как TEXT. +CREATE TABLE IF NOT EXISTS stg.flights ( + flight_id TEXT, + route_no TEXT, + status TEXT, + scheduled_departure TEXT, + scheduled_arrival TEXT, + actual_departure TEXT, + actual_arrival TEXT, + src_created_at_ts TIMESTAMP, + load_dttm TIMESTAMP NOT NULL DEFAULT now(), + batch_id TEXT +) +WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) +-- Ключ распределения: flight_id +-- Обоснование: flight_id — это уникальный идентификатор рейса. +-- Использование flight_id обеспечивает: +-- 1. Равномерное распределение данных по сегментам (flight_id имеет высокую кардинальность) +-- 2. Оптимизацию запросов, которые фильтруют или группируют по flight_id +-- 3. Co-location данных flights и boarding_passes при JOIN по flight_id +DISTRIBUTED BY (flight_id); diff --git a/sql/stg/flights_dq.sql b/sql/stg/flights_dq.sql new file mode 100644 index 0000000..1b3bc1d --- /dev/null +++ b/sql/stg/flights_dq.sql @@ -0,0 +1,95 @@ +-- Проверки качества данных для flights + +DO $$ +DECLARE + v_batch_id TEXT := '{{ run_id }}'::text; + v_prev_ts TIMESTAMP; + v_src_count BIGINT; + v_stg_count BIGINT; + v_dup_count BIGINT; + v_null_count BIGINT; + v_orphan_route_count BIGINT; +BEGIN + -- Опорная метка: максимум src_created_at_ts среди предыдущих батчей + SELECT max(src_created_at_ts) + INTO v_prev_ts + FROM stg.flights + WHERE batch_id <> v_batch_id + OR batch_id IS NULL; + + -- Источник: считаем строки во внешней таблице, которые вошли в окно инкремента + SELECT COUNT(*) + INTO v_src_count + FROM stg.flights_ext + WHERE scheduled_departure > COALESCE(v_prev_ts, TIMESTAMP '1900-01-01 00:00:00'); + + IF v_src_count = 0 THEN + RAISE EXCEPTION + 'В источнике flights_ext нет строк для окна инкремента (scheduled_departure > %).', + COALESCE(v_prev_ts, TIMESTAMP '1900-01-01 00:00:00'); + END IF; + + -- Считаем строки, реально вставленные в stg.flights в этом батче + SELECT COUNT(*) + INTO v_stg_count + FROM stg.flights + WHERE batch_id = v_batch_id; + + IF v_src_count <> v_stg_count THEN + RAISE EXCEPTION + 'DQ FAILED: несовпадение количества строк. Источник: %, STG: %', + v_src_count, + v_stg_count; + END IF; + + -- Проверка на дубликаты flight_id + SELECT COUNT(*) - COUNT(DISTINCT flight_id) + INTO v_dup_count + FROM stg.flights AS f + WHERE f.batch_id = v_batch_id; + + IF v_dup_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: найдены дубликаты flight_id (batch_id=%): %', + v_batch_id, + v_dup_count; + END IF; + + -- Проверка обязательных полей (flight_id, route_no, status, scheduled_departure) + SELECT COUNT(*) + INTO v_null_count + FROM stg.flights AS f + WHERE f.batch_id = v_batch_id + AND (f.flight_id IS NULL OR f.flight_id = '' + OR f.route_no IS NULL OR f.route_no = '' + OR f.status IS NULL OR f.status = '' + OR f.scheduled_departure IS NULL OR f.scheduled_departure = ''); + + IF v_null_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: найдены строки с NULL в обязательных полях (batch_id=%): %', + v_batch_id, + v_null_count; + END IF; + + -- Проверка ссылочной целостности: все flights должны иметь соответствующие routes + SELECT COUNT(*) + INTO v_orphan_route_count + FROM stg.flights AS f + LEFT JOIN stg.routes AS r ON f.route_no = r.route_no + WHERE f.batch_id = v_batch_id + AND r.route_no IS NULL; + + IF v_orphan_route_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: найдены flights без соответствующих routes (batch_id=%): %', + v_batch_id, + v_orphan_route_count; + END IF; + + RAISE NOTICE + 'DQ PASSED: flights ок (batch_id=%): source=% stg=%', + v_batch_id, + v_src_count, + v_stg_count; +END $$; diff --git a/sql/stg/flights_load.sql b/sql/stg/flights_load.sql new file mode 100644 index 0000000..5a6188e --- /dev/null +++ b/sql/stg/flights_load.sql @@ -0,0 +1,48 @@ +-- Загрузка инкремента из stg.flights_ext в stg.flights. +-- Окно инкремента определяется по src_created_at_ts: +-- берём строки, где scheduled_departure больше максимального src_created_at_ts +-- среди "старых" батчей; верхняя граница по дате не используется. + +-- CTE для определения максимальной даты загрузки предыдущего батча +WITH max_batch_ts AS ( + SELECT COALESCE(MAX(src_created_at_ts), TIMESTAMP '1900-01-01 00:00:00') AS max_ts + FROM stg.flights + WHERE batch_id <> '{{ run_id }}'::text + OR batch_id IS NULL +) +INSERT INTO stg.flights ( + flight_id, + route_no, + status, + scheduled_departure, + scheduled_arrival, + actual_departure, + actual_arrival, + src_created_at_ts, + load_dttm, + batch_id +) +SELECT + ext.flight_id::text, + ext.route_no::text, + ext.status::text, + ext.scheduled_departure::text, + ext.scheduled_arrival::text, + ext.actual_departure::text, + ext.actual_arrival::text, + ext.scheduled_departure::timestamp, + now(), + '{{ run_id }}'::text +FROM stg.flights_ext AS ext +CROSS JOIN max_batch_ts AS mb +WHERE ext.scheduled_departure > mb.max_ts +AND NOT EXISTS ( + SELECT 1 + FROM stg.flights AS f + WHERE f.batch_id = '{{ run_id }}'::text + AND f.flight_id = ext.flight_id::text +); + +-- Обновляем статистику для оптимизатора Greenplum +-- Это критично для корректной работы оптимизатора и выбора оптимального плана выполнения +ANALYZE stg.flights; diff --git a/sql/stg/routes_ddl.sql b/sql/stg/routes_ddl.sql new file mode 100644 index 0000000..c3c93c4 --- /dev/null +++ b/sql/stg/routes_ddl.sql @@ -0,0 +1,44 @@ +-- DDL для слоя STG по таблице routes (справочник). +-- Используется как из общего скрипта ddl_gp.sql (через \i), +-- так и может выполняться отдельно при изменении схемы. + +-- Схема stg для сырого слоя DWH. +CREATE SCHEMA IF NOT EXISTS stg; + +-- Внешняя таблица в схеме stg для чтения данных из bookings.routes через PXF. +DROP EXTERNAL TABLE IF EXISTS stg.routes_ext; +CREATE EXTERNAL TABLE stg.routes_ext ( + route_no TEXT, + validity TSTZRANGE, + departure_airport TEXT, + arrival_airport TEXT, + airplane_code TEXT, + days_of_week INTEGER[], + scheduled_time TIME WITHOUT TIME ZONE, + duration INTERVAL +) +LOCATION ('pxf://bookings.routes?PROFILE=JDBC&SERVER=bookings-db') +FORMAT 'CUSTOM' (formatter='pxfwritable_import'); + +-- Внутренняя таблица stg.routes — сырой слой, все бизнес-колонки как TEXT. +CREATE TABLE IF NOT EXISTS stg.routes ( + route_no TEXT, + validity TEXT, + departure_airport TEXT, + arrival_airport TEXT, + airplane_code TEXT, + days_of_week TEXT, + scheduled_time TEXT, + duration TEXT, + src_created_at_ts TIMESTAMP, + load_dttm TIMESTAMP NOT NULL DEFAULT now(), + batch_id TEXT +) +WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) +-- Ключ распределения: route_no +-- Обоснование: route_no — это уникальный идентификатор маршрута. +-- Использование route_no обеспечивает: +-- 1. Равномерное распределение данных по сегментам (route_no имеет высокую кардинальность) +-- 2. Co-location данных routes с airports и airplanes при JOIN +-- 3. Оптимизацию запросов, которые фильтруют или группируют по route_no +DISTRIBUTED BY (route_no); diff --git a/sql/stg/routes_dq.sql b/sql/stg/routes_dq.sql new file mode 100644 index 0000000..8c31452 --- /dev/null +++ b/sql/stg/routes_dq.sql @@ -0,0 +1,116 @@ +-- Проверки качества данных для routes (справочник) + +DO $$ +DECLARE + v_batch_id TEXT := '{{ run_id }}'::text; + v_src_count BIGINT; + v_stg_count BIGINT; + v_dup_count BIGINT; + v_null_count BIGINT; + v_orphan_airports_count BIGINT; + v_orphan_airplanes_count BIGINT; +BEGIN + -- Источник: считаем все строки во внешней таблице + SELECT COUNT(*) + INTO v_src_count + FROM stg.routes_ext; + + IF v_src_count = 0 THEN + RAISE EXCEPTION + 'В источнике routes_ext нет строк.'; + END IF; + + -- Считаем строки, реально вставленные в stg.routes в этом батче + SELECT COUNT(*) + INTO v_stg_count + FROM stg.routes + WHERE batch_id = v_batch_id; + + IF v_src_count <> v_stg_count THEN + RAISE EXCEPTION + 'DQ FAILED: несовпадение количества строк. Источник: %, STG: %', + v_src_count, + v_stg_count; + END IF; + + -- Проверка на дубликаты составного ключа (route_no, validity) + SELECT COUNT(*) - COUNT(DISTINCT route_no || '|' || validity) + INTO v_dup_count + FROM stg.routes AS r + WHERE r.batch_id = v_batch_id; + + IF v_dup_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: найдены дубликаты (route_no, validity) (batch_id=%): %', + v_batch_id, + v_dup_count; + END IF; + + -- Проверка обязательных полей: route_no, departure_airport, arrival_airport, airplane_code + SELECT COUNT(*) + INTO v_null_count + FROM stg.routes AS r + WHERE r.batch_id = v_batch_id + AND (r.route_no IS NULL OR r.route_no = '' + OR r.departure_airport IS NULL OR r.departure_airport = '' + OR r.arrival_airport IS NULL OR r.arrival_airport = '' + OR r.airplane_code IS NULL OR r.airplane_code = ''); + + IF v_null_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: найдены строки с NULL в обязательных полях (route_no, departure_airport, arrival_airport, airplane_code) (batch_id=%): %', + v_batch_id, + v_null_count; + END IF; + + -- Проверка ссылочной целостности: departure_airport должен существовать в airports + SELECT COUNT(*) + INTO v_orphan_airports_count + FROM stg.routes AS r + LEFT JOIN stg.airports AS da ON r.departure_airport = da.airport_code + WHERE r.batch_id = v_batch_id + AND da.airport_code IS NULL; + + IF v_orphan_airports_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: найдены routes с несуществующим departure_airport в airports (batch_id=%): %', + v_batch_id, + v_orphan_airports_count; + END IF; + + -- Проверка ссылочной целостности: arrival_airport должен существовать в airports + SELECT COUNT(*) + INTO v_orphan_airports_count + FROM stg.routes AS r + LEFT JOIN stg.airports AS aa ON r.arrival_airport = aa.airport_code + WHERE r.batch_id = v_batch_id + AND aa.airport_code IS NULL; + + IF v_orphan_airports_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: найдены routes с несуществующим arrival_airport в airports (batch_id=%): %', + v_batch_id, + v_orphan_airports_count; + END IF; + + -- Проверка ссылочной целостности: airplane_code должен существовать в airplanes + SELECT COUNT(*) + INTO v_orphan_airplanes_count + FROM stg.routes AS r + LEFT JOIN stg.airplanes AS a ON r.airplane_code = a.airplane_code + WHERE r.batch_id = v_batch_id + AND a.airplane_code IS NULL; + + IF v_orphan_airplanes_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: найдены routes с несуществующим airplane_code в airplanes (batch_id=%): %', + v_batch_id, + v_orphan_airplanes_count; + END IF; + + RAISE NOTICE + 'DQ PASSED: routes ок (batch_id=%): source=% stg=%', + v_batch_id, + v_src_count, + v_stg_count; +END $$; diff --git a/sql/stg/routes_load.sql b/sql/stg/routes_load.sql new file mode 100644 index 0000000..77c2077 --- /dev/null +++ b/sql/stg/routes_load.sql @@ -0,0 +1,41 @@ +-- Загрузка всех строк из stg.routes_ext в stg.routes (full load). +-- Используем batch_id для отслеживания загрузки. + +INSERT INTO stg.routes ( + route_no, + validity, + departure_airport, + arrival_airport, + airplane_code, + days_of_week, + scheduled_time, + duration, + src_created_at_ts, + load_dttm, + batch_id +) +SELECT + ext.route_no::text, + ext.validity::text, + ext.departure_airport::text, + ext.arrival_airport::text, + ext.airplane_code::text, + ext.days_of_week::text, + ext.scheduled_time::text, + ext.duration::text, + now()::timestamp, + now()::timestamp, + '{{ run_id }}'::text +FROM stg.routes_ext AS ext +WHERE NOT EXISTS ( + -- Защита от дублей в рамках одного batch_id по составному ключу (route_no, validity) + SELECT 1 + FROM stg.routes AS r + WHERE r.batch_id = '{{ run_id }}'::text + AND r.route_no = ext.route_no::text + AND r.validity = ext.validity::text +); + +-- Обновляем статистику для оптимизатора Greenplum +-- Это критично для корректной работы оптимизатора и выбора оптимального плана выполнения +ANALYZE stg.routes; diff --git a/sql/stg/seats_ddl.sql b/sql/stg/seats_ddl.sql new file mode 100644 index 0000000..4b6ead7 --- /dev/null +++ b/sql/stg/seats_ddl.sql @@ -0,0 +1,34 @@ +-- DDL для слоя STG по таблице seats (справочник). +-- Используется как из общего скрипта ddl_gp.sql (через \i), +-- так и может выполняться отдельно при изменении схемы. + +-- Схема stg для сырого слоя DWH. +CREATE SCHEMA IF NOT EXISTS stg; + +-- Внешняя таблица в схеме stg для чтения данных из bookings.seats через PXF. +DROP EXTERNAL TABLE IF EXISTS stg.seats_ext; +CREATE EXTERNAL TABLE stg.seats_ext ( + airplane_code TEXT, + seat_no TEXT, + fare_conditions TEXT +) +LOCATION ('pxf://bookings.seats?PROFILE=JDBC&SERVER=bookings-db') +FORMAT 'CUSTOM' (formatter='pxfwritable_import'); + +-- Внутренняя таблица stg.seats — сырой слой, все бизнес-колонки как TEXT. +CREATE TABLE IF NOT EXISTS stg.seats ( + airplane_code TEXT, + seat_no TEXT, + fare_conditions TEXT, + src_created_at_ts TIMESTAMP, + load_dttm TIMESTAMP NOT NULL DEFAULT now(), + batch_id TEXT +) +WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) +-- Ключ распределения: airplane_code +-- Обоснование: airplane_code обеспечивает co-location с таблицей airplanes. +-- Использование airplane_code обеспечивает: +-- 1. Co-location данных seats и airplanes при JOIN по airplane_code +-- 2. Группировка мест по самолётам (в одном самолёте обычно много мест) +-- 3. Оптимизацию запросов, которые фильтруют или группируют по airplane_code +DISTRIBUTED BY (airplane_code); diff --git a/sql/stg/seats_dq.sql b/sql/stg/seats_dq.sql new file mode 100644 index 0000000..91d195a --- /dev/null +++ b/sql/stg/seats_dq.sql @@ -0,0 +1,84 @@ +-- Проверки качества данных для seats (справочник) + +DO $$ +DECLARE + v_batch_id TEXT := '{{ run_id }}'::text; + v_src_count BIGINT; + v_stg_count BIGINT; + v_dup_count BIGINT; + v_null_count BIGINT; + v_orphan_airplanes_count BIGINT; +BEGIN + -- Источник: считаем все строки во внешней таблице + SELECT COUNT(*) + INTO v_src_count + FROM stg.seats_ext; + + IF v_src_count = 0 THEN + RAISE EXCEPTION + 'В источнике seats_ext нет строк.'; + END IF; + + -- Считаем строки, реально вставленные в stg.seats в этом батче + SELECT COUNT(*) + INTO v_stg_count + FROM stg.seats + WHERE batch_id = v_batch_id; + + IF v_src_count <> v_stg_count THEN + RAISE EXCEPTION + 'DQ FAILED: несовпадение количества строк. Источник: %, STG: %', + v_src_count, + v_stg_count; + END IF; + + -- Проверка на дубликаты составного ключа (airplane_code, seat_no) + SELECT COUNT(*) - COUNT(DISTINCT airplane_code || '|' || seat_no) + INTO v_dup_count + FROM stg.seats AS s + WHERE s.batch_id = v_batch_id; + + IF v_dup_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: найдены дубликаты (airplane_code, seat_no) (batch_id=%): %', + v_batch_id, + v_dup_count; + END IF; + + -- Проверка обязательных полей: airplane_code, seat_no, fare_conditions + SELECT COUNT(*) + INTO v_null_count + FROM stg.seats AS s + WHERE s.batch_id = v_batch_id + AND (s.airplane_code IS NULL OR s.airplane_code = '' + OR s.seat_no IS NULL OR s.seat_no = '' + OR s.fare_conditions IS NULL OR s.fare_conditions = ''); + + IF v_null_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: найдены строки с NULL в обязательных полях (airplane_code, seat_no, fare_conditions) (batch_id=%): %', + v_batch_id, + v_null_count; + END IF; + + -- Проверка ссылочной целостности: airplane_code должен существовать в airplanes + SELECT COUNT(*) + INTO v_orphan_airplanes_count + FROM stg.seats AS s + LEFT JOIN stg.airplanes AS a ON s.airplane_code = a.airplane_code + WHERE s.batch_id = v_batch_id + AND a.airplane_code IS NULL; + + IF v_orphan_airplanes_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: найдены seats с несуществующим airplane_code в airplanes (batch_id=%): %', + v_batch_id, + v_orphan_airplanes_count; + END IF; + + RAISE NOTICE + 'DQ PASSED: seats ок (batch_id=%): source=% stg=%', + v_batch_id, + v_src_count, + v_stg_count; +END $$; diff --git a/sql/stg/seats_load.sql b/sql/stg/seats_load.sql new file mode 100644 index 0000000..3c530d0 --- /dev/null +++ b/sql/stg/seats_load.sql @@ -0,0 +1,31 @@ +-- Загрузка всех строк из stg.seats_ext в stg.seats (full load). +-- Используем batch_id для отслеживания загрузки. + +INSERT INTO stg.seats ( + airplane_code, + seat_no, + fare_conditions, + src_created_at_ts, + load_dttm, + batch_id +) +SELECT + ext.airplane_code::text, + ext.seat_no::text, + ext.fare_conditions::text, + now()::timestamp, + now()::timestamp, + '{{ run_id }}'::text +FROM stg.seats_ext AS ext +WHERE NOT EXISTS ( + -- Защита от дублей в рамках одного batch_id по составному ключу (airplane_code, seat_no) + SELECT 1 + FROM stg.seats AS s + WHERE s.batch_id = '{{ run_id }}'::text + AND s.airplane_code = ext.airplane_code::text + AND s.seat_no = ext.seat_no::text +); + +-- Обновляем статистику для оптимизатора Greenplum +-- Это критично для корректной работы оптимизатора и выбора оптимального плана выполнения +ANALYZE stg.seats; diff --git a/sql/stg/segments_ddl.sql b/sql/stg/segments_ddl.sql new file mode 100644 index 0000000..402bc05 --- /dev/null +++ b/sql/stg/segments_ddl.sql @@ -0,0 +1,36 @@ +-- DDL для слоя STG по таблице segments. +-- Используется как из общего скрипта ddl_gp.sql (через \i), +-- так и может выполняться отдельно при изменении схемы. + +-- Схема stg для сырого слоя DWH. +CREATE SCHEMA IF NOT EXISTS stg; + +-- Внешняя таблица в схеме stg для чтения данных из bookings.segments через PXF. +DROP EXTERNAL TABLE IF EXISTS stg.segments_ext; +CREATE EXTERNAL TABLE stg.segments_ext ( + ticket_no TEXT, + flight_id TEXT, + fare_conditions TEXT, + price NUMERIC(10,2) +) +LOCATION ('pxf://bookings.segments?PROFILE=JDBC&SERVER=bookings-db') +FORMAT 'CUSTOM' (formatter='pxfwritable_import'); + +-- Внутренняя таблица stg.segments — сырой слой, все бизнес-колонки как TEXT. +CREATE TABLE IF NOT EXISTS stg.segments ( + ticket_no TEXT, + flight_id TEXT, + fare_conditions TEXT, + price TEXT, + src_created_at_ts TIMESTAMP, + load_dttm TIMESTAMP NOT NULL DEFAULT now(), + batch_id TEXT +) +WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) +-- Ключ распределения: ticket_no +-- Обоснование: ticket_no — это основной бизнес-ключ для билетов. +-- Использование ticket_no обеспечивает: +-- 1. Co-location данных segments и tickets при JOIN по ticket_no +-- 2. Co-location данных segments и boarding_passes при JOIN по ticket_no +-- 3. Равномерное распределение данных по сегментам (ticket_no имеет высокую кардинальность) +DISTRIBUTED BY (ticket_no); diff --git a/sql/stg/segments_dq.sql b/sql/stg/segments_dq.sql new file mode 100644 index 0000000..3c00e12 --- /dev/null +++ b/sql/stg/segments_dq.sql @@ -0,0 +1,113 @@ +-- Проверки качества данных для segments + +DO $$ +DECLARE + v_batch_id TEXT := '{{ run_id }}'::text; + v_prev_ts TIMESTAMP; + v_src_count BIGINT; + v_stg_count BIGINT; + v_dup_count BIGINT; + v_null_count BIGINT; + v_orphan_ticket_count BIGINT; + v_orphan_flight_count BIGINT; +BEGIN + -- Опорная метка: максимум src_created_at_ts среди предыдущих батчей + SELECT max(src_created_at_ts) + INTO v_prev_ts + FROM stg.segments + WHERE batch_id <> v_batch_id + OR batch_id IS NULL; + + -- Источник: считаем строки во внешней таблице, которые вошли в окно инкремента + SELECT COUNT(*) + INTO v_src_count + FROM stg.segments_ext AS s + JOIN stg.tickets_ext AS t ON s.ticket_no = t.ticket_no + JOIN stg.bookings_ext AS b ON t.book_ref = b.book_ref + WHERE b.book_date > COALESCE(v_prev_ts, TIMESTAMP '1900-01-01 00:00:00'); + + IF v_src_count = 0 THEN + RAISE EXCEPTION + 'В источнике segments_ext нет строк для окна инкремента (book_date > %).', + COALESCE(v_prev_ts, TIMESTAMP '1900-01-01 00:00:00'); + END IF; + + -- Считаем строки, реально вставленные в stg.segments в этом батче + SELECT COUNT(*) + INTO v_stg_count + FROM stg.segments + WHERE batch_id = v_batch_id; + + IF v_src_count <> v_stg_count THEN + RAISE EXCEPTION + 'DQ FAILED: несовпадение количества строк. Источник: %, STG: %', + v_src_count, + v_stg_count; + END IF; + + -- Проверка на дубликаты (ticket_no, flight_id) + SELECT COUNT(*) - COUNT(DISTINCT ticket_no || '|' || flight_id) + INTO v_dup_count + FROM stg.segments AS s + WHERE s.batch_id = v_batch_id; + + IF v_dup_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: найдены дубликаты (ticket_no, flight_id) (batch_id=%): %', + v_batch_id, + v_dup_count; + END IF; + + -- Проверка обязательных полей (ticket_no, flight_id, fare_conditions, price) + SELECT COUNT(*) + INTO v_null_count + FROM stg.segments AS s + WHERE s.batch_id = v_batch_id + AND (s.ticket_no IS NULL OR s.ticket_no = '' + OR s.flight_id IS NULL OR s.flight_id = '' + OR s.fare_conditions IS NULL OR s.fare_conditions = '' + OR s.price IS NULL OR s.price = ''); + + IF v_null_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: найдены строки с NULL в обязательных полях (batch_id=%): %', + v_batch_id, + v_null_count; + END IF; + + -- Проверка ссылочной целостности: все segments должны иметь соответствующие tickets + SELECT COUNT(*) + INTO v_orphan_ticket_count + FROM stg.segments AS s + LEFT JOIN stg.tickets AS t ON s.ticket_no = t.ticket_no + WHERE s.batch_id = v_batch_id + AND t.ticket_no IS NULL; + + IF v_orphan_ticket_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: найдены segments без соответствующих tickets (batch_id=%): %', + v_batch_id, + v_orphan_ticket_count; + END IF; + + -- Проверка ссылочной целостности: все segments должны иметь соответствующие flights + SELECT COUNT(*) + INTO v_orphan_flight_count + FROM stg.segments AS s + LEFT JOIN stg.flights AS f ON s.flight_id = f.flight_id + WHERE s.batch_id = v_batch_id + AND f.flight_id IS NULL; + + IF v_orphan_flight_count <> 0 THEN + RAISE EXCEPTION + 'DQ FAILED: найдены segments без соответствующих flights (batch_id=%): %', + v_batch_id, + v_orphan_flight_count; + END IF; + + RAISE NOTICE + 'DQ PASSED: segments ок (batch_id=%): source=% stg=%', + v_batch_id, + v_src_count, + v_stg_count; +END $$; diff --git a/sql/stg/segments_load.sql b/sql/stg/segments_load.sql new file mode 100644 index 0000000..96154a5 --- /dev/null +++ b/sql/stg/segments_load.sql @@ -0,0 +1,45 @@ +-- Загрузка инкремента из stg.segments_ext в stg.segments. +-- Инкремент определяется по дате бронирования (book_date из bookings.bookings) +-- через JOIN с таблицей tickets. + +-- CTE для определения максимальной даты загрузки предыдущего батча +WITH max_batch_ts AS ( + SELECT COALESCE(MAX(src_created_at_ts), TIMESTAMP '1900-01-01 00:00:00') AS max_ts + FROM stg.segments + WHERE batch_id <> '{{ run_id }}'::text + OR batch_id IS NULL +) +INSERT INTO stg.segments ( + ticket_no, + flight_id, + fare_conditions, + price, + src_created_at_ts, + load_dttm, + batch_id +) +SELECT + ext.ticket_no, + ext.flight_id, + ext.fare_conditions, + ext.price::text, + b.book_date::timestamp, + now(), + '{{ run_id }}'::text +FROM stg.segments_ext AS ext +JOIN stg.tickets_ext AS t ON ext.ticket_no = t.ticket_no +JOIN stg.bookings_ext AS b ON t.book_ref = b.book_ref +CROSS JOIN max_batch_ts AS mb +WHERE b.book_date > mb.max_ts +AND NOT EXISTS ( + -- Защита от дублей в рамках одного batch_id + SELECT 1 + FROM stg.segments AS s + WHERE s.batch_id = '{{ run_id }}'::text + AND s.ticket_no = ext.ticket_no + AND s.flight_id = ext.flight_id +); + +-- Обновляем статистику для оптимизатора Greenplum +-- Это критично для корректной работы оптимизатора и выбора оптимального плана выполнения +ANALYZE stg.segments; diff --git a/tests/test_dags_smoke.py b/tests/test_dags_smoke.py index 2d2f727..20ee766 100644 --- a/tests/test_dags_smoke.py +++ b/tests/test_dags_smoke.py @@ -70,3 +70,85 @@ def test_csv_to_greenplum_dq_dag_structure(): assert h in s.get_direct_relatives("downstream") assert d in h.get_direct_relatives("downstream") assert q in d.get_direct_relatives("downstream") + + +def test_bookings_stg_ddl_dag_structure(): + """Проверка структуры DAG bookings_stg_ddl.""" + dag = _load_dag("airflow.dags.bookings_stg_ddl") + + expected_tasks = { + "apply_stg_bookings_ddl", + "apply_stg_tickets_ddl", + "apply_stg_airports_ddl", + "apply_stg_airplanes_ddl", + "apply_stg_routes_ddl", + "apply_stg_seats_ddl", + "apply_stg_flights_ddl", + "apply_stg_segments_ddl", + "apply_stg_boarding_passes_ddl", + } + assert expected_tasks.issubset(dag.task_dict.keys()) + + # Проверка линейных зависимостей + # Справочники создаются после bookings/tickets + assert dag.has_task("apply_stg_bookings_ddl") + assert dag.has_task("apply_stg_tickets_ddl") + assert dag.has_task("apply_stg_airports_ddl") + assert dag.has_task("apply_stg_airplanes_ddl") + assert dag.has_task("apply_stg_routes_ddl") + assert dag.has_task("apply_stg_seats_ddl") + assert dag.has_task("apply_stg_flights_ddl") + assert dag.has_task("apply_stg_segments_ddl") + assert dag.has_task("apply_stg_boarding_passes_ddl") + + +def test_bookings_to_gp_stage_dag_structure(): + """Проверка структуры DAG bookings_to_gp_stage.""" + dag = _load_dag("airflow.dags.bookings_to_gp_stage") + + expected_tasks = { + "generate_bookings_day", + "load_bookings_to_stg", + "check_row_counts", + "load_tickets_to_stg", + "check_tickets_dq", + "load_airports_to_stg", + "check_airports_dq", + "load_airplanes_to_stg", + "check_airplanes_dq", + "load_routes_to_stg", + "check_routes_dq", + "load_seats_to_stg", + "check_seats_dq", + "load_flights_to_stg", + "check_flights_dq", + "load_segments_to_stg", + "check_segments_dq", + "load_boarding_passes_to_stg", + "check_boarding_passes_dq", + "finish_summary", + } + assert expected_tasks.issubset(dag.task_dict.keys()) + + # Проверка линейных зависимостей + # bookings/tickets → справочники → транзакции → финальный лог + assert dag.has_task("generate_bookings_day") + assert dag.has_task("load_bookings_to_stg") + assert dag.has_task("check_row_counts") + assert dag.has_task("load_tickets_to_stg") + assert dag.has_task("check_tickets_dq") + assert dag.has_task("load_airports_to_stg") + assert dag.has_task("check_airports_dq") + assert dag.has_task("load_airplanes_to_stg") + assert dag.has_task("check_airplanes_dq") + assert dag.has_task("load_routes_to_stg") + assert dag.has_task("check_routes_dq") + assert dag.has_task("load_seats_to_stg") + assert dag.has_task("check_seats_dq") + assert dag.has_task("load_flights_to_stg") + assert dag.has_task("check_flights_dq") + assert dag.has_task("load_segments_to_stg") + assert dag.has_task("check_segments_dq") + assert dag.has_task("load_boarding_passes_to_stg") + assert dag.has_task("check_boarding_passes_dq") + assert dag.has_task("finish_summary")