Compare commits
5
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
999a8d08bf | ||
|
|
d3f2f42e38 | ||
|
|
a8148c8afa | ||
|
|
f053fe1e35 | ||
|
|
ab97d6a535 |
@@ -1,5 +1,9 @@
|
||||
# airflow-dwh-gp-lab
|
||||
|
||||
> Этот стенд можно использовать самостоятельно, но он также является курсовой работой
|
||||
> плана обучения [Data Engineering Roadmap](https://github.com/dementev-dev/de-roadmap).
|
||||
> Нужна помощь? Обратитесь к автору материалов — ментору [@dementev_dev](https://t.me/dementev_dev).
|
||||
|
||||
Учебный стенд для лабораторных по Data Engineering: **Airflow** оркестрирует загрузку данных из демо‑БД
|
||||
**bookings** (Postgres) в **Greenplum**.
|
||||
|
||||
@@ -95,6 +99,12 @@ SELECT COUNT(*) FROM dm.route_performance;
|
||||
|
||||
Подробнее про логику DAG и проверки — `docs/bookings_to_gp_stage.md`.
|
||||
|
||||
7) Изучите готовую реализацию:
|
||||
|
||||
Весь пайплайн `STG → ODS → DDS → DM` на этой ветке полностью реализован.
|
||||
Переходите к [навигатору по реализации](docs/assignment/README.md) —
|
||||
там описана структура, ключевые паттерны и рекомендуемый порядок изучения.
|
||||
|
||||
## DAG-и в стенде
|
||||
|
||||
Основные (для потока bookings → DWH):
|
||||
|
||||
@@ -44,11 +44,11 @@
|
||||
|
||||
> Делаем пока полный пайплайн работает — можно сверяться с реальными данными.
|
||||
|
||||
- [ ] Создать `docs/assignment/analyst_spec.md`
|
||||
- [ ] Для каждой таблицы-задания: имя, описание, поля, маппинг,
|
||||
- [x] Создать `docs/assignment/analyst_spec.md`
|
||||
- [x] Для каждой таблицы-задания: имя, описание, поля, маппинг,
|
||||
бизнес-правила, тип SCD, гранулярность, distribution key
|
||||
- [ ] Для dim_routes (SCD2): пошаговый алгоритм текстом, формула hashdiff
|
||||
- [ ] Рекомендуемый порядок выполнения
|
||||
- [x] Для dim_routes (SCD2): пошаговый алгоритм текстом, формула hashdiff
|
||||
- [x] Рекомендуемый порядок выполнения
|
||||
|
||||
### Этап 3. Валидационный DAG
|
||||
|
||||
@@ -57,27 +57,26 @@
|
||||
|
||||
> Делаем и тестируем на полных данных, пока ничего не удалено.
|
||||
|
||||
- [ ] Создать `airflow/dags/bookings_validate.py`
|
||||
- [ ] Создать SQL-скрипты в `sql/validate/`
|
||||
- [ ] Таски по слоям: STG, ODS, DDS, DM
|
||||
- [ ] Дружелюбные сообщения об ошибках с подсказками
|
||||
- [ ] Протестировать на работающем стенде
|
||||
- [x] Создать `airflow/dags/bookings_validate.py`
|
||||
- [x] Создать SQL-скрипты в `sql/validate/`
|
||||
- [x] Таски по слоям: STG, ODS, DDS, DM
|
||||
- [x] Дружелюбные сообщения об ошибках с подсказками
|
||||
- [x] Протестировать на работающем стенде
|
||||
|
||||
### Этап 4. Подготовка main и ветка solution
|
||||
|
||||
**Инструмент:** Sonnet — удаление файлов и добавление заглушек
|
||||
по списку из [assignment_design.md](docs/design/assignment_design.md).
|
||||
**Инструмент:** Opus — раскладка по веткам, заглушки, ослабление DQ.
|
||||
|
||||
> Финальный этап: всё готово и протестировано, теперь раскладываем по веткам.
|
||||
> План: `docs/plans/2026-03-12_main-solution-split.md`
|
||||
|
||||
- [ ] Оставить только эталонный срез (sales_report + цепочка)
|
||||
- [ ] Убрать реализации таблиц-заданий (airplanes, seats, routes в STG/ODS;
|
||||
dim_airplanes, dim_passengers, dim_routes в DDS; 4 витрины DM)
|
||||
- [ ] Добавить TODO-маркеры / заглушки в DAG'ах для студенческих тасков
|
||||
- [ ] Обновить `ddl_gp.sql` (убрать `\i` для таблиц-заданий)
|
||||
- [ ] Создать ветку `solution` от main
|
||||
- [ ] Добавить полные реализации всех таблиц-заданий
|
||||
- [ ] Проверить, что всё работает end-to-end на обеих ветках
|
||||
- [x] Смержить chore/bookings-etl → main
|
||||
- [x] Общие правки на main (документация, TODO.md)
|
||||
- [x] Создать ветку solution (снимок полного эталона)
|
||||
- [x] Main-only правки: заглушки, ослабление DQ, адаптация тестов
|
||||
- [x] Очистка docs на main (удалить внутренние документы)
|
||||
- [x] Верификация main (`make test`, `make lint`)
|
||||
- [x] Верификация solution (`make test`)
|
||||
|
||||
---
|
||||
|
||||
|
||||
@@ -0,0 +1,34 @@
|
||||
# План: Онбординг студента + маркетинг (v4)
|
||||
|
||||
> Статус: выполнен (2026-03-13)
|
||||
|
||||
## Контекст
|
||||
|
||||
Ментор опционален. Студент, клонировав main, проходит «Быстрый старт» — и застревает:
|
||||
нет явного «что дальше», ссылка на задания спрятана, validate DAG не объяснён, нет маркетинга.
|
||||
На main есть противоречия: docs и DAG-docstrings говорят «все реализовано»,
|
||||
фактически стоят заглушки `SELECT 1;`.
|
||||
|
||||
## Что было сделано
|
||||
|
||||
### На main (10 файлов)
|
||||
|
||||
1. **README.md** — маркетинг-баннер, шаг 6 (только эталонные таблицы), шаг 7 → задания.
|
||||
2. **docs/assignment/README.md** — полный гид студента (эталон → ТЗ → заглушки → validate).
|
||||
3. **docs/assignment/analyst_spec.md** — DAG-интеграция: «таски уже подключены, DAG менять не нужно».
|
||||
4. **docs/bookings_to_gp_dm.md** — пометки заглушек, адаптация секций проверки.
|
||||
5. **docs/bookings_to_gp_ods.md** — пометки заглушек airplanes/seats.
|
||||
6. **docs/bookings_to_gp_dds.md** — пометки заглушек dim_routes/dim_passengers/dim_airplanes.
|
||||
7-10. **4 DAG docstrings** — эталон vs задания (заглушки).
|
||||
|
||||
### На solution (follow-up)
|
||||
|
||||
- README.md — маркетинг-баннер + шаг 7 (solution-specific: «Изучите готовую реализацию»).
|
||||
- docs/assignment/README.md — навигационный гид по всем паттернам (ODS → DDS → DM).
|
||||
- Этот план архивирован в docs/archive/.
|
||||
|
||||
## Верификация
|
||||
|
||||
- `make test` — passed на обеих ветках.
|
||||
- `make lint` — clean на обеих ветках.
|
||||
- Grep «заглушки/SELECT 1» в README.md и docs/assignment/README.md на solution — не найдено.
|
||||
@@ -1,17 +1,50 @@
|
||||
# Учебные задания
|
||||
|
||||
> Это ветка `solution` — здесь всё уже реализовано.
|
||||
> Используйте её как справочник, а задания выполняйте на ветке `main`.
|
||||
|
||||
## Навигатор по реализации
|
||||
|
||||
На этой ветке полностью работает пайплайн `STG → ODS → DDS → DM`.
|
||||
Ниже — рекомендуемый порядок изучения: от простых паттернов к сложным.
|
||||
|
||||
### ODS: загрузка из STG
|
||||
|
||||
| Паттерн | Файл | Что посмотреть |
|
||||
|---------|------|----------------|
|
||||
| TRUNCATE + INSERT (snapshot-справочник) | `sql/ods/airports_load.sql` | Фильтрация по `stg_batch_id`, дедупликация через `ROW_NUMBER()` |
|
||||
| SCD1 UPSERT (транзакционные данные) | `sql/ods/bookings_load.sql` | TEMP TABLE → UPDATE (`IS DISTINCT FROM`) → INSERT, HWM |
|
||||
|
||||
### DDS: измерения и факт
|
||||
|
||||
| Паттерн | Файл | Что посмотреть |
|
||||
|---------|------|----------------|
|
||||
| SCD1 (простое измерение) | `sql/dds/dim_airports_load.sql` | UPSERT с суррогатным ключом |
|
||||
| SCD1 с обогащением | `sql/dds/dim_airplanes_load.sql` | JOIN с `ods.seats` для `total_seats` |
|
||||
| SCD2 (версионирование) | `sql/dds/dim_routes_load.sql` | `hashdiff`, `valid_from`/`valid_to`, денормализация |
|
||||
| Факт с point-in-time JOIN | `sql/dds/fact_flight_sales_load.sql` | SCD2 lookup, защитные LEFT JOIN |
|
||||
|
||||
### DM: витрины
|
||||
|
||||
| Паттерн | Файл | Что посмотреть |
|
||||
|---------|------|----------------|
|
||||
| UPSERT с HWM по датам | `sql/dm/sales_report_load.sql` | Инкрементальная агрегация фактов |
|
||||
| Full Rebuild (AO Column) | `sql/dm/route_performance_load.sql` | TRUNCATE + INSERT, агрегация по SCD2 `route_bk` |
|
||||
| UNION ALL (dual-role dimension) | `sql/dm/airport_traffic_load.sql` | Разворот факта: departure + arrival |
|
||||
| Двухуровневая агрегация | `sql/dm/monthly_overview_load.sql` | Честный `avg_load_factor` через 2 уровня |
|
||||
| Метод затронутых ключей | `sql/dm/passenger_loyalty_load.sql` | HWM по пассажирам, пересчёт полной истории |
|
||||
|
||||
## Техническое задание
|
||||
|
||||
Основной документ: **[analyst_spec.md](analyst_spec.md)** — ТЗ от аналитика
|
||||
с описанием всех таблиц, маппингами, бизнес-правилами и подсказками.
|
||||
|
||||
## Эталон для изучения
|
||||
## Описания DAG-ов
|
||||
|
||||
Перед началом работы изучите эталонный срез (витрина `dm.sales_report`
|
||||
и вся её цепочка STG → ODS → DDS → DM):
|
||||
|
||||
- SQL-скрипты: `sql/stg/`, `sql/ods/`, `sql/dds/`, `sql/dm/`
|
||||
- DAG-файлы: `airflow/dags/`
|
||||
- [STG](../bookings_to_gp_stage.md) ·
|
||||
[ODS](../bookings_to_gp_ods.md) ·
|
||||
[DDS](../bookings_to_gp_dds.md) ·
|
||||
[DM](../bookings_to_gp_dm.md)
|
||||
- [Порядок запуска DAG-ов](../reference/dag_execution_order.md)
|
||||
|
||||
## Для менторов
|
||||
|
||||
@@ -123,8 +123,12 @@ DQ-проверки факта отражают эту логику:
|
||||
> В боевых системах для этого используют «строку-заглушку» (unknown member, SK = 0)
|
||||
> и отдельный процесс backfill.
|
||||
|
||||
**Point-in-time join для SCD2.**
|
||||
Для `dim_routes` используется привязка по дате рейса:
|
||||
**Два пути lookup для аэропортов и маршрутов.**
|
||||
Аэропорты (`departure_airport_sk`, `arrival_airport_sk`) разрешаются через `ods.routes` →
|
||||
`dim_airports`. Аэропорты вылета/прилёта одинаковы во всех версиях маршрута, поэтому
|
||||
point-in-time логика не нужна — безопасно брать актуальную версию из ODS.
|
||||
|
||||
`route_sk` и `airplane_sk` разрешаются через point-in-time join с SCD2 `dim_routes`:
|
||||
|
||||
```sql
|
||||
LEFT JOIN dds.dim_routes AS rte
|
||||
|
||||
@@ -168,13 +168,14 @@ DISTRIBUTED BY (ticket_no)
|
||||
|
||||
### 3.8. Политика NULL FK в факте
|
||||
|
||||
FK суррогатные ключи разделены на три группы:
|
||||
FK суррогатные ключи разделены на четыре группы по источнику lookup:
|
||||
|
||||
| Группа | FK | NULL допустим? | Причина |
|
||||
|--------|-----|---------------|---------|
|
||||
| **Обязательные** | `tariff_sk`, `passenger_sk` | Нет | Данные всегда есть в ODS (segments, tickets). NULL = баг загрузки. |
|
||||
| **Зависят от маршрута** | `route_sk`, `departure_airport_sk`, `arrival_airport_sk`, `airplane_sk` | Нет (в норме) | Маршрут должен быть в ODS. NULL = аномалия данных, DQ предупреждает. |
|
||||
| **Зависят от расписания** | `calendar_sk` | Допустим (редко) | `scheduled_departure` может быть NULL в ODS. DQ считает и логирует, но не фейлит. |
|
||||
| Группа | FK | Lookup-путь | NULL допустим? | Причина |
|
||||
|--------|-----|-------------|---------------|---------|
|
||||
| **Обязательные** | `tariff_sk`, `passenger_sk` | `dim_tariffs` / `dim_passengers` (SCD1) | Нет | Данные всегда есть в ODS. NULL = баг загрузки. |
|
||||
| **Эталонный справочник** | `departure_airport_sk`, `arrival_airport_sk` | `ods.routes` → `dim_airports` | Нет (в норме) | Аэропорты разрешаются через `ods.routes` (эталон), не через `dim_routes`. NULL = аномалия данных, DQ проверяет с порогом 1%. |
|
||||
| **Point-in-time SCD2** | `route_sk`, `airplane_sk` | `dim_routes` (SCD2 lookup) | Нет (в норме) | Point-in-time привязка к версии маршрута на дату рейса. NULL = аномалия, DQ предупреждает. |
|
||||
| **Расписание** | `calendar_sk` | `dim_calendar` | Допустим (редко) | `scheduled_departure` может быть NULL в ODS. DQ считает и логирует. |
|
||||
|
||||
DQ-проверки явно контролируют каждую группу (см. секцию 6).
|
||||
|
||||
@@ -459,15 +460,27 @@ WITH fact_src AS (
|
||||
JOIN ods.tickets AS tkt ON tkt.ticket_no = seg.ticket_no
|
||||
JOIN ods.bookings AS bkg ON bkg.book_ref = tkt.book_ref
|
||||
JOIN ods.flights AS flt ON flt.flight_id = seg.flight_id
|
||||
-- Airport lookup через ods.routes (эталонный справочник).
|
||||
-- Аэропорты вылета/прилёта одинаковы во всех версиях маршрута —
|
||||
-- безопасно брать из ODS без point-in-time логики.
|
||||
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
|
||||
-- SCD2 point-in-time: версия маршрута, актуальная на дату вылета
|
||||
LEFT JOIN dds.dim_routes AS rte
|
||||
ON rte.route_bk = flt.route_no
|
||||
AND flt.scheduled_departure::DATE >= rte.valid_from
|
||||
AND (rte.valid_to IS NULL OR flt.scheduled_departure::DATE < rte.valid_to)
|
||||
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
|
||||
LEFT JOIN dds.dim_airports AS arr ON arr.airport_bk = rte.arrival_airport
|
||||
LEFT JOIN dds.dim_airplanes AS ap ON ap.airplane_bk = rte.airplane_code
|
||||
LEFT JOIN dds.dim_airports AS dep ON dep.airport_bk = ods_rte.departure_airport -- через ods.routes
|
||||
LEFT JOIN dds.dim_airports AS arr ON arr.airport_bk = ods_rte.arrival_airport -- через ods.routes
|
||||
LEFT JOIN dds.dim_airplanes AS ap ON ap.airplane_bk = rte.airplane_code -- через dim_routes
|
||||
LEFT JOIN dds.dim_tariffs AS tar ON tar.fare_conditions = seg.fare_conditions
|
||||
LEFT JOIN dds.dim_passengers AS pax ON pax.passenger_bk = tkt.passenger_id
|
||||
LEFT JOIN ods.boarding_passes AS bp
|
||||
@@ -497,8 +510,9 @@ ANALYZE dds.fact_flight_sales;
|
||||
### 5.5. Модель историчности факта
|
||||
|
||||
Dimension SK фиксируются **при INSERT** и не перезаписываются:
|
||||
- `route_sk` — версия маршрута на дату `scheduled_departure` (point-in-time SCD2 lookup);
|
||||
- `departure_airport_sk`, `arrival_airport_sk`, `airplane_sk` — из той же версии маршрута;
|
||||
- `departure_airport_sk`, `arrival_airport_sk` — из `ods.routes` → `dim_airports` (эталонный справочник; аэропорты одинаковы для всех версий маршрута);
|
||||
- `route_sk` — версия маршрута на дату `scheduled_departure` (point-in-time SCD2 lookup через `dim_routes`);
|
||||
- `airplane_sk` — из той же point-in-time версии `dim_routes`;
|
||||
- `calendar_sk`, `tariff_sk`, `passenger_sk` — из текущих SCD1-измерений на момент INSERT.
|
||||
|
||||
UPDATE факта обновляет только **мутабельные поля**: `is_boarded`, `seat_no`, `price`
|
||||
|
||||
@@ -115,6 +115,7 @@ graph LR
|
||||
ODS_Tickets -->|book_ref, passenger_id| FACT_Sales
|
||||
ODS_Bookings -->|book_date| FACT_Sales
|
||||
ODS_Flights -->|schedule, route_no| FACT_Sales
|
||||
ODS_Routes -->|dep/arr airports| FACT_Sales
|
||||
ODS_Boarding -->|LEFT JOIN seat_no| FACT_Sales
|
||||
|
||||
%% Dimensions to Fact
|
||||
@@ -220,12 +221,12 @@ SK генерация: `MAX(sk) + ROW_NUMBER()` (безопасно при `conc
|
||||
|
||||
`dds.fact_flight_sales` — зерно: 1 строка = 1 сегмент билета (`ticket_no` + `flight_id`).
|
||||
|
||||
| FK | Источник | Примечание |
|
||||
|----|----------|------------|
|
||||
| FK | Lookup-путь | Примечание |
|
||||
|----|-------------|------------|
|
||||
| `calendar_sk` | `dim_calendar` | Дата вылета |
|
||||
| `departure_airport_sk` | `dim_airports` | Аэропорт вылета |
|
||||
| `arrival_airport_sk` | `dim_airports` | Аэропорт прилёта |
|
||||
| `airplane_sk` | `dim_airplanes` | Самолёт |
|
||||
| `departure_airport_sk` | `ods.routes` → `dim_airports` | Аэропорт вылета (эталонный справочник) |
|
||||
| `arrival_airport_sk` | `ods.routes` → `dim_airports` | Аэропорт прилёта (эталонный справочник) |
|
||||
| `airplane_sk` | `dim_routes` → `dim_airplanes` | Самолёт (point-in-time SCD2) |
|
||||
| `tariff_sk` | `dim_tariffs` | Тариф |
|
||||
| `passenger_sk` | `dim_passengers` | Пассажир |
|
||||
| `route_sk` | `dim_routes` | Версия маршрута (SCD2, point-in-time) |
|
||||
|
||||
Reference in New Issue
Block a user