Генерация dds слоя по ТЗ - без тестов
This commit is contained in:
@@ -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);
|
||||
@@ -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 $$;
|
||||
@@ -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;
|
||||
@@ -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);
|
||||
@@ -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 $$;
|
||||
@@ -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;
|
||||
@@ -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);
|
||||
@@ -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 $$;
|
||||
@@ -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;
|
||||
@@ -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);
|
||||
@@ -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 $$;
|
||||
@@ -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;
|
||||
@@ -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);
|
||||
@@ -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 $$;
|
||||
@@ -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;
|
||||
@@ -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);
|
||||
@@ -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 $$;
|
||||
@@ -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;
|
||||
@@ -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);
|
||||
@@ -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 $$;
|
||||
@@ -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;
|
||||
Reference in New Issue
Block a user