Генерация dds слоя по ТЗ - без тестов

This commit is contained in:
2026-01-18 11:04:46 +03:00
parent 15081384f4
commit 9c033b39aa
26 changed files with 1566 additions and 5 deletions
+60 -1
View File
@@ -38,4 +38,63 @@ with DAG(
sql="stg/tickets_ddl.sql", 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,
]
)
+106 -1
View File
@@ -94,11 +94,116 @@ with DAG(
sql="stg/tickets_dq.sql", 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. Финальный лог/сводка # 6. Финальный лог/сводка
finish_summary = PythonOperator( finish_summary = PythonOperator(
task_id="finish_summary", task_id="finish_summary",
python_callable=_finish_summary, python_callable=_finish_summary,
) )
# Сначала загружаются и проверяются bookings и tickets
generate_bookings_day >> load_bookings_to_stg >> check_row_counts 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
+124 -3
View File
@@ -1,6 +1,6 @@
# Схема БД DWH (Bookings → Greenplum) # Схема БД DWH (Bookings → Greenplum)
> **Статус:** Проект в разработке. Реализован только STG слой (частично: bookings, tickets). > **Статус:** Проект в разработке. Реализован STG слой полностью (все 9 таблиц).
## Обзор ## Обзор
@@ -28,7 +28,7 @@
| Слой | Статус | Реализовано | | Слой | Статус | Реализовано |
|------|--------|-------------| |------|--------|-------------|
| **Source** | ✅ Готово | Демо-БД bookings (Postgres) | | **Source** | ✅ Готово | Демо-БД bookings (Postgres) |
| **STG** | ⚠️ В процессе | 2 из 9 таблиц (bookings, tickets) | | **STG** | ✅ Готово | 9 из 9 таблиц (bookings, tickets, airports, airplanes, routes, seats, flights, segments, boarding_passes) |
| **ODS** | ❌ Не реализован | Планируется | | **ODS** | ❌ Не реализован | Планируется |
| **DDS** | ❌ Не реализован | Планируется | | **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) ## Полная схема потоков данных (Data Lineage)
```mermaid ```mermaid
@@ -299,7 +420,7 @@ graph LR
## TODO ## TODO
- [ ] Реализовать STG слой полностью (все 9 таблиц) - [x] Реализовать STG слой полностью (все 9 таблиц)
- [ ] Реализовать ODS слой - [ ] Реализовать ODS слой
- [ ] Реализовать DDS слой (измерения и факт) - [ ] Реализовать DDS слой (измерения и факт)
- [ ] Создать DAG для загрузки ODS - [ ] Создать DAG для загрузки ODS
+11
View File
@@ -26,3 +26,14 @@ FORMAT 'CUSTOM' (formatter='pxfwritable_import');
-- Здесь подключаем их через psql \i, чтобы сохранить единый входной скрипт. -- Здесь подключаем их через psql \i, чтобы сохранить единый входной скрипт.
\i stg/bookings_ddl.sql \i stg/bookings_ddl.sql
\i stg/tickets_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
+37
View File
@@ -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);
+67
View File
@@ -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 $$;
+32
View File
@@ -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;
+40
View File
@@ -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);
+69
View File
@@ -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 $$;
+36
View File
@@ -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;
+38
View File
@@ -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);
+99
View File
@@ -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 $$;
+36
View File
@@ -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;
+42
View File
@@ -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);
+95
View File
@@ -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 $$;
+48
View File
@@ -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;
+44
View File
@@ -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);
+116
View File
@@ -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 $$;
+41
View File
@@ -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;
+34
View File
@@ -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);
+84
View File
@@ -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 $$;
+31
View File
@@ -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;
+36
View File
@@ -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);
+113
View File
@@ -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 $$;
+45
View File
@@ -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;
+82
View File
@@ -70,3 +70,85 @@ def test_csv_to_greenplum_dq_dag_structure():
assert h in s.get_direct_relatives("downstream") assert h in s.get_direct_relatives("downstream")
assert d in h.get_direct_relatives("downstream") assert d in h.get_direct_relatives("downstream")
assert q in d.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")