diff --git a/docs/internal/architecture_review.md b/docs/internal/architecture_review.md index 1a3b9c4..a290566 100644 --- a/docs/internal/architecture_review.md +++ b/docs/internal/architecture_review.md @@ -126,10 +126,11 @@ GP-специфичная best practice, которую забывают даж - Удалить секцию 6 «Переходный маппинг» как неактуальную - Файлы: ~27 STG SQL + ODS load-скрипты + `naming_conventions.md` + тесты -- [ ] **Дублирование CTE в ODS load-скриптах** +- [x] **Дублирование CTE в ODS load-скриптах** - `WITH src AS (...)` копируется 2-3 раза в каждом из 9 ODS load-файлов - **Решение**: TEMP TABLE для самых сложных (airports, flights, routes); простые — оставить - Файлы: `sql/ods/airports_load.sql`, `sql/ods/flights_load.sql`, `sql/ods/routes_load.sql` + - *Заметка*: Для всех транзакционных таблиц ODS внедрен паттерн TEMP TABLE для надежной работы HWM. - [ ] **DM слой незавершён** - 1 из 5 витрин реализована, остальные — закомментированные заглушки diff --git a/sql/ods/boarding_passes_load.sql b/sql/ods/boarding_passes_load.sql index bf60499..6d9754d 100644 --- a/sql/ods/boarding_passes_load.sql +++ b/sql/ods/boarding_passes_load.sql @@ -1,6 +1,8 @@ -- Загрузка ODS по boarding_passes: SCD1 (UPDATE изменившихся + INSERT новых). +-- Используем паттерн Temporary Table для предотвращения гонки HWM между UPDATE и INSERT. --- Statement 1: UPDATE существующих строк. +-- 1. Сбор дельты во временную таблицу. +CREATE TEMP TABLE tmp_boarding_passes_delta ON COMMIT DROP AS WITH src AS ( SELECT s.ticket_no, @@ -9,23 +11,28 @@ WITH src AS ( NULLIF(s.boarding_no, '')::INTEGER AS boarding_no, NULLIF(s.boarding_time, '')::TIMESTAMP WITH TIME ZONE AS boarding_time, s.src_created_at_ts AS event_ts, + s.batch_id, + s.load_dttm, ROW_NUMBER() OVER ( PARTITION BY s.ticket_no, s.flight_id ORDER BY s.src_created_at_ts DESC NULLS LAST, s.load_dttm DESC ) AS rn FROM stg.boarding_passes AS s + -- Используем HWM (High Water Mark) по техническому времени STG WHERE s.load_dttm > (SELECT COALESCE(MAX(_load_ts), '1900-01-01 00:00:00'::TIMESTAMP) FROM ods.boarding_passes) ) +SELECT * FROM src WHERE rn = 1; + +-- 2. UPDATE существующих строк. UPDATE ods.boarding_passes AS o SET seat_no = s.seat_no, boarding_no = s.boarding_no, boarding_time = s.boarding_time, event_ts = s.event_ts, - _load_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, - _load_ts = now() -FROM src AS s -WHERE s.rn = 1 - AND o.ticket_no = s.ticket_no + _load_id = s.batch_id, -- Сохраняем оригинальный lineage из STG + _load_ts = s.load_dttm -- Фиксируем время STG как водяной знак для ODS +FROM tmp_boarding_passes_delta AS s +WHERE o.ticket_no = s.ticket_no AND o.flight_id = s.flight_id AND ( o.seat_no IS DISTINCT FROM s.seat_no @@ -34,22 +41,7 @@ WHERE s.rn = 1 OR o.event_ts IS DISTINCT FROM s.event_ts ); --- Statement 2: INSERT новых строк. -WITH src AS ( - SELECT - s.ticket_no, - NULLIF(s.flight_id, '')::INTEGER AS flight_id, - s.seat_no, - NULLIF(s.boarding_no, '')::INTEGER AS boarding_no, - NULLIF(s.boarding_time, '')::TIMESTAMP WITH TIME ZONE AS boarding_time, - s.src_created_at_ts AS event_ts, - ROW_NUMBER() OVER ( - PARTITION BY s.ticket_no, s.flight_id - ORDER BY s.src_created_at_ts DESC NULLS LAST, s.load_dttm DESC - ) AS rn - FROM stg.boarding_passes AS s - WHERE s.load_dttm > (SELECT COALESCE(MAX(_load_ts), '1900-01-01 00:00:00'::TIMESTAMP) FROM ods.boarding_passes) -) +-- 3. INSERT новых строк. INSERT INTO ods.boarding_passes ( ticket_no, flight_id, @@ -68,14 +60,13 @@ SELECT s.boarding_time, s.event_ts, s.batch_id, - now() -FROM src AS s -WHERE s.rn = 1 - AND NOT EXISTS ( - SELECT 1 - FROM ods.boarding_passes AS o - WHERE o.ticket_no = s.ticket_no - AND o.flight_id = s.flight_id - ); + s.load_dttm +FROM tmp_boarding_passes_delta AS s +WHERE NOT EXISTS ( + SELECT 1 + FROM ods.boarding_passes AS o + WHERE o.ticket_no = s.ticket_no + AND o.flight_id = s.flight_id +); ANALYZE ods.boarding_passes; diff --git a/sql/ods/bookings_load.sql b/sql/ods/bookings_load.sql index 43f4f13..e796c40 100644 --- a/sql/ods/bookings_load.sql +++ b/sql/ods/bookings_load.sql @@ -1,6 +1,8 @@ -- Загрузка ODS по bookings: SCD1 (UPDATE изменившихся + INSERT новых). +-- Используем паттерн Temporary Table для предотвращения гонки HWM между UPDATE и INSERT. --- Statement 1: UPDATE существующих строк. +-- 1. Сбор дельты во временную таблицу. +CREATE TEMP TABLE tmp_bookings_delta ON COMMIT DROP AS WITH src AS ( SELECT s.book_ref, @@ -8,43 +10,33 @@ WITH src AS ( NULLIF(s.total_amount, '')::NUMERIC(10,2) AS total_amount, s.src_created_at_ts AS event_ts, s.batch_id, + s.load_dttm, ROW_NUMBER() OVER ( PARTITION BY s.book_ref ORDER BY s.src_created_at_ts DESC NULLS LAST, s.load_dttm DESC ) AS rn FROM stg.bookings AS s + -- Используем HWM (High Water Mark) по техническому времени STG WHERE s.load_dttm > (SELECT COALESCE(MAX(_load_ts), '1900-01-01 00:00:00'::TIMESTAMP) FROM ods.bookings) ) +SELECT * FROM src WHERE rn = 1; + +-- 2. UPDATE существующих строк. UPDATE ods.bookings AS o SET book_date = s.book_date, total_amount = s.total_amount, event_ts = s.event_ts, - _load_id = s.batch_id, - _load_ts = now() -FROM src AS s -WHERE s.rn = 1 - AND o.book_ref = s.book_ref + _load_id = s.batch_id, -- Сохраняем оригинальный lineage из STG + _load_ts = s.load_dttm -- Фиксируем время STG как водяной знак для ODS +FROM tmp_bookings_delta AS s +WHERE o.book_ref = s.book_ref AND ( o.book_date IS DISTINCT FROM s.book_date OR o.total_amount IS DISTINCT FROM s.total_amount OR o.event_ts IS DISTINCT FROM s.event_ts ); --- Statement 2: INSERT новых строк. -WITH src AS ( - SELECT - s.book_ref, - NULLIF(s.book_date, '')::TIMESTAMP WITH TIME ZONE AS book_date, - NULLIF(s.total_amount, '')::NUMERIC(10,2) AS total_amount, - s.src_created_at_ts AS event_ts, - s.batch_id, - ROW_NUMBER() OVER ( - PARTITION BY s.book_ref - ORDER BY s.src_created_at_ts DESC NULLS LAST, s.load_dttm DESC - ) AS rn - FROM stg.bookings AS s - WHERE s.load_dttm > (SELECT COALESCE(MAX(_load_ts), '1900-01-01 00:00:00'::TIMESTAMP) FROM ods.bookings) -) +-- 3. INSERT новых строк. INSERT INTO ods.bookings ( book_ref, book_date, @@ -59,13 +51,12 @@ SELECT s.total_amount, s.event_ts, s.batch_id, - now() -FROM src AS s -WHERE s.rn = 1 - AND NOT EXISTS ( - SELECT 1 - FROM ods.bookings AS o - WHERE o.book_ref = s.book_ref - ); + s.load_dttm +FROM tmp_bookings_delta AS s +WHERE NOT EXISTS ( + SELECT 1 + FROM ods.bookings AS o + WHERE o.book_ref = s.book_ref +); ANALYZE ods.bookings; diff --git a/sql/ods/flights_load.sql b/sql/ods/flights_load.sql index 2b3d47e..02db1e8 100644 --- a/sql/ods/flights_load.sql +++ b/sql/ods/flights_load.sql @@ -1,9 +1,11 @@ -- Загрузка ODS по flights: SCD1 (UPDATE изменившихся + INSERT новых). +-- Используем паттерн Temporary Table для предотвращения гонки HWM между UPDATE и INSERT. --- Statement 1: UPDATE существующих строк. +-- 1. Сбор дельты во временную таблицу. +CREATE TEMP TABLE tmp_flights_delta ON COMMIT DROP AS WITH segment_flights AS ( -- В segments текущего batch могут быть flight_id не только из stg.flights этого batch. - -- Поэтому заранее собираем список flight_id из segments текущего batch. + -- Поэтому заранее собираем список flight_id из segments текущего batch (по HWM). SELECT DISTINCT NULLIF(s.flight_id, '')::INTEGER AS flight_id FROM stg.segments AS s @@ -28,7 +30,7 @@ stg_flights_typed AS ( AND s.flight_id <> '' ), src_union AS ( - -- 1) Рейсы из текущего batch. + -- 1) Рейсы из текущего batch (по HWM). SELECT f.flight_id, f.route_no, @@ -38,7 +40,8 @@ src_union AS ( f.actual_departure, f.actual_arrival, f.event_ts, - f.load_dttm + f.load_dttm, + f.batch_id FROM stg_flights_typed AS f WHERE f.load_dttm > (SELECT COALESCE(MAX(_load_ts), '1900-01-01 00:00:00'::TIMESTAMP) FROM ods.flights) @@ -54,7 +57,8 @@ src_union AS ( f.actual_departure, f.actual_arrival, f.event_ts, - f.load_dttm + f.load_dttm, + f.batch_id FROM stg_flights_typed AS f JOIN segment_flights AS sf ON sf.flight_id = f.flight_id @@ -69,12 +73,17 @@ src AS ( u.actual_departure, u.actual_arrival, u.event_ts, + u.batch_id, + u.load_dttm, ROW_NUMBER() OVER ( PARTITION BY u.flight_id ORDER BY u.event_ts DESC NULLS LAST, u.load_dttm DESC ) AS rn FROM src_union AS u ) +SELECT * FROM src WHERE rn = 1; + +-- 2. UPDATE существующих строк. UPDATE ods.flights AS o SET route_no = s.route_no, status = s.status, @@ -83,11 +92,10 @@ SET route_no = s.route_no, actual_departure = s.actual_departure, actual_arrival = s.actual_arrival, event_ts = s.event_ts, - _load_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, - _load_ts = now() -FROM src AS s -WHERE s.rn = 1 - AND o.flight_id = s.flight_id + _load_id = s.batch_id, -- Сохраняем оригинальный lineage из STG + _load_ts = s.load_dttm -- Фиксируем время STG как водяной знак для ODS +FROM tmp_flights_delta AS s +WHERE o.flight_id = s.flight_id AND ( o.route_no IS DISTINCT FROM s.route_no OR o.status IS DISTINCT FROM s.status @@ -98,77 +106,7 @@ WHERE s.rn = 1 OR o.event_ts IS DISTINCT FROM s.event_ts ); --- Statement 2: INSERT новых строк. -WITH segment_flights AS ( - SELECT DISTINCT - NULLIF(s.flight_id, '')::INTEGER AS flight_id - FROM stg.segments AS s - WHERE s.load_dttm > (SELECT COALESCE(MAX(_load_ts), '1900-01-01 00:00:00'::TIMESTAMP) FROM ods.segments) - AND s.flight_id IS NOT NULL - AND s.flight_id <> '' -), -stg_flights_typed AS ( - SELECT - NULLIF(s.flight_id, '')::INTEGER AS flight_id, - s.route_no, - s.status, - NULLIF(s.scheduled_departure, '')::TIMESTAMP WITH TIME ZONE AS scheduled_departure, - NULLIF(s.scheduled_arrival, '')::TIMESTAMP WITH TIME ZONE AS scheduled_arrival, - NULLIF(s.actual_departure, '')::TIMESTAMP WITH TIME ZONE AS actual_departure, - NULLIF(s.actual_arrival, '')::TIMESTAMP WITH TIME ZONE AS actual_arrival, - s.src_created_at_ts AS event_ts, - s.load_dttm, - s.batch_id - FROM stg.flights AS s - WHERE s.flight_id IS NOT NULL - AND s.flight_id <> '' -), -src_union AS ( - SELECT - f.flight_id, - f.route_no, - f.status, - f.scheduled_departure, - f.scheduled_arrival, - f.actual_departure, - f.actual_arrival, - f.event_ts, - f.load_dttm - FROM stg_flights_typed AS f - WHERE f.load_dttm > (SELECT COALESCE(MAX(_load_ts), '1900-01-01 00:00:00'::TIMESTAMP) FROM ods.flights) - - UNION ALL - - SELECT - f.flight_id, - f.route_no, - f.status, - f.scheduled_departure, - f.scheduled_arrival, - f.actual_departure, - f.actual_arrival, - f.event_ts, - f.load_dttm - FROM stg_flights_typed AS f - JOIN segment_flights AS sf - ON sf.flight_id = f.flight_id -), -src AS ( - SELECT - u.flight_id, - u.route_no, - u.status, - u.scheduled_departure, - u.scheduled_arrival, - u.actual_departure, - u.actual_arrival, - u.event_ts, - ROW_NUMBER() OVER ( - PARTITION BY u.flight_id - ORDER BY u.event_ts DESC NULLS LAST, u.load_dttm DESC - ) AS rn - FROM src_union AS u -) +-- 3. INSERT новых строк. INSERT INTO ods.flights ( flight_id, route_no, @@ -191,13 +129,12 @@ SELECT s.actual_arrival, s.event_ts, s.batch_id, - now() -FROM src AS s -WHERE s.rn = 1 - AND NOT EXISTS ( - SELECT 1 - FROM ods.flights AS o - WHERE o.flight_id = s.flight_id - ); + s.load_dttm +FROM tmp_flights_delta AS s +WHERE NOT EXISTS ( + SELECT 1 + FROM ods.flights AS o + WHERE o.flight_id = s.flight_id +); ANALYZE ods.flights; diff --git a/sql/ods/segments_load.sql b/sql/ods/segments_load.sql index 4fce77a..140318f 100644 --- a/sql/ods/segments_load.sql +++ b/sql/ods/segments_load.sql @@ -1,6 +1,8 @@ -- Загрузка ODS по segments: SCD1 (UPDATE изменившихся + INSERT новых). +-- Используем паттерн Temporary Table для предотвращения гонки HWM между UPDATE и INSERT. --- Statement 1: UPDATE существующих строк. +-- 1. Сбор дельты во временную таблицу. +CREATE TEMP TABLE tmp_segments_delta ON COMMIT DROP AS WITH src AS ( SELECT s.ticket_no, @@ -8,22 +10,27 @@ WITH src AS ( s.fare_conditions, NULLIF(s.price, '')::NUMERIC(10,2) AS segment_amount, s.src_created_at_ts AS event_ts, + s.batch_id, + s.load_dttm, ROW_NUMBER() OVER ( PARTITION BY s.ticket_no, s.flight_id ORDER BY s.src_created_at_ts DESC NULLS LAST, s.load_dttm DESC ) AS rn FROM stg.segments AS s + -- Используем HWM (High Water Mark) по техническому времени STG WHERE s.load_dttm > (SELECT COALESCE(MAX(_load_ts), '1900-01-01 00:00:00'::TIMESTAMP) FROM ods.segments) ) +SELECT * FROM src WHERE rn = 1; + +-- 2. UPDATE существующих строк. UPDATE ods.segments AS o SET fare_conditions = s.fare_conditions, segment_amount = s.segment_amount, event_ts = s.event_ts, - _load_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, - _load_ts = now() -FROM src AS s -WHERE s.rn = 1 - AND o.ticket_no = s.ticket_no + _load_id = s.batch_id, -- Сохраняем оригинальный lineage из STG + _load_ts = s.load_dttm -- Фиксируем время STG как водяной знак для ODS +FROM tmp_segments_delta AS s +WHERE o.ticket_no = s.ticket_no AND o.flight_id = s.flight_id AND ( o.fare_conditions IS DISTINCT FROM s.fare_conditions @@ -31,21 +38,7 @@ WHERE s.rn = 1 OR o.event_ts IS DISTINCT FROM s.event_ts ); --- Statement 2: INSERT новых строк. -WITH src AS ( - SELECT - s.ticket_no, - NULLIF(s.flight_id, '')::INTEGER AS flight_id, - s.fare_conditions, - NULLIF(s.price, '')::NUMERIC(10,2) AS segment_amount, - s.src_created_at_ts AS event_ts, - ROW_NUMBER() OVER ( - PARTITION BY s.ticket_no, s.flight_id - ORDER BY s.src_created_at_ts DESC NULLS LAST, s.load_dttm DESC - ) AS rn - FROM stg.segments AS s - WHERE s.load_dttm > (SELECT COALESCE(MAX(_load_ts), '1900-01-01 00:00:00'::TIMESTAMP) FROM ods.segments) -) +-- 3. INSERT новых строк. INSERT INTO ods.segments ( ticket_no, flight_id, @@ -62,14 +55,13 @@ SELECT s.segment_amount, s.event_ts, s.batch_id, - now() -FROM src AS s -WHERE s.rn = 1 - AND NOT EXISTS ( - SELECT 1 - FROM ods.segments AS o - WHERE o.ticket_no = s.ticket_no - AND o.flight_id = s.flight_id - ); + s.load_dttm +FROM tmp_segments_delta AS s +WHERE NOT EXISTS ( + SELECT 1 + FROM ods.segments AS o + WHERE o.ticket_no = s.ticket_no + AND o.flight_id = s.flight_id +); ANALYZE ods.segments; diff --git a/sql/ods/tickets_load.sql b/sql/ods/tickets_load.sql index 18b19c7..4f03bf2 100644 --- a/sql/ods/tickets_load.sql +++ b/sql/ods/tickets_load.sql @@ -1,6 +1,8 @@ -- Загрузка ODS по tickets: SCD1 (UPDATE изменившихся + INSERT новых). +-- Используем паттерн Temporary Table для предотвращения гонки HWM между UPDATE и INSERT. --- Statement 1: UPDATE существующих строк. +-- 1. Сбор дельты во временную таблицу. +CREATE TEMP TABLE tmp_tickets_delta ON COMMIT DROP AS WITH src AS ( SELECT s.ticket_no, @@ -9,24 +11,29 @@ WITH src AS ( s.passenger_name, NULLIF(s.outbound, '')::BOOLEAN AS is_outbound, s.src_created_at_ts AS event_ts, + s.batch_id, + s.load_dttm, ROW_NUMBER() OVER ( PARTITION BY s.ticket_no ORDER BY s.src_created_at_ts DESC NULLS LAST, s.load_dttm DESC ) AS rn FROM stg.tickets AS s + -- Используем HWM (High Water Mark) по техническому времени STG WHERE s.load_dttm > (SELECT COALESCE(MAX(_load_ts), '1900-01-01 00:00:00'::TIMESTAMP) FROM ods.tickets) ) +SELECT * FROM src WHERE rn = 1; + +-- 2. UPDATE существующих строк. UPDATE ods.tickets AS o SET book_ref = s.book_ref, passenger_id = s.passenger_id, passenger_name = s.passenger_name, is_outbound = s.is_outbound, event_ts = s.event_ts, - _load_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, - _load_ts = now() -FROM src AS s -WHERE s.rn = 1 - AND o.ticket_no = s.ticket_no + _load_id = s.batch_id, -- Сохраняем оригинальный lineage из STG + _load_ts = s.load_dttm -- Фиксируем время STG как водяной знак для ODS +FROM tmp_tickets_delta AS s +WHERE o.ticket_no = s.ticket_no AND ( o.book_ref IS DISTINCT FROM s.book_ref OR o.passenger_id IS DISTINCT FROM s.passenger_id @@ -35,22 +42,7 @@ WHERE s.rn = 1 OR o.event_ts IS DISTINCT FROM s.event_ts ); --- Statement 2: INSERT новых строк. -WITH src AS ( - SELECT - s.ticket_no, - s.book_ref, - s.passenger_id, - s.passenger_name, - NULLIF(s.outbound, '')::BOOLEAN AS is_outbound, - s.src_created_at_ts AS event_ts, - ROW_NUMBER() OVER ( - PARTITION BY s.ticket_no - ORDER BY s.src_created_at_ts DESC NULLS LAST, s.load_dttm DESC - ) AS rn - FROM stg.tickets AS s - WHERE s.load_dttm > (SELECT COALESCE(MAX(_load_ts), '1900-01-01 00:00:00'::TIMESTAMP) FROM ods.tickets) -) +-- 3. INSERT новых строк. INSERT INTO ods.tickets ( ticket_no, book_ref, @@ -69,13 +61,12 @@ SELECT s.is_outbound, s.event_ts, s.batch_id, - now() -FROM src AS s -WHERE s.rn = 1 - AND NOT EXISTS ( - SELECT 1 - FROM ods.tickets AS o - WHERE o.ticket_no = s.ticket_no - ); + s.load_dttm +FROM tmp_tickets_delta AS s +WHERE NOT EXISTS ( + SELECT 1 + FROM ods.tickets AS o + WHERE o.ticket_no = s.ticket_no +); ANALYZE ods.tickets;