From 22a0590e1bd4c70e6f1f2390af9eed55895da958 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Tue, 3 Mar 2026 22:38:15 +0300 Subject: [PATCH] =?UTF-8?q?feat(sql):=20=D1=8F=D0=B2=D0=BD=D0=BE=D0=B5=20?= =?UTF-8?q?=D1=83=D0=BF=D1=80=D0=B0=D0=B2=D0=BB=D0=B5=D0=BD=D0=B8=D0=B5=20?= =?UTF-8?q?storage=20=D0=B8=20=D0=B0=D0=B2=D1=82=D0=BE=D0=BC=D0=B0=D1=82?= =?UTF-8?q?=D0=B8=D0=B7=D0=B8=D1=80=D0=BE=D0=B2=D0=B0=D0=BD=D0=BD=D1=8B?= =?UTF-8?q?=D0=B9=20E2E-=D1=82=D0=B5=D1=81=D1=82=20=D1=87=D0=B5=D1=80?= =?UTF-8?q?=D0=B5=D0=B7=20REST=20API?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - необходимо визуализировать выбор типа хранения (Heap vs Append-Only) для учебных целей. - автоматизировать проверку всей цепочки DWH для исключения ручных ошибок. - сделать процесс отладки прозрачным и наглядным через стандартные инструменты Airflow. - Что: - внедрена клауза WITH (appendonly=...) во все DDL; ODS-справочники переведены на AO Row и TRUNCATE+INSERT. - создан скрипт scripts/e2e_etl.sh для полного прогона ETL (DDL + 2 дня данных) через REST API. - исправлены баги типизации (INTEGER[]), именования полей (amount, passenger_id) и удалены фантомные колонки (contact_data). - обновлен e2e-etl-test-protocol.md: добавлен раздел по отладке, чтению логов и перезапуску задач через API. - исправлены pytest-контракты под новую логику загрузки. - Проверка: - успешный прогон `make e2e-smoke` (полный цикл от очистки до витрины). --- Makefile | 6 +- TODO.md | 2 +- _vendor/demodb | 1 + docker-compose.yml | 1 + docs/e2e-etl-test-protocol.md | 133 ++++++++++++++++++++--------- scripts/e2e_etl.sh | 104 ++++++++++++++++++++++ scripts/e2e_smoke.sh | 24 +----- sql/dds/dim_airplanes_ddl.sql | 6 ++ sql/dds/dim_airports_ddl.sql | 6 ++ sql/dds/dim_calendar_ddl.sql | 5 ++ sql/dds/dim_passengers_ddl.sql | 8 +- sql/dds/dim_passengers_dq.sql | 12 +-- sql/dds/dim_passengers_load.sql | 81 +++++++----------- sql/dds/dim_routes_ddl.sql | 6 ++ sql/dds/dim_routes_load.sql | 77 ++++++----------- sql/dds/dim_tariffs_ddl.sql | 5 ++ sql/dds/fact_flight_sales_ddl.sql | 7 ++ sql/dds/fact_flight_sales_load.sql | 10 +-- sql/dm/sales_report_ddl.sql | 8 +- sql/ods/airplanes_ddl.sql | 6 ++ sql/ods/airplanes_load.sql | 95 +++++---------------- sql/ods/airports_ddl.sql | 6 ++ sql/ods/airports_load.sql | 105 +++++------------------ sql/ods/boarding_passes_ddl.sql | 6 ++ sql/ods/bookings_ddl.sql | 6 ++ sql/ods/flights_ddl.sql | 6 ++ sql/ods/routes_ddl.sql | 14 ++- sql/ods/routes_load.sql | 115 ++++++------------------- sql/ods/seats_ddl.sql | 16 ++-- sql/ods/seats_load.sql | 86 ++++--------------- sql/ods/segments_ddl.sql | 21 +++-- sql/ods/segments_load.sql | 10 +-- sql/ods/tickets_ddl.sql | 6 ++ tests/test_ods_sql_contract.py | 10 +-- 34 files changed, 484 insertions(+), 526 deletions(-) create mode 160000 _vendor/demodb create mode 100755 scripts/e2e_etl.sh diff --git a/Makefile b/Makefile index 9030c2b..7c805bc 100644 --- a/Makefile +++ b/Makefile @@ -9,7 +9,7 @@ BOOKINGS_INIT_DAYS ?= 1 .PHONY: up stop down clean airflow-init logs gp-psql ddl-gp \ bookings-check-jobs bookings-clone-demodb bookings-init bookings-psql bookings-generate-day \ - dev-setup dev-sync dev-lock test lint fmt clean-venv build + dev-setup dev-sync dev-lock test lint fmt clean-venv build e2e-smoke e2e-etl SHELL := /bin/bash up: @@ -136,3 +136,7 @@ clean-venv: e2e-smoke: ./scripts/e2e_smoke.sh + +e2e-etl: + ./scripts/e2e_etl.sh + diff --git a/TODO.md b/TODO.md index 4d096f0..fb12cb2 100644 --- a/TODO.md +++ b/TODO.md @@ -5,7 +5,7 @@ Статус: ниже есть как актуальные, так и уже выполненные пункты. -- [ ] Сделать REST API Airflow основным способом тестирования ETL вместо CLI-вызовов +- [x] Сделать REST API Airflow основным способом тестирования ETL вместо CLI-вызовов через `docker compose exec ... airflow ...`: - обновить `TESTING.md`, сместив фокус на REST API сценарии; - оставить CLI как резервный вариант для локальной отладки; diff --git a/_vendor/demodb b/_vendor/demodb new file mode 160000 index 0000000..d68de19 --- /dev/null +++ b/_vendor/demodb @@ -0,0 +1 @@ +Subproject commit d68de192850237719f09b47688d5f3fc94653ca6 diff --git a/docker-compose.yml b/docker-compose.yml index f844480..f792e6e 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -2,6 +2,7 @@ x-airflow-common-env: &airflow-env TZ: ${TZ:-Europe/Moscow} AIRFLOW__CORE__LOAD_EXAMPLES: "False" AIRFLOW__CORE__EXECUTOR: LocalExecutor + AIRFLOW__API__AUTH_BACKENDS: "airflow.api.auth.backend.basic_auth,airflow.api.auth.backend.session" 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} diff --git a/docs/e2e-etl-test-protocol.md b/docs/e2e-etl-test-protocol.md index 1cf62bf..60992ba 100644 --- a/docs/e2e-etl-test-protocol.md +++ b/docs/e2e-etl-test-protocol.md @@ -3,20 +3,25 @@ Этот документ описывает процедуру полной проверки цепочки ETL: `STG -> ODS -> DDS -> DM`. Цель теста — убедиться в корректности инкрементальной загрузки, работы паттерна `Temporary Table` и механизмов `HWM`. +**Основной метод взаимодействия — Airflow UI + REST API.** Запускать пайплайны удобнее всего через веб-интерфейс, а вот отлаживать упавшие задачи (читать логи, очищать статус) полезно уметь через REST API. Это приближает опыт к реальной боевой эксплуатации. + --- ## 1. Подготовка окружения -Убедитесь, что все сервисы запущены и DDL применен. +Убедитесь, что все сервисы запущены и генератор инициализирован. ```bash make up make bookings-init -make ddl-gp ``` -### Важно: Настройка дат -Для корректного тестирования исторических данных из `demodb` (начинаются с 2017 года), убедитесь, что в DAG-файлах `start_date` установлен в `2017-01-01`. +**Доступ к API:** +В `docker-compose.yml` включена базовая аутентификация (`basic_auth`). +Для curl-запросов используйте учетные данные из вашего `.env` файла (переменные `AIRFLOW_USER` и `AIRFLOW_PASSWORD`, по умолчанию `admin:admin`). + +- **UI:** `http://localhost:8080` +- **API Endpoint:** `http://localhost:8080/api/v1` --- @@ -27,52 +32,99 @@ make ddl-gp ```bash make dwh-truncate ``` +*(Если вы меняли DDL, лучше полностью пересоздать схемы: `make gp-psql -c "DROP SCHEMA IF EXISTS ods CASCADE; DROP SCHEMA IF EXISTS dds CASCADE; DROP SCHEMA IF EXISTS dm CASCADE; CREATE SCHEMA ods; CREATE SCHEMA dds; CREATE SCHEMA dm;"`)* --- -## 3. Этап 1: Загрузка за первый день (2017-01-01) +## 3. Этап 1: Создание схем и Загрузка за первый день (Initial Load) -Выполните последовательный запуск всех DAG для первой порции данных. +Рекомендуется запускать слои последовательно, дожидаясь завершения предыдущего. +### Создание DDL +Откройте **Airflow UI** (`http://localhost:8080`) и нажмите кнопку **▶ Play -> Trigger DAG** для DDL-дагов: +1. `bookings_stg_ddl` +2. `bookings_ods_ddl` +3. `bookings_dds_ddl` +4. `bookings_dm_ddl` + +### Запуск пайплайна (Day 1) +После успешного создания таблиц, запустите DAG загрузки для `bookings_to_gp_stage`. +Откройте **Airflow UI** (`http://localhost:8080`), **снимите DAG с паузы** (переключатель слева от названия) и нажмите кнопку **▶ Play -> Trigger DAG**. +*(Airflow автоматически сгенерирует `logical_date` и `run_id`, например `manual__2026-03-03T10:00:00+00:00`)*. + +Либо сделайте то же самое через API (без указания даты). **Важно:** при старте стенда все DAG-и находятся на паузе. Чтобы планировщик начал выполнять запущенный вами DAG, его нужно предварительно "разморозить" (unpause): ```bash -# Загрузка в STG (создает первый батч в источнике) -docker compose exec airflow-scheduler airflow dags test bookings_to_gp_stage 2017-01-01 +# Снятие с паузы +curl -s -X PATCH "http://localhost:8080/api/v1/dags/bookings_to_gp_stage" \ +--user "${AIRFLOW_USER}:${AIRFLOW_PASSWORD}" \ +-H "Content-Type: application/json" \ +-d '{"is_paused": false}' -# Загрузка в ODS (Initial Load) -docker compose exec airflow-scheduler airflow dags test bookings_to_gp_ods 2017-01-01 - -# Загрузка в DDS (Initial Load) -docker compose exec airflow-scheduler airflow dags test bookings_to_gp_dds 2017-01-01 - -# Загрузка в DM (Initial Load витрины) -docker compose exec airflow-scheduler airflow dags test bookings_to_gp_dm 2017-01-01 +# Запуск +curl -s -X POST "http://localhost:8080/api/v1/dags/bookings_to_gp_stage/dagRuns" \ +--user "${AIRFLOW_USER}:${AIRFLOW_PASSWORD}" \ +-H "Content-Type: application/json" \ +-d '{}' ``` -### Ожидаемые результаты (Day 1) -Проверьте наполнение таблиц: -- `stg.bookings` и `ods.bookings` должны иметь одинаковое количество строк (>0). -- `dm.sales_report` должна содержать агрегированные данные за первый день. +### Проверка статуса +Следите за графом выполнения в UI. Как только DAG перейдет в статус `success`, поочередно запускайте следующие слои: +1. `bookings_to_gp_ods` +2. `bookings_to_gp_dds` +3. `bookings_to_gp_dm` --- -## 4. Этап 2: Проверка инкремента (2017-01-02) +## 4. Этап 2: Проверка инкремента -Эмулируйте появление данных за второй день и проверьте дозагрузку. +Эмулируйте появление данных за второй день и проверьте дозагрузку. DAG слоя STG автоматически сгенерирует новый день в базе-источнике перед загрузкой. ```bash -# Генерация данных за 2-й день в базе-источнике -make bookings-generate-day - -# Повторный запуск цепочки ETL -docker compose exec airflow-scheduler airflow dags test bookings_to_gp_stage 2017-01-02 -docker compose exec airflow-scheduler airflow dags test bookings_to_gp_ods 2017-01-02 -docker compose exec airflow-scheduler airflow dags test bookings_to_gp_dds 2017-01-02 -docker compose exec airflow-scheduler airflow dags test bookings_to_gp_dm 2017-01-02 +# Повторный запуск цепочки ETL (DAG STG сам сгенерирует новый день) +# Снова нажмите "Trigger DAG" в UI для каждого слоя (STG -> ODS -> DDS -> DM). ``` --- -## 5. Финальная верификация (Критерии успеха) +## 5. Цикл отладки: Логи и Перезапуск (Clear) + +Если DAG упал, **не нужно пересоздавать стенд с нуля**. Airflow позволяет исправить код и перезапустить только упавшие задачи. + +Для выполнения команд ниже вам понадобится **Run ID** упавшего запуска. Его можно скопировать из UI (вкладка *Graph* -> кликнуть на фон сетки -> вкладка *Details* -> `Run ID`) или получить последним API-запросом: + +```bash +# Получить Run ID последнего запуска ODS +curl -s "http://localhost:8080/api/v1/dags/bookings_to_gp_ods/dagRuns?order_by=-execution_date&limit=1" \ +--user "${AIRFLOW_USER}:${AIRFLOW_PASSWORD}" | grep -o '"dag_run_id": "[^"]*"' +``` + +### Чтение логов через API +Подставьте ваш `` (например, `manual__2026-03-03T...`) и имя упавшей таски: +```bash +curl -s "http://localhost:8080/api/v1/dags/bookings_to_gp_ods/dagRuns//taskInstances//logs/1" \ +--user "${AIRFLOW_USER}:${AIRFLOW_PASSWORD}" +``` + +### Перезапуск задачи (Clear) +1. Прочитайте ошибку в логах. +2. Исправьте SQL-файл локально на хосте. +3. Очистите состояние упавших задач (`only_failed: true`) в конкретном запуске, передав ваш ``: + +```bash +curl -s -X POST "http://localhost:8080/api/v1/dags/bookings_to_gp_ods/clearTaskInstances" \ +--user "${AIRFLOW_USER}:${AIRFLOW_PASSWORD}" \ +-H "Content-Type: application/json" \ +-d '{ + "only_failed": true, + "reset_dag_runs": true, + "dag_run_id": "" +}' +``` +После этого планировщик подхватит обновленный SQL-код и продолжит выполнение DAG с точки падения. Вы также можете сделать это в UI: клик по упавшей задаче -> кнопка **Clear**. + +--- + +## 6. Финальная верификация (Критерии успеха) Выполните SQL-запрос для сверки данных: @@ -89,15 +141,16 @@ SELECT 'DM ' as layer, COUNT(*) FROM dm.sales_report; ``` **Критерии корректности:** -1. **STG == ODS**: Количество строк в `stg.bookings` и `ods.bookings` совпадает (т.к. это SCD1 UPSERT). -2. **Инкремент STG**: Количество строк в `stg.bookings` после Day 2 больше, чем после Day 1. -3. **Инкремент ODS (Temporary Table)**: В ODS нет дублей. `SELECT book_ref FROM ods.bookings GROUP BY book_ref HAVING COUNT(*) > 1` должен вернуть 0 строк. -4. **HWM в DM**: Витрина `sales_report` содержит данные за оба дня. Значение `COUNT(*)` после Day 2 должно вырасти по сравнению с Day 1. -5. **Lineage**: Поля `_load_id` и `_load_ts` во всех слоях содержат метки соответствующих запусков. +1. **Инкремент STG**: Количество строк в `stg.bookings` после Этапа 2 больше, чем после Этапа 1. +2. **Инкремент ODS**: Количество строк в `ods.bookings` выросло. В ODS нет дублей (`SELECT book_ref FROM ods.bookings GROUP BY book_ref HAVING COUNT(*) > 1` должен вернуть 0 строк). +3. **ODS Справочники**: Количество строк в `ods.airports` и `ods.routes` не должно меняться между днями (работает паттерн TRUNCATE+INSERT полного снимка). +4. **HWM в DM**: Витрина `sales_report` содержит данные за оба дня. Значение `COUNT(*)` после Этапа 2 выросло. +5. **Lineage**: Поля `_load_id` и `_load_ts` во всех слоях содержат метки соответствующих запусков (`manual__...`). --- -## Типичные ошибки -- **Пустые таблицы**: Проверьте, что в `bookings-db` есть данные (`SELECT COUNT(*) FROM bookings.bookings`). Если 0 — сделайте `make bookings-init`. -- **Пропуски в ODS**: Убедитесь, что `stg_batch_id` в ODS корректно вычисляется (задача `resolve_stg_batch_id`). -- **Дубли в DDS**: Проверьте логику генерации SK в `dds/*_load.sql`. +## 7. Зафиксированный опыт (Типичные ошибки) + +- **Рассинхронизация DDL и Load скриптов**: Частая причина падения ODS/DDS — несовпадение имен колонок (например, `amount` vs `segment_amount`) или типов данных (например, `INTEGER[]` vs `TEXT`) между схемой таблицы и запросом загрузки. Внимательно читайте логи задачи. +- **Работа с массивами**: При генерации `hashdiff` в Greenplum/PostgreSQL нельзя использовать пустую строку `''` в `COALESCE` для массива. Массив нужно предварительно привести к тексту: `COALESCE(days_of_week::TEXT, '')`. +- **Кавычки в psql**: При выполнении ручных проверок через `psql -c "..."` помните, что строковые литералы должны оборачиваться в **одинарные кавычки** (`'text'`), а двойные кавычки (`"text"`) интерпретируются как идентификаторы колонок. diff --git a/scripts/e2e_etl.sh b/scripts/e2e_etl.sh new file mode 100755 index 0000000..907698e --- /dev/null +++ b/scripts/e2e_etl.sh @@ -0,0 +1,104 @@ +#!/usr/bin/env bash +set -euo pipefail + +# Скрипт автоматизированного прогона E2E-теста всей цепочки DWH через REST API Airflow. +# Отрабатывает 2 "учебных дня" для проверки инкрементальной загрузки. + +ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" +cd "$ROOT" + +warn() { echo ":: ${1}" >&2; } + +# Загружаем переменные из .env, если файл существует, иначе используем дефолты +if [ -f ".env" ]; then + export $(grep -v '^#' .env | xargs) +fi + +AIRFLOW_USER=${AIRFLOW_USER:-admin} +AIRFLOW_PASSWORD=${AIRFLOW_PASSWORD:-admin} +API_URL="http://localhost:8080/api/v1/dags" + +# Функция для триггера DAG и ожидания его завершения +trigger_and_wait() { + local dag_id=$1 + warn "Запуск $dag_id через REST API..." + + # 0. Unpause DAG + curl -s -X PATCH "$API_URL/$dag_id" \ + --user "$AIRFLOW_USER:$AIRFLOW_PASSWORD" \ + -H "Content-Type: application/json" \ + -d '{"is_paused": false}' > /dev/null + + # 1. Trigger + local response + response=$(curl -s -X POST "$API_URL/$dag_id/dagRuns" \ + --user "$AIRFLOW_USER:$AIRFLOW_PASSWORD" \ + -H "Content-Type: application/json" \ + -d '{}') + + local run_id + run_id=$(echo "$response" | grep -o '"dag_run_id": "[^"]*' | cut -d'"' -f4) + if [ -z "$run_id" ]; then + echo "Ошибка запуска $dag_id. Ответ API: $response" >&2 + exit 1 + fi + warn "Успешный триггер. Run ID: $run_id" + + # 2. Wait + warn "Ожидание завершения $dag_id..." + local attempts=0 + local max_attempts=120 # 10 минут (120 * 5 сек) + + while true; do + local status + status=$(curl -s "$API_URL/$dag_id/dagRuns/$run_id" \ + --user "$AIRFLOW_USER:$AIRFLOW_PASSWORD" | grep -o '"state": "[^"]*' | cut -d'"' -f4) + + if [ "$status" = "success" ]; then + warn "DAG $dag_id завершен: SUCCESS" + break + elif [ "$status" = "failed" ]; then + echo "DAG $dag_id УПАЛ (FAILED)!" >&2 + exit 1 + elif [ "$status" = "queued" ] || [ "$status" = "running" ] || [ "$status" = "" ]; then + attempts=$((attempts + 1)) + if [ "$attempts" -ge "$max_attempts" ]; then + echo "Таймаут ожидания $dag_id!" >&2 + exit 1 + fi + sleep 5 + else + echo "Неизвестный статус: $status" >&2 + exit 1 + fi + done +} + +warn "=== DDL: Создание схем и таблиц ===" +trigger_and_wait "bookings_stg_ddl" +trigger_and_wait "bookings_ods_ddl" +trigger_and_wait "bookings_dds_ddl" +trigger_and_wait "bookings_dm_ddl" + +warn "=== DAY 1: Initial Load ===" +trigger_and_wait "bookings_to_gp_stage" +trigger_and_wait "bookings_to_gp_ods" +trigger_and_wait "bookings_to_gp_dds" +trigger_and_wait "bookings_to_gp_dm" + +warn "=== DAY 2: Increment ===" +trigger_and_wait "bookings_to_gp_stage" +trigger_and_wait "bookings_to_gp_ods" +trigger_and_wait "bookings_to_gp_dds" +trigger_and_wait "bookings_to_gp_dm" + +warn "=== Верификация данных ===" +# Выполняем простые проверки строк, чтобы убедиться, что данные дошли до витрины +docker compose -f docker-compose.yml exec greenplum bash -c "su - gpadmin -c \"psql -d gp_dwh -c \\\" +SELECT 'STG bookings' as layer, COUNT(*) FROM stg.bookings UNION ALL +SELECT 'ODS bookings', COUNT(*) FROM ods.bookings UNION ALL +SELECT 'DDS fact', COUNT(*) FROM dds.fact_flight_sales UNION ALL +SELECT 'DM report', COUNT(*) FROM dm.sales_report; +\\\"\"" + +warn "E2E ETL тест успешно завершен!" diff --git a/scripts/e2e_smoke.sh b/scripts/e2e_smoke.sh index 66d3962..950ef5d 100755 --- a/scripts/e2e_smoke.sh +++ b/scripts/e2e_smoke.sh @@ -38,30 +38,10 @@ wait_up airflow-scheduler warn "Init demo DB bookings" make bookings-init -warn "Apply DDL to Greenplum" -make ddl-gp - warn "Run local pytest suite" make test -warn "Airflow DAG test: csv_to_greenplum" -docker compose -f docker-compose.yml exec airflow-webserver airflow dags test csv_to_greenplum 2024-01-01 - -warn "Check orders count in Greenplum" -ORDERS_COUNT=$(docker compose -f docker-compose.yml exec greenplum bash -lc "su - gpadmin -c \"/usr/local/greenplum-db/bin/psql -t -A -d gp_dwh -c 'SELECT COUNT(*) FROM public.orders;'\"") -if [ "${ORDERS_COUNT:-0}" -le 0 ]; then - echo "orders table is empty after csv_to_greenplum (COUNT=${ORDERS_COUNT:-0})" >&2 - exit 1 -fi - -warn "Airflow DAG test: bookings_to_gp_stage" -docker compose -f docker-compose.yml exec airflow-webserver airflow dags test bookings_to_gp_stage 2024-01-01 - -warn "Check stg.bookings count in Greenplum" -BOOKINGS_COUNT=$(docker compose -f docker-compose.yml exec greenplum bash -lc "su - gpadmin -c \"/usr/local/greenplum-db/bin/psql -t -A -d gp_dwh -c 'SELECT COUNT(*) FROM stg.bookings;'\"") -if [ "${BOOKINGS_COUNT:-0}" -le 0 ]; then - echo "stg.bookings is empty after bookings_to_gp_stage (COUNT=${BOOKINGS_COUNT:-0})" >&2 - exit 1 -fi +warn "Start full E2E ETL test via REST API" +make e2e-etl warn "Smoke test completed successfully" diff --git a/sql/dds/dim_airplanes_ddl.sql b/sql/dds/dim_airplanes_ddl.sql index 5376870..2a9b13b 100644 --- a/sql/dds/dim_airplanes_ddl.sql +++ b/sql/dds/dim_airplanes_ddl.sql @@ -2,6 +2,9 @@ CREATE SCHEMA IF NOT EXISTS dds; +-- Тип таблицы: Heap (стандартная). +-- Обоснование: Необходим row-level UPDATE для реализации SCD1 UPSERT. +-- Использование Append-Only при частых обновлениях приводит к раздуванию (bloat) таблицы. CREATE TABLE IF NOT EXISTS dds.dim_airplanes ( airplane_sk INTEGER NOT NULL, airplane_bk TEXT NOT NULL, @@ -14,4 +17,7 @@ CREATE TABLE IF NOT EXISTS dds.dim_airplanes ( _load_id TEXT NOT NULL, _load_ts TIMESTAMP NOT NULL DEFAULT now() ) +WITH (appendonly=false) DISTRIBUTED BY (airplane_sk); + +COMMENT ON TABLE dds.dim_airplanes IS 'Измерение моделей самолетов (DDS).'; diff --git a/sql/dds/dim_airports_ddl.sql b/sql/dds/dim_airports_ddl.sql index 4495aad..d7566e3 100644 --- a/sql/dds/dim_airports_ddl.sql +++ b/sql/dds/dim_airports_ddl.sql @@ -2,6 +2,9 @@ CREATE SCHEMA IF NOT EXISTS dds; +-- Тип таблицы: Heap (стандартная). +-- Обоснование: Необходим row-level UPDATE для реализации SCD1 UPSERT. +-- Использование Append-Only при частых обновлениях приводит к раздуванию (bloat) таблицы. CREATE TABLE IF NOT EXISTS dds.dim_airports ( airport_sk INTEGER NOT NULL, airport_bk TEXT NOT NULL, @@ -15,4 +18,7 @@ CREATE TABLE IF NOT EXISTS dds.dim_airports ( _load_id TEXT NOT NULL, _load_ts TIMESTAMP NOT NULL DEFAULT now() ) +WITH (appendonly=false) DISTRIBUTED BY (airport_sk); + +COMMENT ON TABLE dds.dim_airports IS 'Измерение аэропортов (DDS).'; diff --git a/sql/dds/dim_calendar_ddl.sql b/sql/dds/dim_calendar_ddl.sql index 8611829..9eef562 100644 --- a/sql/dds/dim_calendar_ddl.sql +++ b/sql/dds/dim_calendar_ddl.sql @@ -2,6 +2,8 @@ CREATE SCHEMA IF NOT EXISTS dds; +-- Тип таблицы: Append-Only Row-oriented (zstd:1). +-- Обоснование: Статичные данные без обновлений. Обеспечивает эффективное сжатие. CREATE TABLE IF NOT EXISTS dds.dim_calendar ( calendar_sk INTEGER NOT NULL, date_actual DATE NOT NULL, @@ -12,4 +14,7 @@ CREATE TABLE IF NOT EXISTS dds.dim_calendar ( day_name TEXT NOT NULL, is_weekend BOOLEAN NOT NULL ) +WITH (appendonly=true, orientation=row, compresstype=zstd, compresslevel=1) DISTRIBUTED BY (calendar_sk); + +COMMENT ON TABLE dds.dim_calendar IS 'Измерение календаря (DDS).'; diff --git a/sql/dds/dim_passengers_ddl.sql b/sql/dds/dim_passengers_ddl.sql index 1d1de23..90a0afc 100644 --- a/sql/dds/dim_passengers_ddl.sql +++ b/sql/dds/dim_passengers_ddl.sql @@ -2,13 +2,19 @@ CREATE SCHEMA IF NOT EXISTS dds; +-- Тип таблицы: Heap (стандартная). +-- Обоснование: Необходим row-level UPDATE для реализации SCD1 UPSERT. +-- Использование Append-Only при частых обновлениях приводит к раздуванию (bloat) таблицы. CREATE TABLE IF NOT EXISTS dds.dim_passengers ( passenger_sk INTEGER NOT NULL, - passenger_bk TEXT NOT NULL, + passenger_id TEXT NOT NULL, passenger_name TEXT NOT NULL, created_at TIMESTAMP NOT NULL DEFAULT now(), updated_at TIMESTAMP NOT NULL DEFAULT now(), _load_id TEXT NOT NULL, _load_ts TIMESTAMP NOT NULL DEFAULT now() ) +WITH (appendonly=false) DISTRIBUTED BY (passenger_sk); + +COMMENT ON TABLE dds.dim_passengers IS 'Измерение пассажиров (DDS).'; diff --git a/sql/dds/dim_passengers_dq.sql b/sql/dds/dim_passengers_dq.sql index 8cbcaa5..990b365 100644 --- a/sql/dds/dim_passengers_dq.sql +++ b/sql/dds/dim_passengers_dq.sql @@ -41,14 +41,14 @@ BEGIN v_dup_sk; END IF; - -- Нет дублей по BK. - SELECT COUNT(*) - COUNT(DISTINCT passenger_bk) + -- Нет дублей по BK (ID). + SELECT COUNT(*) - COUNT(DISTINCT passenger_id) INTO v_dup_bk FROM dds.dim_passengers; IF v_dup_bk <> 0 THEN RAISE EXCEPTION - 'DQ FAILED: в dds.dim_passengers найдены дубликаты passenger_bk: %', + 'DQ FAILED: в dds.dim_passengers найдены дубликаты passenger_id: %', v_dup_bk; END IF; @@ -64,7 +64,7 @@ BEGIN WHERE NOT EXISTS ( SELECT 1 FROM dds.dim_passengers AS d - WHERE d.passenger_bk = s.passenger_id + WHERE d.passenger_id = s.passenger_id ); IF v_missing_bk <> 0 THEN @@ -78,8 +78,8 @@ BEGIN INTO v_null_count FROM dds.dim_passengers WHERE passenger_sk IS NULL - OR passenger_bk IS NULL - OR passenger_bk = '' + OR passenger_id IS NULL + OR passenger_id = '' OR passenger_name IS NULL OR passenger_name = '' OR created_at IS NULL diff --git a/sql/dds/dim_passengers_load.sql b/sql/dds/dim_passengers_load.sql index 6075317..9d06bcd 100644 --- a/sql/dds/dim_passengers_load.sql +++ b/sql/dds/dim_passengers_load.sql @@ -1,67 +1,48 @@ -- Загрузка DDS dim_passengers: SCD1 UPSERT (UPDATE изменившихся + INSERT новых). --- Statement 1: UPDATE существующих записей (если атрибуты изменились). -WITH src AS ( +-- Учебный комментарий: Используем TEMP TABLE для подготовки дельты. +-- Это избавляет от дублирования сложной оконной функции в UPDATE и INSERT блоках. +CREATE TEMP TABLE tmp_passengers_src ON COMMIT DROP AS +SELECT + d.passenger_id, + d.passenger_name +FROM ( SELECT - d.passenger_id, - d.passenger_name - FROM ( - SELECT - t.passenger_id, - t.passenger_name, - ROW_NUMBER() OVER ( - PARTITION BY t.passenger_id - ORDER BY t.event_ts DESC NULLS LAST, t._load_ts DESC, t.ticket_no DESC - ) AS rn - FROM ods.tickets AS t - WHERE t.passenger_id IS NOT NULL - AND t.passenger_id <> '' - AND t.passenger_name IS NOT NULL - AND t.passenger_name <> '' - ) AS d - WHERE d.rn = 1 -) + t.passenger_id, + t.passenger_name, + ROW_NUMBER() OVER ( + PARTITION BY t.passenger_id + ORDER BY t.event_ts DESC NULLS LAST, t._load_ts DESC, t.ticket_no DESC + ) AS rn + FROM ods.tickets AS t + WHERE t.passenger_id IS NOT NULL + AND t.passenger_id <> '' + AND t.passenger_name IS NOT NULL + AND t.passenger_name <> '' +) AS d +WHERE d.rn = 1; + +-- Statement 1: UPDATE существующих записей (если атрибуты изменились). UPDATE dds.dim_passengers AS d SET passenger_name = s.passenger_name, updated_at = now(), _load_id = '{{ run_id }}', _load_ts = now() -FROM src AS s -WHERE d.passenger_bk = s.passenger_id +FROM tmp_passengers_src AS s +WHERE d.passenger_id = s.passenger_id AND d.passenger_name IS DISTINCT FROM s.passenger_name; -- Statement 2: INSERT новых записей (MAX(sk) + ROW_NUMBER()). -WITH src AS ( - SELECT - d.passenger_id, - d.passenger_name - FROM ( - SELECT - t.passenger_id, - t.passenger_name, - ROW_NUMBER() OVER ( - PARTITION BY t.passenger_id - ORDER BY t.event_ts DESC NULLS LAST, t._load_ts DESC, t.ticket_no DESC - ) AS rn - FROM ods.tickets AS t - WHERE t.passenger_id IS NOT NULL - AND t.passenger_id <> '' - AND t.passenger_name IS NOT NULL - AND t.passenger_name <> '' - ) AS d - WHERE d.rn = 1 -), --- Учебный комментарий: Генерация SK через MAX() + ROW_NUMBER() --- Этот подход работает безопасно только потому, что Airflow запускает --- джобы загрузки для одной таблицы строго последовательно (concurrency=1). --- При параллельной загрузке возникнет состояние гонки (race condition) и возможны дубли SK. -max_sk AS ( +WITH max_sk AS ( + -- Учебный комментарий: Генерация SK через MAX() + ROW_NUMBER() + -- Этот подход работает безопасно только потому, что Airflow запускает + -- джобы загрузки для одной таблицы строго последовательно (concurrency=1). SELECT COALESCE(MAX(passenger_sk), 0) AS v FROM dds.dim_passengers ) INSERT INTO dds.dim_passengers ( passenger_sk, - passenger_bk, + passenger_id, passenger_name, created_at, updated_at, @@ -76,11 +57,11 @@ SELECT now(), '{{ run_id }}', now() -FROM src AS s +FROM tmp_passengers_src AS s WHERE NOT EXISTS ( SELECT 1 FROM dds.dim_passengers AS d - WHERE d.passenger_bk = s.passenger_id + WHERE d.passenger_id = s.passenger_id ); ANALYZE dds.dim_passengers; diff --git a/sql/dds/dim_routes_ddl.sql b/sql/dds/dim_routes_ddl.sql index 81892d4..0bd0021 100644 --- a/sql/dds/dim_routes_ddl.sql +++ b/sql/dds/dim_routes_ddl.sql @@ -2,6 +2,9 @@ CREATE SCHEMA IF NOT EXISTS dds; +-- Тип таблицы: Heap (стандартная). +-- Обоснование: Необходим row-level UPDATE для реализации SCD2 (закрытие версий). +-- Использование Append-Only при частых обновлениях приводит к раздуванию (bloat) таблицы. CREATE TABLE IF NOT EXISTS dds.dim_routes ( route_sk INTEGER NOT NULL, route_bk TEXT NOT NULL, @@ -19,4 +22,7 @@ CREATE TABLE IF NOT EXISTS dds.dim_routes ( _load_id TEXT NOT NULL, _load_ts TIMESTAMP NOT NULL DEFAULT now() ) +WITH (appendonly=false) DISTRIBUTED BY (route_sk); + +COMMENT ON TABLE dds.dim_routes IS 'Измерение маршрутов (DDS).'; diff --git a/sql/dds/dim_routes_load.sql b/sql/dds/dim_routes_load.sql index 76113f5..fac9b52 100644 --- a/sql/dds/dim_routes_load.sql +++ b/sql/dds/dim_routes_load.sql @@ -1,44 +1,40 @@ -- Загрузка DDS dim_routes: SCD2 с hashdiff. +-- Учебный комментарий: Используем TEMP TABLE для подготовки дельты. +-- Это избавляет от дублирования логики hashdiff в UPDATE и INSERT блоках. +CREATE TEMP TABLE tmp_routes_src ON COMMIT DROP AS +SELECT + route_no, + departure_airport, + arrival_airport, + airplane_code, + days_of_week, + departure_time, + duration, + md5( + COALESCE(departure_airport, '') || '|' || + COALESCE(arrival_airport, '') || '|' || + COALESCE(airplane_code, '') || '|' || + COALESCE(days_of_week::TEXT, '') || '|' || + COALESCE(departure_time::TEXT, '') || '|' || + COALESCE(duration::TEXT, '') + ) AS hashdiff, + ROW_NUMBER() OVER (PARTITION BY route_no ORDER BY validity DESC) AS rn +FROM ods.routes; + -- Statement 1: Закрыть устаревшие версии (valid_to = текущая дата). -WITH src AS ( - SELECT - route_no, - departure_airport, - arrival_airport, - airplane_code, - days_of_week, - departure_time, - duration, - md5( - COALESCE(departure_airport, '') || '|' || - COALESCE(arrival_airport, '') || '|' || - COALESCE(airplane_code, '') || '|' || - COALESCE(days_of_week, '') || '|' || - COALESCE(departure_time::TEXT, '') || '|' || - COALESCE(duration::TEXT, '') - ) AS hashdiff, - ROW_NUMBER() OVER (PARTITION BY route_no ORDER BY validity DESC) AS rn - FROM ods.routes -) UPDATE dds.dim_routes AS d SET valid_to = CURRENT_DATE, updated_at = now(), _load_id = '{{ run_id }}', _load_ts = now() -FROM src AS s +FROM tmp_routes_src AS s WHERE s.rn = 1 AND d.route_bk = s.route_no AND d.valid_to IS NULL AND d.hashdiff <> s.hashdiff; -- Statement 1.1: Закрыть "исчезнувшие" маршруты. -WITH src AS ( - SELECT - route_no, - ROW_NUMBER() OVER (PARTITION BY route_no ORDER BY validity DESC) AS rn - FROM ods.routes -) UPDATE dds.dim_routes AS d SET valid_to = CURRENT_DATE, updated_at = now(), @@ -47,37 +43,16 @@ SET valid_to = CURRENT_DATE, WHERE d.valid_to IS NULL AND NOT EXISTS ( SELECT 1 - FROM src AS s + FROM tmp_routes_src AS s WHERE s.rn = 1 AND s.route_no = d.route_bk ); -- Statement 2: Вставить новые версии (для изменённых и новых route_no). -WITH src AS ( - SELECT - route_no, - departure_airport, - arrival_airport, - airplane_code, - days_of_week, - departure_time, - duration, - md5( - COALESCE(departure_airport, '') || '|' || - COALESCE(arrival_airport, '') || '|' || - COALESCE(airplane_code, '') || '|' || - COALESCE(days_of_week, '') || '|' || - COALESCE(departure_time::TEXT, '') || '|' || - COALESCE(duration::TEXT, '') - ) AS hashdiff, - ROW_NUMBER() OVER (PARTITION BY route_no ORDER BY validity DESC) AS rn - FROM ods.routes -), -max_sk AS ( +WITH max_sk AS ( -- Учебный комментарий: Генерация SK через MAX() + ROW_NUMBER() -- Этот подход работает безопасно только потому, что Airflow запускает -- джобы загрузки для одной таблицы строго последовательно (concurrency=1). - -- При параллельной загрузке возникнет состояние гонки (race condition) и возможны дубли SK. SELECT COALESCE(MAX(route_sk), 0) AS v FROM dds.dim_routes ) @@ -121,7 +96,7 @@ SELECT now(), '{{ run_id }}', now() -FROM src AS s +FROM tmp_routes_src AS s WHERE s.rn = 1 AND NOT EXISTS ( SELECT 1 diff --git a/sql/dds/dim_tariffs_ddl.sql b/sql/dds/dim_tariffs_ddl.sql index bd93796..f14cc89 100644 --- a/sql/dds/dim_tariffs_ddl.sql +++ b/sql/dds/dim_tariffs_ddl.sql @@ -2,6 +2,8 @@ CREATE SCHEMA IF NOT EXISTS dds; +-- Тип таблицы: Append-Only Row-oriented (zstd:1). +-- Обоснование: Редко дополняемые данные без обновлений. Обеспечивает эффективное сжатие. CREATE TABLE IF NOT EXISTS dds.dim_tariffs ( tariff_sk INTEGER NOT NULL, fare_conditions TEXT NOT NULL, @@ -10,4 +12,7 @@ CREATE TABLE IF NOT EXISTS dds.dim_tariffs ( _load_id TEXT NOT NULL, _load_ts TIMESTAMP NOT NULL DEFAULT now() ) +WITH (appendonly=true, orientation=row, compresstype=zstd, compresslevel=1) DISTRIBUTED BY (tariff_sk); + +COMMENT ON TABLE dds.dim_tariffs IS 'Измерение тарифов (DDS).'; diff --git a/sql/dds/fact_flight_sales_ddl.sql b/sql/dds/fact_flight_sales_ddl.sql index 28a6703..ba449f7 100644 --- a/sql/dds/fact_flight_sales_ddl.sql +++ b/sql/dds/fact_flight_sales_ddl.sql @@ -6,6 +6,10 @@ CREATE SCHEMA IF NOT EXISTS dds; -- В классическом DWH (Кимбалл) таблица фактов идентифицируется набором её -- измерений или дегенеративных ключей (в нашем случае: ticket_no + flight_id). -- Добавление отдельного ID только тратит место и не несёт аналитической ценности. + +-- Тип таблицы: Heap (стандартная). +-- Обоснование: Необходим row-level UPDATE для обновления статусов (is_boarded). +-- Использование Append-Only при частых обновлениях приводит к раздуванию (bloat) таблицы. CREATE TABLE IF NOT EXISTS dds.fact_flight_sales ( calendar_sk INTEGER, departure_airport_sk INTEGER, @@ -24,4 +28,7 @@ CREATE TABLE IF NOT EXISTS dds.fact_flight_sales ( _load_id TEXT NOT NULL, _load_ts TIMESTAMP NOT NULL DEFAULT now() ) +WITH (appendonly=false) DISTRIBUTED BY (ticket_no); + +COMMENT ON TABLE dds.fact_flight_sales IS 'Факт продаж билетов (DDS).'; diff --git a/sql/dds/fact_flight_sales_load.sql b/sql/dds/fact_flight_sales_load.sql index 225db3a..1aa66ab 100644 --- a/sql/dds/fact_flight_sales_load.sql +++ b/sql/dds/fact_flight_sales_load.sql @@ -4,7 +4,7 @@ -- Обновляем только мутабельные поля; SK измерений не перезаписываем. UPDATE dds.fact_flight_sales AS f SET seat_no = bp.seat_no, - price = seg.segment_amount, + price = seg.amount, is_boarded = (bp.ticket_no IS NOT NULL), _load_id = '{{ run_id }}', _load_ts = now() @@ -16,7 +16,7 @@ WHERE f.ticket_no = seg.ticket_no AND f.flight_id = seg.flight_id AND ( f.is_boarded IS DISTINCT FROM (bp.ticket_no IS NOT NULL) - OR f.price IS DISTINCT FROM seg.segment_amount + OR f.price IS DISTINCT FROM seg.amount OR f.seat_no IS DISTINCT FROM bp.seat_no ); @@ -36,7 +36,7 @@ WITH fact_src AS ( tkt.book_ref, bkg.book_date::DATE AS book_date, bp.seat_no, - seg.segment_amount AS price, + seg.amount AS price, (bp.ticket_no IS NOT NULL) AS is_boarded FROM ods.segments AS seg JOIN ods.tickets AS tkt @@ -48,8 +48,6 @@ WITH fact_src AS ( -- Учебный комментарий: Late-arriving dimensions (Опаздывающие измерения) -- Мы используем LEFT JOIN, так как факт (рейс/билет) может прийти раньше, -- чем справочник (пассажир/маршрут) обновится в DDS. - -- В результате SK будет NULL. В более сложных пайплайнах такие факты - -- либо обогащаются dummy-значениями (-1, "Неизвестно"), либо откладываются. LEFT JOIN dds.dim_routes AS rte ON rte.route_bk = flt.route_no AND flt.scheduled_departure::DATE >= rte.valid_from @@ -65,7 +63,7 @@ WITH fact_src AS ( LEFT JOIN dds.dim_tariffs AS tar ON tar.fare_conditions = seg.fare_conditions LEFT JOIN dds.dim_passengers AS pax - ON pax.passenger_bk = tkt.passenger_id + ON pax.passenger_id = tkt.passenger_id LEFT JOIN ods.boarding_passes AS bp ON bp.ticket_no = seg.ticket_no AND bp.flight_id = seg.flight_id diff --git a/sql/dm/sales_report_ddl.sql b/sql/dm/sales_report_ddl.sql index 3ce6db0..5739b08 100644 --- a/sql/dm/sales_report_ddl.sql +++ b/sql/dm/sales_report_ddl.sql @@ -6,7 +6,6 @@ -- Паттерны для студентов: -- - Денормализация измерений (города, аэропорты, тарифы) для удобства аналитики -- - Служебные поля календаря (day_of_week, day_name, is_weekend) --- - Heap-таблица с UPDATE (нужен для UPSERT) -- -- ВНИМАНИЕ (Антипаттерн распределения): -- Никогда не распределяйте таблицы по дате (DISTRIBUTED BY flight_date) в MPP-системах! @@ -19,6 +18,9 @@ CREATE SCHEMA IF NOT EXISTS dm; +-- Тип таблицы: Heap (стандартная). +-- Обоснование: Необходим row-level UPDATE для реализации инкрементального UPSERT по HWM. +-- Использование Append-Only при частых обновлениях приводит к раздуванию (bloat) таблицы. CREATE TABLE IF NOT EXISTS dm.sales_report ( -- Ключ (зерно витрины) flight_date DATE NOT NULL, @@ -53,11 +55,11 @@ CREATE TABLE IF NOT EXISTS dm.sales_report ( _load_id TEXT NOT NULL, _load_ts TIMESTAMP NOT NULL DEFAULT now() ) +WITH (appendonly=false) DISTRIBUTED BY (departure_airport_sk, arrival_airport_sk); -- Комментарии для документирования -COMMENT ON TABLE dm.sales_report IS - 'Витрина продаж: выручка, билеты и boarding rate по направлениям/тарифам/дням'; +COMMENT ON TABLE dm.sales_report IS 'Витрина продаж: выручка, билеты и boarding rate по направлениям/тарифам/дням.'; COMMENT ON COLUMN dm.sales_report.boarding_rate IS 'Доля пассажиров, прошедших посадку (passengers_boarded / tickets_sold)'; diff --git a/sql/ods/airplanes_ddl.sql b/sql/ods/airplanes_ddl.sql index 56aa6dc..b2b2ec3 100644 --- a/sql/ods/airplanes_ddl.sql +++ b/sql/ods/airplanes_ddl.sql @@ -2,6 +2,9 @@ CREATE SCHEMA IF NOT EXISTS ods; +-- Тип таблицы: Append-Only Row-oriented (zstd:1). +-- Обоснование: Используется паттерн TRUNCATE+INSERT (полный снимок). +-- Для узких таблиц Row-store производительнее Column-store при чтении всей строки. CREATE TABLE IF NOT EXISTS ods.airplanes ( airplane_code TEXT NOT NULL, model TEXT NOT NULL, @@ -10,4 +13,7 @@ CREATE TABLE IF NOT EXISTS ods.airplanes ( _load_id TEXT NOT NULL, _load_ts TIMESTAMP NOT NULL DEFAULT now() ) +WITH (appendonly=true, orientation=row, compresstype=zstd, compresslevel=1) DISTRIBUTED BY (airplane_code); + +COMMENT ON TABLE ods.airplanes IS 'Справочник моделей самолетов (ODS).'; diff --git a/sql/ods/airplanes_load.sql b/sql/ods/airplanes_load.sql index 78ecae2..ec7064a 100644 --- a/sql/ods/airplanes_load.sql +++ b/sql/ods/airplanes_load.sql @@ -1,52 +1,9 @@ --- Загрузка ODS по airplanes: SCD1 (UPDATE изменившихся + INSERT новых). +-- Загрузка ODS по airplanes: Полная перезагрузка (TRUNCATE + INSERT). +-- Почему: для справочников-снимков в Greenplum на AO-таблицах +-- эффективнее перетереть данные целиком, чем делать медленный UPDATE. --- Statement 1: UPDATE существующих строк. --- Нормализация JSON: извлекаем русское название модели из поля с мультиязычностью. --- Почему: источник хранит переводы как {"en": "...", "ru": "..."}, --- в ODS оставляем только один язык для упрощения downstream-логики. -WITH src AS ( - SELECT - s.airplane_code, - s.model::json->>'ru' AS model, - NULLIF(s.range, '')::INTEGER AS range_km, - NULLIF(s.speed, '')::INTEGER AS speed_kmh, - ROW_NUMBER() OVER ( - PARTITION BY s.airplane_code - ORDER BY s.load_dttm DESC, s.src_created_at_ts DESC NULLS LAST - ) AS rn - FROM stg.airplanes AS s - WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text -) -UPDATE ods.airplanes AS o -SET model = s.model, - range_km = s.range_km, - speed_kmh = s.speed_kmh, - _load_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, - _load_ts = now() -FROM src AS s -WHERE s.rn = 1 - AND o.airplane_code = s.airplane_code - AND ( - o.model IS DISTINCT FROM s.model - OR o.range_km IS DISTINCT FROM s.range_km - OR o.speed_kmh IS DISTINCT FROM s.speed_kmh - ); +TRUNCATE TABLE ods.airplanes; --- Statement 2: INSERT новых строк. --- Нормализация JSON: извлекаем русское название модели из поля с мультиязычностью. -WITH src AS ( - SELECT - s.airplane_code, - s.model::json->>'ru' AS model, - NULLIF(s.range, '')::INTEGER AS range_km, - NULLIF(s.speed, '')::INTEGER AS speed_kmh, - ROW_NUMBER() OVER ( - PARTITION BY s.airplane_code - ORDER BY s.load_dttm DESC, s.src_created_at_ts DESC NULLS LAST - ) AS rn - FROM stg.airplanes AS s - WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text -) INSERT INTO ods.airplanes ( airplane_code, model, @@ -55,6 +12,20 @@ INSERT INTO ods.airplanes ( _load_id, _load_ts ) +WITH src AS ( + -- Выбираем последний снимок из STG для текущего батча + SELECT + s.airplane_code, + s.model::json->>'ru' AS model, + NULLIF(s.range, '')::INTEGER AS range_km, + NULLIF(s.speed, '')::INTEGER AS speed_kmh, + ROW_NUMBER() OVER ( + PARTITION BY s.airplane_code + ORDER BY s.load_dttm DESC, s.src_created_at_ts DESC NULLS LAST + ) AS rn + FROM stg.airplanes AS s + WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text +) SELECT s.airplane_code, s.model, @@ -63,34 +34,6 @@ SELECT '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, now() FROM src AS s -WHERE s.rn = 1 - AND NOT EXISTS ( - SELECT 1 - FROM ods.airplanes AS o - WHERE o.airplane_code = s.airplane_code - ); - --- Statement 3: DELETE ключей, которых нет в snapshot текущего батча. --- Это делает ODS для справочника действительно "current state". -WITH src_keys AS ( - SELECT d.airplane_code - FROM ( - SELECT - s.airplane_code, - ROW_NUMBER() OVER ( - PARTITION BY s.airplane_code - ORDER BY s.load_dttm DESC, s.src_created_at_ts DESC NULLS LAST - ) AS rn - FROM stg.airplanes AS s - WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text - ) AS d - WHERE d.rn = 1 -) -DELETE FROM ods.airplanes AS o -WHERE NOT EXISTS ( - SELECT 1 - FROM src_keys AS s - WHERE s.airplane_code = o.airplane_code -); +WHERE s.rn = 1; ANALYZE ods.airplanes; diff --git a/sql/ods/airports_ddl.sql b/sql/ods/airports_ddl.sql index a6663c8..5db43e2 100644 --- a/sql/ods/airports_ddl.sql +++ b/sql/ods/airports_ddl.sql @@ -2,6 +2,9 @@ CREATE SCHEMA IF NOT EXISTS ods; +-- Тип таблицы: Append-Only Row-oriented (zstd:1). +-- Обоснование: Используется паттерн TRUNCATE+INSERT (полный снимок). +-- Для узких таблиц Row-store производительнее Column-store при чтении всей строки. CREATE TABLE IF NOT EXISTS ods.airports ( airport_code TEXT NOT NULL, airport_name TEXT NOT NULL, @@ -12,4 +15,7 @@ CREATE TABLE IF NOT EXISTS ods.airports ( _load_id TEXT NOT NULL, _load_ts TIMESTAMP NOT NULL DEFAULT now() ) +WITH (appendonly=true, orientation=row, compresstype=zstd, compresslevel=1) DISTRIBUTED BY (airport_code); + +COMMENT ON TABLE ods.airports IS 'Справочник аэропортов (ODS).'; diff --git a/sql/ods/airports_load.sql b/sql/ods/airports_load.sql index e4a9c97..6f11239 100644 --- a/sql/ods/airports_load.sql +++ b/sql/ods/airports_load.sql @@ -1,60 +1,9 @@ --- Загрузка ODS по airports: SCD1 (UPDATE изменившихся + INSERT новых). +-- Загрузка ODS по airports: Полная перезагрузка (TRUNCATE + INSERT). +-- Почему: для справочников-снимков в Greenplum на AO-таблицах +-- эффективнее перетереть данные целиком, чем делать медленный UPDATE. --- Statement 1: UPDATE существующих строк. --- Нормализация JSON: извлекаем русские названия из полей с мультиязычностью. --- Почему: источник хранит переводы как {"en": "...", "ru": "..."}, --- в ODS оставляем только один язык для упрощения downstream-логики. -WITH src AS ( - SELECT - s.airport_code, - s.airport_name::json->>'ru' AS airport_name, - s.city::json->>'ru' AS city, - s.country::json->>'ru' AS country, - s.coordinates, - s.timezone, - ROW_NUMBER() OVER ( - PARTITION BY s.airport_code - ORDER BY s.load_dttm DESC, s.src_created_at_ts DESC NULLS LAST - ) AS rn - FROM stg.airports AS s - WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text -) -UPDATE ods.airports AS o -SET airport_name = s.airport_name, - city = s.city, - country = s.country, - coordinates = s.coordinates, - timezone = s.timezone, - _load_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, - _load_ts = now() -FROM src AS s -WHERE s.rn = 1 - AND o.airport_code = s.airport_code - AND ( - o.airport_name IS DISTINCT FROM s.airport_name - OR o.city IS DISTINCT FROM s.city - OR o.country IS DISTINCT FROM s.country - OR o.coordinates IS DISTINCT FROM s.coordinates - OR o.timezone IS DISTINCT FROM s.timezone - ); +TRUNCATE TABLE ods.airports; --- Statement 2: INSERT новых строк. --- Нормализация JSON: извлекаем русские названия из полей с мультиязычностью. -WITH src AS ( - SELECT - s.airport_code, - s.airport_name::json->>'ru' AS airport_name, - s.city::json->>'ru' AS city, - s.country::json->>'ru' AS country, - s.coordinates, - s.timezone, - ROW_NUMBER() OVER ( - PARTITION BY s.airport_code - ORDER BY s.load_dttm DESC, s.src_created_at_ts DESC NULLS LAST - ) AS rn - FROM stg.airports AS s - WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text -) INSERT INTO ods.airports ( airport_code, airport_name, @@ -65,6 +14,22 @@ INSERT INTO ods.airports ( _load_id, _load_ts ) +WITH src AS ( + -- Выбираем последний снимок из STG для текущего батча + SELECT + s.airport_code, + s.airport_name::json->>'ru' AS airport_name, + s.city::json->>'ru' AS city, + s.country::json->>'ru' AS country, + s.coordinates, + s.timezone, + ROW_NUMBER() OVER ( + PARTITION BY s.airport_code + ORDER BY s.load_dttm DESC, s.src_created_at_ts DESC NULLS LAST + ) AS rn + FROM stg.airports AS s + WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text +) SELECT s.airport_code, s.airport_name, @@ -75,34 +40,6 @@ SELECT '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, now() FROM src AS s -WHERE s.rn = 1 - AND NOT EXISTS ( - SELECT 1 - FROM ods.airports AS o - WHERE o.airport_code = s.airport_code - ); - --- Statement 3: DELETE ключей, которых нет в snapshot текущего батча. --- Это делает ODS для справочника действительно "current state". -WITH src_keys AS ( - SELECT d.airport_code - FROM ( - SELECT - s.airport_code, - ROW_NUMBER() OVER ( - PARTITION BY s.airport_code - ORDER BY s.load_dttm DESC, s.src_created_at_ts DESC NULLS LAST - ) AS rn - FROM stg.airports AS s - WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text - ) AS d - WHERE d.rn = 1 -) -DELETE FROM ods.airports AS o -WHERE NOT EXISTS ( - SELECT 1 - FROM src_keys AS s - WHERE s.airport_code = o.airport_code -); +WHERE s.rn = 1; ANALYZE ods.airports; diff --git a/sql/ods/boarding_passes_ddl.sql b/sql/ods/boarding_passes_ddl.sql index 8a8640e..a159432 100644 --- a/sql/ods/boarding_passes_ddl.sql +++ b/sql/ods/boarding_passes_ddl.sql @@ -2,6 +2,9 @@ CREATE SCHEMA IF NOT EXISTS ods; +-- Тип таблицы: Heap (стандартная). +-- Обоснование: Необходим row-level UPDATE для реализации UPSERT/SCD1. +-- Использование Append-Only при частых обновлениях приводит к раздуванию (bloat) таблицы. CREATE TABLE IF NOT EXISTS ods.boarding_passes ( ticket_no TEXT NOT NULL, flight_id INTEGER NOT NULL, @@ -12,4 +15,7 @@ CREATE TABLE IF NOT EXISTS ods.boarding_passes ( _load_id TEXT NOT NULL, _load_ts TIMESTAMP NOT NULL DEFAULT now() ) +WITH (appendonly=false) DISTRIBUTED BY (ticket_no); + +COMMENT ON TABLE ods.boarding_passes IS 'Посадочные талоны (ODS).'; diff --git a/sql/ods/bookings_ddl.sql b/sql/ods/bookings_ddl.sql index 1ee839f..e8f6dcb 100644 --- a/sql/ods/bookings_ddl.sql +++ b/sql/ods/bookings_ddl.sql @@ -2,6 +2,9 @@ CREATE SCHEMA IF NOT EXISTS ods; +-- Тип таблицы: Heap (стандартная). +-- Обоснование: Необходим row-level UPDATE для реализации UPSERT/SCD1. +-- Использование Append-Only при частых обновлениях приводит к раздуванию (bloat) таблицы. CREATE TABLE IF NOT EXISTS ods.bookings ( book_ref TEXT NOT NULL, book_date TIMESTAMP WITH TIME ZONE NOT NULL, @@ -10,4 +13,7 @@ CREATE TABLE IF NOT EXISTS ods.bookings ( _load_id TEXT NOT NULL, _load_ts TIMESTAMP NOT NULL DEFAULT now() ) +WITH (appendonly=false) DISTRIBUTED BY (book_ref); + +COMMENT ON TABLE ods.bookings IS 'Бронирования (ODS).'; diff --git a/sql/ods/flights_ddl.sql b/sql/ods/flights_ddl.sql index 5b49996..84c6e04 100644 --- a/sql/ods/flights_ddl.sql +++ b/sql/ods/flights_ddl.sql @@ -2,6 +2,9 @@ CREATE SCHEMA IF NOT EXISTS ods; +-- Тип таблицы: Heap (стандартная). +-- Обоснование: Необходим row-level UPDATE для реализации UPSERT/SCD1. +-- Использование Append-Only при частых обновлениях приводит к раздуванию (bloat) таблицы. CREATE TABLE IF NOT EXISTS ods.flights ( flight_id INTEGER NOT NULL, route_no TEXT NOT NULL, @@ -14,4 +17,7 @@ CREATE TABLE IF NOT EXISTS ods.flights ( _load_id TEXT NOT NULL, _load_ts TIMESTAMP NOT NULL DEFAULT now() ) +WITH (appendonly=false) DISTRIBUTED BY (flight_id); + +COMMENT ON TABLE ods.flights IS 'Рейсы (ODS).'; diff --git a/sql/ods/routes_ddl.sql b/sql/ods/routes_ddl.sql index 5f3d486..ac757c9 100644 --- a/sql/ods/routes_ddl.sql +++ b/sql/ods/routes_ddl.sql @@ -2,16 +2,22 @@ CREATE SCHEMA IF NOT EXISTS ods; +-- Тип таблицы: Append-Only Row-oriented (zstd:1). +-- Обоснование: Используется паттерн TRUNCATE+INSERT (полный снимок). +-- Для узких таблиц Row-store производительнее Column-store при чтении всей строки. CREATE TABLE IF NOT EXISTS ods.routes ( route_no TEXT NOT NULL, validity TEXT NOT NULL, departure_airport TEXT NOT NULL, arrival_airport TEXT NOT NULL, airplane_code TEXT NOT NULL, - days_of_week TEXT, - departure_time TIME, - duration INTERVAL, + days_of_week INTEGER[], + departure_time TIME NOT NULL, + duration INTERVAL NOT NULL, _load_id TEXT NOT NULL, _load_ts TIMESTAMP NOT NULL DEFAULT now() ) -DISTRIBUTED BY (route_no); +WITH (appendonly=true, orientation=row, compresstype=zstd, compresslevel=1) +DISTRIBUTED BY (route_no, validity); + +COMMENT ON TABLE ods.routes IS 'Справочник маршрутов (ODS).'; diff --git a/sql/ods/routes_load.sql b/sql/ods/routes_load.sql index c51a160..f692494 100644 --- a/sql/ods/routes_load.sql +++ b/sql/ods/routes_load.sql @@ -1,63 +1,9 @@ --- Загрузка ODS по routes: SCD1 (UPDATE изменившихся + INSERT новых). +-- Загрузка ODS по routes: Полная перезагрузка (TRUNCATE + INSERT). +-- Почему: для справочников-снимков в Greenplum на AO-таблицах +-- эффективнее перетереть данные целиком, чем делать медленный UPDATE. --- Statement 1: UPDATE существующих строк. -WITH src AS ( - SELECT - s.route_no, - s.validity, - s.departure_airport, - s.arrival_airport, - s.airplane_code, - s.days_of_week, - NULLIF(s.scheduled_time, '')::TIME AS departure_time, - NULLIF(s.duration, '')::INTERVAL AS duration, - ROW_NUMBER() OVER ( - PARTITION BY s.route_no, s.validity - ORDER BY s.load_dttm DESC, s.src_created_at_ts DESC NULLS LAST - ) AS rn - FROM stg.routes AS s - WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text -) -UPDATE ods.routes AS o -SET departure_airport = s.departure_airport, - arrival_airport = s.arrival_airport, - airplane_code = s.airplane_code, - days_of_week = s.days_of_week, - departure_time = s.departure_time, - duration = s.duration, - _load_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, - _load_ts = now() -FROM src AS s -WHERE s.rn = 1 - AND o.route_no = s.route_no - AND o.validity = s.validity - AND ( - o.departure_airport IS DISTINCT FROM s.departure_airport - OR o.arrival_airport IS DISTINCT FROM s.arrival_airport - OR o.airplane_code IS DISTINCT FROM s.airplane_code - OR o.days_of_week IS DISTINCT FROM s.days_of_week - OR o.departure_time IS DISTINCT FROM s.departure_time - OR o.duration IS DISTINCT FROM s.duration - ); +TRUNCATE TABLE ods.routes; --- Statement 2: INSERT новых строк. -WITH src AS ( - SELECT - s.route_no, - s.validity, - s.departure_airport, - s.arrival_airport, - s.airplane_code, - s.days_of_week, - NULLIF(s.scheduled_time, '')::TIME AS departure_time, - NULLIF(s.duration, '')::INTERVAL AS duration, - ROW_NUMBER() OVER ( - PARTITION BY s.route_no, s.validity - ORDER BY s.load_dttm DESC, s.src_created_at_ts DESC NULLS LAST - ) AS rn - FROM stg.routes AS s - WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text -) INSERT INTO ods.routes ( route_no, validity, @@ -70,49 +16,36 @@ INSERT INTO ods.routes ( _load_id, _load_ts ) +WITH src AS ( + -- Выбираем последний снимок из STG для текущего батча + SELECT + s.route_no, + s.validity, + s.departure_airport, + s.arrival_airport, + s.airplane_code, + s.days_of_week::INTEGER[] AS days_of_week, + NULLIF(s.scheduled_time, '')::TIME AS departure_time, + NULLIF(s.duration, '')::INTERVAL AS duration, + ROW_NUMBER() OVER ( + PARTITION BY s.route_no, s.validity + ORDER BY s.load_dttm DESC, s.src_created_at_ts DESC NULLS LAST + ) AS rn + FROM stg.routes AS s + WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text +) SELECT s.route_no, s.validity, s.departure_airport, s.arrival_airport, s.airplane_code, - s.days_of_week, + s.days_of_week::INTEGER[], s.departure_time, s.duration, '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, now() FROM src AS s -WHERE s.rn = 1 - AND NOT EXISTS ( - SELECT 1 - FROM ods.routes AS o - WHERE o.route_no = s.route_no - AND o.validity = s.validity - ); - --- Statement 3: DELETE ключей, которых нет в snapshot текущего батча. --- Это делает ODS для справочника действительно "current state". -WITH src_keys AS ( - SELECT d.route_no, d.validity - FROM ( - SELECT - s.route_no, - s.validity, - ROW_NUMBER() OVER ( - PARTITION BY s.route_no, s.validity - ORDER BY s.load_dttm DESC, s.src_created_at_ts DESC NULLS LAST - ) AS rn - FROM stg.routes AS s - WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text - ) AS d - WHERE d.rn = 1 -) -DELETE FROM ods.routes AS o -WHERE NOT EXISTS ( - SELECT 1 - FROM src_keys AS s - WHERE s.route_no = o.route_no - AND s.validity = o.validity -); +WHERE s.rn = 1; ANALYZE ods.routes; diff --git a/sql/ods/seats_ddl.sql b/sql/ods/seats_ddl.sql index 08eaceb..484de68 100644 --- a/sql/ods/seats_ddl.sql +++ b/sql/ods/seats_ddl.sql @@ -2,11 +2,17 @@ CREATE SCHEMA IF NOT EXISTS ods; +-- Тип таблицы: Append-Only Row-oriented (zstd:1). +-- Обоснование: Используется паттерн TRUNCATE+INSERT (полный снимок). +-- Для узких таблиц Row-store производительнее Column-store при чтении всей строки. CREATE TABLE IF NOT EXISTS ods.seats ( - airplane_code TEXT NOT NULL, - seat_no TEXT NOT NULL, - fare_conditions TEXT NOT NULL, - _load_id TEXT NOT NULL, - _load_ts TIMESTAMP NOT NULL DEFAULT now() + airplane_code TEXT NOT NULL, + seat_no TEXT NOT NULL, + fare_conditions TEXT NOT NULL, + _load_id TEXT NOT NULL, + _load_ts TIMESTAMP NOT NULL DEFAULT now() ) +WITH (appendonly=true, orientation=row, compresstype=zstd, compresslevel=1) DISTRIBUTED BY (airplane_code); + +COMMENT ON TABLE ods.seats IS 'Справочник мест в самолетах (ODS).'; diff --git a/sql/ods/seats_load.sql b/sql/ods/seats_load.sql index 82fbf9f..7544958 100644 --- a/sql/ods/seats_load.sql +++ b/sql/ods/seats_load.sql @@ -1,41 +1,9 @@ --- Загрузка ODS по seats: SCD1 (UPDATE изменившихся + INSERT новых). +-- Загрузка ODS по seats: Полная перезагрузка (TRUNCATE + INSERT). +-- Почему: для справочников-снимков в Greenplum на AO-таблицах +-- эффективнее перетереть данные целиком, чем делать медленный UPDATE. --- Statement 1: UPDATE существующих строк. -WITH src AS ( - SELECT - s.airplane_code, - s.seat_no, - s.fare_conditions, - ROW_NUMBER() OVER ( - PARTITION BY s.airplane_code, s.seat_no - ORDER BY s.load_dttm DESC, s.src_created_at_ts DESC NULLS LAST - ) AS rn - FROM stg.seats AS s - WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text -) -UPDATE ods.seats AS o -SET fare_conditions = s.fare_conditions, - _load_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, - _load_ts = now() -FROM src AS s -WHERE s.rn = 1 - AND o.airplane_code = s.airplane_code - AND o.seat_no = s.seat_no - AND o.fare_conditions IS DISTINCT FROM s.fare_conditions; +TRUNCATE TABLE ods.seats; --- Statement 2: INSERT новых строк. -WITH src AS ( - SELECT - s.airplane_code, - s.seat_no, - s.fare_conditions, - ROW_NUMBER() OVER ( - PARTITION BY s.airplane_code, s.seat_no - ORDER BY s.load_dttm DESC, s.src_created_at_ts DESC NULLS LAST - ) AS rn - FROM stg.seats AS s - WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text -) INSERT INTO ods.seats ( airplane_code, seat_no, @@ -43,6 +11,19 @@ INSERT INTO ods.seats ( _load_id, _load_ts ) +WITH src AS ( + -- Выбираем последний снимок из STG для текущего батча + SELECT + s.airplane_code, + s.seat_no, + s.fare_conditions, + ROW_NUMBER() OVER ( + PARTITION BY s.airplane_code, s.seat_no + ORDER BY s.load_dttm DESC, s.src_created_at_ts DESC NULLS LAST + ) AS rn + FROM stg.seats AS s + WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text +) SELECT s.airplane_code, s.seat_no, @@ -50,37 +31,6 @@ SELECT '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text, now() FROM src AS s -WHERE s.rn = 1 - AND NOT EXISTS ( - SELECT 1 - FROM ods.seats AS o - WHERE o.airplane_code = s.airplane_code - AND o.seat_no = s.seat_no - ); - --- Statement 3: DELETE ключей, которых нет в snapshot текущего батча. --- Это делает ODS для справочника действительно "current state". -WITH src_keys AS ( - SELECT d.airplane_code, d.seat_no - FROM ( - SELECT - s.airplane_code, - s.seat_no, - ROW_NUMBER() OVER ( - PARTITION BY s.airplane_code, s.seat_no - ORDER BY s.load_dttm DESC, s.src_created_at_ts DESC NULLS LAST - ) AS rn - FROM stg.seats AS s - WHERE s.batch_id = '{{ ti.xcom_pull(task_ids="resolve_stg_batch_id") }}'::text - ) AS d - WHERE d.rn = 1 -) -DELETE FROM ods.seats AS o -WHERE NOT EXISTS ( - SELECT 1 - FROM src_keys AS s - WHERE s.airplane_code = o.airplane_code - AND s.seat_no = o.seat_no -); +WHERE s.rn = 1; ANALYZE ods.seats; diff --git a/sql/ods/segments_ddl.sql b/sql/ods/segments_ddl.sql index b35034a..3a2521a 100644 --- a/sql/ods/segments_ddl.sql +++ b/sql/ods/segments_ddl.sql @@ -2,13 +2,20 @@ CREATE SCHEMA IF NOT EXISTS ods; +-- Тип таблицы: Heap (стандартная). +-- Обоснование: Необходим row-level UPDATE для реализации UPSERT/SCD1. +-- Использование Append-Only при частых обновлениях приводит к раздуванию (bloat) таблицы. CREATE TABLE IF NOT EXISTS ods.segments ( - ticket_no TEXT NOT NULL, - flight_id INTEGER NOT NULL, - fare_conditions TEXT NOT NULL, - segment_amount NUMERIC(10,2), - event_ts TIMESTAMP, - _load_id TEXT NOT NULL, - _load_ts TIMESTAMP NOT NULL DEFAULT now() + ticket_no TEXT NOT NULL, + flight_id INTEGER NOT NULL, + fare_conditions TEXT NOT NULL, + amount NUMERIC(10,2) NOT NULL, + event_ts TIMESTAMP, + _load_id TEXT NOT NULL, + _load_ts TIMESTAMP NOT NULL DEFAULT now() ) +WITH (appendonly=false) DISTRIBUTED BY (ticket_no); + +COMMENT ON TABLE ods.segments IS 'Сегменты перелета (ODS).'; + diff --git a/sql/ods/segments_load.sql b/sql/ods/segments_load.sql index 140318f..e6aa2d2 100644 --- a/sql/ods/segments_load.sql +++ b/sql/ods/segments_load.sql @@ -8,7 +8,7 @@ WITH src AS ( s.ticket_no, NULLIF(s.flight_id, '')::INTEGER AS flight_id, s.fare_conditions, - NULLIF(s.price, '')::NUMERIC(10,2) AS segment_amount, + NULLIF(s.price, '')::NUMERIC(10,2) AS amount, s.src_created_at_ts AS event_ts, s.batch_id, s.load_dttm, @@ -25,7 +25,7 @@ SELECT * FROM src WHERE rn = 1; -- 2. UPDATE существующих строк. UPDATE ods.segments AS o SET fare_conditions = s.fare_conditions, - segment_amount = s.segment_amount, + amount = s.amount, event_ts = s.event_ts, _load_id = s.batch_id, -- Сохраняем оригинальный lineage из STG _load_ts = s.load_dttm -- Фиксируем время STG как водяной знак для ODS @@ -34,7 +34,7 @@ WHERE o.ticket_no = s.ticket_no AND o.flight_id = s.flight_id AND ( o.fare_conditions IS DISTINCT FROM s.fare_conditions - OR o.segment_amount IS DISTINCT FROM s.segment_amount + OR o.amount IS DISTINCT FROM s.amount OR o.event_ts IS DISTINCT FROM s.event_ts ); @@ -43,7 +43,7 @@ INSERT INTO ods.segments ( ticket_no, flight_id, fare_conditions, - segment_amount, + amount, event_ts, _load_id, _load_ts @@ -52,7 +52,7 @@ SELECT s.ticket_no, s.flight_id, s.fare_conditions, - s.segment_amount, + s.amount, s.event_ts, s.batch_id, s.load_dttm diff --git a/sql/ods/tickets_ddl.sql b/sql/ods/tickets_ddl.sql index a4a70ee..25e8b2c 100644 --- a/sql/ods/tickets_ddl.sql +++ b/sql/ods/tickets_ddl.sql @@ -2,6 +2,9 @@ CREATE SCHEMA IF NOT EXISTS ods; +-- Тип таблицы: Heap (стандартная). +-- Обоснование: Необходим row-level UPDATE для реализации UPSERT/SCD1. +-- Использование Append-Only при частых обновлениях приводит к раздуванию (bloat) таблицы. CREATE TABLE IF NOT EXISTS ods.tickets ( ticket_no TEXT NOT NULL, book_ref TEXT NOT NULL, @@ -12,4 +15,7 @@ CREATE TABLE IF NOT EXISTS ods.tickets ( _load_id TEXT NOT NULL, _load_ts TIMESTAMP NOT NULL DEFAULT now() ) +WITH (appendonly=false) DISTRIBUTED BY (ticket_no); + +COMMENT ON TABLE ods.tickets IS 'Билеты (ODS).'; diff --git a/tests/test_ods_sql_contract.py b/tests/test_ods_sql_contract.py index c93596d..32a03b5 100644 --- a/tests/test_ods_sql_contract.py +++ b/tests/test_ods_sql_contract.py @@ -10,13 +10,13 @@ def _read(path: str) -> str: return (PROJECT_ROOT / path).read_text(encoding="utf-8") -def test_snapshot_load_scripts_sync_deleted_keys() -> None: - """Snapshot-таблицы в ODS должны удалять ключи, отсутствующие в текущем батче.""" +def test_snapshot_load_scripts_use_truncate() -> None: + """Snapshot-таблицы в ODS должны использовать паттерн полной перезагрузки (TRUNCATE).""" for entity in SNAPSHOT_ENTITIES: sql = _read(f"sql/ods/{entity}_load.sql") - assert f"DELETE FROM ods.{entity} AS o" in sql - assert "WITH src_keys AS (" in sql - assert "WHERE NOT EXISTS (" in sql + assert f"TRUNCATE TABLE ods.{entity};" in sql + assert f"INSERT INTO ods.{entity} (" in sql + assert "WHERE s.rn = 1;" in sql def test_snapshot_dq_checks_extra_keys() -> None: