diff --git a/docs/assignment/analyst_spec.md b/docs/assignment/analyst_spec.md index 998c8ed..42769db 100644 --- a/docs/assignment/analyst_spec.md +++ b/docs/assignment/analyst_spec.md @@ -5,8 +5,9 @@ > опираясь на эталонный срез (витрина `dm.sales_report` и вся её цепочка). > > **Эталон для изучения:** -> - STG: `bookings`, `tickets`, `flights`, `segments`, `boarding_passes`, `airports` -> - ODS: те же таблицы +> - STG: весь слой — `bookings`, `tickets`, `flights`, `segments`, `boarding_passes`, +> `airports`, `airplanes`, `seats`, `routes` (реализован как образец) +> - ODS: `bookings`, `tickets`, `segments`, `flights`, `boarding_passes`, `airports`, `routes` > - DDS: `dim_airports`, `dim_tariffs`, `dim_calendar`, `fact_flight_sales` > - DM: `sales_report` > @@ -16,18 +17,20 @@ ## Рекомендуемый порядок выполнения -1. **STG** — `airplanes`, `seats`, `routes` (разминка по аналогии) -2. **ODS** — `airplanes`, `seats`, `routes` (закрепление UPSERT / TRUNCATE+INSERT) -3. **DDS** — `dim_airplanes`, `dim_passengers` (SCD1 — новые измерения) -4. **DDS** — `dim_routes` (SCD2 — ключевой вызов курсовой) -5. **DM** — `airport_traffic` (простая витрина, похожа на `sales_report`) -6. **DM** — `route_performance` (Full Rebuild, работа с SCD2-измерением) -7. **DM** — `monthly_overview` (двухуровневая агрегация) -8. **DM** — `passenger_loyalty` (самая сложная, пересчёт истории) +> **Примечание:** STG-слой полностью реализован в эталоне (все 9 таблиц), `ods.routes` +> тоже эталонный. Ваше задание начинается с ODS: `airplanes` и `seats`. + +1. **ODS** — `airplanes`, `seats` (практика TRUNCATE+INSERT по аналогии с `ods.airports`) +2. **DDS** — `dim_airplanes`, `dim_passengers` (SCD1 — новые измерения) +3. **DDS** — `dim_routes` (SCD2 — ключевой вызов курсовой) +4. **DM** — `airport_traffic` (простая витрина, похожа на `sales_report`) +5. **DM** — `route_performance` (Full Rebuild, работа с SCD2-измерением) +6. **DM** — `monthly_overview` (двухуровневая агрегация) +7. **DM** — `passenger_loyalty` (самая сложная, пересчёт истории) Порядок выстроен от простого к сложному. Каждый шаг опирается на опыт предыдущего. Вы можете двигаться в ином порядке, но убедитесь, что зависимости -слоёв соблюдены (STG → ODS → DDS → DM). +слоёв соблюдены (ODS → DDS → DM). --- @@ -46,106 +49,40 @@ --- -## Часть 1. STG-слой (Staging) +## Почему STG-слой уже реализован -> **Цель:** Скопировать данные из источника (bookings-db) в Greenplum «как есть», -> сохраняя все поля в текстовом виде. Преобразование типов — задача ODS. -> -> **Аналог для изучения:** `sql/stg/airports_ddl.sql`, `sql/stg/airports_load.sql` +STG-слой (все 9 таблиц) и `ods.routes` полностью реализованы в эталоне. Причина: -### 1.1. stg.airplanes +Эталонный пайплайн использует **batch resolver** — механизм, который находит +согласованный набор данных во всех четырёх snapshot-таблицах (`stg.airports`, +`stg.airplanes`, `stg.routes`, `stg.seats`) по одному `_load_id`. Если любая +из этих таблиц пуста — batch resolver не найдёт общего батча, и весь ODS pipeline +не запустится. -**Описание:** Справочник моделей воздушных судов. Содержит технические -характеристики каждой модели: дальность полёта и крейсерскую скорость. +Это означает: без заполненного STG невозможна загрузка ODS и всей последующей +цепочки (DDS, DM). Поэтому весь STG реализован как эталон — чтобы пайплайн +работал с первого запуска. -**Источник:** `bookings.airplanes_data` (через PXF) +`ods.routes` тоже эталонный: факт `fact_flight_sales` использует его для резолвинга +аэропортов вылета/прилёта, независимо от студенческого `dds.dim_routes`. -| Поле | Тип | Описание | Маппинг из источника | -|------|-----|----------|----------------------| -| `airplane_code` | TEXT | Код модели самолёта (бизнес-ключ) | `airplanes_data.airplane_code::TEXT` | -| `model` | TEXT | Полное наименование модели (JSON в источнике) | `airplanes_data.model::TEXT` | -| `range` | TEXT | Дальность полёта, км | `airplanes_data.range::TEXT` | -| `speed` | TEXT | Крейсерская скорость, км/ч | `airplanes_data.speed::TEXT` | -| `event_ts` | TIMESTAMP | Момент загрузки (`now()`). У snapshot-справочников нет бизнес-события с точным временем, поэтому `event_ts` заполняется при загрузке | `now()` | -| `_load_ts` | TIMESTAMP NOT NULL | Момент загрузки в STG | `now()` | -| `_load_id` | TEXT NOT NULL | Идентификатор батча загрузки | `'{{ run_id }}'` | - -**Тип историзации:** Нет (накопительный snapshot). Каждый запуск добавляет -новый батч по `_load_id`; прошлые батчи остаются в таблице. Не используйте TRUNCATE. - -**Стратегия загрузки:** Полный снимок (Full Snapshot). Справочник маленький — -загружаем целиком каждый раз. Идемпотентность — через проверку `_load_id`. - -**Distribution Key:** `airplane_code` +**Что вам делать:** Изучите эталонные скрипты как образец — именно так написан +«боевой» код загрузки: +- `sql/stg/airports_load.sql` — инкрементальная загрузка (HWM) +- `sql/stg/airplanes_load.sql` — full snapshot с проверкой `_load_id` +- `sql/ods/airports_load.sql` — TRUNCATE+INSERT из STG --- -### 1.2. stg.seats - -**Описание:** Карта посадочных мест для каждой модели самолёта. -Каждая строка — одно конкретное место в конкретной модели. - -**Источник:** `bookings.seats` (через PXF) - -| Поле | Тип | Описание | Маппинг из источника | -|------|-----|----------|----------------------| -| `airplane_code` | TEXT | Код модели самолёта (FK → airplanes) | `seats.airplane_code::TEXT` | -| `seat_no` | TEXT | Номер места (напр. «1A», «12C») | `seats.seat_no::TEXT` | -| `fare_conditions` | TEXT | Класс обслуживания (`Economy`, `Business`, `Comfort`) | `seats.fare_conditions::TEXT` | -| `event_ts` | TIMESTAMP | Момент загрузки (`now()`). Snapshot-справочник — бизнес-событие отсутствует | `now()` | -| `_load_ts` | TIMESTAMP NOT NULL | Момент загрузки в STG | `now()` | -| `_load_id` | TEXT NOT NULL | Идентификатор батча загрузки | `'{{ run_id }}'` | - -**Тип историзации:** Нет (накопительный snapshot). Каждый запуск добавляет -новый батч по `_load_id`; прошлые батчи остаются в таблице. Не используйте TRUNCATE. - -**Стратегия загрузки:** Полный снимок. Идемпотентность — через -`airplane_code` + `seat_no` + `_load_id`. - -**Distribution Key:** `airplane_code` - ---- - -### 1.3. stg.routes - -**Описание:** Справочник авиамаршрутов. Маршрут — регулярный рейс между двумя -аэропортами на определённом типе самолёта с фиксированным расписанием. - -**Источник:** `bookings.routes` (через PXF) - -| Поле | Тип | Описание | Маппинг из источника | -|------|-----|----------|----------------------| -| `route_no` | TEXT | Номер маршрута (бизнес-ключ, часть 1) | `routes.route_no::TEXT` | -| `validity` | TEXT | Период действия маршрута (бизнес-ключ, часть 2) | `routes.validity::TEXT` | -| `departure_airport` | TEXT | Код аэропорта вылета (FK → airports) | `routes.departure_airport::TEXT` | -| `arrival_airport` | TEXT | Код аэропорта прилёта (FK → airports) | `routes.arrival_airport::TEXT` | -| `airplane_code` | TEXT | Код модели самолёта (FK → airplanes) | `routes.airplane_code::TEXT` | -| `days_of_week` | TEXT | Дни недели выполнения рейса | `routes.days_of_week::TEXT` | -| `scheduled_time` | TEXT | Время вылета по расписанию | `routes.scheduled_time::TEXT` | -| `duration` | TEXT | Плановая продолжительность полёта | `routes.duration::TEXT` | -| `event_ts` | TIMESTAMP | Момент загрузки (`now()`). Snapshot-справочник — бизнес-событие отсутствует | `now()` | -| `_load_ts` | TIMESTAMP NOT NULL | Момент загрузки в STG | `now()` | -| `_load_id` | TEXT NOT NULL | Идентификатор батча загрузки | `'{{ run_id }}'` | - -**Тип историзации:** Нет (накопительный snapshot). Каждый запуск добавляет -новый батч по `_load_id`; прошлые батчи остаются в таблице. Не используйте TRUNCATE. - -**Стратегия загрузки:** Полный снимок. Идемпотентность — по составному ключу -`route_no` + `validity` + `_load_id`. - -**Distribution Key:** `route_no` - ---- - -## Часть 2. ODS-слой (Operational Data Store) +## Часть 1. ODS-слой (Operational Data Store) > **Цель:** Привести данные из STG к целевым типам, очистить, дедуплицировать. -> Справочники (airplanes, seats, routes) загружаются стратегией TRUNCATE + INSERT -> из последнего согласованного батча STG. +> Справочники (`airplanes`, `seats`) загружаются стратегией TRUNCATE + INSERT +> из последнего согласованного батча STG. (`ods.routes` реализован в эталоне.) > > **Аналог для изучения:** `sql/ods/airports_ddl.sql`, `sql/ods/airports_load.sql` -### 2.1. ods.airplanes +### 1.1. ods.airplanes **Описание:** Очищенный справочник моделей воздушных судов с правильными типами. @@ -175,7 +112,7 @@ --- -### 2.2. ods.seats +### 1.2. ods.seats **Описание:** Карта посадочных мест с корректными типами. @@ -203,41 +140,9 @@ --- -### 2.3. ods.routes - -**Описание:** Справочник маршрутов с правильными типами данных. - -**Источник:** `stg.routes` - -| Поле | Тип | Описание | Маппинг из STG | -|------|-----|----------|----------------| -| `route_no` | TEXT NOT NULL | Номер маршрута (PK, часть 1) | `route_no` | -| `validity` | TEXT NOT NULL | Период действия (PK, часть 2) | `validity` | -| `departure_airport` | TEXT NOT NULL | Код аэропорта вылета | `departure_airport` | -| `arrival_airport` | TEXT NOT NULL | Код аэропорта прилёта | `arrival_airport` | -| `airplane_code` | TEXT NOT NULL | Код модели самолёта | `airplane_code` | -| `days_of_week` | INTEGER[] | Дни недели (массив) | `days_of_week` — преобразовать TEXT в `INTEGER[]` | -| `departure_time` | TIME NOT NULL | Время вылета | `scheduled_time::TIME` | -| `duration` | INTERVAL NOT NULL | Длительность полёта | `duration::INTERVAL` | -| `_load_id` | TEXT NOT NULL | Идентификатор батча | `'{{ run_id }}'` | -| `_load_ts` | TIMESTAMP NOT NULL | Момент загрузки | `now()` | - -**Тип историзации:** Нет (текущее состояние справочника, TRUNCATE + INSERT). - -**Стратегия загрузки:** TRUNCATE + INSERT. - -**Бизнес-правила:** -- Составной PK: `(route_no, validity)`. -- Преобразование `days_of_week` из текста в массив целых чисел (`INTEGER[]`). -- Преобразование `scheduled_time` → `TIME`, `duration` → `INTERVAL`. - -**Тип хранения Greenplum:** Append-Only Row. - -**Distribution Key:** `(route_no, validity)` - --- -## Часть 3. DDS-слой (Detailed Data Store) — Измерения +## Часть 2. DDS-слой (Detailed Data Store) — Измерения > **Цель:** Построить измерения звёздной схемы (Star Schema) с суррогатными > ключами. SCD1-измерения обновляют атрибуты «на месте». SCD2-измерение @@ -245,7 +150,7 @@ > > **Аналог для SCD1:** `sql/dds/dim_airports_ddl.sql`, `sql/dds/dim_airports_load.sql` -### 3.1. dds.dim_airplanes (SCD1) +### 2.1. dds.dim_airplanes (SCD1) **Описание:** Измерение моделей самолётов. Содержит технические характеристики и рассчитанное общее количество мест (обогащение из `ods.seats`). @@ -284,7 +189,7 @@ --- -### 3.2. dds.dim_passengers (SCD1) +### 2.2. dds.dim_passengers (SCD1) **Описание:** Измерение пассажиров. Извлекается из таблицы билетов — каждый уникальный `passenger_id` становится строкой измерения. @@ -322,7 +227,7 @@ --- -### 3.3. dds.dim_routes (SCD2) +### 2.3. dds.dim_routes (SCD2) > **Это ключевой вызов курсовой.** Реализация SCD Type 2 — обязательный навык > для Data Engineer. Ниже — алгоритм текстом; SQL вы пишете самостоятельно. @@ -406,14 +311,14 @@ md5(concat_ws('|', --- -## Часть 4. DM-слой (Data Marts) — Витрины +## Часть 3. DM-слой (Data Marts) — Витрины > **Цель:** Построить аналитические витрины поверх DDS. > Каждая витрина отвечает на конкретный бизнес-вопрос. > > **Аналог для изучения:** `sql/dm/sales_report_ddl.sql`, `sql/dm/sales_report_load.sql` -### 4.1. dm.airport_traffic +### 3.1. dm.airport_traffic **Бизнес-вопрос:** «Какой пассажиропоток и выручка у каждого аэропорта по дням?» @@ -469,7 +374,7 @@ md5(concat_ws('|', --- -### 4.2. dm.route_performance +### 3.2. dm.route_performance **Бизнес-вопрос:** «Какие маршруты самые эффективные? Где высокий load factor, а где теряем пассажиров?» @@ -532,7 +437,7 @@ md5(concat_ws('|', --- -### 4.3. dm.monthly_overview +### 3.3. dm.monthly_overview **Бизнес-вопрос:** «Какова помесячная динамика: рейсы, выручка, load factor в разрезе типов самолётов?» @@ -594,7 +499,7 @@ md5(concat_ws('|', --- -### 4.4. dm.passenger_loyalty +### 3.4. dm.passenger_loyalty **Бизнес-вопрос:** «Кто наши самые лояльные пассажиры? Сколько они летают, тратят, какой класс предпочитают?» @@ -660,6 +565,34 @@ md5(concat_ws('|', --- +## Пересчёт факта после реализации измерений + +После того, как вы реализуете все DDS-измерения (`dim_airplanes`, `dim_passengers`, +`dim_routes`), нужно пересчитать факт — он загружался ещё без ваших SK: + +1. Запустите загрузку своих измерений (dim_airplanes, dim_passengers, dim_routes) +2. Выполните `TRUNCATE dds.fact_flight_sales;` +3. Перезапустите загрузку факта (DAG или только таск `load_dds_fact_flight_sales`) +4. Перезапустите DM-витрины + +```sql +-- Проверка: все SK заполнены +SELECT + COUNT(*) AS total, + COUNT(departure_airport_sk) AS has_dep_sk, + COUNT(arrival_airport_sk) AS has_arr_sk, + COUNT(route_sk) AS has_route_sk, + COUNT(passenger_sk) AS has_passenger_sk, + COUNT(airplane_sk) AS has_airplane_sk +FROM dds.fact_flight_sales; +``` + +Это стандартная практика при **late-arriving dimensions** (опаздывающих измерениях): +факт загрузился раньше, чем были готовы справочники, поэтому SK заполнены не были. +После пересчёта все SK проставятся корректно. + +--- + ## Валидация После реализации каждого слоя запустите валидационный DAG `bookings_validate` diff --git a/docs/design/assignment_design.md b/docs/design/assignment_design.md index 5c410b5..13d18db 100644 --- a/docs/design/assignment_design.md +++ b/docs/design/assignment_design.md @@ -20,15 +20,14 @@ |------|-----------------------------------------------------------------| | DM | `sales_report` | | DDS | `fact_flight_sales`, `dim_airports` (SCD1), `dim_tariffs` (SCD1), `dim_calendar` | -| ODS | `bookings`, `tickets`, `segments`, `flights`, `boarding_passes`, `airports` | -| STG | `bookings`, `tickets`, `segments`, `flights`, `boarding_passes`, `airports` | +| ODS | `bookings`, `tickets`, `segments`, `flights`, `boarding_passes`, `airports`, `routes` | +| STG | весь слой: `bookings`, `tickets`, `segments`, `flights`, `boarding_passes`, `airports`, `airplanes`, `seats`, `routes` | ### Задание студенту | Слой | Таблицы | Что нового для студента | |------|-------------------------------------------------------------------|--------------------------------------------------| -| STG | `airplanes`, `seats`, `routes` | Практика по аналогии с эталоном | -| ODS | `airplanes`, `seats`, `routes` | Практика SCD1 UPSERT по аналогии | +| ODS | `airplanes`, `seats` | Практика TRUNCATE+INSERT по аналогии с эталоном | | DDS | `dim_airplanes` (SCD1), `dim_passengers` (SCD1), `dim_routes` (SCD2) | **SCD2 — ключевой вызов курсовой** | | DM | `airport_traffic`, `monthly_overview`, `route_performance`, `passenger_loyalty` | Разная сложность (от простой к сложной) | @@ -38,14 +37,13 @@ Студенту рекомендуется (но не обязательно) двигаться в таком порядке: -1. **STG** (airplanes, seats, routes) — разминка, по аналогии -2. **ODS** (airplanes, seats, routes) — закрепление UPSERT -3. **DDS** dim_airplanes, dim_passengers (SCD1) — новые измерения -4. **DDS** dim_routes (**SCD2**) — ключевой вызов -5. **DM** airport_traffic — простая витрина, похожа на sales_report -6. **DM** route_performance — TRUNCATE+INSERT, SCD2-агрегация по BK -7. **DM** monthly_overview — двухуровневая агрегация -8. **DM** passenger_loyalty — самая сложная, пересчёт истории +1. **ODS** (airplanes, seats) — практика TRUNCATE+INSERT +2. **DDS** dim_airplanes, dim_passengers (SCD1) — новые измерения +3. **DDS** dim_routes (**SCD2**) — ключевой вызов +4. **DM** airport_traffic — простая витрина, похожа на sales_report +5. **DM** route_performance — TRUNCATE+INSERT, SCD2-агрегация по BK +6. **DM** monthly_overview — двухуровневая агрегация +7. **DM** passenger_loyalty — самая сложная, пересчёт истории Порядок выстроен от простого к сложному. Каждый шаг опирается на опыт предыдущего. @@ -83,14 +81,9 @@ ``` bookings_validate -├── validate_stg -│ ├── check_stg_airplanes_exists (таблица создана, >0 строк) -│ ├── check_stg_seats_exists -│ └── check_stg_routes_exists ├── validate_ods │ ├── check_ods_airplanes_rowcount (ODS >= STG по кол-ву уникальных BK) │ ├── check_ods_seats_rowcount -│ ├── check_ods_routes_rowcount │ └── check_ods_no_null_pks (PK not null) ├── validate_dds │ ├── check_dim_airplanes_exists diff --git a/docs/design/db_schema.md b/docs/design/db_schema.md index abeac40..ffee180 100644 --- a/docs/design/db_schema.md +++ b/docs/design/db_schema.md @@ -1,6 +1,10 @@ # Схема БД DWH (Bookings → Greenplum) -> **Статус:** Все слои реализованы (STG, ODS, DDS, DM). +> **Статус:** Все слои реализованы в ветке `solution` (STG, ODS, DDS, DM). +> На ветке `main` студенческие таблицы ODS (`airplanes`, `seats`), DDS-измерения +> (`dim_airplanes`, `dim_passengers`, `dim_routes`) и студенческие DM-витрины +> (`airport_traffic`, `route_performance`, `monthly_overview`, `passenger_loyalty`) +> — заглушки (`SELECT 1;`). Данные появятся после реализации заданий. Архитектура хранилища данных (DWH) для учебного проекта Airflow + Greenplum. Источник — демо-БД `bookings` (Postgres). Документ даёт цельный взгляд «сверху»; @@ -250,7 +254,9 @@ Degenerate keys: `book_ref`, `ticket_no`, `flight_id`, `book_date`, `seat_no`. ## Ключевые договорённости - **Нейминг полей**: [`naming_conventions.md`](naming_conventions.md) -- **DQ-проверки**: SQL-скрипты с `RAISE EXCEPTION` (не отдельный DQ-слой) +- **DQ-проверки**: SQL-скрипты с `RAISE EXCEPTION` (не отдельный DQ-слой). + На ветке `main` студенческие DQ-скрипты (`airplanes_dq.sql`, `seats_dq.sql`, + student dims/DM) — заглушки (`SELECT 1;`) без проверок. - **Инкремент STG**: для `tickets` опорная дата — из `bookings.book_date` - **Point-in-time JOIN**: факт ↔ `dim_routes` по `[valid_from, valid_to)` - **Суррогатные ключи**: `MAX(sk) + ROW_NUMBER()` (не SERIAL — GP-специфика) diff --git a/sql/dds/fact_flight_sales_load.sql b/sql/dds/fact_flight_sales_load.sql index 1aa66ab..d61f357 100644 --- a/sql/dds/fact_flight_sales_load.sql +++ b/sql/dds/fact_flight_sales_load.sql @@ -48,6 +48,21 @@ WITH fact_src AS ( -- Учебный комментарий: Late-arriving dimensions (Опаздывающие измерения) -- Мы используем LEFT JOIN, так как факт (рейс/билет) может прийти раньше, -- чем справочник (пассажир/маршрут) обновится в DDS. + -- + -- Airport lookup через ods.routes (эталонный справочник), а не через dds.dim_routes + -- (студенческое задание SCD2). Это архитектурное решение: эталонный пайплайн работает + -- независимо от студенческого кода. Аэропорты вылета/прилёта одинаковы во всех + -- версиях одного route_no — безопасно брать из ODS без point-in-time логики. + -- ROW_NUMBER по validity DESC: выбираем актуальную версию расписания маршрута. + LEFT JOIN ( + SELECT route_no, departure_airport, arrival_airport + FROM ( + SELECT route_no, departure_airport, arrival_airport, + ROW_NUMBER() OVER (PARTITION BY route_no ORDER BY validity DESC) AS rn + FROM ods.routes + ) ranked + WHERE rn = 1 + ) AS ods_rte ON ods_rte.route_no = flt.route_no LEFT JOIN dds.dim_routes AS rte ON rte.route_bk = flt.route_no AND flt.scheduled_departure::DATE >= rte.valid_from @@ -55,11 +70,11 @@ WITH fact_src AS ( LEFT JOIN dds.dim_calendar AS cal ON cal.date_actual = flt.scheduled_departure::DATE LEFT JOIN dds.dim_airports AS dep - ON dep.airport_bk = rte.departure_airport + ON dep.airport_bk = ods_rte.departure_airport -- через ods.routes (эталон) LEFT JOIN dds.dim_airports AS arr - ON arr.airport_bk = rte.arrival_airport + ON arr.airport_bk = ods_rte.arrival_airport -- через ods.routes (эталон) LEFT JOIN dds.dim_airplanes AS ap - ON ap.airplane_bk = rte.airplane_code + ON ap.airplane_bk = rte.airplane_code -- через dim_routes (point-in-time, как прежде) LEFT JOIN dds.dim_tariffs AS tar ON tar.fare_conditions = seg.fare_conditions LEFT JOIN dds.dim_passengers AS pax