diff --git a/README.md b/README.md index 85fb6e1..cc31b5e 100644 --- a/README.md +++ b/README.md @@ -13,8 +13,8 @@ - работы с Airflow и Greenplum. В курсовой у нас один источник данных — демо‑БД **bookings**. В стенде уже есть готовые учебные примеры -загрузки **bookings → stg**, **stg -> ods** и **ods -> dds** в Greenplum, чтобы вы могли сфокусироваться -на DWH‑части (ODS/DDS/DM) и не тратить время на инфраструктуру. +загрузки **bookings → STG → ODS → DDS → DM** в Greenplum, чтобы вы могли сфокусироваться +на DWH‑части и не тратить время на инфраструктуру. ## Что внутри @@ -43,7 +43,7 @@ Про PXF и технические детали стенда: [docs/stack.md](docs/stack.md). -## Быстрый старт (основной сценарий: bookings -> stg -> ods -> dds) +## Быстрый старт (основной сценарий: bookings → STG → ODS → DDS → DM) 1) Скопируйте настройки: @@ -65,18 +65,19 @@ make bookings-init # быстрое восстановление из s # make bookings-generate # альтернатива: полная генерация с нуля (занимает часы) ``` -4) Подготовьте STG/ODS/DDS-объекты в Greenplum (выберите один вариант): +4) Подготовьте объекты DWH в Greenplum (выберите один вариант): -- Учебный вариант: в Airflow UI запустите DAG `bookings_stg_ddl`, затем `bookings_ods_ddl`, затем `bookings_dds_ddl`; -- Технический шорткат: `make ddl-gp` (применяет DDL для STG, ODS и DDS разом вручную). +- Учебный вариант: в Airflow UI запустите DAG'и по порядку — `bookings_stg_ddl`, `bookings_ods_ddl`, `bookings_dds_ddl`, `bookings_dm_ddl`; +- Технический шорткат: `make ddl-gp` (применяет DDL для STG, ODS, DDS и DM разом). -5) Запустите основной DAG `bookings_to_gp_stage`. +5) Запустите DAG'и загрузки данных по порядку: -6) Запустите DAG `bookings_to_gp_ods`. +- `bookings_to_gp_stage` — загрузка в STG +- `bookings_to_gp_ods` — загрузка в ODS +- `bookings_to_gp_dds` — загрузка в DDS +- `bookings_to_gp_dm` — загрузка в DM (витрины) -7) Запустите DAG `bookings_to_gp_dds`. - -8) Проверьте результат в Greenplum: +6) Проверьте результат в Greenplum: ```bash make gp-psql @@ -88,6 +89,8 @@ SELECT COUNT(*) FROM ods.bookings; SELECT COUNT(*) FROM ods.tickets; SELECT COUNT(*) FROM dds.dim_routes; SELECT COUNT(*) FROM dds.fact_flight_sales; +SELECT COUNT(*) FROM dm.sales_report; +SELECT COUNT(*) FROM dm.route_performance; ``` Подробнее про логику DAG и проверки — `docs/bookings_to_gp_stage.md`. @@ -103,6 +106,8 @@ SELECT COUNT(*) FROM dds.fact_flight_sales; - `bookings_to_gp_ods` — загружает данные из STG в ODS (SCD1 UPSERT) и выполняет DQ‑проверки. - `bookings_dds_ddl` — создаёт/обновляет DDS-таблицы (`dim_*`, `fact_flight_sales`) по домену bookings. - `bookings_to_gp_dds` — загружает данные из ODS в DDS (SCD1/SCD2 + факт) и выполняет DQ‑проверки. +- `bookings_dm_ddl` — создаёт/обновляет DM-витрины (sales_report, route_performance, passenger_loyalty, airport_traffic, monthly_overview). +- `bookings_to_gp_dm` — загружает данные из DDS в DM (агрегированные витрины). ## Полезные команды @@ -111,7 +116,7 @@ make up # поднять стек make logs # логи airflow-webserver и airflow-scheduler make gp-psql # psql в Greenplum make bookings-psql # psql в демо-БД bookings (Postgres) -make ddl-gp # применить DDL STG+ODS+DDS к Greenplum вручную (вместо DDL-DAG) +make ddl-gp # применить DDL STG+ODS+DDS+DM к Greenplum вручную (вместо DDL-DAG) make down # остановить и удалить контейнеры/сети (volumes сохраняются) make clean # полный reset: удалить контейнеры/сети и volumes (данные будут потеряны) ``` @@ -138,7 +143,7 @@ make clean # полный reset: удалить контейнер ## Документация -- Учебные задания: `educational-tasks.md`). +- Учебные задания: `educational-tasks.md`. - План тестирования/проверок и негативные кейсы: `TESTING.md`. - Дополнительные заметки и технические детали: `docs/README.md`. - Детали по ODS DAG: `docs/bookings_to_gp_ods.md`. @@ -150,7 +155,7 @@ make clean # полный reset: удалить контейнер |----------|---------| | Airflow UI не открывается | Дождитесь сообщения `Listening at: http://0.0.0.0:8080` в логах (`make logs`) | | `database "demo" does not exist` в bookings‑DAG | Вы сделали reset с удалением volumes (`make clean` / `docker compose down -v`). Запустите `make bookings-init` (быстрое восстановление из дампа, ~18 сек) и повторите DAG. | -| Ошибка `bookings.jobs должен быть >= 1` | Проверьте значение `BOOKINGS_JOBS` в `.env` — оно должно быть целым числом >= 1. По умолчанию `BOOKINGS_JOBS=2` (параллельная генерация через dblink). При `BOOKINGS_JOBS=1` генерация идёт синхронно. | +| Ошибка `bookings.jobs должен быть >= 1` | Проверьте значение `BOOKINGS_JOBS` в `.env` — оно должно быть целым числом >= 1. По умолчанию `BOOKINGS_JOBS=1` (синхронная генерация). При `BOOKINGS_JOBS>1` генерация идёт параллельно через dblink. | | Ошибка подключения к Greenplum | Убедитесь, что контейнер `greenplum` имеет статус `healthy` (`docker compose ps`) | | `protocol "pxf" does not exist` | Перезапустите `greenplum` и повторите `bookings_stg_ddl`/`make ddl-gp` — расширение `pxf` создаётся автоматически при старте контейнера. | | DAG `bookings_to_gp_stage` ругается на отсутствующие таблицы stg | Запустите `bookings_stg_ddl` (или выполните `make ddl-gp`), затем повторите запуск | diff --git a/TODO.md b/TODO.md index 091ebb0..acff723 100644 --- a/TODO.md +++ b/TODO.md @@ -27,12 +27,12 @@ **Инструмент:** Opus (глубокий анализ кода и контекста проекта) + ручное тестирование (make up, запуск DAG'ов, проверка данных). -- [ ] Протестировать полный ETL-цикл с нуля +- [x] Протестировать полный ETL-цикл с нуля (make up → bookings-init (восстановление из дампа) → STG → ODS → DDS → DM) -- [ ] Прогнать инкремент (bookings-generate-day → повторный запуск DAG'ов) -- [ ] Почистить код эталонного среза -- [ ] Актуализировать README и документацию -- [ ] Убедиться, что `make test` и `make lint` проходят +- [x] Прогнать инкремент (bookings-generate-day → повторный запуск DAG'ов) +- [x] Почистить код эталонного среза +- [x] Актуализировать README и документацию +- [x] Убедиться, что `make test` и `make lint` проходят - [ ] Проверить, что стенд поднимается на чистой машине ### Этап 2. Подготовка main diff --git a/airflow/dags/bookings_dm_ddl.py b/airflow/dags/bookings_dm_ddl.py index 7d9b7b7..6af5223 100644 --- a/airflow/dags/bookings_dm_ddl.py +++ b/airflow/dags/bookings_dm_ddl.py @@ -67,4 +67,10 @@ with DAG( ) # Линейная цепочка - apply_dm_sales_report_ddl >> apply_dm_route_performance_ddl >> apply_dm_passenger_loyalty_ddl >> apply_dm_airport_traffic_ddl >> apply_dm_monthly_overview_ddl + ( + apply_dm_sales_report_ddl + >> apply_dm_route_performance_ddl + >> apply_dm_passenger_loyalty_ddl + >> apply_dm_airport_traffic_ddl + >> apply_dm_monthly_overview_ddl + ) diff --git a/airflow/dags/bookings_to_gp_dm.py b/airflow/dags/bookings_to_gp_dm.py index 0297594..c2cfb57 100644 --- a/airflow/dags/bookings_to_gp_dm.py +++ b/airflow/dags/bookings_to_gp_dm.py @@ -131,7 +131,7 @@ with DAG( load_dm_route_performance, load_dm_passenger_loyalty, load_dm_airport_traffic, - load_dm_monthly_overview + load_dm_monthly_overview, ] # Связываем dq с finish diff --git a/airflow/dags/bookings_to_gp_ods.py b/airflow/dags/bookings_to_gp_ods.py index f5637da..1acd63a 100644 --- a/airflow/dags/bookings_to_gp_ods.py +++ b/airflow/dags/bookings_to_gp_ods.py @@ -35,17 +35,17 @@ def _resolve_stg_batch_id(**context) -> str: Возвращает stg_batch_id из dag_run.conf или вычисляет последний согласованный батч. ВНИМАНИЕ: Это значение используется ТОЛЬКО для загрузки snapshot-справочников - (airports, airplanes, routes, seats). - + (airports, airplanes, routes, seats). + Для инкрементальных транзакционных таблиц (bookings, tickets, flights, segments, - boarding_passes) этот батч НЕ используется. Вместо этого они грузят все новые + boarding_passes) этот батч НЕ используется. Вместо этого они грузят все новые записи по HWM: WHERE load_dttm > (SELECT MAX(_load_ts) FROM ods.table). - Это сделано для того, чтобы не потерять инкременты, если STG-DAG запускался + Это сделано для того, чтобы не потерять инкременты, если STG-DAG запускался несколько раз до запуска ODS-DAG'а. Зачем нужна согласованность (INTERSECT по всем справочникам)? Чтобы ODS загружал только те данные, для которых уже приехали ВСЕ связанные - справочники. Это защищает от рассинхрона данных, когда часть измерений в текущем + справочники. Это защищает от рассинхрона данных, когда часть измерений в текущем батче обновилась, а часть — упала или не доехала, что могло бы привести к потере ссылочной целостности при сборке витрин. """ diff --git a/docs/stack.md b/docs/stack.md index 0272f42..afb944b 100644 --- a/docs/stack.md +++ b/docs/stack.md @@ -132,9 +132,4 @@ make fmt - `BOOKINGS_DB_PORT` — внешний порт (по умолчанию `5434`) - `BOOKINGS_START_DATE` — стартовая дата модельного времени - `BOOKINGS_INIT_DAYS` — сколько дней генерировать при `make bookings-generate` (генерация с нуля) -- `BOOKINGS_JOBS` — число джобов генератора (по умолчанию `2`) - -### CSV pipeline (побочный пример) - -- `CSV_DIR` — путь к каталогу с CSV внутри контейнеров Airflow (по умолчанию `/opt/airflow/data`) -- `CSV_ROWS` — количество строк, генерируемых DAG (по умолчанию `1000`) +- `BOOKINGS_JOBS` — число джобов генератора (по умолчанию `1`) diff --git a/sql/dm/airport_traffic_load.sql b/sql/dm/airport_traffic_load.sql index 2303b8d..214113e 100644 --- a/sql/dm/airport_traffic_load.sql +++ b/sql/dm/airport_traffic_load.sql @@ -127,3 +127,5 @@ WHERE NOT EXISTS ( WHERE tgt.traffic_date = src.traffic_date AND tgt.airport_sk = src.airport_sk ); + +ANALYZE dm.airport_traffic; diff --git a/sql/dm/monthly_overview_load.sql b/sql/dm/monthly_overview_load.sql index f17dd08..7a0aa89 100644 --- a/sql/dm/monthly_overview_load.sql +++ b/sql/dm/monthly_overview_load.sql @@ -163,3 +163,5 @@ WHERE NOT EXISTS ( AND tgt.month_actual = src.month_actual AND tgt.airplane_sk = src.airplane_sk ); + +ANALYZE dm.monthly_overview; diff --git a/sql/dm/passenger_loyalty_load.sql b/sql/dm/passenger_loyalty_load.sql index 273030f..fb8e3ad 100644 --- a/sql/dm/passenger_loyalty_load.sql +++ b/sql/dm/passenger_loyalty_load.sql @@ -55,7 +55,7 @@ fare_modes AS ( SELECT bm.*, fm.favorite_fare_conditions, - p.passenger_id AS passenger_bk, -- Исправлено: в dim_passengers BK называется passenger_id + p.passenger_id AS passenger_bk, p.passenger_name, (bm.last_flight_date - bm.first_flight_date) AS days_as_customer FROM base_metrics bm @@ -127,3 +127,5 @@ WHERE NOT EXISTS ( SELECT 1 FROM dm.passenger_loyalty AS tgt WHERE tgt.passenger_sk = src.passenger_sk ); + +ANALYZE dm.passenger_loyalty; diff --git a/sql/dm/route_performance_load.sql b/sql/dm/route_performance_load.sql index 64e9353..c8fd3c3 100644 --- a/sql/dm/route_performance_load.sql +++ b/sql/dm/route_performance_load.sql @@ -79,3 +79,5 @@ SELECT '{{ run_id }}' AS _load_id FROM tmp_route_metrics m JOIN dds.dim_routes r_curr ON m.route_bk = r_curr.route_bk AND r_curr.valid_to IS NULL; + +ANALYZE dm.route_performance; diff --git a/sql/dm/sales_report_load.sql b/sql/dm/sales_report_load.sql index 88fe9ce..8f8a6f8 100644 --- a/sql/dm/sales_report_load.sql +++ b/sql/dm/sales_report_load.sql @@ -162,3 +162,5 @@ WHERE NOT EXISTS ( AND tgt.arrival_airport_sk = src.arrival_airport_sk AND tgt.tariff_sk = src.tariff_sk ); + +ANALYZE dm.sales_report; diff --git a/tests/test_dags_smoke.py b/tests/test_dags_smoke.py index f94a597..6ade3ef 100644 --- a/tests/test_dags_smoke.py +++ b/tests/test_dags_smoke.py @@ -355,10 +355,18 @@ def test_bookings_dm_ddl_dag_structure(): assert expected_tasks.issubset(dag.task_dict.keys()) # Проверяем линейную цепочку - _assert_direct_edge(dag, "apply_dm_sales_report_ddl", "apply_dm_route_performance_ddl") - _assert_direct_edge(dag, "apply_dm_route_performance_ddl", "apply_dm_passenger_loyalty_ddl") - _assert_direct_edge(dag, "apply_dm_passenger_loyalty_ddl", "apply_dm_airport_traffic_ddl") - _assert_direct_edge(dag, "apply_dm_airport_traffic_ddl", "apply_dm_monthly_overview_ddl") + _assert_direct_edge( + dag, "apply_dm_sales_report_ddl", "apply_dm_route_performance_ddl" + ) + _assert_direct_edge( + dag, "apply_dm_route_performance_ddl", "apply_dm_passenger_loyalty_ddl" + ) + _assert_direct_edge( + dag, "apply_dm_passenger_loyalty_ddl", "apply_dm_airport_traffic_ddl" + ) + _assert_direct_edge( + dag, "apply_dm_airport_traffic_ddl", "apply_dm_monthly_overview_ddl" + ) def test_bookings_to_gp_dm_dag_structure():