From 3de4db8aa656f1a67f07ac5e9a781821aeec8b1b Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Mon, 9 Mar 2026 20:57:58 +0300 Subject: [PATCH] =?UTF-8?q?docs(all):=20=D1=80=D0=B5=D0=B2=D0=B8=D0=B7?= =?UTF-8?q?=D0=B8=D1=8F=20=D0=B4=D0=BE=D0=BA=D1=83=D0=BC=D0=B5=D0=BD=D1=82?= =?UTF-8?q?=D0=B0=D1=86=D0=B8=D0=B8=20=D0=BF=D0=B5=D1=80=D0=B5=D0=B4=20?= =?UTF-8?q?=D0=BC=D0=B5=D1=80=D0=B6=D0=B5=D0=BC=20=D0=B2=20main?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - ветка содержала устаревшие ссылки, артефакты CSV-пайплайна и метки «черновик» для полностью реализованных слоёв STG→ODS→DDS→DM. - Что: - AGENTS.md: заменена фраза «в будущем» на перечисление реальных слоёв ODS/DDS/DM. - TESTING.md: удалены две строки про каталог data/ (артефакт CSV-пайплайна). - README.md: список документации заменён на кликабельные markdown-ссылки, добавлены STG и DM DAG. - docs/README.md: добавлен DM DAG в «Быстрый путь», DM design в «Технические детали»; убраны метки «(черновик)». - docs/internal/bookings_stg_design.md: убран заголовок «черновик», исправлены описания слоёв и DDL. - docs/internal/PRD.md: битая ссылка на analyst_spec.md заменена текстом с пометкой TODO. - TODO.md: ссылка на plans/ обновлена на docs/internal/bookings_db_issues.md. - docs/bookings_to_gp_dm.md: создан новый документ по аналогии с DDS doc (5 витрин, граф, DQ, ошибки). - plans/ и docs/chore/: каталоги удалены (планы выполнены, история сохранена в git). - Проверка: - make lint && make test — прошло чисто. - grep -n "черновик|data/|в будущем|plans/" — пустой результат. --- AGENTS.md | 4 +- README.md | 12 +- TESTING.md | 2 - TODO.md | 2 +- docs/README.md | 10 +- docs/bookings_to_gp_dm.md | 79 +++ docs/chore/bookings-etl.md | 466 -------------- docs/internal/PRD.md | 2 +- docs/internal/bookings_stg_design.md | 12 +- plans/bookings-demodb-bugfix-plan.md | 100 --- plans/dockerfile-improvements.md | 489 --------------- plans/greenplum-pxf-custom-image-plan.md | 142 ----- plans/stg_layer_implementation_plan.md | 744 ----------------------- 13 files changed, 102 insertions(+), 1962 deletions(-) create mode 100644 docs/bookings_to_gp_dm.md delete mode 100644 docs/chore/bookings-etl.md delete mode 100644 plans/bookings-demodb-bugfix-plan.md delete mode 100644 plans/dockerfile-improvements.md delete mode 100644 plans/greenplum-pxf-custom-image-plan.md delete mode 100644 plans/stg_layer_implementation_plan.md diff --git a/AGENTS.md b/AGENTS.md index b86b3de..4ffc6c0 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -39,7 +39,9 @@ - В каталоге `sql/` придерживаемся слоёв DWH: - `sql/src/` — скрипты, работающие с исходными системами (например, `bookings_generate_day_if_missing.sql`); - `sql/stg/` — скрипты для стейджинга (`bookings_ddl.sql`, `bookings_load.sql`, `bookings_dq.sql`); - - в будущем можно добавить `sql/ods/`, `sql/dds/`, `sql/dm/` по мере роста стенда. + - `sql/ods/` — скрипты ODS (операционное хранилище); + - `sql/dds/` — скрипты DDS (детальное хранилище, star schema); + - `sql/dm/` — скрипты DM (витрины / data marts). - Нейминг служебных полей и SCD-полей фиксирован в `docs/internal/naming_conventions.md` (единый источник для всех новых слоёв). - Именование файлов: `{объект}_{роль}.sql`, где: - `объект` — логическое имя сущности (`bookings`, `orders`, и т.п.); diff --git a/README.md b/README.md index cc31b5e..07f59ac 100644 --- a/README.md +++ b/README.md @@ -143,11 +143,13 @@ make clean # полный reset: удалить контейнер ## Документация -- Учебные задания: `educational-tasks.md`. -- План тестирования/проверок и негативные кейсы: `TESTING.md`. -- Дополнительные заметки и технические детали: `docs/README.md`. -- Детали по ODS DAG: `docs/bookings_to_gp_ods.md`. -- Детали по DDS DAG: `docs/bookings_to_gp_dds.md`. +- [Учебные задания](educational-tasks.md) +- [План тестирования/проверок и негативные кейсы](TESTING.md) +- [Дополнительные заметки и технические детали](docs/README.md) +- [Детали по STG DAG](docs/bookings_to_gp_stage.md) +- [Детали по ODS DAG](docs/bookings_to_gp_ods.md) +- [Детали по DDS DAG](docs/bookings_to_gp_dds.md) +- [Детали по DM DAG](docs/bookings_to_gp_dm.md) ## Типичные проблемы и решения diff --git a/TESTING.md b/TESTING.md index 80f9480..67747f4 100644 --- a/TESTING.md +++ b/TESTING.md @@ -69,12 +69,10 @@ ## 8. Снятие метрик и мониторинг - Контейнеры: `docker compose ps`, `docker stats` (по желанию). - Логи задач: в Airflow UI → конкретный таск → Log. -- Хостовые CSV: каталог `data/` (можно открыть любой файл и убедиться в структуре). ## 9. Завершение работы - `make down` — выключает сервисы и удаляет контейнеры/сети (volumes сохраняются). - Полный сброс данных (удаляет volumes): `make clean`. -- При необходимости сохранить данные: скопировать CSV из `data/` и сделать дампы до `make clean`. ## Текущий статус (пример успешного прогона) - `uv run pytest -q` — 14 passed, 9 smoke-тестов DAG пропущены (Airflow не установлен в venv). diff --git a/TODO.md b/TODO.md index 45e07d5..f5d8432 100644 --- a/TODO.md +++ b/TODO.md @@ -125,7 +125,7 @@ - патчи `bookings/patches/engine_jobs1_sync.patch` и `bookings/patches/install_drop_if_exists.patch` падают при применении (hunk failed / garbage in patch); - из‑за этого DAG `bookings_to_gp_stage` валится на проверках (источник пустой). - - план: `plans/bookings-demodb-bugfix-plan.md` + - детали: `docs/internal/bookings_db_issues.md` - [x] Добавить раздел «Благодарности» в `README.md`: - явно поблагодарить Postgres Pro за демо‑БД bookings (репозиторий `postgrespro/demodb`); diff --git a/docs/README.md b/docs/README.md index fc12a28..970374c 100644 --- a/docs/README.md +++ b/docs/README.md @@ -10,13 +10,15 @@ - [Главный учебный DAG: bookings → stg](bookings_to_gp_stage.md) - [Учебный DAG: stg -> ods](bookings_to_gp_ods.md) - [Учебный DAG: ods -> dds](bookings_to_gp_dds.md) +- [Учебный DAG: dds -> dm](bookings_to_gp_dm.md) ## Технические детали (опционально) - [Как устроен Docker-стенд (образы, Connections, переменные окружения)](stack.md) - [Единые конвенции нейминга DWH (служебные поля и SCD)](internal/naming_conventions.md) - [PXF в этом проекте (проектная реализация)](internal/pxf_bookings.md) -- [Дизайн stg для bookings (черновик)](internal/bookings_stg_design.md) -- [Дизайн ods для bookings (черновик)](internal/bookings_ods_design.md) -- [Дизайн dds для bookings (черновик)](internal/bookings_dds_design.md) -- [Про время/UTC в bookings (черновик)](internal/bookings_tz.md) +- [Дизайн-документ STG](internal/bookings_stg_design.md) +- [Дизайн-документ ODS](internal/bookings_ods_design.md) +- [Дизайн-документ DDS](internal/bookings_dds_design.md) +- [Дизайн-документ DM](internal/bookings_dm_design.md) +- [Про время/UTC в bookings](internal/bookings_tz.md) diff --git a/docs/bookings_to_gp_dm.md b/docs/bookings_to_gp_dm.md new file mode 100644 index 0000000..07c2c66 --- /dev/null +++ b/docs/bookings_to_gp_dm.md @@ -0,0 +1,79 @@ +# DAG `bookings_to_gp_dm`: `dds` -> `dm` в Greenplum + +Этот DAG — учебный пример загрузки слоя **DM** (Data Mart / витрины) из текущего состояния **DDS**. +Логика: все 5 витрин загружаются параллельно, для каждой — пара `load -> dq`. + +## Что делает DAG + +- Загружает витрины DM параллельно (паттерны загрузки разные — учебная демонстрация выбора стратегии): + - `dm.sales_report` — UPSERT по датам; DQ проверяет только строки текущего `run_id` (`_load_id`); + - `dm.route_performance` — Full Rebuild (TRUNCATE + INSERT): таблица маленькая, дельту считать дороже; + - `dm.passenger_loyalty` — инкрементальный UPSERT по «затронутым ключам» (HWM по `_load_ts`): пересчитываем агрегаты только для пассажиров с новыми фактами; + - `dm.airport_traffic` — инкрементальный UPSERT по датам (HWM по `_load_ts`); + - `dm.monthly_overview` — инкрементальный UPSERT по месяцам (HWM по `_load_ts`). +- Для каждой витрины выполняет пару задач `load -> dq`. + +## Что должно быть готово перед запуском + +1) Стенд поднят: + +```bash +make up +``` + +2) STG, ODS и DDS уже загружены: + +- выполнены DAG-и `bookings_to_gp_stage`, `bookings_to_gp_ods`, `bookings_to_gp_dds`; +- DDL-объекты созданы (`bookings_dm_ddl` или `make ddl-gp`). + +## Как запустить + +1) Откройте Airflow UI: http://localhost:8080. +2) Если запускаете DM впервые — выполните `bookings_dm_ddl`. +3) Запустите `bookings_to_gp_dm`. + +## Граф зависимостей + +Все 5 веток запускаются параллельно от `start_dm`, затем сходятся в `finish_dm_summary`: + +``` +start_dm +├── load_dm_sales_report -> dq_dm_sales_report -> finish_dm_summary +├── load_dm_route_performance -> dq_dm_route_performance -> finish_dm_summary +├── load_dm_passenger_loyalty -> dq_dm_passenger_loyalty -> finish_dm_summary +├── load_dm_airport_traffic -> dq_dm_airport_traffic -> finish_dm_summary +└── load_dm_monthly_overview -> dq_dm_monthly_overview -> finish_dm_summary +``` + +## Как проверить результат + +```bash +make gp-psql +``` + +```sql +SELECT COUNT(*) FROM dm.sales_report; +SELECT COUNT(*) FROM dm.route_performance; +SELECT COUNT(*) FROM dm.passenger_loyalty; +SELECT COUNT(*) FROM dm.airport_traffic; +SELECT COUNT(*) FROM dm.monthly_overview; + +-- Проверка инварианта sales_report: посаженных не больше, чем продано +SELECT COUNT(*) FROM dm.sales_report WHERE tickets_sold < passengers_boarded; + +-- Проверка route_performance: нет дублей по бизнес-ключу +SELECT route_bk, COUNT(*) FROM dm.route_performance GROUP BY route_bk HAVING COUNT(*) > 1; +``` + +Ожидаемо: все витрины непусты, инварианты соблюдены, дублей нет. + +## Типичные ошибки + +- `relation "dm..." does not exist`: + - не применён DM DDL (`bookings_dm_ddl` или `make ddl-gp`). +- DQ падает на `sales_report` по `boarding_rate`: + - проверьте, что DDS загрузился корректно (`dds.fact_flight_sales` непуста). +- DQ падает на `passenger_loyalty` с ошибкой FK: + - проверьте, что `dds.dim_passengers` содержит всех пассажиров из факта. +- `finish_dm_summary` не выполняется: + - одна из DQ-задач упала; найдите в логах Airflow задачу с ошибкой и исправьте. diff --git a/docs/chore/bookings-etl.md b/docs/chore/bookings-etl.md deleted file mode 100644 index e910dcf..0000000 --- a/docs/chore/bookings-etl.md +++ /dev/null @@ -1,466 +0,0 @@ -# План ETL для загрузки `bookings.tickets` в STG слой - -> Архивный документ: это рабочий план, который использовался при разработке. -> Актуальная реализация потока — DAG `bookings_to_gp_stage` и SQL в `sql/stg/`. - -**Ветка:** `chore/bookings-etl` -**Цель:** Добавить загрузку таблицы `tickets` в STG слой Greenplum по аналогии с `bookings` - -## 1. Выбор таблицы и обоснование - -**Выбранная таблица:** `bookings.tickets` - -**Почему `tickets`:** -- Аналитическая ценность: билеты нужны для анализа выручки, загрузки рейсов, пассажиропотока -- Связь с существующим потоком: таблица связана с `bookings.bookings` через `book_ref` -- Инкрементальная природа: билеты создаются вместе с бронированием → понятная логика -- Простая структура: без сложных типов данных (хорошо для обучения) - -## 2. Структура исходной таблицы - -Исходная таблица в `bookings-db` (Postgres): - -| Колонка | Тип | Описание | -|---------|-----|----------| -| `ticket_no` | text (PK) | Уникальный номер билета | -| `book_ref` | text (FK) | Ссылка на бронирование (`bookings.book_ref`) | -| `passenger_id` | text | Идентификатор пассажира | -| `passenger_name` | text | Имя пассажира | -| `outbound` | boolean | Направление рейса (прямой/обратный) | - -**Особенности:** -- Количество записей: примерно в 1.5 раза больше, чем бронирований -- Один `book_ref` может иметь несколько `ticket_no` -- В исходной таблице нет явной временной колонки → используем дату из связанного бронирования - -## 3. Проектирование STG-таблицы в Greenplum - -### Внутренняя таблица `stg.tickets`: - -```sql --- sql/stg/tickets_ddl.sql - -CREATE TABLE IF NOT EXISTS stg.tickets ( - -- Бизнес-атрибуты (из источника) - ticket_no TEXT NOT NULL, -- номер билета - book_ref TEXT NOT NULL, -- ссылка на бронирование - passenger_id TEXT, -- идентификатор пассажира - passenger_name TEXT, -- имя пассажира - outbound TEXT, -- флаг направления (в источнике boolean; в STG храним как TEXT) - - -- Технические атрибуты - src_created_at_ts TIMESTAMP, -- временная метка источника (для инкремента) - load_dttm TIMESTAMP NOT NULL DEFAULT now(), -- время загрузки в Greenplum - batch_id TEXT -- идентификатор батча (из Airflow run_id) -) -WITH (appendonly=true, orientation=row, compresstype=zstd, compresslevel=1) -DISTRIBUTED BY (book_ref); -- распределение по ключу связи с bookings -``` - -**Решения по проектированию:** -- **DISTRIBUTED BY (book_ref):** типовой джойн `tickets → bookings` идёт по `book_ref`, так меньше motion в MPP -- **boolean → TEXT:** в сыром STG храним бизнес-колонки как TEXT (для обучения и минимизации кастов на входе) -- **src_created_at_ts:** временная метка из даты связанного бронирования (см. раздел 7) - -## 4. Проектирование внешней таблицы (PXF) - -Используем существующий PXF-профиль для Postgres (как в `stg.bookings_ext`): - -```sql --- sql/stg/tickets_ddl.sql (продолжение) - -DROP EXTERNAL TABLE IF EXISTS stg.tickets_ext; - -CREATE EXTERNAL TABLE stg.tickets_ext ( - ticket_no TEXT, - book_ref TEXT, - passenger_id TEXT, - passenger_name TEXT, - outbound TEXT -) -LOCATION ('pxf://bookings.tickets?PROFILE=JDBC&SERVER=bookings-db') -FORMAT 'CUSTOM' (formatter='pxfwritable_import'); -``` - -**Решения:** -- Только бизнес-атрибуты во внешней таблице (без тех.колонок) -- PXF-профиль `JDBC` (как в `stg.bookings_ext`) — более стабильный вариант -- PXF-сервер настроен как `bookings-db` в конфигурации - -## 5. DDL-скрипт (создание таблиц) - -Полный файл `sql/stg/tickets_ddl.sql`: - -```sql --- DDL для слоя STG по таблице tickets. --- Используется как из общего скрипта ddl_gp.sql (через \i), --- так и может выполняться отдельно при изменении схемы. - --- Схема stg для сырого слоя DWH. -CREATE SCHEMA IF NOT EXISTS stg; - --- Внешняя таблица в схеме stg для чтения данных из bookings.tickets через PXF. -DROP EXTERNAL TABLE IF EXISTS stg.tickets_ext; - -CREATE EXTERNAL TABLE stg.tickets_ext ( - ticket_no TEXT, - book_ref TEXT, - passenger_id TEXT, - passenger_name TEXT, - outbound TEXT -) -LOCATION ('pxf://bookings.tickets?PROFILE=JDBC&SERVER=bookings-db') -FORMAT 'CUSTOM' (formatter='pxfwritable_import'); - --- Внутренняя таблица stg.tickets — сырой слой, все бизнес-колонки как TEXT. -CREATE TABLE IF NOT EXISTS stg.tickets ( - ticket_no TEXT NOT NULL, - book_ref TEXT NOT NULL, - passenger_id TEXT, - passenger_name TEXT, - outbound TEXT, - src_created_at_ts TIMESTAMP, - load_dttm TIMESTAMP NOT NULL DEFAULT now(), - batch_id TEXT -) -WITH (appendonly=true, orientation=row, compresstype=zstd, compresslevel=1) --- Распределяем по book_ref, чтобы джойны tickets → bookings по book_ref были без motion. -DISTRIBUTED BY (book_ref); - --- На случай, если таблица уже была создана раньше с другим ключом распределения. -ALTER TABLE IF EXISTS stg.tickets SET DISTRIBUTED BY (book_ref); -``` - -## 6. LOAD-скрипт (загрузка инкремента) - -### Логика инкремента - -**Проблема:** В `bookings.tickets` нет явной временной колонки. - -**Решение:** Используем дату бронирования из связанной таблицы `bookings.bookings`: -1. Связь через `tickets.book_ref = bookings.book_ref` -2. Временная колонка: `bookings.book_date` -3. Фильтр инкремента: `book_date > max(src_created_at_ts)` - -### Скрипт `sql/stg/tickets_load.sql`: - -```sql --- Загрузка инкремента из stg.tickets_ext в stg.tickets --- Инкремент определяется по дате бронирования (book_date из bookings.bookings) - -INSERT INTO stg.tickets ( - ticket_no, - book_ref, - passenger_id, - passenger_name, - outbound, - src_created_at_ts, - load_dttm, - batch_id -) -SELECT - ext.ticket_no, - ext.book_ref, - ext.passenger_id, - ext.passenger_name, - ext.outbound, - b.book_date::timestamp, -- временная метка из бронирования - now(), - '{{ run_id }}'::text -FROM stg.tickets_ext AS ext -JOIN stg.bookings_ext AS b ON ext.book_ref = b.book_ref -WHERE b.book_date > COALESCE( - ( - SELECT max(src_created_at_ts) - FROM stg.tickets - WHERE batch_id <> '{{ run_id }}'::text - OR batch_id IS NULL - ), - TIMESTAMP '1900-01-01 00:00:00' -) -AND NOT EXISTS ( - -- Защита от дублей: ticket_no в источнике уникален, и в stg его не дублируем. - SELECT 1 - FROM stg.tickets AS t - WHERE t.ticket_no = ext.ticket_no -); -``` - -**Объяснение логики:** -1. Из `stg.tickets_ext` берём все билеты -2. Прямой JOIN с `stg.bookings_ext` по `book_ref` — это даёт `book_date` из бронирования -3. Фильтр по `book_date > max(src_created_at_ts)` — берём только новые билеты -4. `NOT EXISTS` — защита от повторной загрузки того же билета в текущем батче - -**Важное примечание:** Используем только внешние таблицы (`stg.tickets_ext` и `stg.bookings_ext`), так как прямой доступ к `bookings.bookings` через PXF невозможен. - -## 7. DQ-проверки (качество данных) - -Скрипт `sql/stg/tickets_dq.sql`: - -```sql --- Проверки качества данных для tickets - -DO $$ -DECLARE - v_batch_id TEXT := '{{ run_id }}'::text; - v_prev_ts TIMESTAMP; - v_source_count BIGINT; - v_stg_count BIGINT; - v_orphan_count BIGINT; - v_null_count BIGINT; -BEGIN - -- Опорная метка: максимум src_created_at_ts среди предыдущих батчей - SELECT max(src_created_at_ts) - INTO v_prev_ts - FROM stg.tickets - WHERE batch_id <> v_batch_id - OR batch_id IS NULL; - - -- Источник: считаем строки в том же окне инкремента, что и загрузка - SELECT COUNT(*) - INTO v_source_count - FROM stg.tickets_ext AS t - JOIN stg.bookings_ext AS b ON t.book_ref = b.book_ref - WHERE b.book_date > COALESCE(v_prev_ts, TIMESTAMP '1900-01-01 00:00:00'); - - IF v_source_count = 0 THEN - RAISE EXCEPTION - 'В источнике tickets_ext нет строк для окна инкремента (book_date > %). Проверьте генерацию данных (таск generate_bookings_day).', - COALESCE(v_prev_ts, TIMESTAMP '1900-01-01 00:00:00'); - END IF; - - -- STG: считаем строки текущего батча - SELECT COUNT(*) - INTO v_stg_count - FROM stg.tickets - WHERE batch_id = v_batch_id; - - IF v_source_count <> v_stg_count THEN - RAISE EXCEPTION - 'DQ FAILED: несовпадение количества билетов. Источник: %, STG (batch_id=%): %', - v_source_count, - v_batch_id, - v_stg_count; - END IF; - - -- Ссылочная целостность: tickets должны иметь соответствующие bookings в STG - SELECT COUNT(*) - INTO v_orphan_count - FROM stg.tickets AS t - WHERE t.batch_id = v_batch_id - AND NOT EXISTS ( - SELECT 1 - FROM stg.bookings AS b - WHERE b.book_ref = t.book_ref - ); - - IF v_orphan_count <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: найдены tickets без соответствующих bookings (batch_id=%): %', - v_batch_id, - v_orphan_count; - END IF; - - -- Обязательные поля - SELECT COUNT(*) - INTO v_null_count - FROM stg.tickets AS t - WHERE t.batch_id = v_batch_id - AND (t.ticket_no IS NULL OR t.book_ref IS NULL); - - IF v_null_count <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: найдены tickets с NULL в обязательных полях (batch_id=%): %', - v_batch_id, - v_null_count; - END IF; -END $$; -``` - -## 8. Интеграция с существующими DAG - -### 8.1. Создание DDL через `bookings_stg_ddl.py` - -**Почему расширяем существующий DDL DAG:** -- Уже есть инфраструктура для создания `stg.bookings_ext` и `stg.bookings` -- Единый DAG для создания всех STG-объектов -- Минимальные изменения → проще для новичков - -**Изменения в `airflow/dags/bookings_stg_ddl.py`:** - -Добавить задачу после создания bookings DDL: - -```python -# После существующих задач: - -apply_stg_tickets_ddl = PostgresOperator( - task_id="apply_stg_tickets_ddl", - postgres_conn_id=GREENPLUM_CONN_ID, - sql="stg/tickets_ddl.sql", -) - -# Обновляем связи задач -apply_stg_bookings_ddl >> apply_stg_tickets_ddl -``` - -### 8.2. Загрузка данных через `bookings_to_gp_stage.py` - -**Почему расширяем существующий загрузочный DAG:** -- `generate_bookings_day` уже есть -- Минимальные изменения → проще для новичков -- Единый поток данных (bookings + tickets за один запуск) - -### Изменения в `airflow/dags/bookings_to_gp_stage.py`: - -Добавить задачи после загрузки `bookings`: - -```python -# После существующих задач: - -# Загрузка билетов -load_tickets_to_stg = PostgresOperator( - task_id="load_tickets_to_stg", - postgres_conn_id=GREENPLUM_CONN_ID, - sql="stg/tickets_load.sql", -) - -# DQ-проверки билетов -check_tickets_dq = PostgresOperator( - task_id="check_tickets_dq", - postgres_conn_id=GREENPLUM_CONN_ID, - sql="stg/tickets_dq.sql", -) - -# Обновляем связи задач -check_row_counts >> load_tickets_to_stg >> check_tickets_dq >> finish_summary -``` - -**Фактические изменения в DAG:** - -Обновлён `description` DAG и добавлены две новые задачи: -- `load_tickets_to_stg` — загружает билеты через PXF -- `check_tickets_dq` — проверяет качество данных билетов - -Порядок выполнения: -``` -generate_bookings_day - → load_bookings_to_stg - → check_row_counts - → load_tickets_to_stg - → check_tickets_dq - → finish_summary -``` - -## 9. Тестирование - -### План проверки: - -**1. Подготовка окружения:** -```bash -make up # поднять стенд -make bookings-init # инициализировать демо-БД -``` - -**2. Создание таблиц:** -```bash -# Вариант 1: через Airflow UI (предпочтительно) -# Запустите DAG `bookings_stg_ddl` в Airflow UI - -# Вариант 2: через make-команду -make ddl-gp - -# Вариант 3: напрямую через psql -make gp-psql -``` -```sql --- внутри psql (если выбрали вариант 3): -\i sql/stg/tickets_ddl.sql -\dt stg.* -``` - -**3. Первый запуск DAG:** -- Запустить DAG `bookings_to_gp_stage` в Airflow UI -- Проверить успешность всех задач - -**4. Проверка данных:** -```sql --- Количество билетов -SELECT COUNT(*) FROM stg.tickets; - --- Проверка батчей -SELECT batch_id, COUNT(*) -FROM stg.tickets -GROUP BY batch_id; - --- Выборка данных -SELECT * FROM stg.tickets -ORDER BY load_dttm DESC -LIMIT 10; -``` - -**5. Инкрементальная загрузка:** -```bash -make bookings-generate-day # сгенерировать новый день -``` -- Повторный запуск DAG -- Проверить, что добавились только новые билеты - -**6. Проверка связей с bookings:** -```sql --- Все билеты должны иметь соответствующие бронирования -SELECT COUNT(*) -FROM stg.tickets t -LEFT JOIN stg.bookings b ON t.book_ref = b.book_ref -WHERE b.book_ref IS NULL; --- Ожидаемое значение: 0 -``` - -## 10. Порядок реализации - -1. ✅ Создать файл `sql/stg/tickets_ddl.sql` -2. ✅ Создать файл `sql/stg/tickets_load.sql` -3. ✅ Создать файл `sql/stg/tickets_dq.sql` -4. ✅ Обновить `airflow/dags/bookings_to_gp_stage.py` (добавить задачи tickets) -5. ✅ Обновить `airflow/dags/bookings_stg_ddl.py` (добавить создание tickets DDL) -6. ✅ Обновить `sql/ddl_gp.sql` (подключить tickets_ddl.sql) -7. ✅ Локальное тестирование (раздел 9) -8. ✅ Проверка через Airflow UI -9. ✅ Проверка идемпотентности DDL -10. ✅ Проверка инкрементальной загрузки -11. ✅ Обновить `README.md` (добавить tickets в список STG-таблиц) - -## 11. Сопутствующие изменения - -### Изменения в коде: - -**sql/ddl_gp.sql:** -- ✅ Добавлено подключение `sql/stg/tickets_ddl.sql` - -**educational-tasks.md:** -- ✅ Добавлен раздел 2.3 по анализу структуры `tickets` в STG -- ✅ Добавлен раздел 3.3 по анализу данных `bookings + tickets` - -### Необходимые изменения в документации: - -**README.md:** -- ✅ Добавить `tickets` в список STG-таблиц -- ✅ Обновить описание DAG `bookings_to_gp_stage` (упомянуть загрузку билетов) - -## 12. Результаты тестирования - -### Первый запуск (полная загрузка): -- **Загружено:** 182436 билетов -- **Бронирования:** 45730 -- **Связи:** все билеты имеют соответствующие бронирования (0 orphan tickets) -- **DQ-проверки:** все пройдены успешно - -### Второй запуск (инкрементальная загрузка): -- **Загружено:** новые билеты (количество зависит от сгенерированных данных) -- **Всего в таблице:** сумма всех батчей -- **Уникальность:** все ticket_no уникальные (нет дубликатов) -- **DQ-проверки:** все пройдены успешно - -### Идемпотентность DDL: -- **Первый запуск DDL:** таблицы созданы, данные не затронуты -- **Второй запуск DDL:** данные не пропали, таблицы существуют (CREATE TABLE IF NOT EXISTS) diff --git a/docs/internal/PRD.md b/docs/internal/PRD.md index 6e6efc6..7db5630 100644 --- a/docs/internal/PRD.md +++ b/docs/internal/PRD.md @@ -143,7 +143,7 @@ pgmeta (Postgres 16) ─────────────> Airflow (webserver 1. **Эталонный вертикальный срез** — полностью реализованная цепочка `sales_report` и все её источники вниз по слоям (STG → ODS → DDS → DM). 2. **ТЗ от аналитика** — описание остальных таблиц - ([analyst_spec.md](../assignment/analyst_spec.md)). + (analyst_spec.md — будет создан на Этапе 3, см. TODO.md). 3. **Частично готовый DAG** — студент добавляет свои таски по аналогии. 4. **Валидационный DAG** — студент запускает для самоконтроля. diff --git a/docs/internal/bookings_stg_design.md b/docs/internal/bookings_stg_design.md index 62fe546..7e1750d 100644 --- a/docs/internal/bookings_stg_design.md +++ b/docs/internal/bookings_stg_design.md @@ -1,12 +1,10 @@ -# Дизайн STG для bookings в Greenplum (черновик) - -_Внутренний документ для учебного стенда. Перед итоговой сдачей можно объединить с основной документацией._ +# Дизайн STG для bookings в Greenplum ## 1. Цель и общий контур - Источник: Postgres в контейнере `bookings-db`, база `demo`, таблица `bookings.bookings` (см. `docs/internal/bookings_tz.md`). - Цель: показываем путь данных от операционной БД до сырого слоя DWH в Greenplum. -- В этом документе описываем только часть `src (bookings-db) → STG (Greenplum)`. Слои ODS/DDS/DM студент проектирует сам по статье про моделирование DWH. +- В этом документе описываем часть `src (bookings-db) → STG (Greenplum)`. STG — входной слой; далее данные обрабатываются в ODS → DDS → DM (см. соответствующие design-документы). Примечание: в текущей версии стенда слой `stg` содержит не только `bookings`, но и остальные таблицы потока (`tickets`, `airports`, `airplanes`, `routes`, `seats`, `flights`, `segments`, `boarding_passes`). @@ -16,7 +14,7 @@ _Внутренний документ для учебного стенда. П - `src`: оперативная система (`bookings-db`, схема `bookings`). - `stg`: сырой слой в Greenplum, максимально близкий к источнику, без бизнес‑логики. -- дальше по заданию менти можно строить `ods`/`dds`/`dm` поверх STG. +- далее данные обрабатываются в слоях: `ods` (см. [bookings_ods_design.md](bookings_ods_design.md)), `dds` (см. [bookings_dds_design.md](bookings_dds_design.md)), `dm` (см. [bookings_dm_design.md](bookings_dm_design.md)). ## 2. Схема и таблицы в Greenplum @@ -37,7 +35,7 @@ _Внутренний документ для учебного стенда. П В текущей реализации аналогично созданы внешние таблицы `*_ext` и для остальных сущностей (см. `sql/stg/*_ddl.sql`). -DDL будет добавлен в `sql/ddl_gp.sql` в блоке DDL для Greenplum (примерно по шаблону из `docs/internal/pxf_bookings.md`), с `LOCATION ('pxf://bookings.bookings?PROFILE=JDBC&SERVER=bookings-db')`. +DDL определён в `sql/stg/bookings_ddl.sql` и подключается из `sql/ddl_gp.sql` через `\i` (применяется через `make ddl-gp`). ### 2.3. Внутренняя таблица STG @@ -133,7 +131,7 @@ DDL будет добавлен в `sql/ddl_gp.sql` в блоке DDL для Gre - `docs/internal/pxf_bookings.md` — детали настройки PXF и внешней таблицы для чтения из `bookings-db`. - `sql/stg/bookings_ddl.sql` — DDL для схемы `stg` и таблиц `stg.bookings_ext` / `stg.bookings` (подключается из `sql/ddl_gp.sql` и применяется через `make ddl-gp`). -Дальнейшая модель DWH (слои ODS/DDS/DM, факт/измерения, SCD) должна быть спроектирована студентом по статье о моделировании данных, используя `stg.bookings` как входной слой. +Дальнейшая обработка данных описана в design-документах ODS/DDS/DM (см. раздел 5). ## 6. Требования к читаемости и комментариям diff --git a/plans/bookings-demodb-bugfix-plan.md b/plans/bookings-demodb-bugfix-plan.md deleted file mode 100644 index 0813ac1..0000000 --- a/plans/bookings-demodb-bugfix-plan.md +++ /dev/null @@ -1,100 +0,0 @@ -# План исправления: генератор demodb (bookings) остаётся пустым - -## Контекст - -В стенде используется демобаза bookings из репозитория `postgrespro/demodb`, закреплённая на коммите `866e56f7fe54596a1d2a88f5f32f4aa3b2698121` (см. `DEMODB_COMMIT` в `Makefile`). - -Инициализация источника для ETL выполняется двумя способами: -- `make bookings-init` — быстрое восстановление из seed-дампа (~18 сек, рекомендуется для студентов); -- `make bookings-generate` — полная генерация с нуля (для разработчиков): - - клонирует demodb в `bookings/demodb/`; - - пытается применить патчи из `bookings/patches/`; - - запускает `install.sql` в контейнере `bookings-db`; - - выставляет GUC-параметры (`gen.connstr`, `bookings.start_date/init_days/jobs`); - - запускает `/bookings/generate_next_day.sql` (должен сгенерировать минимум 1 день данных). - -## Симптомы (как в TODO) - -- после `make bookings-generate` таблица `bookings.bookings` остаётся пустой; -- патчи `bookings/patches/engine_jobs1_sync.patch` и `bookings/patches/install_drop_if_exists.patch` падают при применении; -- из‑за этого DAG `bookings_to_gp_stage` валится на проверках (источник пустой). - -## Предварительный диагноз (что уже видно) - -1) `engine_jobs1_sync.patch` не является валидным unified diff (в hunk’ах нет номеров строк вида `@@ -N,M +N,M @@`), поэтому `patch` отвечает: -`patch: **** Only garbage was found in the patch input.` - -2) `install_drop_if_exists.patch` устарел относительно закреплённого коммита demodb: в `install.sql` уже есть `DROP DATABASE IF EXISTS demo;`, поэтому hunk “не находится” и патч не накатывается. - -3) Ошибки патча сейчас замаскированы в `Makefile` через `|| true`, поэтому `make bookings-generate` может завершаться “успешно”, хотя критичные правки в demodb не применились. - -## Цель фикса - -- `make bookings-init` (восстановление из дампа) и `make bookings-generate` (генерация с нуля) воспроизводимо создают и наполняют `demo.bookings.bookings` (>0 строк). -- Если патчи не применяются — процесс останавливается с понятным сообщением, что делать дальше. -- Патчи соответствуют закреплённому коммиту demodb и применяются идемпотентно. - -## План диагностики (чтобы быстро подтвердить проблему) - -1) Чистое воспроизведение: -- `make clean` -- `rm -rf bookings/demodb` -- `make bookings-generate` - -2) Проверка данных: -- `make bookings-psql` -- выполнить: - - `SELECT COUNT(*) FROM bookings.bookings;` - - `SELECT min(book_date), max(book_date) FROM bookings.bookings;` - -3) Проверка патчей (без изменения файлов): -- `patch -d bookings/demodb -p1 --dry-run < bookings/patches/engine_jobs1_sync.patch` -- `patch -d bookings/demodb -p1 --dry-run < bookings/patches/install_drop_if_exists.patch` - -Ожидаемо: сейчас dry-run показывает “garbage in patch” и/или “Hunk FAILED”. - -## План решения - -### Шаг 1. Пересобрать патчи под закреплённый демо‑коммит - -Собираем патчи через `git diff`, чтобы получился корректный unified diff. - -1) `bookings/patches/engine_jobs1_sync.patch`: -- Добавить/подтвердить 2 изменения в `engine.sql`: - - `busy()` игнорирует текущий backend: `AND pid <> pg_backend_pid()`. - - `continue()` при `jobs = 1` выполняет `process_queue(end_date)` синхронно и пишет заметный маркер в лог (`Job 1 (local): ok`), иначе — оставляет текущую логику через `dblink`. - -2) `bookings/patches/install_drop_if_exists.patch`: -- Поменять строку (в актуальном `install.sql`): - - было: `DROP DATABASE IF EXISTS demo;` - - стало: `DROP DATABASE IF EXISTS demo WITH (FORCE);` - -### Шаг 2. Сделать `make bookings-generate` fail-fast на проблемах с патчами - -В `Makefile`: -- убрать `|| true` у применения патчей; -- при ошибке патча — завершать `make` с ненулевым кодом и короткой подсказкой: - - “удалите `bookings/demodb` и повторите `make bookings-generate`”, - - “если не помогло — проверьте, что `DEMODB_COMMIT` не менялся и патчи собраны под него”. - -### Шаг 3. Добавить “защиту от тихого пустого результата” - -После запуска `/bookings/generate_next_day.sql` (в `Makefile` или внутри SQL): -- выполнить проверку `COUNT(*)` по `bookings.bookings`; -- если 0 — завершаться ошибкой с подсказкой, куда смотреть (патчи/логи генератора). - -Цель: чтобы проблема не уезжала дальше в DAG’и и DQ‑проверки, а ловилась сразу при init. - -## Проверка (критерии готовности) - -- `make clean && rm -rf bookings/demodb && make bookings-generate` завершается без ошибок. -- `make bookings-psql` → `SELECT COUNT(*) FROM bookings.bookings;` возвращает `> 0` (для обоих способов: `bookings-init` и `bookings-generate`). -- `make bookings-generate-day` добавляет следующий день: - - `max(book_date)` сдвигается на +1 сутки. -- `./scripts/e2e_smoke.sh` проходит до проверки `stg.bookings` (или хотя бы DAG `bookings_to_gp_stage` перестаёт падать на “источник пустой”). - -## Откат (если нужно быстро вернуть стенд в рабочее состояние) - -- Временно отключить применение патчей в `Makefile` и явно предупреждать, что генерация может быть нестабильной (нежелательно для студентов). -- Или зафиксировать альтернативный `DEMODB_COMMIT`, под который уже готовы патчи (делать только вместе с обновлением документации и проверкой, что генерация стабильна). - diff --git a/plans/dockerfile-improvements.md b/plans/dockerfile-improvements.md deleted file mode 100644 index 5903725..0000000 --- a/plans/dockerfile-improvements.md +++ /dev/null @@ -1,489 +0,0 @@ -# План улучшения Dockerfile и docker-compose.yml - -## Обзор - -Документ описывает план улучшения Dockerfile для Airflow и его интеграции с docker-compose.yml на основе анализа best practices. - -## Согласованные изменения - -✅ Переименовать `Dockerfile` → `Dockerfile.airflow` -✅ Обновить `build: .` → `build: { context: ., dockerfile: Dockerfile.airflow }` -✅ Добавить YAML anchors для устранения дублирования конфигурации -✅ Добавить healthcheck для airflow-webserver -✅ Исправить расположение requirements.txt в Dockerfile -✅ Добавить LABEL в Dockerfile -✅ Добавить `USER airflow` после установки зависимостей -✅ Добавить проверку `pip check` -✅ Переименовать контейнеры (`gp_airflow_web` → `gp_airflow_webserver`, `gp_airflow_sch` → `gp_airflow_scheduler`) - ---- - -## Часть 1: Изменения в Dockerfile - -### Текущее состояние (Dockerfile) -```dockerfile -FROM apache/airflow:2.9.2 - -COPY airflow/requirements.txt /requirements.txt -RUN pip install --no-cache-dir -r /requirements.txt -``` - -### Новое состояние (Dockerfile.airflow) -```dockerfile -# Apache Airflow с дополнительными зависимостями для Greenplum -FROM apache/airflow:2.9.2 - -LABEL maintainer="your-email@example.com" -LABEL description="Airflow with Greenplum and Pandas dependencies" -LABEL version="1.0" - -# Копируем requirements в стандартное расположение -COPY airflow/requirements.txt /opt/airflow/requirements.txt - -# Устанавливаем зависимости -RUN pip install --no-cache-dir -r /opt/airflow/requirements.txt \ - && pip check \ - && rm -rf /tmp/pip-* - -# Переключаемся на пользователя airflow -USER airflow - -WORKDIR /opt/airflow -``` - -### Обоснование изменений: - -1. **LABEL** - стандартная практика для документирования образов -2. **`/opt/airflow/requirements.txt`** - стандартное расположение для Airflow -3. **`pip check`** - проверка совместимости установленных пакетов -4. **`USER airflow`** - безопасность и соответствие best practices -5. **`WORKDIR /opt/airflow`** - явное указание рабочей директории - ---- - -## Часть 2: Изменения в docker-compose.yml - -### Основные изменения: - -1. **Добавить YAML anchors** для устранения дублирования -2. **Обновить build** для всех Airflow сервисов -3. **Добавить healthcheck** для airflow-webserver -4. **Добавить комментарии** для пояснения сокращений в именах контейнеров - -### Структура YAML anchors: - -```yaml -x-airflow-common-env: &airflow-env - TZ: ${TZ:-Europe/Moscow} - AIRFLOW__CORE__LOAD_EXAMPLES: "False" - AIRFLOW__DATABASE__SQL_ALCHEMY_CONN: postgresql+psycopg2://${PG_USER}:${PG_PASSWORD}@pgmeta:5432/${PG_DB} - AIRFLOW__WEBSERVER__SECRET_KEY: ${AIRFLOW__WEBSERVER__SECRET_KEY} - AIRFLOW_CONN_GREENPLUM_CONN: postgresql://${GP_USER}:${GP_PASSWORD}@greenplum:${GP_PORT:-5432}/${GP_DB} - AIRFLOW_CONN_BOOKINGS_DB: postgresql://${BOOKINGS_DB_USER}:${BOOKINGS_DB_PASSWORD}@bookings-db:5432/demo - -x-airflow-common-volumes: &airflow-volumes - - ./airflow/dags:/opt/airflow/dags - - ./sql:/sql:ro - - airflow_data:/opt/airflow/data - -x-airflow-common-depends: &airflow-depends - pgmeta: - condition: service_healthy - greenplum: - condition: service_healthy - airflow-init: - condition: service_completed_successfully -``` - -### Изменения для airflow-webserver: - -```yaml -airflow-webserver: - build: - context: . - dockerfile: Dockerfile.airflow - image: airflow-custom:latest - container_name: gp_airflow_webserver - env_file: .env - environment: - <<: *airflow-env - command: airflow webserver - ports: - - "8080:8080" - volumes: - <<: *airflow-volumes - # requirements.txt монтируется как volume для удобства разработки - # При изменении зависимостей не требуется пересборка образа - - ./airflow/requirements.txt:/opt/airflow/requirements.txt - healthcheck: - test: ["CMD", "curl", "-f", "http://localhost:8080/health"] - interval: 30s - timeout: 10s - retries: 5 - start_period: 40s - depends_on: - <<: *airflow-depends -``` - -### Изменения для airflow-scheduler: - -```yaml -airflow-scheduler: - build: - context: . - dockerfile: Dockerfile.airflow - image: airflow-custom:latest - container_name: gp_airflow_scheduler - env_file: .env - environment: - <<: *airflow-env - command: airflow scheduler - volumes: - <<: *airflow-volumes - - ./airflow/requirements.txt:/opt/airflow/requirements.txt - depends_on: - <<: *airflow-depends -``` - -### Изменения для airflow-init: - -```yaml -airflow-init: - build: - context: . - dockerfile: Dockerfile.airflow - image: airflow-custom:latest - # Имя контейнера закомментировано (одноразовый сервис) - user: "0" - env_file: .env - environment: - <<: *airflow-env - volumes: - <<: *airflow-volumes - command: > - bash -lc " - set -e; - mkdir -p /opt/airflow/data && chown -R airflow:root /opt/airflow/data; - for i in {1..30}; do - su -s /bin/bash airflow -c \"PATH='/home/airflow/.local/bin:$${PATH}' airflow db migrate\" && break || echo 'waiting for pgmeta' && sleep 3; - done; - su -s /bin/bash airflow -c \"PATH='/home/airflow/.local/bin:$${PATH}' airflow users create --username ${AIRFLOW_USER} --password ${AIRFLOW_PASSWORD} --firstname Admin --lastname User --role Admin --email admin@example.org\" || true - " - depends_on: - pgmeta: - condition: service_healthy -``` - ---- - -## Часть 3: Влияние переименования контейнеров - -### Важно: Переименование контейнеров - -При переименовании контейнеров (`gp_airflow_web` → `gp_airflow_webserver`, `gp_airflow_sch` → `gp_airflow_scheduler`) необходимо учитывать: - -1. **Имена сервисов в docker-compose.yml остаются без изменений** - - `airflow-webserver`, `airflow-scheduler`, `airflow-init` - это имена сервисов - - Они используются в командах `docker-compose logs`, `docker-compose restart` и т.д. - - Эти имена НЕ меняются - -2. **Имена контейнеров меняются** - - `gp_airflow_web` → `gp_airflow_webserver` - - `gp_airflow_sch` → `gp_airflow_scheduler` - - Они используются в командах `docker exec`, `docker inspect`, `docker logs` - -3. **Проверка использования старых имён контейнеров** - ```bash - # Поиск упоминаний старых имён в скриптах - grep -r "gp_airflow_web" . - grep -r "gp_airflow_sch" . - ``` - -4. **Обновление Makefile (если используется)** - - Проверить команды, которые используют старые имена контейнеров - - Пример: `make logs` может использовать `docker-compose logs airflow-webserver` (это корректно) - -5. **Обновление документации** - - Обновить все упоминания имён контейнеров в README.md, TESTING.md - - Обновить примеры команд в документации - ---- - -## Часть 4: План тестирования изменений - -### Этап 1: Подготовка окружения - -1. **Остановить текущие контейнеры** - ```bash - make down - ``` - -2. **Удалить старые образы (опционально)** - ```bash - docker rmi airflow-custom:latest - ``` - -3. **Проверить наличие .env файла** - ```bash - ls -la .env - # Если отсутствует, скопировать из .env.example - cp .env.example .env - ``` - -### Этап 2: Сборка новых образов - -1. **Собрать образы с новым Dockerfile** - ```bash - docker-compose build - ``` - -2. **Проверить успешность сборки** - ```bash - docker images | grep airflow-custom - ``` - -3. **Проверить наличие LABEL** - ```bash - docker inspect airflow-custom:latest | grep -A 10 "Labels" - ``` - -### Этап 3: Запуск сервисов - -1. **Запустить весь стек** - ```bash - make up - ``` - -2. **Проверить статус контейнеров** - ```bash - docker-compose ps - ``` - -3. **Проверить логи инициализации** - ```bash - docker-compose logs airflow-init - ``` - -### Этап 4: Проверка healthcheck - -1. **Проверить healthcheck для airflow-webserver** - ```bash - docker inspect gp_airflow_webserver | grep -A 20 "Health" - ``` - -2. **Дождаться healthy статуса** - ```bash - watch -n 2 'docker inspect --format="{{.State.Health.Status}}" gp_airflow_webserver' - ``` - -### Этап 5: Функциональное тестирование - -1. **Проверить доступ к Airflow UI** - - Открыть http://localhost:8080 - - Проверить авторизацию (использовать креды из .env) - - Проверить наличие DAG'ов в списке - -2. **Проверить запуск DAG'ов** - - Выбрать любой DAG (например, `bookings_to_gp_stage`) - - Запустить вручную через UI - - Проверить успешность выполнения задач - -3. **Проверить подключения к БД** - - Проверить Admin → Connections - - Убедиться, что `greenplum_conn` и `bookings_db` доступны - - Проверить Test Connection для каждого подключения - -4. **Проверить логи scheduler** - ```bash - docker-compose logs airflow-scheduler | tail -50 - ``` - -### Этап 6: Проверка использования старых имён контейнеров - -1. **Поиск упоминаний старых имён в проекте** - ```bash - grep -r "gp_airflow_web" . --exclude-dir=.git --exclude-dir=__pycache__ - grep -r "gp_airflow_sch" . --exclude-dir=.git --exclude-dir=__pycache__ - ``` - -2. **Проверка Makefile** - ```bash - grep -E "gp_airflow_web|gp_airflow_sch" Makefile - ``` - -3. **Проверка документации** - ```bash - grep -r "gp_airflow_web" README.md TESTING.md AGENTS.md docs/ - grep -r "gp_airflow_sch" README.md TESTING.md AGENTS.md docs/ - ``` - -4. **Обновление найденных упоминаний** - - Заменить `gp_airflow_web` на `gp_airflow_webserver` - - Заменить `gp_airflow_sch` на `gp_airflow_scheduler` - - Обновить примеры команд в документации - -### Этап 7: Тестирование требований - -1. **Проверить установленные пакеты** - ```bash - docker exec gp_airflow_webserver pip list | grep -E "psycopg2|pandas" - ``` - -2. **Проверить совместимость пакетов** - ```bash - docker exec gp_airflow_webserver pip check - ``` - -3. **Проверить пользователя внутри контейнера** - ```bash - docker exec gp_airflow_webserver whoami - # Ожидаемый результат: airflow - ``` - -### Этап 8: Тестирование изменений requirements.txt - -1. **Добавить тестовый пакет в requirements.txt** - ```bash - echo "requests==2.31.0" >> airflow/requirements.txt - ``` - -2. **Перезапустить контейнеры** - ```bash - docker-compose restart airflow-webserver airflow-scheduler - ``` - -3. **Проверить установку пакета** - ```bash - docker exec gp_airflow_webserver pip list | grep requests - ``` - -4. **Удалить тестовый пакет** - ```bash - # Вернуть исходный requirements.txt - git checkout airflow/requirements.txt - ``` - -### Этап 9: Регрессионное тестирование - -1. **Запустить существующие тесты** - ```bash - make test - ``` - -2. **Проверить smoke-тесты DAG** - ```bash - uv run pytest tests/test_dags_smoke.py -v - ``` - -3. **Проверить тесты helpers** - ```bash - uv run pytest tests/test_greenplum_helpers.py -v - ``` - -### Этап 10: Проверка после перезапуска - -1. **Полный перезапуск стека** - ```bash - make down - make up - ``` - -2. **Проверить сохранность данных** - - Проверить наличие DAG'ов в UI - - Проверить историю запусков DAG'ов - - Проверить подключения к БД - -3. **Проверить логи на наличие ошибок** - ```bash - docker-compose logs airflow-webserver | grep -i error - docker-compose logs airflow-scheduler | grep -i error - # Обратите внимание: имена сервисов не изменились, только имена контейнеров - ``` - ---- - -## Часть 5: Критерии успеха - -### Функциональные требования: - -- ✅ Все контейнеры успешно запускаются -- ✅ Airflow UI доступен на http://localhost:8080 -- ✅ Авторизация работает корректно -- ✅ Все DAG'ы отображаются в UI -- ✅ DAG'и успешно выполняются -- ✅ Подключения к Greenplum и bookings-db работают -- ✅ Healthcheck для airflow-webserver работает корректно - -### Технические требования: - -- ✅ Образ собирается без ошибок -- ✅ LABEL присутствуют в образе -- ✅ Пакеты устанавливаются от пользователя airflow -- ✅ `pip check` не возвращает ошибок -- ✅ requirements.txt монтируется как volume -- ✅ YAML anchors работают корректно -- ✅ Нет дублирования конфигурации - -### Требования к совместимости: - -- ✅ Существующие тесты проходят успешно -- ✅ Данные сохраняются после перезапуска -- ✅ История запусков DAG'ов сохраняется -- ✅ Подключения к БД работают как раньше - ---- - -## Часть 6: Откат изменений - -Если изменения вызывают проблемы, план отката: - -1. **Остановить контейнеры** - ```bash - make down - ``` - -2. **Вернуть старые файлы** - ```bash - git checkout Dockerfile docker-compose.yml - ``` - -3. **Удалить новый образ** - ```bash - docker rmi airflow-custom:latest - ``` - -4. **Перезапустить стек** - ```bash - make up - ``` - ---- - -## Часть 7: Документация - -После успешного внедрения изменений необходимо обновить: - -1. **README.md** - добавить информацию о новом Dockerfile -2. **AGENTS.md** - обновить инструкции для агентов -3. **Комментарии в docker-compose.yml** - добавить пояснения к YAML anchors - ---- - -## Резюме - -План включает: -- 2 файла для изменения: `Dockerfile` → `Dockerfile.airflow`, `docker-compose.yml` -- 10 этапов тестирования -- 12 функциональных критериев успеха -- План отката на случай проблем -- Переименование контейнеров для улучшения читаемости -- Проверка использования старых имён контейнеров в проекте - -### Важные замечания: - -1. **Имена сервисов НЕ меняются** - `airflow-webserver`, `airflow-scheduler`, `airflow-init` -2. **Имена контейнеров меняются** - `gp_airflow_web` → `gp_airflow_webserver`, `gp_airflow_sch` → `gp_airflow_scheduler` -3. **Необходимо проверить** использование старых имён контейнеров в скриптах и документации -4. **Makefile команды** используют имена сервисов, поэтому они продолжат работать без изменений - -Все изменения следуют best practices для Docker и Airflow, сохраняют обратную совместимость и не требуют изменений в DAG'ах. diff --git a/plans/greenplum-pxf-custom-image-plan.md b/plans/greenplum-pxf-custom-image-plan.md deleted file mode 100644 index 10ce36a..0000000 --- a/plans/greenplum-pxf-custom-image-plan.md +++ /dev/null @@ -1,142 +0,0 @@ -# План работ: свой образ Greenplum с интегрированным PXF - -## Контекст и проблема - -Иногда (и у нас воспроизводится стабильно) контейнер `greenplum` падает при старте с ошибкой: - -`chown: changing ownership of '/docker-entrypoint-initdb.d/10_pxf_bookings.sh': Read-only file system` - -Причина: init-скрипт `pxf/init/10_pxf_bookings.sh` проброшен в контейнер как bind-mount `:ro`, а entrypoint базового образа пытается сделать `chown` файлов в `/docker-entrypoint-initdb.d/`. - -## Цели - -- Greenplum стабильно стартует после любого `restart/up` (без “иногда не стартует”). -- PXF готов к работе **после каждого запуска** контейнера. -- Не меняем права/владельца файлов в репозитории на хосте (никаких “файл стал root’ом / не редактируется”). -- Для студентов всё остаётся простым: `make build` (явная сборка) + `make up` (поднимает стенд; если образа нет — Docker Compose соберёт сам). - -## Выбранный подход (высокоуровневый дизайн) - -1) Делаем свой образ Greenplum на базе `woblerr/greenplum:6.27.1`. - -2) Встраиваем в образ “seed” для PXF: - - JDBC-драйвер (JAR), - - конфиг сервера `bookings-db` (`jdbc-site.xml`), - - скрипт “ensure”, который идемпотентно гарантирует, что файлы лежат в `PXF_BASE` (на persistent volume). - -3) Запускаем “ensure” **на каждом старте контейнера** через wrapper-entrypoint, а затем передаём управление оригинальному entrypoint базового образа. - -4) `PXF_BASE` по умолчанию остаётся на volume (`/data/pxf`), чтобы настройки переживали рестарты. - -## План работ (по шагам) - -### Шаг 1. Разведка базового образа - -- Проверить, где находится оригинальный entrypoint и как он запускается (путь, параметры, пользователь). -- Понять, как `GREENPLUM_PXF_ENABLE=true` влияет на старт (чтобы wrapper не ломал поведение). - -Результат: фиксируем в README “как устроен старт” (1–2 абзаца). - -### Шаг 2. Новый Dockerfile для Greenplum - -- Добавить `Dockerfile.greenplum`: - - `FROM woblerr/greenplum:6.27.1` - - `COPY` seed-артефакты в образ (например, в `/opt/pxf-seed/...`) - - `COPY` wrapper-entrypoint в образ - - настроить права/владельца внутри образа так, чтобы старт был без ошибок - -Результат: образ собирается локально через `make build` и/или автоматически через `make up`. - -### Шаг 3. Wrapper-entrypoint (каждый старт) - -- Добавить скрипт entrypoint-обёртки (например, `greenplum/entrypoint-wrapper.sh` или `pxf/entrypoint-wrapper.sh`): - - на старте вызывает `ensure`-скрипт; - - затем делает `exec` оригинального entrypoint базового образа с теми же аргументами. - -Важно: wrapper не должен “перехватывать” логику инициализации кластера — только добавлять шаг подготовки PXF. - -### Шаг 4. Переписать текущий init-скрипт в “ensure” (идемпотентный) - -- Превратить `pxf/init/10_pxf_bookings.sh` в скрипт, который можно безопасно выполнять на каждом запуске: - - не опираться на `~/.bashrc`; - - `PXF_BASE` вычислять через env (`PXF_BASE`, `GREENPLUM_DATA_DIRECTORY`, fallback `/data/pxf`); - - seed-копирование делать “если файла нет”; - - добавить понятные логи (что сделано / что пропущено); - - `pxf cluster sync`: - - выполнять, только если команда доступна, - - не валить контейнер при ошибке (но писать предупреждение). - -Дополнительно (опционально, но полезно для стенда): -- env-переключатель `PXF_SEED_OVERWRITE=1` — принудительно перезаписывать конфиг из образа в volume (для обновлений без удаления volume). - -### Шаг 5. Обновить `docker-compose.yml` - -- Для сервиса `greenplum` перейти на `build:` (и при желании оставить `image:` как тег). -- Убрать bind-mount’ы PXF (jar/config/init-скрипт), т.к. теперь всё в образе. -- Оставить `greenplum_data:/data` и `./sql:/sql:ro`. -- Исправить healthcheck Greenplum (сейчас конструкция вида `... || echo 1` делает healthcheck “вечно успешным”): - - healthcheck должен возвращать ненулевой код, если БД не готова; - - добавить `start_period`, чтобы не ловить ложные падения на холодном старте. - -### Шаг 6. Обновить Makefile и README - -- `Makefile`: - - убедиться, что `make build` собирает также Greenplum-образ (если введём build для сервиса); - - оставить `make up` как есть (Compose сам соберёт образ, если его нет). -- `README.md`: - - зачем свой образ (устойчивость, права на хосте, меньше mount’ов), - - как пересобрать образ, - - как обновить PXF-конфиг (через rebuild + `PXF_SEED_OVERWRITE=1` или через очистку volume). - -## Проверка и критерии готовности - -- `docker compose up -d greenplum` → контейнер остаётся `Up`, не падает. -- `docker compose restart greenplum` повторить 10–20 раз → без падений. -- Healthcheck Greenplum становится `healthy` (не “вечно healthy” и не “вечно starting”). -- Поднятие всего стека (`make up`) приводит к старту Airflow (scheduler/webserver), т.к. `depends_on: condition: service_healthy` начинает работать корректно. - -## Риски и как их снизить - -- **Стартап может замедлиться**, если `pxf cluster sync` делать каждый раз: поэтому скрипт должен быть быстрым, а sync — не фатальным при ошибках. -- **Обновление конфигов**: так как PXF_BASE на volume, изменения в образе сами не перетрут файлы — поэтому нужен `PXF_SEED_OVERWRITE=1` или понятная инструкция “как обновить”. - -## Откат - -- Вернуться к использованию `image: woblerr/greenplum:6.27.1` в `docker-compose.yml`. -- Вернуть mount’ы PXF, если нужно (но это вернёт риск с `:ro`). - ---- - -## Статус (реализовано) - -- Добавлен кастомный образ Greenplum: `Dockerfile.greenplum` (seed PXF + startup wrapper). -- Ensure‑скрипт перенесён в образ и стал идемпотентным: `pxf/init/10_pxf_bookings.sh`. -- Стартовый скрипт контейнера включает ensure и создаёт `EXTENSION pxf`: `pxf/init/start_greenplum_with_pxf.sh`. -- В `docker-compose.yml`: - - `greenplum` собирается через `build: Dockerfile.greenplum`; - - добавлен `hostname: gpdbsne`; - - убраны PXF bind-mount’ы (jar/config/init); - - healthcheck ждёт не только GPDB, но и готовность PXF. -- Обновлены инструкции: `README.md`, `.env.example`. -- Проблема с генератором demodb (пустая `bookings.bookings`) зафиксирована в `TODO.md`. - -## Проверка (как воспроизвести) - -Команды для ручной проверки: - -- Пересобрать и перезапустить Greenplum: - - `make build` - - `docker compose up -d --force-recreate greenplum` -- 3–10 рестартов: - - `docker compose restart greenplum` - - дождаться `healthy` в `docker compose ps` -- Проверить PXF: - - `docker compose exec greenplum bash -lc "su - gpadmin -c '/usr/local/pxf/bin/pxf cluster status'"` -- Проверить DAG, который использует PXF: - - `docker exec -i gp_airflow_scheduler airflow dags test bookings_stg_ddl` - -Текущий результат: - -- `bookings_stg_ddl` проходит (PXF и `protocol pxf` доступны). -- `bookings_to_gp_stage` падает не из-за PXF, а из-за пустого источника - (`demo.bookings.bookings` = 0 строк). Это отдельная задача (см. `TODO.md`). diff --git a/plans/stg_layer_implementation_plan.md b/plans/stg_layer_implementation_plan.md deleted file mode 100644 index 518d4ad..0000000 --- a/plans/stg_layer_implementation_plan.md +++ /dev/null @@ -1,744 +0,0 @@ -# План реализации STG слоя целиком - -> **Статус:** План готов к реализации -> **Дата:** 2026-01-17 -> **Автор:** Architect Mode - -## Обзор задачи - -Согласно [`docs/internal/db_schema.md`](../docs/internal/db_schema.md), STG слой реализован частично (2 из 9 таблиц: bookings, tickets). Необходимо реализовать оставшиеся 7 таблиц. - -## Стратегия загрузки данных - -| Тип таблиц | Стратегия | Обоснование | -|-----------|-----------|-------------| -| Справочники (airports, airplanes, routes, seats) | **Full load** | Маленький объём (<10K строк), простота реализации | -| Транзакции (flights, segments) | **Инкремент** | Больший объём; выбираем максимально естественное опорное поле | -| Транзакции (boarding_passes) | **Full snapshot** | В источнике строки создаются и обновляются со временем, простого инкремента без усложнений нет | - -> Важно: под **Full load** в STG подразумеваем «сняли слепок и дописали в append-only таблицу с `batch_id`», а не `TRUNCATE + INSERT`. Это даёт простую идемпотентность (по `batch_id`) и сохраняет историю загрузок. - -### Оценка размера справочников - -| Справочник | Примерный размер | Оценка | -|-------------|------------------|---------| -| **airports** | ~5K-6K аэропортов | **Маленький** | -| **airplanes** | ~10 моделей самолётов | **Крошечный** | -| **seats** | ~1700-2000 записей | **Маленький** | -| **routes** | Ожидается несколько тысяч | **Маленький/Средний** | - -**Вывод:** Все справочники очень маленькие (до нескольких тысяч строк). Даже если routes будет 5000-10000 строк - это всё равно минимальный объём для Greenplum. - -## Список таблиц для реализации - -| Таблица источника | Таблица STG | Тип данных | Стратегия загрузки | Опорное поле для инкремента | -|-------------------|-------------|------------|-------------------|---------------------------| -| `bookings.airports_data` | `stg.airports` | Справочник | Full | - | -| `bookings.airplanes_data` | `stg.airplanes` | Справочник | Full | - | -| `bookings.routes` | `stg.routes` | Справочник | Full | - | -| `bookings.seats` | `stg.seats` | Справочник | Full | - | -| `bookings.flights` | `stg.flights` | Транзакции | Инкремент | `scheduled_departure` | -| `bookings.segments` | `stg.segments` | Транзакции | Инкремент | `book_date` (через tickets) | -| `bookings.boarding_passes` | `stg.boarding_passes` | Транзакции | Full (snapshot) | - | - -## Паттерн реализации (на основе bookings/tickets) - -Для каждой таблицы создаются 3 файла: - -1. **`sql/stg/{table}_ddl.sql`** - DDL для внешней и внутренней таблиц -2. **`sql/stg/{table}_load.sql`** - Загрузка (full или инкремент) -3. **`sql/stg/{table}_dq.sql`** - Проверки качества данных - -### Общая структура DDL файла - -```sql --- DDL для слоя STG по таблице {table}. --- Используется как из общего скрипта ddl_gp.sql (через \i), --- так и может выполняться отдельно при изменении схемы. - --- Схема stg для сырого слоя DWH. -CREATE SCHEMA IF NOT EXISTS stg; - --- Внешняя таблица в схеме stg для чтения данных из bookings.{table} через PXF. -DROP EXTERNAL TABLE IF EXISTS stg.{table}_ext; -CREATE EXTERNAL TABLE stg.{table}_ext ( - -- поля из источника -) -LOCATION ('pxf://bookings.{table}?PROFILE=JDBC&SERVER=bookings-db') -FORMAT 'CUSTOM' (formatter='pxfwritable_import'); - --- Внутренняя таблица stg.{table} — сырой слой, все бизнес-колонки как TEXT. -CREATE TABLE IF NOT EXISTS stg.{table} ( - -- бизнес-колонки как TEXT - src_created_at_ts TIMESTAMP, - load_dttm TIMESTAMP NOT NULL DEFAULT now(), - batch_id TEXT -) -WITH (appendonly=true, orientation=row, compresstype=zstd, compresslevel=1) -DISTRIBUTED BY ({distribution_key}); -``` - -Примечание по PXF/JDBC типам: для «сложных» типов Postgres (например, `jsonb`, `point`, массивы, `tstzrange`) -чаще всего проще и надёжнее объявлять колонки во внешней таблице как `TEXT`, чтобы избежать несовместимостей -драйвера/маппинга типов. Внутренний STG всё равно хранит бизнес-поля как `TEXT`. - -### Общая структура LOAD файла (Full load для справочников) - -```sql --- Загрузка всех строк из stg.{table}_ext в stg.{table}. --- Используем batch_id для отслеживания загрузки. - -INSERT INTO stg.{table} ( - -- бизнес-колонки - src_created_at_ts, - load_dttm, - batch_id -) -SELECT - ext.{field}::text, - now()::timestamp, - '{{ run_id }}'::text -FROM stg.{table}_ext AS ext -WHERE NOT EXISTS ( - -- Защита от дублей в рамках одного batch_id - SELECT 1 - FROM stg.{table} AS t - WHERE t.batch_id = '{{ run_id }}'::text - AND t.{pk} = ext.{pk}::text -); - --- Обновляем статистику для оптимизатора Greenplum -ANALYZE stg.{table}; -``` - -### Общая структура LOAD файла (Инкремент для транзакций) - -```sql --- Загрузка инкремента из stg.{table}_ext в stg.{table}. --- Окно инкремента определяется по src_created_at_ts: --- берём строки, где {increment_field} больше максимального src_created_at_ts --- среди "старых" батчей; верхняя граница по дате не используется. - --- CTE для определения максимальной даты загрузки предыдущего батча -WITH max_batch_ts AS ( - SELECT COALESCE(MAX(src_created_at_ts), TIMESTAMP '1900-01-01 00:00:00') AS max_ts - FROM stg.{table} - WHERE batch_id <> '{{ run_id }}'::text - OR batch_id IS NULL -) -INSERT INTO stg.{table} ( - -- бизнес-колонки - src_created_at_ts, - load_dttm, - batch_id -) -SELECT - ext.{field}::text, - ext.{increment_field}::timestamp, - now(), - '{{ run_id }}'::text -FROM stg.{table}_ext AS ext -CROSS JOIN max_batch_ts AS mb -WHERE ext.{increment_field} > mb.max_ts -AND NOT EXISTS ( - SELECT 1 - FROM stg.{table} AS t - WHERE t.batch_id = '{{ run_id }}'::text - AND t.{pk} = ext.{pk}::text -); - --- Обновляем статистику для оптимизатора Greenplum -ANALYZE stg.{table}; -``` - -### Общая структура DQ файла - -Для **инкрементальных** таблиц сравниваем окно инкремента (по `src_created_at_ts`) между источником и STG. -Для **full snapshot** таблиц (справочники и `boarding_passes`) обычно достаточно сравнить общее количество строк в источнике с количеством строк, загруженных в текущий `batch_id`, плюс проверить дубликаты/NULL/ссылочную целостность. - -```sql --- Проверки качества данных для {table} - -DO $$ -DECLARE - v_batch_id TEXT := '{{ run_id }}'::text; - v_prev_ts TIMESTAMP; - v_src_count BIGINT; - v_stg_count BIGINT; - v_dup_count BIGINT; - v_null_count BIGINT; - -- другие переменные для специфических проверок -BEGIN - -- Опорная метка: максимум src_created_at_ts среди предыдущих батчей - SELECT max(src_created_at_ts) - INTO v_prev_ts - FROM stg.{table} - WHERE batch_id <> v_batch_id - OR batch_id IS NULL; - - -- Источник: считаем строки во внешней таблице, которые вошли в окно инкремента - SELECT COUNT(*) - INTO v_src_count - FROM stg.{table}_ext - WHERE {increment_field} > COALESCE(v_prev_ts, TIMESTAMP '1900-01-01 00:00:00'); - - IF v_src_count = 0 THEN - RAISE EXCEPTION - 'В источнике {table}_ext нет строк для окна инкремента.'; - END IF; - - -- Считаем строки, реально вставленные в stg.{table} в этом батче - SELECT COUNT(*) - INTO v_stg_count - FROM stg.{table} - WHERE batch_id = v_batch_id; - - IF v_src_count <> v_stg_count THEN - RAISE EXCEPTION - 'DQ FAILED: несовпадение количества строк. Источник: %, STG: %', - v_src_count, - v_stg_count; - END IF; - - -- Проверка на дубликаты первичного ключа - SELECT COUNT(*) - COUNT(DISTINCT {pk}) - INTO v_dup_count - FROM stg.{table} AS t - WHERE t.batch_id = v_batch_id; - - IF v_dup_count <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: найдены дубликаты {pk} (batch_id=%): %', - v_batch_id, - v_dup_count; - END IF; - - -- Проверка обязательных полей - SELECT COUNT(*) - INTO v_null_count - FROM stg.{table} AS t - WHERE t.batch_id = v_batch_id - AND (t.{required_field} IS NULL OR t.{required_field} = ''); - - IF v_null_count <> 0 THEN - RAISE EXCEPTION - 'DQ FAILED: найдены строки с NULL в обязательных полях (batch_id=%): %', - v_batch_id, - v_null_count; - END IF; - - RAISE NOTICE - 'DQ PASSED: {table} ок (batch_id=%): source=% stg=%', - v_batch_id, - v_src_count, - v_stg_count; -END $$; -``` - -## Детали реализации по таблицам - -### 1. airports (справочник, full load) - -**Внешняя таблица**: `stg.airports_ext` -- Поля: `airport_code`, `airport_name` (JSONB), `city` (JSONB), `country` (JSONB), `coordinates`, `timezone` -- PXF: `pxf://bookings.airports_data?PROFILE=JDBC&SERVER=bookings-db` - -**Внутренняя таблица**: `stg.airports` -- Бизнес-колонки как TEXT: - - `airport_code TEXT` - - `airport_name TEXT` - - `city TEXT` - - `country TEXT` - - `coordinates TEXT` - - `timezone TEXT` -- Тех.колонки: `src_created_at_ts`, `load_dttm`, `batch_id` -- Распределение: `DISTRIBUTED BY (airport_code)` -- Обоснование: `airport_code` — это уникальный идентификатор аэропорта - -**Загрузка**: Full (все строки при каждом запуске) - -**DQ проверки**: -- Count между источником и STG -- Дубликаты `airport_code` -- NULL обязательных полей (airport_code, airport_name, city, timezone) - -### 2. airplanes (справочник, full load) - -**Внешняя таблица**: `stg.airplanes_ext` -- Поля: `airplane_code`, `model` (JSONB), `range`, `speed` -- PXF: `pxf://bookings.airplanes_data?PROFILE=JDBC&SERVER=bookings-db` - -**Внутренняя таблица**: `stg.airplanes` -- Бизнес-колонки как TEXT: - - `airplane_code TEXT` - - `model TEXT` - - `range TEXT` - - `speed TEXT` -- Тех.колонки: `src_created_at_ts`, `load_dttm`, `batch_id` -- Распределение: `DISTRIBUTED BY (airplane_code)` -- Обоснование: `airplane_code` — это уникальный идентификатор самолёта - -**Загрузка**: Full - -**DQ проверки**: -- Count между источником и STG -- Дубликаты `airplane_code` -- NULL обязательных полей (airplane_code, model) - -### 3. routes (справочник, full load) - -**Внешняя таблица**: `stg.routes_ext` -- Поля: `route_no`, `validity` (tstzrange), `departure_airport`, `arrival_airport`, `airplane_code`, `days_of_week` (int[]), `scheduled_time`, `duration` -- PXF: `pxf://bookings.routes?PROFILE=JDBC&SERVER=bookings-db` - -**Внутренняя таблица**: `stg.routes` -- Бизнес-колонки как TEXT: - - `route_no TEXT` - - `validity TEXT` - - `departure_airport TEXT` - - `arrival_airport TEXT` - - `airplane_code TEXT` - - `days_of_week TEXT` - - `scheduled_time TEXT` - - `duration TEXT` -- Тех.колонки: `src_created_at_ts`, `load_dttm`, `batch_id` -- Распределение: `DISTRIBUTED BY (route_no)` -- Обоснование: `route_no` — логический идентификатор маршрута; он нужен для JOIN с flights по `route_no` - -**Загрузка**: Full - -**DQ проверки**: -- Count между источником и STG -- Дубликаты `(route_no, validity)` -- NULL обязательных полей (route_no, departure_airport, arrival_airport, airplane_code) -- Ссылочная целостность на airports (departure_airport, arrival_airport) -- Ссылочная целостность на airplanes (airplane_code) - -### 4. seats (справочник, full load) - -**Внешняя таблица**: `stg.seats_ext` -- Поля: `airplane_code`, `seat_no`, `fare_conditions` -- PXF: `pxf://bookings.seats?PROFILE=JDBC&SERVER=bookings-db` - -**Внутренняя таблица**: `stg.seats` -- Бизнес-колонки как TEXT: - - `airplane_code TEXT` - - `seat_no TEXT` - - `fare_conditions TEXT` -- Тех.колонки: `src_created_at_ts`, `load_dttm`, `batch_id` -- Распределение: `DISTRIBUTED BY (airplane_code)` -- Обоснование: co-location с airplanes для оптимизации JOIN - -**Загрузка**: Full - -**DQ проверки**: -- Count между источником и STG -- Дубликаты `(airplane_code, seat_no)` -- NULL обязательных полей (airplane_code, seat_no, fare_conditions) -- Ссылочная целостность на airplanes (airplane_code) - -### 5. flights (транзакции, инкремент) - -**Внешняя таблица**: `stg.flights_ext` -- Поля: `flight_id`, `route_no`, `status`, `scheduled_departure`, `scheduled_arrival`, `actual_departure`, `actual_arrival` -- PXF: `pxf://bookings.flights?PROFILE=JDBC&SERVER=bookings-db` - -**Внутренняя таблица**: `stg.flights` -- Бизнес-колонки как TEXT: - - `flight_id TEXT` - - `route_no TEXT` - - `status TEXT` - - `scheduled_departure TEXT` - - `scheduled_arrival TEXT` - - `actual_departure TEXT` - - `actual_arrival TEXT` -- Тех.колонки: `src_created_at_ts`, `load_dttm`, `batch_id` -- `src_created_at_ts` = `scheduled_departure` -- Распределение: `DISTRIBUTED BY (flight_id)` -- Обоснование: `flight_id` — это уникальный идентификатор рейса - -**Загрузка**: Инкремент по `scheduled_departure` - -**DQ проверки**: -- Count между источником и STG -- Дубликаты `flight_id` -- NULL обязательных полей (flight_id, route_no, status, scheduled_departure) -- Ссылочная целостность на routes (route_no) - -> Примечание: `flights.status/actual_*` в источнике могут меняться со временем. Для учебного STG можно принять допущение "insert-only" (снимаем слепок на момент загрузки), либо усложнить и перезагружать скользящее окно по датам вылета. - -### 6. segments (транзакции, инкремент) - -**Внешняя таблица**: `stg.segments_ext` -- Поля: `ticket_no`, `flight_id`, `fare_conditions`, `price` -- PXF: `pxf://bookings.segments?PROFILE=JDBC&SERVER=bookings-db` - -**Внутренняя таблица**: `stg.segments` -- Бизнес-колонки как TEXT: - - `ticket_no TEXT` - - `flight_id TEXT` - - `fare_conditions TEXT` - - `price TEXT` -- Тех.колонки: `src_created_at_ts`, `load_dttm`, `batch_id` -- `src_created_at_ts` = берётся из `bookings.book_date` через JOIN с tickets -- Распределение: `DISTRIBUTED BY (ticket_no)` -- Обоснование: co-location с tickets для оптимизации JOIN - -**Загрузка**: Инкремент по `book_date` (как в tickets) - -**DQ проверки**: -- Count между источником и STG -- Дубликаты `(ticket_no, flight_id)` -- NULL обязательных полей (ticket_no, flight_id, fare_conditions, price) -- Ссылочная целостность на tickets (ticket_no) -- Ссылочная целостность на flights (flight_id) - -### 7. boarding_passes (транзакции, full snapshot) - -**Внешняя таблица**: `stg.boarding_passes_ext` -- Поля: `ticket_no`, `flight_id`, `seat_no`, `boarding_no`, `boarding_time` -- PXF: `pxf://bookings.boarding_passes?PROFILE=JDBC&SERVER=bookings-db` - -**Внутренняя таблица**: `stg.boarding_passes` -- Бизнес-колонки как TEXT: - - `ticket_no TEXT` - - `flight_id TEXT` - - `seat_no TEXT` - - `boarding_no TEXT` - - `boarding_time TEXT` -- Тех.колонки: `src_created_at_ts`, `load_dttm`, `batch_id` -- `src_created_at_ts` = `now()` (в этой таблице нет удобного поля для инкремента, потому что строки могут создаваться и обновляться со временем) -- Распределение: `DISTRIBUTED BY (ticket_no)` -- Обоснование: co-location с tickets/segments для оптимизации JOIN - -**Загрузка**: Full snapshot (все строки при каждом запуске) - -**DQ проверки**: -- Count между источником и STG -- Дубликаты `(ticket_no, flight_id)` -- NULL обязательных полей (ticket_no, flight_id) -- Ссылочная целостность на tickets (ticket_no) -- Ссылочная целостность на segments (ticket_no, flight_id) - -> Примечание: в источнике `boarding_passes` строки сначала создаются при CHECK-IN (без `boarding_time`), а потом обновляются при BOARDING. Поэтому инкремент "по времени" без усложнений будет пропускать часть событий и/или изменения. Для учебного стенда самый стабильный вариант — снимать полный слепок. - -## Обновление существующих DAG - -### `airflow/dags/bookings_stg_ddl.py` - -Добавить задачи для создания DDL новых таблиц: - -```python -apply_stg_airports_ddl = PostgresOperator( - task_id="apply_stg_airports_ddl", - postgres_conn_id=GREENPLUM_CONN_ID, - sql="stg/airports_ddl.sql", -) - -apply_stg_airplanes_ddl = PostgresOperator( - task_id="apply_stg_airplanes_ddl", - postgres_conn_id=GREENPLUM_CONN_ID, - sql="stg/airplanes_ddl.sql", -) - -apply_stg_routes_ddl = PostgresOperator( - task_id="apply_stg_routes_ddl", - postgres_conn_id=GREENPLUM_CONN_ID, - sql="stg/routes_ddl.sql", -) - -apply_stg_seats_ddl = PostgresOperator( - task_id="apply_stg_seats_ddl", - postgres_conn_id=GREENPLUM_CONN_ID, - sql="stg/seats_ddl.sql", -) - -apply_stg_flights_ddl = PostgresOperator( - task_id="apply_stg_flights_ddl", - postgres_conn_id=GREENPLUM_CONN_ID, - sql="stg/flights_ddl.sql", -) - -apply_stg_segments_ddl = PostgresOperator( - task_id="apply_stg_segments_ddl", - postgres_conn_id=GREENPLUM_CONN_ID, - sql="stg/segments_ddl.sql", -) - -apply_stg_boarding_passes_ddl = PostgresOperator( - task_id="apply_stg_boarding_passes_ddl", - postgres_conn_id=GREENPLUM_CONN_ID, - sql="stg/boarding_passes_ddl.sql", -) -``` - -Зависимости: -- Сначала создаются справочники (airports, airplanes, routes, seats) -- Затем транзакционные таблицы (flights, segments, boarding_passes) - -### `airflow/dags/bookings_to_gp_stage.py` - -Добавить задачи для загрузки новых таблиц: - -```python -# Загрузка справочников (full load) -load_airports_to_stg = PostgresOperator( - task_id="load_airports_to_stg", - postgres_conn_id=GREENPLUM_CONN_ID, - sql="stg/airports_load.sql", -) - -check_airports_dq = PostgresOperator( - task_id="check_airports_dq", - postgres_conn_id=GREENPLUM_CONN_ID, - sql="stg/airports_dq.sql", -) - -load_airplanes_to_stg = PostgresOperator( - task_id="load_airplanes_to_stg", - postgres_conn_id=GREENPLUM_CONN_ID, - sql="stg/airplanes_load.sql", -) - -check_airplanes_dq = PostgresOperator( - task_id="check_airplanes_dq", - postgres_conn_id=GREENPLUM_CONN_ID, - sql="stg/airplanes_dq.sql", -) - -load_routes_to_stg = PostgresOperator( - task_id="load_routes_to_stg", - postgres_conn_id=GREENPLUM_CONN_ID, - sql="stg/routes_load.sql", -) - -check_routes_dq = PostgresOperator( - task_id="check_routes_dq", - postgres_conn_id=GREENPLUM_CONN_ID, - sql="stg/routes_dq.sql", -) - -load_seats_to_stg = PostgresOperator( - task_id="load_seats_to_stg", - postgres_conn_id=GREENPLUM_CONN_ID, - sql="stg/seats_load.sql", -) - -check_seats_dq = PostgresOperator( - task_id="check_seats_dq", - postgres_conn_id=GREENPLUM_CONN_ID, - sql="stg/seats_dq.sql", -) - -# Загрузка транзакций (инкремент) -load_flights_to_stg = PostgresOperator( - task_id="load_flights_to_stg", - postgres_conn_id=GREENPLUM_CONN_ID, - sql="stg/flights_load.sql", -) - -check_flights_dq = PostgresOperator( - task_id="check_flights_dq", - postgres_conn_id=GREENPLUM_CONN_ID, - sql="stg/flights_dq.sql", -) - -load_segments_to_stg = PostgresOperator( - task_id="load_segments_to_stg", - postgres_conn_id=GREENPLUM_CONN_ID, - sql="stg/segments_load.sql", -) - -check_segments_dq = PostgresOperator( - task_id="check_segments_dq", - postgres_conn_id=GREENPLUM_CONN_ID, - sql="stg/segments_dq.sql", -) - -load_boarding_passes_to_stg = PostgresOperator( - task_id="load_boarding_passes_to_stg", - postgres_conn_id=GREENPLUM_CONN_ID, - sql="stg/boarding_passes_load.sql", -) - -check_boarding_passes_dq = PostgresOperator( - task_id="check_boarding_passes_dq", - postgres_conn_id=GREENPLUM_CONN_ID, - sql="stg/boarding_passes_dq.sql", -) -``` - -Зависимости: -- Сначала загружаются и проверяются bookings и tickets (уже есть) -- Затем загружаются справочники (airports, airplanes, routes, seats) -- Затем загружаются транзакции (flights, segments, boarding_passes) -- В конце финальный лог - -### Обновление `sql/ddl_gp.sql` - -Добавить подключение новых DDL файлов: - -```sql --- DDL для слоя stg по таблицам bookings и tickets вынесены в отдельные файлы. --- Здесь подключаем их через psql \i, чтобы сохранить единый входной скрипт. -\i stg/bookings_ddl.sql -\i stg/tickets_ddl.sql - --- DDL для новых таблиц STG слоя -\i stg/airports_ddl.sql -\i stg/airplanes_ddl.sql -\i stg/routes_ddl.sql -\i stg/seats_ddl.sql -\i stg/flights_ddl.sql -\i stg/segments_ddl.sql -\i stg/boarding_passes_ddl.sql -``` - -### Добавление тестов - -Обновить `tests/test_dags_smoke.py` для проверки структуры обновлённых DAG: - -```python -def test_bookings_stg_ddl_dag_structure(): - dag = _load_dag("airflow.dags.bookings_stg_ddl") - - expected_tasks = { - "apply_stg_bookings_ddl", - "apply_stg_tickets_ddl", - "apply_stg_airports_ddl", - "apply_stg_airplanes_ddl", - "apply_stg_routes_ddl", - "apply_stg_seats_ddl", - "apply_stg_flights_ddl", - "apply_stg_segments_ddl", - "apply_stg_boarding_passes_ddl", - } - assert expected_tasks.issubset(dag.task_dict.keys()) - - # Проверка линейных зависимостей - # ... (проверка зависимостей между задачами) -``` - -```python -def test_bookings_to_gp_stage_dag_structure(): - dag = _load_dag("airflow.dags.bookings_to_gp_stage") - - expected_tasks = { - "generate_bookings_day", - "load_bookings_to_stg", - "check_row_counts", - "load_tickets_to_stg", - "check_tickets_dq", - "load_airports_to_stg", - "check_airports_dq", - "load_airplanes_to_stg", - "check_airplanes_dq", - "load_routes_to_stg", - "check_routes_dq", - "load_seats_to_stg", - "check_seats_dq", - "load_flights_to_stg", - "check_flights_dq", - "load_segments_to_stg", - "check_segments_dq", - "load_boarding_passes_to_stg", - "check_boarding_passes_dq", - "finish_summary", - } - assert expected_tasks.issubset(dag.task_dict.keys()) - - # Проверка линейных зависимостей - # ... (проверка зависимостей между задачами) -``` - -### Обновление документации - -Обновить статус в [`docs/internal/db_schema.md`](../docs/internal/db_schema.md) с "2 из 9" на "9 из 9". - -Добавить описание новых таблиц в документацию. - -## Диаграмма потока данных STG слоя - -```mermaid -graph TB - subgraph Source[Source: bookings-db] - B1[airports_data] - B2[airplanes_data] - B3[routes] - B4[seats] - B5[flights] - B6[segments] - B7[boarding_passes] - end - - subgraph STG[STG Layer: Greenplum] - S1[stg.airports] - S2[stg.airplanes] - S3[stg.routes] - S4[stg.seats] - S5[stg.flights] - S6[stg.segments] - S7[stg.boarding_passes] - end - - B1 --> S1 - B2 --> S2 - B3 --> S3 - B4 --> S4 - B5 --> S5 - B6 --> S6 - B7 --> S7 -``` - -## Чек-лист реализации - -- [ ] Создать файлы DDL для новых таблиц (7 файлов) - - [ ] `sql/stg/airports_ddl.sql` - - [ ] `sql/stg/airplanes_ddl.sql` - - [ ] `sql/stg/routes_ddl.sql` - - [ ] `sql/stg/seats_ddl.sql` - - [ ] `sql/stg/flights_ddl.sql` - - [ ] `sql/stg/segments_ddl.sql` - - [ ] `sql/stg/boarding_passes_ddl.sql` -- [ ] Создать файлы LOAD для новых таблиц (7 файлов) - - [ ] `sql/stg/airports_load.sql` - - [ ] `sql/stg/airplanes_load.sql` - - [ ] `sql/stg/routes_load.sql` - - [ ] `sql/stg/seats_load.sql` - - [ ] `sql/stg/flights_load.sql` - - [ ] `sql/stg/segments_load.sql` - - [ ] `sql/stg/boarding_passes_load.sql` -- [ ] Создать файлы DQ для новых таблиц (7 файлов) - - [ ] `sql/stg/airports_dq.sql` - - [ ] `sql/stg/airplanes_dq.sql` - - [ ] `sql/stg/routes_dq.sql` - - [ ] `sql/stg/seats_dq.sql` - - [ ] `sql/stg/flights_dq.sql` - - [ ] `sql/stg/segments_dq.sql` - - [ ] `sql/stg/boarding_passes_dq.sql` -- [ ] Обновить DAG `bookings_stg_ddl.py` -- [ ] Обновить DAG `bookings_to_gp_stage.py` -- [ ] Обновить `sql/ddl_gp.sql` -- [ ] Добавить тесты для новых DAG в `tests/test_dags_smoke.py` -- [ ] Обновить документацию `docs/internal/db_schema.md` -- [ ] Провести тестирование реализации - -## Примечания для реализации - -1. **Именование файлов**: Использовать `{table}_ddl.sql`, `{table}_load.sql`, `{table}_dq.sql` -2. **Ключи распределения**: Выбирать ключи с высокой кардинальностью для равномерного распределения -3. **Co-location**: Использовать одинаковые ключи распределения для связанных таблиц (tickets, segments, boarding_passes по ticket_no) -4. **Комментарии**: Добавлять русскоязычные комментарии в SQL-файлы для студентов -5. **DQ проверки**: Все проверки должны падать с `RAISE EXCEPTION` при ошибке -6. **Batch ID**: Использовать `{{ run_id }}` для идентификации батча -7. **Защита от дублей**: Использовать `NOT EXISTS` для предотвращения дублирования в рамках одного batch_id - -## Связанные документы - -- [`docs/internal/db_schema.md`](../docs/internal/db_schema.md) - Общая схема DWH -- [`docs/internal/bookings_stg_design.md`](../docs/internal/bookings_stg_design.md) - Детальный дизайн STG для bookings -- [`sql/stg/bookings_ddl.sql`](../sql/stg/bookings_ddl.sql) - Образец DDL -- [`sql/stg/bookings_load.sql`](../sql/stg/bookings_load.sql) - Образец LOAD -- [`sql/stg/bookings_dq.sql`](../sql/stg/bookings_dq.sql) - Образец DQ -- [`airflow/dags/bookings_stg_ddl.py`](../airflow/dags/bookings_stg_ddl.py) - Образец DAG DDL -- [`airflow/dags/bookings_to_gp_stage.py`](../airflow/dags/bookings_to_gp_stage.py) - Образец DAG загрузки