docs(all): ревизия документации перед мержем в main

- Зачем:
  - ветка содержала устаревшие ссылки, артефакты 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/" — пустой результат.
This commit is contained in:
2026-03-09 20:59:09 +03:00
parent a207f8d281
commit 3de4db8aa6
13 changed files with 102 additions and 1962 deletions
+3 -1
View File
@@ -39,7 +39,9 @@
- В каталоге `sql/` придерживаемся слоёв DWH: - В каталоге `sql/` придерживаемся слоёв DWH:
- `sql/src/` — скрипты, работающие с исходными системами (например, `bookings_generate_day_if_missing.sql`); - `sql/src/` — скрипты, работающие с исходными системами (например, `bookings_generate_day_if_missing.sql`);
- `sql/stg/` — скрипты для стейджинга (`bookings_ddl.sql`, `bookings_load.sql`, `bookings_dq.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` (единый источник для всех новых слоёв). - Нейминг служебных полей и SCD-полей фиксирован в `docs/internal/naming_conventions.md` (единый источник для всех новых слоёв).
- Именование файлов: `{объект}_{роль}.sql`, где: - Именование файлов: `{объект}_{роль}.sql`, где:
- `объект` — логическое имя сущности (`bookings`, `orders`, и т.п.); - `объект` — логическое имя сущности (`bookings`, `orders`, и т.п.);
+7 -5
View File
@@ -143,11 +143,13 @@ make clean # полный reset: удалить контейнер
## Документация ## Документация
- Учебные задания: `educational-tasks.md`. - [Учебные задания](educational-tasks.md)
- План тестирования/проверок и негативные кейсы: `TESTING.md`. - [План тестирования/проверок и негативные кейсы](TESTING.md)
- Дополнительные заметки и технические детали: `docs/README.md`. - [Дополнительные заметки и технические детали](docs/README.md)
- Детали по ODS DAG: `docs/bookings_to_gp_ods.md`. - [Детали по STG DAG](docs/bookings_to_gp_stage.md)
- Детали по DDS DAG: `docs/bookings_to_gp_dds.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)
## Типичные проблемы и решения ## Типичные проблемы и решения
-2
View File
@@ -69,12 +69,10 @@
## 8. Снятие метрик и мониторинг ## 8. Снятие метрик и мониторинг
- Контейнеры: `docker compose ps`, `docker stats` (по желанию). - Контейнеры: `docker compose ps`, `docker stats` (по желанию).
- Логи задач: в Airflow UI → конкретный таск → Log. - Логи задач: в Airflow UI → конкретный таск → Log.
- Хостовые CSV: каталог `data/` (можно открыть любой файл и убедиться в структуре).
## 9. Завершение работы ## 9. Завершение работы
- `make down` — выключает сервисы и удаляет контейнеры/сети (volumes сохраняются). - `make down` — выключает сервисы и удаляет контейнеры/сети (volumes сохраняются).
- Полный сброс данных (удаляет volumes): `make clean`. - Полный сброс данных (удаляет volumes): `make clean`.
- При необходимости сохранить данные: скопировать CSV из `data/` и сделать дампы до `make clean`.
## Текущий статус (пример успешного прогона) ## Текущий статус (пример успешного прогона)
- `uv run pytest -q` — 14 passed, 9 smoke-тестов DAG пропущены (Airflow не установлен в venv). - `uv run pytest -q` — 14 passed, 9 smoke-тестов DAG пропущены (Airflow не установлен в venv).
+1 -1
View File
@@ -125,7 +125,7 @@
- патчи `bookings/patches/engine_jobs1_sync.patch` и `bookings/patches/install_drop_if_exists.patch` - патчи `bookings/patches/engine_jobs1_sync.patch` и `bookings/patches/install_drop_if_exists.patch`
падают при применении (hunk failed / garbage in patch); падают при применении (hunk failed / garbage in patch);
- из‑за этого DAG `bookings_to_gp_stage` валится на проверках (источник пустой). - из‑за этого DAG `bookings_to_gp_stage` валится на проверках (источник пустой).
- план: `plans/bookings-demodb-bugfix-plan.md` - детали: `docs/internal/bookings_db_issues.md`
- [x] Добавить раздел «Благодарности» в `README.md`: - [x] Добавить раздел «Благодарности» в `README.md`:
- явно поблагодарить Postgres Pro за демо‑БД bookings (репозиторий `postgrespro/demodb`); - явно поблагодарить Postgres Pro за демо‑БД bookings (репозиторий `postgrespro/demodb`);
+6 -4
View File
@@ -10,13 +10,15 @@
- [Главный учебный DAG: bookings → stg](bookings_to_gp_stage.md) - [Главный учебный DAG: bookings → stg](bookings_to_gp_stage.md)
- [Учебный DAG: stg -> ods](bookings_to_gp_ods.md) - [Учебный DAG: stg -> ods](bookings_to_gp_ods.md)
- [Учебный DAG: ods -> dds](bookings_to_gp_dds.md) - [Учебный DAG: ods -> dds](bookings_to_gp_dds.md)
- [Учебный DAG: dds -> dm](bookings_to_gp_dm.md)
## Технические детали (опционально) ## Технические детали (опционально)
- [Как устроен Docker-стенд (образы, Connections, переменные окружения)](stack.md) - [Как устроен Docker-стенд (образы, Connections, переменные окружения)](stack.md)
- [Единые конвенции нейминга DWH (служебные поля и SCD)](internal/naming_conventions.md) - [Единые конвенции нейминга DWH (служебные поля и SCD)](internal/naming_conventions.md)
- [PXF в этом проекте (проектная реализация)](internal/pxf_bookings.md) - [PXF в этом проекте (проектная реализация)](internal/pxf_bookings.md)
- [Дизайн stg для bookings (черновик)](internal/bookings_stg_design.md) - [Дизайн-документ STG](internal/bookings_stg_design.md)
- [Дизайн ods для bookings (черновик)](internal/bookings_ods_design.md) - [Дизайн-документ ODS](internal/bookings_ods_design.md)
- [Дизайн dds для bookings (черновик)](internal/bookings_dds_design.md) - [Дизайн-документ DDS](internal/bookings_dds_design.md)
- [Про время/UTC в bookings (черновик)](internal/bookings_tz.md) - [Дизайн-документ DM](internal/bookings_dm_design.md)
- [Про время/UTC в bookings](internal/bookings_tz.md)
+79
View File
@@ -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 задачу с ошибкой и исправьте.
-466
View File
@@ -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)
+1 -1
View File
@@ -143,7 +143,7 @@ pgmeta (Postgres 16) ─────────────> Airflow (webserver
1. **Эталонный вертикальный срез** — полностью реализованная цепочка 1. **Эталонный вертикальный срез** — полностью реализованная цепочка
`sales_report` и все её источники вниз по слоям (STG → ODS → DDS → DM). `sales_report` и все её источники вниз по слоям (STG → ODS → DDS → DM).
2. **ТЗ от аналитика** — описание остальных таблиц 2. **ТЗ от аналитика** — описание остальных таблиц
([analyst_spec.md](../assignment/analyst_spec.md)). (analyst_spec.md — будет создан на Этапе 3, см. TODO.md).
3. **Частично готовый DAG** — студент добавляет свои таски по аналогии. 3. **Частично готовый DAG** — студент добавляет свои таски по аналогии.
4. **Валидационный DAG** — студент запускает для самоконтроля. 4. **Валидационный DAG** — студент запускает для самоконтроля.
+5 -7
View File
@@ -1,12 +1,10 @@
# Дизайн STG для bookings в Greenplum (черновик) # Дизайн STG для bookings в Greenplum
_Внутренний документ для учебного стенда. Перед итоговой сдачей можно объединить с основной документацией._
## 1. Цель и общий контур ## 1. Цель и общий контур
- Источник: Postgres в контейнере `bookings-db`, база `demo`, таблица `bookings.bookings` (см. `docs/internal/bookings_tz.md`). - Источник: Postgres в контейнере `bookings-db`, база `demo`, таблица `bookings.bookings` (см. `docs/internal/bookings_tz.md`).
- Цель: показываем путь данных от операционной БД до сырого слоя DWH в Greenplum. - Цель: показываем путь данных от операционной БД до сырого слоя DWH в Greenplum.
- В этом документе описываем только часть `src (bookings-db) → STG (Greenplum)`. Слои ODS/DDS/DM студент проектирует сам по статье про моделирование DWH. - В этом документе описываем часть `src (bookings-db) → STG (Greenplum)`. STG — входной слой; далее данные обрабатываются в ODS → DDS → DM (см. соответствующие design-документы).
Примечание: в текущей версии стенда слой `stg` содержит не только `bookings`, но и остальные таблицы потока Примечание: в текущей версии стенда слой `stg` содержит не только `bookings`, но и остальные таблицы потока
(`tickets`, `airports`, `airplanes`, `routes`, `seats`, `flights`, `segments`, `boarding_passes`). (`tickets`, `airports`, `airplanes`, `routes`, `seats`, `flights`, `segments`, `boarding_passes`).
@@ -16,7 +14,7 @@ _Внутренний документ для учебного стенда. П
- `src`: оперативная система (`bookings-db`, схема `bookings`). - `src`: оперативная система (`bookings-db`, схема `bookings`).
- `stg`: сырой слой в Greenplum, максимально близкий к источнику, без бизнес‑логики. - `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 ## 2. Схема и таблицы в Greenplum
@@ -37,7 +35,7 @@ _Внутренний документ для учебного стенда. П
В текущей реализации аналогично созданы внешние таблицы `*_ext` и для остальных сущностей (см. `sql/stg/*_ddl.sql`). В текущей реализации аналогично созданы внешние таблицы `*_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 ### 2.3. Внутренняя таблица STG
@@ -133,7 +131,7 @@ DDL будет добавлен в `sql/ddl_gp.sql` в блоке DDL для Gre
- `docs/internal/pxf_bookings.md` — детали настройки PXF и внешней таблицы для чтения из `bookings-db`. - `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`). - `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. Требования к читаемости и комментариям ## 6. Требования к читаемости и комментариям
-100
View File
@@ -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`, под который уже готовы патчи (делать только вместе с обновлением документации и проверкой, что генерация стабильна).
-489
View File
@@ -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'ах.
-142
View File
@@ -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`
- 310 рестартов:
- `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`).
-744
View File
@@ -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 загрузки