docs(dwh): зафиксированы конвенции нейминга и обновлён план ODS
- Зачем: - нужен единый стандарт именования полей, чтобы новые слои не расходились с учебными материалами. - Что: - добавлен единый документ с правилами нейминга `docs/internal/naming_conventions.md`. - полностью переписан `docs/internal/bookings_ods_design.md` в эталонный учебный план ODS (SCD1, батч-контракт, DQ, граф DAG). - добавлены ссылки на стандарт нейминга в `docs/README.md`, `docs/internal/db_schema.md` и `AGENTS.md`. - Проверка: - проверен diff по измененным файлам (`git diff`).
This commit is contained in:
@@ -41,6 +41,7 @@
|
|||||||
- `sql/src/` — скрипты, работающие с исходными системами (например, `bookings_generate_day_if_missing.sql`);
|
- `sql/src/` — скрипты, работающие с исходными системами (например, `bookings_generate_day_if_missing.sql`);
|
||||||
- `sql/stg/` — скрипты для стейджинга (`bookings_ddl.sql`, `bookings_load.sql`, `bookings_dq.sql`);
|
- `sql/stg/` — скрипты для стейджинга (`bookings_ddl.sql`, `bookings_load.sql`, `bookings_dq.sql`);
|
||||||
- в будущем можно добавить `sql/ods/`, `sql/dds/`, `sql/dm/` по мере роста стенда.
|
- в будущем можно добавить `sql/ods/`, `sql/dds/`, `sql/dm/` по мере роста стенда.
|
||||||
|
- Нейминг служебных полей и SCD-полей фиксирован в `docs/internal/naming_conventions.md` (единый источник для всех новых слоёв).
|
||||||
- Именование файлов: `{объект}_{роль}.sql`, где:
|
- Именование файлов: `{объект}_{роль}.sql`, где:
|
||||||
- `объект` — логическое имя сущности (`bookings`, `orders`, и т.п.);
|
- `объект` — логическое имя сущности (`bookings`, `orders`, и т.п.);
|
||||||
- `роль` — `ddl` (создание/изменение объектов), `load` (загрузка/инкремент), `dq` (проверки качества данных) и т.п.
|
- `роль` — `ddl` (создание/изменение объектов), `load` (загрузка/инкремент), `dq` (проверки качества данных) и т.п.
|
||||||
|
|||||||
@@ -12,6 +12,7 @@
|
|||||||
## Технические детали (опционально)
|
## Технические детали (опционально)
|
||||||
|
|
||||||
- [Как устроен Docker-стенд (образы, Connections, переменные окружения)](stack.md)
|
- [Как устроен Docker-стенд (образы, Connections, переменные окружения)](stack.md)
|
||||||
|
- [Единые конвенции нейминга DWH (служебные поля и SCD)](internal/naming_conventions.md)
|
||||||
- [PXF в этом проекте (проектная реализация)](internal/pxf_bookings.md)
|
- [PXF в этом проекте (проектная реализация)](internal/pxf_bookings.md)
|
||||||
- [Дизайн stg для bookings (черновик)](internal/bookings_stg_design.md)
|
- [Дизайн stg для bookings (черновик)](internal/bookings_stg_design.md)
|
||||||
- [Про время/UTC в bookings (черновик)](internal/bookings_tz.md)
|
- [Про время/UTC в bookings (черновик)](internal/bookings_tz.md)
|
||||||
|
|||||||
@@ -1,188 +1,402 @@
|
|||||||
# ODS Layer: план реализации
|
# ODS Layer: эталонный учебный план реализации (v2)
|
||||||
|
|
||||||
## Контекст
|
## Контекст
|
||||||
|
|
||||||
STG-слой уже реализован как учебный эталон: сырые данные из bookings-db грузятся через PXF, хранятся как TEXT, контролируются batch_id.
|
STG-слой уже реализован как учебный эталон:
|
||||||
|
- данные из `bookings-db` читаются через PXF;
|
||||||
|
- в STG бизнес-колонки хранятся как `TEXT`;
|
||||||
|
- загрузка и DQ работают батчами (`batch_id = {{ run_id }}`).
|
||||||
|
|
||||||
**Следующий шаг** — ODS (Operational Data Store): типизированный, очищенный и исторически отслеживаемый слой.
|
Этот документ фиксирует **простую и каноничную** реализацию ODS для менти.
|
||||||
|
|
||||||
**Образовательные цели** реализации:
|
|
||||||
1. Показать переход TEXT → правильные типы данных
|
|
||||||
2. Объяснить SCD Type 2 на конкретных SQL-примерах
|
|
||||||
3. Показать UPSERT-паттерн для транзакционных данных
|
|
||||||
4. Продемонстрировать DQ-проверки уровня ODS (типы, ссылочная целостность)
|
|
||||||
5. Показать параллельный граф DAG с правильными зависимостями
|
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
## Архитектура ODS
|
## 1) Что считаем эталоном для ODS
|
||||||
|
|
||||||
### Отличие ODS от STG
|
### 1.1. Роль ODS в этом стенде
|
||||||
|
|
||||||
| Аспект | STG | ODS |
|
ODS в учебном проекте — это:
|
||||||
|--------|-----|-----|
|
- типизированные и очищенные данные;
|
||||||
| Типы | Всё TEXT | Правильные типы (TIMESTAMP, NUMERIC и т.д.) |
|
- одна актуальная запись на бизнес-ключ;
|
||||||
| Дедупликация | Нет (все инкременты хранятся) | Одна активная запись на бизнес-ключ |
|
- удобный слой для последующей сборки DDS/DM.
|
||||||
| История | Нет | SCD Type 2 для справочников |
|
|
||||||
| Хранилище | appendonly | heap (справочники) + appendonly (транзакции) |
|
|
||||||
| DQ-проверки | Количество строк, NOT NULL по TEXT | Типы, ссылочная целостность, бизнес-правила |
|
|
||||||
|
|
||||||
### Паттерны загрузки
|
### 1.2. Что делаем, что не делаем
|
||||||
|
|
||||||
**Справочники (airports, airplanes, routes, seats) — SCD Type 2:**
|
Делаем в ODS:
|
||||||
- Читаем актуальный снапшот из STG (последний batch)
|
- приведение типов (`TEXT -> TIMESTAMPTZ/NUMERIC/INT/BOOLEAN/...`);
|
||||||
- Закрываем старые версии при изменении атрибутов (is_active=false, dw_end_date)
|
- дедупликацию внутри батча;
|
||||||
- Вставляем новые версии (dw_version++, dw_start_date=CURRENT_DATE)
|
- `UPSERT` (SCD Type 1): обновляем текущую запись при изменении, вставляем новые.
|
||||||
- Вставляем новые записи (dw_version=1)
|
|
||||||
- Хранилище: heap-таблицы (без appendonly), т.к. UPDATE-операции частые
|
|
||||||
|
|
||||||
**Транзакции (bookings, tickets, flights, segments, boarding_passes) — UPSERT:**
|
Не делаем в ODS (в базовом эталоне):
|
||||||
- Обновляем изменившиеся записи (UPDATE)
|
- SCD Type 2 с периодами действия;
|
||||||
- Вставляем новые записи (INSERT WHERE NOT EXISTS)
|
- сложную обработку late-arriving/backdated событий;
|
||||||
- Хранилище: appendonly (как в STG)
|
- отдельный DQ-слой с хранением результатов.
|
||||||
|
|
||||||
|
### 1.3. Где хранится история изменений
|
||||||
|
|
||||||
|
- История «как приходили данные» уже сохраняется в STG (append + `batch_id`).
|
||||||
|
- Историзацию измерений (SCD2) показываем позже в DDS (как в учебной статье `dwh-modeling`).
|
||||||
|
|
||||||
|
Итог: **ODS = текущий слой (current state), простой и понятный**.
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
## Схема ODS-таблиц
|
## 2) Нейминг служебных полей (консистентно с de-roadmap)
|
||||||
|
|
||||||
### Справочники (с SCD Type 2)
|
Источник правил: [`docs/internal/naming_conventions.md`](naming_conventions.md).
|
||||||
|
|
||||||
#### ods.airports
|
В ODS используем такие техполя:
|
||||||
|
|
||||||
|
- `_load_id TEXT NOT NULL` — идентификатор загрузки (берём `stg_batch_id`);
|
||||||
|
- `_load_ts TIMESTAMP NOT NULL DEFAULT now()` — время загрузки в ODS;
|
||||||
|
- `event_ts TIMESTAMP` — время события из источника (если у сущности оно есть).
|
||||||
|
|
||||||
|
### 2.1. Маппинг из текущего STG
|
||||||
|
|
||||||
|
- `stg.batch_id` -> `ods._load_id`
|
||||||
|
- `stg.load_dttm` не переносим 1:1; в ODS пишем собственный `ods._load_ts = now()`
|
||||||
|
- `stg.src_created_at_ts` -> `ods.event_ts` (для транзакционных таблиц)
|
||||||
|
|
||||||
|
### 2.2. Почему так
|
||||||
|
|
||||||
|
- нейминг совпадает с учебной статьёй (`_load_id`, `_load_ts`);
|
||||||
|
- студентам проще переносить паттерн между проектами;
|
||||||
|
- разделяем «когда событие произошло» (`event_ts`) и «когда загрузили в слой» (`_load_ts`).
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## 3) Гранулярность и бизнес-ключи ODS
|
||||||
|
|
||||||
|
| Таблица | Зерно | Бизнес-ключ |
|
||||||
|
|---|---|---|
|
||||||
|
| `ods.airports` | 1 строка = аэропорт | `airport_code` |
|
||||||
|
| `ods.airplanes` | 1 строка = самолёт | `airplane_code` |
|
||||||
|
| `ods.routes` | 1 строка = версия маршрута | `(route_no, validity)` |
|
||||||
|
| `ods.seats` | 1 строка = место в самолёте | `(airplane_code, seat_no)` |
|
||||||
|
| `ods.bookings` | 1 строка = бронирование | `book_ref` |
|
||||||
|
| `ods.tickets` | 1 строка = билет | `ticket_no` |
|
||||||
|
| `ods.flights` | 1 строка = рейс | `flight_id` |
|
||||||
|
| `ods.segments` | 1 строка = сегмент билета | `(ticket_no, flight_id)` |
|
||||||
|
| `ods.boarding_passes` | 1 строка = посадочный на сегмент | `(ticket_no, flight_id)` |
|
||||||
|
|
||||||
|
Критично для эталона:
|
||||||
|
- `routes` — **составной** ключ `(route_no, validity)`;
|
||||||
|
- `boarding_passes` — **составной** ключ `(ticket_no, flight_id)`.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## 4) Схема ODS-таблиц (v1, без SCD2)
|
||||||
|
|
||||||
|
Ниже — учебный минимум колонок. При необходимости можно добавлять бизнес-атрибуты без изменения паттерна загрузки.
|
||||||
|
|
||||||
|
### 4.1. Справочники
|
||||||
|
|
||||||
|
#### `ods.airports`
|
||||||
```sql
|
```sql
|
||||||
airport_code TEXT NOT NULL -- бизнес-ключ
|
airport_code TEXT NOT NULL
|
||||||
airport_name TEXT NOT NULL -- правильный тип (=TEXT, из JSONB в источнике)
|
airport_name TEXT NOT NULL
|
||||||
city TEXT NOT NULL
|
city TEXT NOT NULL
|
||||||
country TEXT NOT NULL
|
country TEXT NOT NULL
|
||||||
coordinates TEXT -- оставляем TEXT (сложный формат, DDS распарсит)
|
coordinates TEXT
|
||||||
timezone TEXT NOT NULL
|
timezone TEXT NOT NULL
|
||||||
-- SCD Type 2
|
_load_id TEXT NOT NULL
|
||||||
dw_start_date DATE NOT NULL DEFAULT CURRENT_DATE
|
_load_ts TIMESTAMP NOT NULL DEFAULT now()
|
||||||
dw_end_date DATE -- NULL = активная запись
|
DISTRIBUTED BY (airport_code)
|
||||||
is_active BOOLEAN NOT NULL DEFAULT true
|
|
||||||
dw_version INTEGER NOT NULL DEFAULT 1
|
|
||||||
-- Технические
|
|
||||||
load_dttm TIMESTAMP DEFAULT now()
|
|
||||||
batch_id TEXT
|
|
||||||
DISTRIBUTED BY (airport_code) -- heap, для эффективных UPDATE
|
|
||||||
```
|
```
|
||||||
|
|
||||||
#### ods.airplanes
|
#### `ods.airplanes`
|
||||||
```sql
|
```sql
|
||||||
airplane_code TEXT NOT NULL -- бизнес-ключ
|
airplane_code TEXT NOT NULL
|
||||||
model TEXT NOT NULL -- из JSONB в источнике
|
model TEXT NOT NULL
|
||||||
range INTEGER -- TEXT → INTEGER (км)
|
range_km INTEGER
|
||||||
speed INTEGER -- TEXT → INTEGER (км/ч)
|
speed_kmh INTEGER
|
||||||
-- SCD Type 2 + технические (аналогично airports)
|
_load_id TEXT NOT NULL
|
||||||
|
_load_ts TIMESTAMP NOT NULL DEFAULT now()
|
||||||
DISTRIBUTED BY (airplane_code)
|
DISTRIBUTED BY (airplane_code)
|
||||||
```
|
```
|
||||||
|
|
||||||
#### ods.routes
|
#### `ods.routes`
|
||||||
```sql
|
```sql
|
||||||
route_no TEXT NOT NULL -- бизнес-ключ
|
route_no TEXT NOT NULL
|
||||||
validity TEXT -- TSTZRANGE → TEXT (сложно парсить, документируем)
|
validity TEXT NOT NULL
|
||||||
departure_airport TEXT NOT NULL -- FK к ods.airports.airport_code
|
departure_airport TEXT NOT NULL
|
||||||
arrival_airport TEXT NOT NULL
|
arrival_airport TEXT NOT NULL
|
||||||
airplane_code TEXT NOT NULL -- FK к ods.airplanes.airplane_code
|
airplane_code TEXT NOT NULL
|
||||||
days_of_week TEXT -- int[] → TEXT (объясняем ограничение STG/PXF)
|
days_of_week TEXT
|
||||||
scheduled_time TIME -- TEXT → TIME (пример кастинга)
|
scheduled_departure_time TIME
|
||||||
duration INTERVAL -- TEXT → INTERVAL (ключевой пример кастинга)
|
scheduled_duration INTERVAL
|
||||||
-- SCD Type 2 + технические
|
_load_id TEXT NOT NULL
|
||||||
|
_load_ts TIMESTAMP NOT NULL DEFAULT now()
|
||||||
DISTRIBUTED BY (route_no)
|
DISTRIBUTED BY (route_no)
|
||||||
```
|
```
|
||||||
|
|
||||||
#### ods.seats
|
#### `ods.seats`
|
||||||
```sql
|
```sql
|
||||||
airplane_code TEXT NOT NULL -- бизнес-ключ (совместный)
|
airplane_code TEXT NOT NULL
|
||||||
seat_no TEXT NOT NULL -- бизнес-ключ (совместный)
|
seat_no TEXT NOT NULL
|
||||||
fare_conditions TEXT NOT NULL -- Economy/Comfort/Business
|
fare_conditions TEXT NOT NULL
|
||||||
-- SCD Type 2 + технические
|
_load_id TEXT NOT NULL
|
||||||
|
_load_ts TIMESTAMP NOT NULL DEFAULT now()
|
||||||
DISTRIBUTED BY (airplane_code)
|
DISTRIBUTED BY (airplane_code)
|
||||||
```
|
```
|
||||||
|
|
||||||
### Транзакционные (UPSERT)
|
### 4.2. Транзакционные
|
||||||
|
|
||||||
#### ods.bookings
|
#### `ods.bookings`
|
||||||
```sql
|
```sql
|
||||||
book_ref TEXT NOT NULL -- бизнес-ключ (CHAR(6))
|
book_ref TEXT NOT NULL
|
||||||
book_date TIMESTAMP WITH TIME ZONE -- TEXT → TIMESTAMPTZ (ключевой пример)
|
book_date TIMESTAMP WITH TIME ZONE NOT NULL
|
||||||
total_amount NUMERIC(10,2) NOT NULL -- TEXT → NUMERIC
|
total_amount NUMERIC(10,2) NOT NULL
|
||||||
-- Технические
|
event_ts TIMESTAMP
|
||||||
src_created_at_ts TIMESTAMP
|
_load_id TEXT NOT NULL
|
||||||
load_dttm TIMESTAMP DEFAULT now()
|
_load_ts TIMESTAMP NOT NULL DEFAULT now()
|
||||||
batch_id TEXT
|
DISTRIBUTED BY (book_ref)
|
||||||
DISTRIBUTED BY (book_ref) -- appendonly, выравнивание со STG
|
|
||||||
```
|
```
|
||||||
|
|
||||||
#### ods.tickets
|
#### `ods.tickets`
|
||||||
```sql
|
```sql
|
||||||
ticket_no TEXT NOT NULL -- бизнес-ключ
|
ticket_no TEXT NOT NULL
|
||||||
book_ref TEXT NOT NULL -- FK к ods.bookings
|
book_ref TEXT NOT NULL
|
||||||
passenger_id TEXT NOT NULL
|
passenger_id TEXT NOT NULL
|
||||||
passenger_name TEXT NOT NULL
|
passenger_name TEXT NOT NULL
|
||||||
outbound BOOLEAN -- TEXT → BOOLEAN (пример нетривиального кастинга)
|
is_outbound BOOLEAN
|
||||||
-- Технические (src_created_at_ts, load_dttm, batch_id)
|
event_ts TIMESTAMP
|
||||||
DISTRIBUTED BY (book_ref) -- совместно с bookings
|
_load_id TEXT NOT NULL
|
||||||
|
_load_ts TIMESTAMP NOT NULL DEFAULT now()
|
||||||
|
DISTRIBUTED BY (book_ref)
|
||||||
```
|
```
|
||||||
|
|
||||||
#### ods.flights
|
#### `ods.flights`
|
||||||
```sql
|
```sql
|
||||||
flight_id INTEGER NOT NULL -- TEXT → INTEGER (бизнес-ключ)
|
flight_id INTEGER NOT NULL
|
||||||
route_no TEXT NOT NULL -- FK к ods.routes
|
route_no TEXT NOT NULL
|
||||||
status TEXT NOT NULL
|
status TEXT NOT NULL
|
||||||
scheduled_departure TIMESTAMP WITH TIME ZONE
|
scheduled_departure TIMESTAMP WITH TIME ZONE
|
||||||
scheduled_arrival TIMESTAMP WITH TIME ZONE
|
scheduled_arrival TIMESTAMP WITH TIME ZONE
|
||||||
actual_departure TIMESTAMP WITH TIME ZONE
|
actual_departure TIMESTAMP WITH TIME ZONE
|
||||||
actual_arrival TIMESTAMP WITH TIME ZONE
|
actual_arrival TIMESTAMP WITH TIME ZONE
|
||||||
-- Технические (src_created_at_ts, load_dttm, batch_id)
|
event_ts TIMESTAMP
|
||||||
|
_load_id TEXT NOT NULL
|
||||||
|
_load_ts TIMESTAMP NOT NULL DEFAULT now()
|
||||||
DISTRIBUTED BY (flight_id)
|
DISTRIBUTED BY (flight_id)
|
||||||
```
|
```
|
||||||
|
|
||||||
#### ods.segments
|
#### `ods.segments`
|
||||||
```sql
|
```sql
|
||||||
ticket_no TEXT NOT NULL -- бизнес-ключ (совместный), FK к ods.tickets
|
ticket_no TEXT NOT NULL
|
||||||
flight_id INTEGER NOT NULL -- TEXT → INTEGER, FK к ods.flights
|
flight_id INTEGER NOT NULL
|
||||||
fare_conditions TEXT NOT NULL
|
fare_conditions TEXT NOT NULL
|
||||||
price NUMERIC(10,2) -- TEXT → NUMERIC (в STG называется 'price', не 'amount')
|
segment_amount NUMERIC(10,2)
|
||||||
-- Технические
|
event_ts TIMESTAMP
|
||||||
DISTRIBUTED BY (ticket_no) -- совместно с boarding_passes
|
_load_id TEXT NOT NULL
|
||||||
|
_load_ts TIMESTAMP NOT NULL DEFAULT now()
|
||||||
|
DISTRIBUTED BY (ticket_no)
|
||||||
```
|
```
|
||||||
|
|
||||||
#### ods.boarding_passes
|
#### `ods.boarding_passes`
|
||||||
```sql
|
```sql
|
||||||
ticket_no TEXT NOT NULL -- бизнес-ключ, FK к ods.tickets
|
ticket_no TEXT NOT NULL
|
||||||
flight_id INTEGER NOT NULL -- TEXT → INTEGER, FK к ods.flights
|
flight_id INTEGER NOT NULL
|
||||||
seat_no TEXT NOT NULL
|
seat_no TEXT NOT NULL
|
||||||
boarding_no INTEGER -- TEXT → INTEGER
|
boarding_no INTEGER
|
||||||
boarding_time TIMESTAMP WITH TIME ZONE -- TEXT → TIMESTAMPTZ
|
boarding_time TIMESTAMP WITH TIME ZONE
|
||||||
-- Технические
|
event_ts TIMESTAMP
|
||||||
DISTRIBUTED BY (ticket_no) -- совместно с segments
|
_load_id TEXT NOT NULL
|
||||||
|
_load_ts TIMESTAMP NOT NULL DEFAULT now()
|
||||||
|
DISTRIBUTED BY (ticket_no)
|
||||||
```
|
```
|
||||||
|
|
||||||
|
Примечание: в учебном варианте не опираемся на физические `PK/FK`-constraint в Greenplum, а проверяем целостность через DQ-скрипты.
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
## Структура файлов
|
## 5) Контракт батча для ODS
|
||||||
|
|
||||||
|
Чтобы ODS был воспроизводимым, в каждом запуске используем **один фиксированный `stg_batch_id`**.
|
||||||
|
|
||||||
|
### 5.1. Источник `stg_batch_id`
|
||||||
|
|
||||||
|
В `bookings_to_gp_ods`:
|
||||||
|
- принимаем `stg_batch_id` из `dag_run.conf`;
|
||||||
|
- если не передан — берём последний из `stg.bookings`;
|
||||||
|
- логируем, какой `stg_batch_id` выбран.
|
||||||
|
|
||||||
|
### 5.2. Как применяем
|
||||||
|
|
||||||
|
Во всех `sql/ods/*_load.sql`:
|
||||||
|
- читаем STG только с `WHERE batch_id = :stg_batch_id`;
|
||||||
|
- пишем в ODS `_load_id = :stg_batch_id`, `_load_ts = now()`.
|
||||||
|
|
||||||
|
Это простой и понятный паттерн: **один запуск ODS = один снимок STG-батча**.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## 6) SQL-паттерны загрузки (SCD1 / UPSERT)
|
||||||
|
|
||||||
|
### 6.1. Шаблон для справочника (пример `airports_load.sql`)
|
||||||
|
|
||||||
|
```sql
|
||||||
|
WITH src AS (
|
||||||
|
SELECT
|
||||||
|
airport_code,
|
||||||
|
airport_name,
|
||||||
|
city,
|
||||||
|
country,
|
||||||
|
coordinates,
|
||||||
|
timezone
|
||||||
|
FROM stg.airports
|
||||||
|
WHERE batch_id = '{{ params.stg_batch_id }}'::text
|
||||||
|
)
|
||||||
|
UPDATE ods.airports AS o
|
||||||
|
SET airport_name = s.airport_name,
|
||||||
|
city = s.city,
|
||||||
|
country = s.country,
|
||||||
|
coordinates = s.coordinates,
|
||||||
|
timezone = s.timezone,
|
||||||
|
_load_id = '{{ params.stg_batch_id }}'::text,
|
||||||
|
_load_ts = now()
|
||||||
|
FROM src AS s
|
||||||
|
WHERE o.airport_code = s.airport_code
|
||||||
|
AND (
|
||||||
|
o.airport_name <> s.airport_name OR
|
||||||
|
o.city <> s.city OR
|
||||||
|
o.country <> s.country OR
|
||||||
|
COALESCE(o.coordinates, '') <> COALESCE(s.coordinates, '') OR
|
||||||
|
o.timezone <> s.timezone
|
||||||
|
);
|
||||||
|
|
||||||
|
INSERT INTO ods.airports (
|
||||||
|
airport_code, airport_name, city, country, coordinates, timezone,
|
||||||
|
_load_id, _load_ts
|
||||||
|
)
|
||||||
|
SELECT
|
||||||
|
s.airport_code, s.airport_name, s.city, s.country, s.coordinates, s.timezone,
|
||||||
|
'{{ params.stg_batch_id }}'::text, now()
|
||||||
|
FROM src s
|
||||||
|
WHERE NOT EXISTS (
|
||||||
|
SELECT 1
|
||||||
|
FROM ods.airports o
|
||||||
|
WHERE o.airport_code = s.airport_code
|
||||||
|
);
|
||||||
|
|
||||||
|
ANALYZE ods.airports;
|
||||||
```
|
```
|
||||||
|
|
||||||
|
### 6.2. Шаблон для транзакции (пример `bookings_load.sql`)
|
||||||
|
|
||||||
|
```sql
|
||||||
|
WITH src AS (
|
||||||
|
SELECT DISTINCT ON (book_ref)
|
||||||
|
book_ref,
|
||||||
|
book_date::TIMESTAMP WITH TIME ZONE AS book_date,
|
||||||
|
total_amount::NUMERIC(10,2) AS total_amount,
|
||||||
|
src_created_at_ts AS event_ts
|
||||||
|
FROM stg.bookings
|
||||||
|
WHERE batch_id = '{{ params.stg_batch_id }}'::text
|
||||||
|
ORDER BY book_ref, src_created_at_ts DESC NULLS LAST, load_dttm DESC
|
||||||
|
)
|
||||||
|
UPDATE ods.bookings AS o
|
||||||
|
SET book_date = s.book_date,
|
||||||
|
total_amount = s.total_amount,
|
||||||
|
event_ts = s.event_ts,
|
||||||
|
_load_id = '{{ params.stg_batch_id }}'::text,
|
||||||
|
_load_ts = now()
|
||||||
|
FROM src AS s
|
||||||
|
WHERE o.book_ref = s.book_ref
|
||||||
|
AND (
|
||||||
|
o.book_date <> s.book_date OR
|
||||||
|
o.total_amount <> s.total_amount
|
||||||
|
);
|
||||||
|
|
||||||
|
INSERT INTO ods.bookings (
|
||||||
|
book_ref, book_date, total_amount, event_ts,
|
||||||
|
_load_id, _load_ts
|
||||||
|
)
|
||||||
|
SELECT
|
||||||
|
s.book_ref, s.book_date, s.total_amount, s.event_ts,
|
||||||
|
'{{ params.stg_batch_id }}'::text, now()
|
||||||
|
FROM src s
|
||||||
|
WHERE NOT EXISTS (
|
||||||
|
SELECT 1
|
||||||
|
FROM ods.bookings o
|
||||||
|
WHERE o.book_ref = s.book_ref
|
||||||
|
);
|
||||||
|
|
||||||
|
ANALYZE ods.bookings;
|
||||||
|
```
|
||||||
|
|
||||||
|
### 6.3. Поведение при пустом батче
|
||||||
|
|
||||||
|
- для инкрементальных таблиц (`bookings`, `tickets`, `flights`, `segments`) пустой батч допустим;
|
||||||
|
- для snapshot-справочников (`airports`, `airplanes`, `routes`, `seats`) пустой батч считаем ошибкой.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## 7) DQ-проверки ODS (минимум, но строго)
|
||||||
|
|
||||||
|
Каждый DQ-скрипт должен:
|
||||||
|
- быть привязан к `stg_batch_id`;
|
||||||
|
- делать `RAISE EXCEPTION` при нарушении;
|
||||||
|
- давать понятную подсказку в тексте ошибки.
|
||||||
|
|
||||||
|
### 7.1. Обязательные проверки
|
||||||
|
|
||||||
|
1. Нет дублей по бизнес-ключу в ODS.
|
||||||
|
|
||||||
|
2. Все ключи из STG текущего батча присутствуют в ODS.
|
||||||
|
|
||||||
|
3. Обязательные поля не `NULL`/не пустые.
|
||||||
|
|
||||||
|
4. Ссылочная целостность в ODS:
|
||||||
|
- `tickets.book_ref -> bookings.book_ref`
|
||||||
|
- `flights.route_no -> routes.route_no`
|
||||||
|
- `segments.ticket_no -> tickets.ticket_no`
|
||||||
|
- `segments.flight_id -> flights.flight_id`
|
||||||
|
- `boarding_passes (ticket_no, flight_id) -> segments (ticket_no, flight_id)`
|
||||||
|
|
||||||
|
### 7.2. Пример проверки покрытия батча
|
||||||
|
|
||||||
|
```sql
|
||||||
|
SELECT COUNT(*)
|
||||||
|
FROM (
|
||||||
|
SELECT DISTINCT book_ref
|
||||||
|
FROM stg.bookings
|
||||||
|
WHERE batch_id = '{{ params.stg_batch_id }}'::text
|
||||||
|
) s
|
||||||
|
WHERE NOT EXISTS (
|
||||||
|
SELECT 1
|
||||||
|
FROM ods.bookings o
|
||||||
|
WHERE o.book_ref = s.book_ref
|
||||||
|
);
|
||||||
|
```
|
||||||
|
|
||||||
|
Ожидаемый результат: `0`.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## 8) Структура файлов
|
||||||
|
|
||||||
|
```text
|
||||||
sql/ods/
|
sql/ods/
|
||||||
├── airports_ddl.sql ← SCD Type 2 схема + DISTRIBUTED BY
|
├── airports_ddl.sql
|
||||||
├── airports_load.sql ← UPDATE (закрыть) + INSERT (новые/изменённые)
|
├── airports_load.sql
|
||||||
├── airports_dq.sql ← нет дублей is_active, все STG-ключи в ODS
|
├── airports_dq.sql
|
||||||
├── airplanes_ddl.sql
|
├── airplanes_ddl.sql
|
||||||
├── airplanes_load.sql
|
├── airplanes_load.sql
|
||||||
├── airplanes_dq.sql
|
├── airplanes_dq.sql
|
||||||
├── routes_ddl.sql ← обратить внимание: scheduled_time TIME, duration INTERVAL
|
├── routes_ddl.sql
|
||||||
├── routes_load.sql
|
├── routes_load.sql
|
||||||
├── routes_dq.sql
|
├── routes_dq.sql
|
||||||
├── seats_ddl.sql
|
├── seats_ddl.sql
|
||||||
├── seats_load.sql
|
├── seats_load.sql
|
||||||
├── seats_dq.sql
|
├── seats_dq.sql
|
||||||
├── bookings_ddl.sql ← TIMESTAMPTZ, NUMERIC
|
├── bookings_ddl.sql
|
||||||
├── bookings_load.sql ← UPDATE изменений + INSERT новых
|
├── bookings_load.sql
|
||||||
├── bookings_dq.sql
|
├── bookings_dq.sql
|
||||||
├── tickets_ddl.sql ← BOOLEAN для outbound
|
├── tickets_ddl.sql
|
||||||
├── tickets_load.sql
|
├── tickets_load.sql
|
||||||
├── tickets_dq.sql
|
├── tickets_dq.sql
|
||||||
├── flights_ddl.sql ← INTEGER для flight_id
|
├── flights_ddl.sql
|
||||||
├── flights_load.sql
|
├── flights_load.sql
|
||||||
├── flights_dq.sql
|
├── flights_dq.sql
|
||||||
├── segments_ddl.sql
|
├── segments_ddl.sql
|
||||||
@@ -192,244 +406,121 @@ sql/ods/
|
|||||||
├── boarding_passes_load.sql
|
├── boarding_passes_load.sql
|
||||||
└── boarding_passes_dq.sql
|
└── boarding_passes_dq.sql
|
||||||
|
|
||||||
|
sql/ddl_gp_ods.sql
|
||||||
|
|
||||||
airflow/dags/
|
airflow/dags/
|
||||||
├── bookings_ods_ddl.py ← аналог bookings_stg_ddl.py (9 PostgresOperator)
|
├── bookings_ods_ddl.py
|
||||||
└── bookings_to_gp_ods.py ← аналог bookings_to_gp_stage.py (параллельный граф)
|
└── bookings_to_gp_ods.py
|
||||||
|
|
||||||
sql/ddl_gp_ods.sql ← мастер-DDL для make ddl-gp-ods (\i на каждый *_ddl.sql)
|
docs/bookings_to_gp_ods.md
|
||||||
|
Makefile (+ ddl-gp-ods)
|
||||||
docs/bookings_to_gp_ods.md ← описание DAG + примеры DQ-запросов для проверки
|
tests/test_dags_smoke.py (+ smoke для 2 новых DAG)
|
||||||
Makefile ← + таргет ddl-gp-ods
|
|
||||||
tests/test_dags_smoke.py ← + smoke-тесты для двух новых DAG
|
|
||||||
```
|
```
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
## DAG-граф bookings_to_gp_ods (параллельный)
|
## 9) DAG `bookings_to_gp_ods`: учебный граф зависимостей
|
||||||
|
|
||||||
```
|
Принцип: у каждой сущности строго `load -> dq`, и только после `dq` разрешаем downstream.
|
||||||
load_ods_airports → dq_ods_airports ─┐
|
|
||||||
├─ load_ods_routes → dq_ods_routes → load_ods_flights → dq_ods_flights ─┐
|
|
||||||
load_ods_airplanes → dq_ods_airplanes ─┘ │
|
|
||||||
└─ load_ods_seats → dq_ods_seats │
|
|
||||||
↓
|
|
||||||
load_ods_bookings → dq_ods_bookings → load_ods_tickets → dq_ods_tickets ──────────────────── load_ods_segments → dq_ods_segments
|
|
||||||
↓
|
|
||||||
load_ods_boarding_passes → dq_ods_boarding_passes
|
|
||||||
|
|
||||||
Все ветки → finish_ods_summary
|
```text
|
||||||
|
load_ods_bookings -> dq_ods_bookings -> load_ods_tickets -> dq_ods_tickets
|
||||||
|
|
||||||
|
├-> load_ods_airports -> dq_ods_airports ─┐
|
||||||
|
├-> load_ods_airplanes -> dq_ods_airplanes ─┼-> load_ods_routes -> dq_ods_routes -> load_ods_flights -> dq_ods_flights
|
||||||
|
└-> └-> load_ods_seats -> dq_ods_seats
|
||||||
|
|
||||||
|
dq_ods_flights + dq_ods_tickets -> load_ods_segments -> dq_ods_segments -> load_ods_boarding_passes -> dq_ods_boarding_passes
|
||||||
|
|
||||||
|
[dq_ods_boarding_passes, dq_ods_seats] -> finish_ods_summary
|
||||||
```
|
```
|
||||||
|
|
||||||
Зависимости:
|
Зависимости:
|
||||||
- `routes` — после `airports` и `airplanes`
|
- `tickets` после `bookings`;
|
||||||
- `seats` — после `airplanes`
|
- `routes` после `airports` и `airplanes`;
|
||||||
- `flights` — после `routes`
|
- `seats` после `airplanes`;
|
||||||
- `tickets` — после `bookings`
|
- `flights` после `routes`;
|
||||||
- `segments` — после `flights` и `tickets`
|
- `segments` после `flights` и `tickets`;
|
||||||
- `boarding_passes` — после `segments`
|
- `boarding_passes` после `segments`.
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
## Ключевые SQL-паттерны для обучения
|
## 10) Порядок реализации
|
||||||
|
|
||||||
### SCD Type 2 (airports_load.sql)
|
1. Подготовить DDL в `sql/ods/*_ddl.sql`.
|
||||||
|
2. Сделать мастер-скрипт `sql/ddl_gp_ods.sql`.
|
||||||
```sql
|
3. Добавить `Makefile`-таргет `ddl-gp-ods`.
|
||||||
-- Шаг 1: Получаем актуальный снапшот из STG (последний batch по load_dttm)
|
4. Создать DAG `bookings_ods_ddl.py`.
|
||||||
WITH latest_stg AS (
|
5. Реализовать `sql/ods/*_load.sql` (SCD1 UPSERT).
|
||||||
SELECT airport_code, airport_name, city, country, coordinates, timezone
|
6. Реализовать `sql/ods/*_dq.sql`.
|
||||||
FROM stg.airports
|
7. Создать DAG `bookings_to_gp_ods.py` (с параметром `stg_batch_id`).
|
||||||
WHERE load_dttm = (SELECT MAX(load_dttm) FROM stg.airports)
|
8. Дописать smoke-тесты DAG в `tests/test_dags_smoke.py`.
|
||||||
)
|
9. Описать запуск и проверки в `docs/bookings_to_gp_ods.md`.
|
||||||
-- Шаг 2: Закрываем изменившиеся версии
|
|
||||||
UPDATE ods.airports AS o
|
|
||||||
SET dw_end_date = CURRENT_DATE - 1,
|
|
||||||
is_active = false
|
|
||||||
FROM latest_stg AS s
|
|
||||||
WHERE o.airport_code = s.airport_code
|
|
||||||
AND o.is_active = true
|
|
||||||
AND (
|
|
||||||
o.airport_name <> s.airport_name OR
|
|
||||||
o.city <> s.city OR
|
|
||||||
o.country <> s.country OR
|
|
||||||
COALESCE(o.coordinates, '') <> COALESCE(s.coordinates, '') OR
|
|
||||||
o.timezone <> s.timezone
|
|
||||||
);
|
|
||||||
|
|
||||||
-- Шаг 3: Вставляем новые записи и новые версии изменённых
|
|
||||||
INSERT INTO ods.airports (
|
|
||||||
airport_code, airport_name, city, country, coordinates, timezone,
|
|
||||||
dw_start_date, dw_end_date, is_active, dw_version,
|
|
||||||
load_dttm, batch_id
|
|
||||||
)
|
|
||||||
SELECT
|
|
||||||
s.airport_code,
|
|
||||||
s.airport_name,
|
|
||||||
s.city,
|
|
||||||
s.country,
|
|
||||||
s.coordinates,
|
|
||||||
s.timezone,
|
|
||||||
CURRENT_DATE,
|
|
||||||
NULL,
|
|
||||||
true,
|
|
||||||
COALESCE(
|
|
||||||
(SELECT MAX(dw_version) FROM ods.airports WHERE airport_code = s.airport_code),
|
|
||||||
0
|
|
||||||
) + 1,
|
|
||||||
now(),
|
|
||||||
'{{ run_id }}'::text
|
|
||||||
FROM latest_stg s
|
|
||||||
WHERE NOT EXISTS (
|
|
||||||
SELECT 1 FROM ods.airports o
|
|
||||||
WHERE o.airport_code = s.airport_code AND o.is_active = true
|
|
||||||
);
|
|
||||||
|
|
||||||
ANALYZE ods.airports;
|
|
||||||
```
|
|
||||||
|
|
||||||
### UPSERT (bookings_load.sql)
|
|
||||||
|
|
||||||
```sql
|
|
||||||
-- Шаг 1: Обновляем изменившиеся записи (SCD Type 1 для транзакций)
|
|
||||||
UPDATE ods.bookings AS o
|
|
||||||
SET book_date = s.book_date::TIMESTAMP WITH TIME ZONE,
|
|
||||||
total_amount = s.total_amount::NUMERIC(10,2),
|
|
||||||
load_dttm = now(),
|
|
||||||
batch_id = '{{ run_id }}'::text
|
|
||||||
FROM (
|
|
||||||
-- Берём последнюю версию каждой записи из STG
|
|
||||||
SELECT DISTINCT ON (book_ref)
|
|
||||||
book_ref, book_date, total_amount, src_created_at_ts
|
|
||||||
FROM stg.bookings
|
|
||||||
ORDER BY book_ref, load_dttm DESC
|
|
||||||
) AS s
|
|
||||||
WHERE o.book_ref = s.book_ref
|
|
||||||
AND (
|
|
||||||
o.book_date <> s.book_date::TIMESTAMP WITH TIME ZONE OR
|
|
||||||
o.total_amount <> s.total_amount::NUMERIC(10,2)
|
|
||||||
);
|
|
||||||
|
|
||||||
-- Шаг 2: Вставляем новые записи
|
|
||||||
INSERT INTO ods.bookings (book_ref, book_date, total_amount, src_created_at_ts, load_dttm, batch_id)
|
|
||||||
SELECT
|
|
||||||
s.book_ref,
|
|
||||||
s.book_date::TIMESTAMP WITH TIME ZONE,
|
|
||||||
s.total_amount::NUMERIC(10,2),
|
|
||||||
s.src_created_at_ts,
|
|
||||||
now(),
|
|
||||||
'{{ run_id }}'::text
|
|
||||||
FROM (
|
|
||||||
SELECT DISTINCT ON (book_ref)
|
|
||||||
book_ref, book_date, total_amount, src_created_at_ts
|
|
||||||
FROM stg.bookings
|
|
||||||
ORDER BY book_ref, load_dttm DESC
|
|
||||||
) AS s
|
|
||||||
WHERE NOT EXISTS (
|
|
||||||
SELECT 1 FROM ods.bookings o WHERE o.book_ref = s.book_ref
|
|
||||||
);
|
|
||||||
|
|
||||||
ANALYZE ods.bookings;
|
|
||||||
```
|
|
||||||
|
|
||||||
### Нетривиальные кастинги (показываем студентам)
|
|
||||||
|
|
||||||
```sql
|
|
||||||
-- duration: '02:35:00' → INTERVAL
|
|
||||||
s.duration::INTERVAL
|
|
||||||
|
|
||||||
-- scheduled_time: 'HH:MM:SS' → TIME
|
|
||||||
s.scheduled_time::TIME
|
|
||||||
|
|
||||||
-- outbound: 'true'/'false' → BOOLEAN
|
|
||||||
s.outbound::BOOLEAN
|
|
||||||
|
|
||||||
-- flight_id: '12345' → INTEGER
|
|
||||||
s.flight_id::INTEGER
|
|
||||||
|
|
||||||
-- boarding_time: '2017-08-13 09:45+03' → TIMESTAMPTZ
|
|
||||||
s.boarding_time::TIMESTAMP WITH TIME ZONE
|
|
||||||
```
|
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
## DQ-проверки ODS
|
## 11) Критерии готовности (Definition of Done)
|
||||||
|
|
||||||
### Справочники (SCD Type 2)
|
Готово, если:
|
||||||
|
|
||||||
|
1. Оба новых DAG парсятся и проходят smoke-тесты (`make test`).
|
||||||
|
2. `make ddl-gp-ods` создаёт объекты без ошибок.
|
||||||
|
3. Для тестового `stg_batch_id` ODS-загрузка завершается успешно.
|
||||||
|
4. Все DQ-задачи зелёные и реально валят DAG при искусственной ошибке.
|
||||||
|
5. В ODS нет дублей по бизнес-ключам.
|
||||||
|
6. Нейминг техполей консистентен с учебной статьёй: `_load_id`, `_load_ts`, `valid_from/valid_to` (последние — когда перейдём к SCD2 в DDS).
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## 12) Как проверять вручную
|
||||||
|
|
||||||
|
```bash
|
||||||
|
make up
|
||||||
|
make ddl-gp
|
||||||
|
# Trigger bookings_to_gp_stage
|
||||||
|
# Получить batch_id из stg.bookings (последний)
|
||||||
|
# Trigger bookings_to_gp_ods с conf: {"stg_batch_id": "<значение>"}
|
||||||
|
make gp-psql
|
||||||
|
```
|
||||||
|
|
||||||
|
Проверочные SQL:
|
||||||
|
|
||||||
```sql
|
```sql
|
||||||
-- 1. Нет дублей активных записей по бизнес-ключу
|
-- 1) Дубликаты в ODS (пример bookings)
|
||||||
SELECT airport_code, COUNT(*)
|
SELECT book_ref, COUNT(*)
|
||||||
FROM ods.airports
|
FROM ods.bookings
|
||||||
WHERE is_active
|
|
||||||
GROUP BY 1
|
GROUP BY 1
|
||||||
HAVING COUNT(*) > 1;
|
HAVING COUNT(*) > 1;
|
||||||
|
|
||||||
-- 2. Все STG-ключи присутствуют в ODS (нет потерянных)
|
-- 2) Покрытие текущего STG-батча в ODS
|
||||||
SELECT COUNT(*)
|
SELECT COUNT(*)
|
||||||
FROM (SELECT DISTINCT airport_code FROM stg.airports) s
|
FROM (
|
||||||
|
SELECT DISTINCT book_ref
|
||||||
|
FROM stg.bookings
|
||||||
|
WHERE batch_id = '<stg_batch_id>'
|
||||||
|
) s
|
||||||
WHERE NOT EXISTS (
|
WHERE NOT EXISTS (
|
||||||
SELECT 1 FROM ods.airports o
|
SELECT 1
|
||||||
WHERE o.airport_code = s.airport_code AND o.is_active = true
|
FROM ods.bookings o
|
||||||
|
WHERE o.book_ref = s.book_ref
|
||||||
);
|
);
|
||||||
|
|
||||||
-- 3. Даты корректны: dw_end_date IS NULL для активных
|
-- 3) Ссылочная целостность tickets -> bookings
|
||||||
SELECT COUNT(*) FROM ods.airports WHERE is_active AND dw_end_date IS NOT NULL;
|
|
||||||
```
|
|
||||||
|
|
||||||
### Транзакционные (UPSERT)
|
|
||||||
|
|
||||||
```sql
|
|
||||||
-- 1. Нет дублей по бизнес-ключу
|
|
||||||
SELECT book_ref, COUNT(*) FROM ods.bookings GROUP BY 1 HAVING COUNT(*) > 1;
|
|
||||||
|
|
||||||
-- 2. Ссылочная целостность: все tickets ссылаются на существующие bookings
|
|
||||||
SELECT COUNT(*) FROM ods.tickets t
|
|
||||||
WHERE NOT EXISTS (SELECT 1 FROM ods.bookings b WHERE b.book_ref = t.book_ref);
|
|
||||||
|
|
||||||
-- 3. Все STG-записи попали в ODS
|
|
||||||
SELECT COUNT(*)
|
SELECT COUNT(*)
|
||||||
FROM (SELECT DISTINCT book_ref FROM stg.bookings) s
|
FROM ods.tickets t
|
||||||
WHERE NOT EXISTS (SELECT 1 FROM ods.bookings o WHERE o.book_ref = s.book_ref);
|
WHERE NOT EXISTS (
|
||||||
|
SELECT 1
|
||||||
|
FROM ods.bookings b
|
||||||
|
WHERE b.book_ref = t.book_ref
|
||||||
|
);
|
||||||
```
|
```
|
||||||
|
|
||||||
|
Ожидаемо: все три запроса возвращают `0` проблемных строк.
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
## Порядок реализации
|
## 13) Что будет следующим шагом
|
||||||
|
|
||||||
1. **`sql/ods/*_ddl.sql`** (9 файлов) — DDL всех таблиц
|
После стабилизации ODS:
|
||||||
2. **`sql/ddl_gp_ods.sql`** — мастер-DDL (собирает все `_ddl.sql` через `\i`)
|
- строим DDS;
|
||||||
3. **`Makefile`** — таргет `ddl-gp-ods`
|
- показываем SCD2 на измерениях DDS (`valid_from`/`valid_to`, `created_at`/`updated_at`) по тому же неймингу, который уже знаком студентам из `dwh-modeling`.
|
||||||
4. **`airflow/dags/bookings_ods_ddl.py`** — DDL DAG
|
|
||||||
5. **`sql/ods/*_load.sql`** (9 файлов) — скрипты загрузки (сначала справочники, потом транзакционные)
|
|
||||||
6. **`sql/ods/*_dq.sql`** (9 файлов) — DQ-проверки
|
|
||||||
7. **`airflow/dags/bookings_to_gp_ods.py`** — Load DAG с параллельным графом
|
|
||||||
8. **`tests/test_dags_smoke.py`** — smoke-тесты для новых DAG
|
|
||||||
9. **`docs/bookings_to_gp_ods.md`** — документация для студентов
|
|
||||||
|
|
||||||
---
|
|
||||||
|
|
||||||
## Проверка результата
|
|
||||||
|
|
||||||
```bash
|
|
||||||
make test # smoke-тесты: оба новых DAG парсятся без ошибок
|
|
||||||
make up # поднять стек
|
|
||||||
make ddl-gp-ods # применить ODS DDL
|
|
||||||
# Trigger bookings_to_gp_stage → дождаться завершения
|
|
||||||
# Trigger bookings_to_gp_ods → проверить параллельный граф в UI
|
|
||||||
make gp-psql # проверочные запросы к ods.*
|
|
||||||
```
|
|
||||||
|
|
||||||
Проверочные запросы:
|
|
||||||
|
|
||||||
```sql
|
|
||||||
-- Активные аэропорты = уникальным из STG
|
|
||||||
SELECT COUNT(*) FROM ods.airports WHERE is_active;
|
|
||||||
SELECT COUNT(DISTINCT airport_code) FROM stg.airports;
|
|
||||||
|
|
||||||
-- Бронирования без дублей
|
|
||||||
SELECT COUNT(*) FROM ods.bookings;
|
|
||||||
SELECT COUNT(DISTINCT book_ref) FROM stg.bookings;
|
|
||||||
|
|
||||||
-- При первом запуске все SCD-версии = 1
|
|
||||||
SELECT DISTINCT dw_version FROM ods.airports ORDER BY 1;
|
|
||||||
```
|
|
||||||
|
|||||||
@@ -22,6 +22,7 @@
|
|||||||
- **Даты**: как минимум различаем `book_date` (дата покупки) и `scheduled_departure` (дата/время вылета)
|
- **Даты**: как минимум различаем `book_date` (дата покупки) и `scheduled_departure` (дата/время вылета)
|
||||||
- **Инкремент в STG**: для `tickets` опорная дата берётся из `bookings.book_date`, потому что в `tickets` нет собственного поля времени изменения
|
- **Инкремент в STG**: для `tickets` опорная дата берётся из `bookings.book_date`, потому что в `tickets` нет собственного поля времени изменения
|
||||||
- **DQ-проверки**: проверки качества данных выполняем SQL-скриптами, но **не сохраняем результаты в отдельные таблицы/слой DQ** (при проблемах падаем с понятной ошибкой и останавливаем пайплайн)
|
- **DQ-проверки**: проверки качества данных выполняем SQL-скриптами, но **не сохраняем результаты в отдельные таблицы/слой DQ** (при проблемах падаем с понятной ошибкой и останавливаем пайплайн)
|
||||||
|
- **Нейминг полей**: единый стандарт — в [`docs/internal/naming_conventions.md`](naming_conventions.md)
|
||||||
|
|
||||||
### Статус реализации по слоям
|
### Статус реализации по слоям
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,80 @@
|
|||||||
|
# Конвенции Нейминга DWH (Единый Источник)
|
||||||
|
|
||||||
|
> Статус: активный стандарт для новых реализаций в этом репозитории.
|
||||||
|
>
|
||||||
|
> Основа: учебные материалы `de-roadmap/dwh-modeling`.
|
||||||
|
|
||||||
|
## Зачем этот документ
|
||||||
|
|
||||||
|
Чтобы имена полей не «плыли» между слоями, DAG и SQL-скриптами:
|
||||||
|
- студенты видят один и тот же словарь во всех задачах;
|
||||||
|
- новые реализации (ODS/DDS/DM) не расходятся с тем, как уже учили на `dwh-modeling`;
|
||||||
|
- ревью становится проще: сразу видно, где отклонение от стандарта.
|
||||||
|
|
||||||
|
## 1. Базовые правила
|
||||||
|
|
||||||
|
- Имена колонок: `snake_case`, на английском.
|
||||||
|
- Бизнес-ключи источника не переименовываем без необходимости (`book_ref`, `ticket_no`, `route_no`).
|
||||||
|
- Булевы поля начинаются с `is_` (`is_outbound`, `is_boarded`).
|
||||||
|
- Денежные/количественные поля называем явно (`total_amount`, `segment_amount`, `range_km`).
|
||||||
|
|
||||||
|
## 2. Каноничные служебные поля
|
||||||
|
|
||||||
|
| Поле | Тип (рекомендация) | Смысл | Где применять |
|
||||||
|
|---|---|---|---|
|
||||||
|
| `_load_id` | `TEXT NOT NULL` | Идентификатор загрузки/батча | STG/ODS (новые объекты), при необходимости DDS |
|
||||||
|
| `_load_ts` | `TIMESTAMP NOT NULL` | Когда запись попала в слой | STG/ODS (новые объекты) |
|
||||||
|
| `event_ts` | `TIMESTAMP` | Когда событие произошло в источнике (effective time) | ODS/DDS при событийной природе данных |
|
||||||
|
| `created_at` | `TIMESTAMP NOT NULL` | Когда строка создана в таблице слоя | DDS/DM, где есть lifecycle строки |
|
||||||
|
| `updated_at` | `TIMESTAMP NOT NULL` | Когда строка обновлена в таблице слоя | DDS/DM, где есть UPDATE |
|
||||||
|
| `valid_from` | `DATE NOT NULL` (базовый трек) | Начало действия версии SCD2 | DDS SCD2 |
|
||||||
|
| `valid_to` | `DATE` | Конец действия версии SCD2 (`NULL` = current) | DDS SCD2 |
|
||||||
|
| `hashdiff` | `TEXT NOT NULL` | Хэш атрибутов версии для детекта изменений | DDS SCD2 |
|
||||||
|
|
||||||
|
## 3. Ключи в DDS
|
||||||
|
|
||||||
|
- Бизнес-ключ измерения: суффикс `_bk` (`customer_bk`, `airport_bk`).
|
||||||
|
- Суррогатный ключ измерения: суффикс `_sk` (`customer_sk`, `airport_sk`).
|
||||||
|
|
||||||
|
## 4. Правило времени (важно для обучения)
|
||||||
|
|
||||||
|
- `event_ts` (effective time) и `_load_ts` (load time) — разные сущности, не смешиваем.
|
||||||
|
- Если `event_ts` отсутствует в источнике, используем `_load_ts` как fallback и явно документируем это в SQL/доке.
|
||||||
|
|
||||||
|
## 5. Применение по слоям
|
||||||
|
|
||||||
|
### STG
|
||||||
|
|
||||||
|
- Для уже реализованного `bookings` STG сохраняем текущие legacy-имена ради обратной совместимости:
|
||||||
|
- `src_created_at_ts`
|
||||||
|
- `load_dttm`
|
||||||
|
- `batch_id`
|
||||||
|
- Для новых STG-объектов (новые домены/задачи) используем канон `_load_id`, `_load_ts` (и `event_ts`, если нужно).
|
||||||
|
|
||||||
|
### ODS
|
||||||
|
|
||||||
|
- В новых реализациях используем канон:
|
||||||
|
- `_load_id`, `_load_ts`, `event_ts`.
|
||||||
|
- Базовый эталон ODS в этом стенде: SCD Type 1 (current state + UPSERT).
|
||||||
|
|
||||||
|
### DDS
|
||||||
|
|
||||||
|
- Для SCD2 используем:
|
||||||
|
- `valid_from`, `valid_to`, `hashdiff`, `created_at`, `updated_at`.
|
||||||
|
- Интервалы считаем как `[valid_from, valid_to)`, current-версия: `valid_to IS NULL`.
|
||||||
|
|
||||||
|
## 6. Переходный маппинг legacy -> канон
|
||||||
|
|
||||||
|
| Legacy (текущий bookings STG) | Канон |
|
||||||
|
|---|---|
|
||||||
|
| `batch_id` | `_load_id` |
|
||||||
|
| `load_dttm` | `_load_ts` |
|
||||||
|
| `src_created_at_ts` | `event_ts` |
|
||||||
|
|
||||||
|
Примечание: это логический маппинг для новых слоёв. Массовое переименование существующего STG не требуется.
|
||||||
|
|
||||||
|
## 7. Что проверяем в ревью
|
||||||
|
|
||||||
|
- Нет новых техполей-синнонимов вроде `loaded_at`, `ingested_at`, `batch_key`, если уже есть канон.
|
||||||
|
- Нет смешивания `event_ts` и `_load_ts` в одном смысле.
|
||||||
|
- В SCD2 не используются альтернативы `dw_start_date/dw_end_date`, если в проекте принят `valid_from/valid_to`.
|
||||||
Reference in New Issue
Block a user