Files
ddadmin 5972bcc2d8 refactor(all): унифицированы служебные поля STG — переход на канонический нейминг
- Зачем:
  - STG-слой использовал legacy-имена (batch_id, load_dttm, src_created_at_ts),
    тогда как ODS/DDS/DM уже работали с каноном (_load_id, _load_ts, event_ts).
    Студент видел разные имена для одного понятия — это убрано.
- Что:
  - переименованы колонки в 9 STG DDL: batch_id→_load_id, load_dttm→_load_ts,
    src_created_at_ts→event_ts; добавлен NOT NULL для _load_id во всех таблицах.
  - обновлены 9 STG Load, 9 STG DQ, 9 ODS Load, 9 ODS DQ (INSERT/SELECT/WHERE).
  - обновлены DAG-файлы bookings_to_gp_stage.py и bookings_to_gp_ods.py
    (встроенный SQL резолвера, комментарии; Python-идентификаторы не тронуты).
  - обновлены тесты и ~15 документов (naming_conventions, PRD, db_schema,
    design-docs, qa-plan, README, TESTING и др.).
- Проверка:
  - grep -rn 'load_dttm\|src_created_at_ts' sql/ airflow/ tests/ — 0 совпадений.
  - make test — 4 passed.
  - e2e-etl: day1 прошёл полностью, day2 стартовал без ошибок.
2026-03-10 00:07:08 +03:00

141 lines
4.5 KiB
SQL
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
-- Загрузка ODS по flights: SCD1 (UPDATE изменившихся + INSERT новых).
-- Используем паттерн Temporary Table для предотвращения гонки HWM между UPDATE и INSERT.
-- 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 (по HWM).
SELECT DISTINCT
NULLIF(s.flight_id, '')::INTEGER AS flight_id
FROM stg.segments AS s
WHERE s._load_ts > (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.event_ts,
s._load_ts,
s._load_id
FROM stg.flights AS s
WHERE s.flight_id IS NOT NULL
AND s.flight_id <> ''
),
src_union AS (
-- 1) Рейсы из текущего batch (по HWM).
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_ts,
f._load_id
FROM stg_flights_typed AS f
WHERE f._load_ts > (SELECT COALESCE(MAX(_load_ts), '1900-01-01 00:00:00'::TIMESTAMP) FROM ods.flights)
UNION ALL
-- 2) Добираем из истории STG рейсы, на которые ссылаются segments текущего batch.
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_ts,
f._load_id
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,
u._load_id,
u._load_ts,
ROW_NUMBER() OVER (
PARTITION BY u.flight_id
ORDER BY u.event_ts DESC NULLS LAST, u._load_ts 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,
scheduled_departure = s.scheduled_departure,
scheduled_arrival = s.scheduled_arrival,
actual_departure = s.actual_departure,
actual_arrival = s.actual_arrival,
event_ts = s.event_ts,
_load_id = s._load_id, -- Сохраняем оригинальный lineage из STG
_load_ts = s._load_ts -- Фиксируем время 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
OR o.scheduled_departure IS DISTINCT FROM s.scheduled_departure
OR o.scheduled_arrival IS DISTINCT FROM s.scheduled_arrival
OR o.actual_departure IS DISTINCT FROM s.actual_departure
OR o.actual_arrival IS DISTINCT FROM s.actual_arrival
OR o.event_ts IS DISTINCT FROM s.event_ts
);
-- 3. INSERT новых строк.
INSERT INTO ods.flights (
flight_id,
route_no,
status,
scheduled_departure,
scheduled_arrival,
actual_departure,
actual_arrival,
event_ts,
_load_id,
_load_ts
)
SELECT
s.flight_id,
s.route_no,
s.status,
s.scheduled_departure,
s.scheduled_arrival,
s.actual_departure,
s.actual_arrival,
s.event_ts,
s._load_id,
s._load_ts
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;