docs(dds): добавлены комментарии и исправлены ошибки в DWH
- Зачем: - закрыты задачи P0 и P1 из ревью архитектуры для повышения понятности стенда для студентов. - Что: - исправлен distribution key для airport_traffic в дизайн-документе. - добавлены комментарии о генерации SK и отсутствии SK в фактах. - создан документ docs/dag_execution_order.md с описанием порядка запуска DAG-ов. - объяснена логика late-arriving dimensions и batch resolver. - добавлена legacy-пометка для хелпера greenplum.py. - Проверка: - визуальная проверка добавленных комментариев и новых файлов.
This commit is contained in:
@@ -34,8 +34,11 @@ def _resolve_stg_batch_id(**context) -> str:
|
|||||||
"""
|
"""
|
||||||
Возвращает stg_batch_id из dag_run.conf или вычисляет последний согласованный батч.
|
Возвращает stg_batch_id из dag_run.conf или вычисляет последний согласованный батч.
|
||||||
|
|
||||||
Согласованным считаем batch_id, который есть во всех snapshot-справочниках STG:
|
Зачем нужна согласованность (INTERSECT по всем справочникам)?
|
||||||
airports, airplanes, routes, seats.
|
Чтобы ODS загружал только те данные, для которых уже приехали ВСЕ связанные
|
||||||
|
справочники. Это защищает от рассинхрона данных, когда часть измерений в текущем
|
||||||
|
батче обновилась, а часть — упала или не доехала, что могло бы привести к
|
||||||
|
потере ссылочной целостности при сборке витрин.
|
||||||
"""
|
"""
|
||||||
conf = context["dag_run"].conf or {}
|
conf = context["dag_run"].conf or {}
|
||||||
stg_batch_id = conf.get("stg_batch_id")
|
stg_batch_id = conf.get("stg_batch_id")
|
||||||
|
|||||||
@@ -1,5 +1,12 @@
|
|||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
|
"""
|
||||||
|
LEGACY: Вспомогательные функции для прямого подключения к Greenplum через psycopg2.
|
||||||
|
Внимание: этот модуль оставлен только для поддержки базового CSV-пайплайна.
|
||||||
|
В новых DAG (ODS/DDS/DM) используйте встроенный в Airflow PostgresOperator
|
||||||
|
и штатные механизмы XCom.
|
||||||
|
"""
|
||||||
|
|
||||||
import logging
|
import logging
|
||||||
import os
|
import os
|
||||||
from typing import List, Sequence, Tuple
|
from typing import List, Sequence, Tuple
|
||||||
|
|||||||
@@ -0,0 +1,33 @@
|
|||||||
|
# Порядок запуска DAG (Cross-DAG Dependencies)
|
||||||
|
|
||||||
|
В этом стенде пайплайны разделены на несколько DAG-ов по слоям DWH (STG, ODS, DDS, DM).
|
||||||
|
Они настроены с `schedule=None`, так как это учебный проект.
|
||||||
|
|
||||||
|
Чтобы данные корректно прошли от источника до витрин, запускать DAG-и нужно в определённом порядке.
|
||||||
|
|
||||||
|
## 1. DDL-скрипты (выполняются один раз)
|
||||||
|
|
||||||
|
Для создания структуры таблиц в аналитических слоях:
|
||||||
|
1. Запустите `bookings_dds_ddl` — создаст таблицы для измерений и фактов в слое DDS.
|
||||||
|
2. Запустите `bookings_dm_ddl` — создаст таблицы витрин в слое DM.
|
||||||
|
|
||||||
|
*(Слои STG и ODS создаются при старте стенда через `make ddl-gp` или могут быть пересозданы соответствующими DDL-скриптами).*
|
||||||
|
|
||||||
|
## 2. Ежедневная загрузка (ETL)
|
||||||
|
|
||||||
|
Для прогрузки новой порции данных (или полного перерасчёта) соблюдайте следующую цепочку:
|
||||||
|
|
||||||
|
1. **`bookings_to_gp_stage`**
|
||||||
|
- Извлекает новые данные из демо-БД PostgreSQL и сохраняет их в `stg`-схему в Greenplum.
|
||||||
|
- Генерирует `stg_batch_id` для текущей загрузки.
|
||||||
|
2. **`bookings_to_gp_ods`**
|
||||||
|
- Берёт последний согласованный `stg_batch_id` из STG-слоя.
|
||||||
|
- Выполняет нормализацию и SCD1-UPSERT в слой ODS.
|
||||||
|
3. **`bookings_to_gp_dds`**
|
||||||
|
- Читает очищенные данные из ODS.
|
||||||
|
- Обновляет измерения (SCD1, SCD2) и инкрементально догружает новые рейсы в таблицу фактов `dds.fact_flight_sales`.
|
||||||
|
4. **`bookings_to_gp_dm`**
|
||||||
|
- Читает новые факты из DDS.
|
||||||
|
- Обновляет агрегированные витрины (использует HWM-инкрементальность по `_load_ts` или полный перерасчёт).
|
||||||
|
|
||||||
|
> **💡 Архитектурная заметка:** В реальном production-окружении (Airflow) эти связи между DAG-ами обычно настраиваются автоматически через `TriggerDagRunOperator`, `ExternalTaskSensor` или механизмы Data-Aware Scheduling (Datasets/Data Assets). В учебных целях мы оставили их ручными, чтобы вы могли проинспектировать каждый слой после его загрузки.
|
||||||
@@ -53,7 +53,7 @@ GP-специфичная best practice, которую забывают даж
|
|||||||
|
|
||||||
### P0: Фактическая ошибка (исправить до показа студентам)
|
### P0: Фактическая ошибка (исправить до показа студентам)
|
||||||
|
|
||||||
- [ ] **Противоречие в distribution key для airport_traffic**
|
- [x] **Противоречие в distribution key для airport_traffic**
|
||||||
- `bookings_dm_design.md` (строка 182): `DISTRIBUTED BY (traffic_date)`
|
- `bookings_dm_design.md` (строка 182): `DISTRIBUTED BY (traffic_date)`
|
||||||
- `sales_report_ddl.sql`: явно объясняет, почему distribution by date — антипаттерн
|
- `sales_report_ddl.sql`: явно объясняет, почему distribution by date — антипаттерн
|
||||||
- **Нужно**: исправить на `DISTRIBUTED BY (airport_sk)` в дизайн-документе
|
- **Нужно**: исправить на `DISTRIBUTED BY (airport_sk)` в дизайн-документе
|
||||||
@@ -61,27 +61,27 @@ GP-специфичная best practice, которую забывают даж
|
|||||||
|
|
||||||
### P1: Высокий эффект, минимум усилий (комментарии и документация)
|
### P1: Высокий эффект, минимум усилий (комментарии и документация)
|
||||||
|
|
||||||
- [ ] **Нет объяснения «почему не SERIAL» в генерации SK**
|
- [x] **Нет объяснения «почему не SERIAL» в генерации SK**
|
||||||
- `MAX(sk) + ROW_NUMBER()` корректен для GP, но студент на PostgreSQL/Snowflake будет использовать `IDENTITY`/`SEQUENCE`
|
- `MAX(sk) + ROW_NUMBER()` корректен для GP, но студент на PostgreSQL/Snowflake будет использовать `IDENTITY`/`SEQUENCE`
|
||||||
- **Нужно**: 4-строчный комментарий в `sql/dds/dim_airports_load.sql`
|
- **Нужно**: 4-строчный комментарий в `sql/dds/dim_airports_load.sql`
|
||||||
|
|
||||||
- [ ] **Факт без суррогатного ключа — не объяснено «почему»**
|
- [x] **Факт без суррогатного ключа — не объяснено «почему»**
|
||||||
- Натуральный (ticket_no, flight_id) как grain — правильное Kimball-моделирование
|
- Натуральный (ticket_no, flight_id) как grain — правильное Kimball-моделирование
|
||||||
- **Нужно**: комментарий в `sql/dds/fact_flight_sales_ddl.sql`
|
- **Нужно**: комментарий в `sql/dds/fact_flight_sales_ddl.sql`
|
||||||
|
|
||||||
- [ ] **Нет упоминания cross-DAG зависимостей**
|
- [x] **Нет упоминания cross-DAG зависимостей**
|
||||||
- STG, ODS, DDS, DM — отдельные DAG-и с `schedule=None`, студент может не понять порядок
|
- STG, ODS, DDS, DM — отдельные DAG-и с `schedule=None`, студент может не понять порядок
|
||||||
- **Нужно**: комментарий в docstring каждого DAG или `docs/dag_execution_order.md`
|
- **Нужно**: комментарий в docstring каждого DAG или `docs/dag_execution_order.md`
|
||||||
|
|
||||||
- [ ] **Late-arriving dimensions не упомянуты**
|
- [x] **Late-arriving dimensions не упомянуты**
|
||||||
- Факт делает LEFT JOIN → `passenger_sk = NULL` при опоздании; нет механизма исправления
|
- Факт делает LEFT JOIN → `passenger_sk = NULL` при опоздании; нет механизма исправления
|
||||||
- **Нужно**: комментарий в `sql/dds/fact_flight_sales_load.sql` у LEFT JOIN-ов
|
- **Нужно**: комментарий в `sql/dds/fact_flight_sales_load.sql` у LEFT JOIN-ов
|
||||||
|
|
||||||
- [ ] **Batch resolver недообъяснён**
|
- [x] **Batch resolver недообъяснён**
|
||||||
- `_resolve_stg_batch_id` с INTERSECT по 4 таблицам — нет комментария **зачем** нужна согласованность
|
- `_resolve_stg_batch_id` с INTERSECT по 4 таблицам — нет комментария **зачем** нужна согласованность
|
||||||
- **Нужно**: комментарий в `airflow/dags/bookings_to_gp_ods.py` перед SQL-запросом
|
- **Нужно**: комментарий в `airflow/dags/bookings_to_gp_ods.py` перед SQL-запросом
|
||||||
|
|
||||||
- [ ] **`helpers/greenplum.py` без пометки «legacy»**
|
- [x] **`helpers/greenplum.py` без пометки «legacy»**
|
||||||
- Использует прямой psycopg2 + ENV — противоречит PostgresOperator-подходу
|
- Использует прямой psycopg2 + ENV — противоречит PostgresOperator-подходу
|
||||||
- **Нужно**: docstring «LEGACY: только для CSV-пайплайна» в `airflow/dags/helpers/greenplum.py`
|
- **Нужно**: docstring «LEGACY: только для CSV-пайплайна» в `airflow/dags/helpers/greenplum.py`
|
||||||
|
|
||||||
|
|||||||
@@ -178,7 +178,7 @@ FROM traffic GROUP BY ...
|
|||||||
|
|
||||||
**Загрузка**: Инкрементальный UPSERT по HWM (`_load_ts`). Из дельты фактов определяем затронутые `(traffic_date, airport_sk)`, пересчитываем агрегаты только для них.
|
**Загрузка**: Инкрементальный UPSERT по HWM (`_load_ts`). Из дельты фактов определяем затронутые `(traffic_date, airport_sk)`, пересчитываем агрегаты только для них.
|
||||||
|
|
||||||
**Хранение**: `DISTRIBUTED BY (traffic_date)`, heap
|
**Хранение**: `DISTRIBUTED BY (airport_sk)`, heap
|
||||||
|
|
||||||
**Учит**: HWM-инкрементальность, TEMP TABLE для однократной агрегации, dual-role dimension join (UNION ALL), conditional aggregation (CASE WHEN + SUM), паттерн "unpivot → aggregate"
|
**Учит**: HWM-инкрементальность, TEMP TABLE для однократной агрегации, dual-role dimension join (UNION ALL), conditional aggregation (CASE WHEN + SUM), паттерн "unpivot → aggregate"
|
||||||
|
|
||||||
|
|||||||
@@ -22,6 +22,9 @@ WHERE d.airport_bk = s.airport_code
|
|||||||
|
|
||||||
-- Statement 2: INSERT новых записей (MAX(sk) + ROW_NUMBER()).
|
-- Statement 2: INSERT новых записей (MAX(sk) + ROW_NUMBER()).
|
||||||
-- Учебный комментарий: Генерация SK через MAX() + ROW_NUMBER()
|
-- Учебный комментарий: Генерация SK через MAX() + ROW_NUMBER()
|
||||||
|
-- Почему не SERIAL/IDENTITY? В MPP-базах данных (как Greenplum) sequence
|
||||||
|
-- работают через мастер-узел и могут стать узким местом при массовой вставке.
|
||||||
|
-- Паттерн MAX() + ROW_NUMBER() генерирует ключи распределённо на сегментах.
|
||||||
-- Этот подход работает безопасно только потому, что Airflow запускает
|
-- Этот подход работает безопасно только потому, что Airflow запускает
|
||||||
-- джобы загрузки для одной таблицы строго последовательно (concurrency=1).
|
-- джобы загрузки для одной таблицы строго последовательно (concurrency=1).
|
||||||
-- При параллельной загрузке возникнет состояние гонки (race condition) и возможны дубли SK.
|
-- При параллельной загрузке возникнет состояние гонки (race condition) и возможны дубли SK.
|
||||||
|
|||||||
@@ -2,6 +2,10 @@
|
|||||||
|
|
||||||
CREATE SCHEMA IF NOT EXISTS dds;
|
CREATE SCHEMA IF NOT EXISTS dds;
|
||||||
|
|
||||||
|
-- Учебный комментарий: Почему у факта нет своего суррогатного ключа (fact_sk)?
|
||||||
|
-- В классическом DWH (Кимбалл) таблица фактов идентифицируется набором её
|
||||||
|
-- измерений или дегенеративных ключей (в нашем случае: ticket_no + flight_id).
|
||||||
|
-- Добавление отдельного ID только тратит место и не несёт аналитической ценности.
|
||||||
CREATE TABLE IF NOT EXISTS dds.fact_flight_sales (
|
CREATE TABLE IF NOT EXISTS dds.fact_flight_sales (
|
||||||
calendar_sk INTEGER,
|
calendar_sk INTEGER,
|
||||||
departure_airport_sk INTEGER,
|
departure_airport_sk INTEGER,
|
||||||
|
|||||||
@@ -45,6 +45,11 @@ WITH fact_src AS (
|
|||||||
ON bkg.book_ref = tkt.book_ref
|
ON bkg.book_ref = tkt.book_ref
|
||||||
JOIN ods.flights AS flt
|
JOIN ods.flights AS flt
|
||||||
ON flt.flight_id = seg.flight_id
|
ON flt.flight_id = seg.flight_id
|
||||||
|
-- Учебный комментарий: Late-arriving dimensions (Опаздывающие измерения)
|
||||||
|
-- Мы используем LEFT JOIN, так как факт (рейс/билет) может прийти раньше,
|
||||||
|
-- чем справочник (пассажир/маршрут) обновится в DDS.
|
||||||
|
-- В результате SK будет NULL. В более сложных пайплайнах такие факты
|
||||||
|
-- либо обогащаются dummy-значениями (-1, "Неизвестно"), либо откладываются.
|
||||||
LEFT JOIN dds.dim_routes AS rte
|
LEFT JOIN dds.dim_routes AS rte
|
||||||
ON rte.route_bk = flt.route_no
|
ON rte.route_bk = flt.route_no
|
||||||
AND flt.scheduled_departure::DATE >= rte.valid_from
|
AND flt.scheduled_departure::DATE >= rte.valid_from
|
||||||
|
|||||||
Reference in New Issue
Block a user