diff --git a/Makefile b/Makefile index 858994b..9030c2b 100644 --- a/Makefile +++ b/Makefile @@ -36,6 +36,9 @@ logs: gp-psql: docker compose -f docker-compose.yml exec greenplum bash -c "su - gpadmin -c '/usr/local/greenplum-db/bin/psql -p 5432 -d gp_dwh'" +dwh-truncate: + docker compose -f docker-compose.yml exec greenplum bash -c "su - gpadmin -c 'cd /sql && /usr/local/greenplum-db/bin/psql -d gp_dwh -f truncate_gp.sql'" + ddl-gp: docker compose -f docker-compose.yml exec greenplum bash -c "su - gpadmin -c 'cd /sql && /usr/local/greenplum-db/bin/psql -d gp_dwh -f ddl_gp.sql'" diff --git a/airflow/dags/bookings_dds_ddl.py b/airflow/dags/bookings_dds_ddl.py index f1e00b7..93f9d9a 100644 --- a/airflow/dags/bookings_dds_ddl.py +++ b/airflow/dags/bookings_dds_ddl.py @@ -21,7 +21,7 @@ default_args = {"owner": "airflow", "retries": 1, "retry_delay": timedelta(secon with DAG( dag_id="bookings_dds_ddl", - start_date=pendulum.datetime(2024, 1, 1, tz="UTC"), + start_date=pendulum.datetime(2017, 1, 1, tz="UTC"), schedule=None, catchup=False, template_searchpath="/sql", diff --git a/airflow/dags/bookings_dm_ddl.py b/airflow/dags/bookings_dm_ddl.py index 45c6b5b..540100e 100644 --- a/airflow/dags/bookings_dm_ddl.py +++ b/airflow/dags/bookings_dm_ddl.py @@ -24,7 +24,7 @@ default_args = {"owner": "airflow", "retries": 1, "retry_delay": timedelta(secon with DAG( dag_id="bookings_dm_ddl", - start_date=pendulum.datetime(2024, 1, 1, tz="UTC"), + start_date=pendulum.datetime(2017, 1, 1, tz="UTC"), schedule=None, catchup=False, template_searchpath="/sql", diff --git a/airflow/dags/bookings_ods_ddl.py b/airflow/dags/bookings_ods_ddl.py index 80dde5b..eb226b9 100644 --- a/airflow/dags/bookings_ods_ddl.py +++ b/airflow/dags/bookings_ods_ddl.py @@ -21,7 +21,7 @@ default_args = {"owner": "airflow", "retries": 1, "retry_delay": timedelta(secon with DAG( dag_id="bookings_ods_ddl", - start_date=pendulum.datetime(2024, 1, 1, tz="UTC"), + start_date=pendulum.datetime(2017, 1, 1, tz="UTC"), schedule=None, catchup=False, template_searchpath="/sql", diff --git a/airflow/dags/bookings_stg_ddl.py b/airflow/dags/bookings_stg_ddl.py index 7e6bc56..43d7bf2 100644 --- a/airflow/dags/bookings_stg_ddl.py +++ b/airflow/dags/bookings_stg_ddl.py @@ -21,7 +21,7 @@ default_args = {"owner": "airflow", "retries": 1, "retry_delay": timedelta(secon with DAG( dag_id="bookings_stg_ddl", - start_date=pendulum.datetime(2024, 1, 1, tz="UTC"), + start_date=pendulum.datetime(2017, 1, 1, tz="UTC"), schedule=None, catchup=False, template_searchpath="/sql", diff --git a/airflow/dags/bookings_to_gp_dds.py b/airflow/dags/bookings_to_gp_dds.py index a194b5f..f09096a 100644 --- a/airflow/dags/bookings_to_gp_dds.py +++ b/airflow/dags/bookings_to_gp_dds.py @@ -37,7 +37,7 @@ def _finish_summary() -> None: with DAG( dag_id="bookings_to_gp_dds", - start_date=pendulum.datetime(2024, 1, 1, tz="UTC"), + start_date=pendulum.datetime(2017, 1, 1, tz="UTC"), schedule=None, catchup=False, max_active_runs=1, diff --git a/airflow/dags/bookings_to_gp_dm.py b/airflow/dags/bookings_to_gp_dm.py index 1683fa9..e35bd2f 100644 --- a/airflow/dags/bookings_to_gp_dm.py +++ b/airflow/dags/bookings_to_gp_dm.py @@ -40,7 +40,7 @@ def _finish_summary() -> None: with DAG( dag_id="bookings_to_gp_dm", - start_date=pendulum.datetime(2024, 1, 1, tz="UTC"), + start_date=pendulum.datetime(2017, 1, 1, tz="UTC"), schedule=None, catchup=False, max_active_runs=1, diff --git a/airflow/dags/bookings_to_gp_ods.py b/airflow/dags/bookings_to_gp_ods.py index 1dfc1fe..f5637da 100644 --- a/airflow/dags/bookings_to_gp_ods.py +++ b/airflow/dags/bookings_to_gp_ods.py @@ -125,7 +125,7 @@ def _finish_summary() -> None: with DAG( dag_id="bookings_to_gp_ods", - start_date=pendulum.datetime(2024, 1, 1, tz="UTC"), + start_date=pendulum.datetime(2017, 1, 1, tz="UTC"), schedule=None, catchup=False, max_active_runs=1, diff --git a/airflow/dags/bookings_to_gp_stage.py b/airflow/dags/bookings_to_gp_stage.py index bcacfd4..c6dc0bf 100644 --- a/airflow/dags/bookings_to_gp_stage.py +++ b/airflow/dags/bookings_to_gp_stage.py @@ -68,7 +68,7 @@ def _finish_summary() -> None: with DAG( dag_id="bookings_to_gp_stage", - start_date=pendulum.datetime(2024, 1, 1, tz="UTC"), + start_date=pendulum.datetime(2017, 1, 1, tz="UTC"), schedule=None, catchup=False, max_active_runs=1, diff --git a/airflow/dags/csv_to_greenplum.py b/airflow/dags/csv_to_greenplum.py index ad5273a..8049b4f 100644 --- a/airflow/dags/csv_to_greenplum.py +++ b/airflow/dags/csv_to_greenplum.py @@ -130,7 +130,7 @@ default_args = {"owner": "airflow", "retries": 1, "retry_delay": timedelta(secon with DAG( dag_id="csv_to_greenplum", - start_date=datetime(2024, 1, 1), + start_date=datetime(2017, 1, 1), schedule=None, catchup=False, default_args=default_args, diff --git a/airflow/dags/csv_to_greenplum_dq.py b/airflow/dags/csv_to_greenplum_dq.py index 02ae600..6893d5f 100644 --- a/airflow/dags/csv_to_greenplum_dq.py +++ b/airflow/dags/csv_to_greenplum_dq.py @@ -52,7 +52,7 @@ default_args = {"owner": "airflow", "retries": 1, "retry_delay": timedelta(secon with DAG( dag_id="csv_to_greenplum_dq", - start_date=datetime(2024, 1, 1), + start_date=datetime(2017, 1, 1), schedule=None, catchup=False, default_args=default_args, diff --git a/airflow/dags/ddl_greenplum_base.py b/airflow/dags/ddl_greenplum_base.py index 6afd76c..cfd7e64 100644 --- a/airflow/dags/ddl_greenplum_base.py +++ b/airflow/dags/ddl_greenplum_base.py @@ -17,7 +17,7 @@ default_args = {"owner": "airflow", "retries": 1, "retry_delay": timedelta(secon with DAG( dag_id="orders_base_ddl", - start_date=datetime(2024, 1, 1), + start_date=datetime(2017, 1, 1), schedule=None, catchup=False, template_searchpath="/sql", diff --git a/docs/e2e-etl-test-protocol.md b/docs/e2e-etl-test-protocol.md new file mode 100644 index 0000000..1cf62bf --- /dev/null +++ b/docs/e2e-etl-test-protocol.md @@ -0,0 +1,103 @@ +# Протокол сквозного (E2E) тестирования ETL + +Этот документ описывает процедуру полной проверки цепочки ETL: `STG -> ODS -> DDS -> DM`. +Цель теста — убедиться в корректности инкрементальной загрузки, работы паттерна `Temporary Table` и механизмов `HWM`. + +--- + +## 1. Подготовка окружения + +Убедитесь, что все сервисы запущены и DDL применен. + +```bash +make up +make bookings-init +make ddl-gp +``` + +### Важно: Настройка дат +Для корректного тестирования исторических данных из `demodb` (начинаются с 2017 года), убедитесь, что в DAG-файлах `start_date` установлен в `2017-01-01`. + +--- + +## 2. Очистка данных (Reset) + +Перед началом теста необходимо полностью очистить все слои DWH. + +```bash +make dwh-truncate +``` + +--- + +## 3. Этап 1: Загрузка за первый день (2017-01-01) + +Выполните последовательный запуск всех DAG для первой порции данных. + +```bash +# Загрузка в STG (создает первый батч в источнике) +docker compose exec airflow-scheduler airflow dags test bookings_to_gp_stage 2017-01-01 + +# Загрузка в 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 +``` + +### Ожидаемые результаты (Day 1) +Проверьте наполнение таблиц: +- `stg.bookings` и `ods.bookings` должны иметь одинаковое количество строк (>0). +- `dm.sales_report` должна содержать агрегированные данные за первый день. + +--- + +## 4. Этап 2: Проверка инкремента (2017-01-02) + +Эмулируйте появление данных за второй день и проверьте дозагрузку. + +```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 +``` + +--- + +## 5. Финальная верификация (Критерии успеха) + +Выполните SQL-запрос для сверки данных: + +```bash +make gp-psql -c " +SELECT 'STG' as layer, COUNT(*) FROM stg.bookings +UNION ALL +SELECT 'ODS' as layer, COUNT(*) FROM ods.bookings +UNION ALL +SELECT 'DDS' as layer, COUNT(*) FROM dds.fact_flight_sales +UNION ALL +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` во всех слоях содержат метки соответствующих запусков. + +--- + +## Типичные ошибки +- **Пустые таблицы**: Проверьте, что в `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`. diff --git a/sql/truncate_gp.sql b/sql/truncate_gp.sql new file mode 100644 index 0000000..77594c4 --- /dev/null +++ b/sql/truncate_gp.sql @@ -0,0 +1,36 @@ +-- Скрипт для полной очистки всех слоев DWH (STG, ODS, DDS, DM). +-- Используется для сброса состояния перед проведением E2E-тестов. + +-- 1. Слой STG (Сырые данные) +TRUNCATE stg.bookings, + stg.tickets, + stg.airports, + stg.airplanes, + stg.routes, + stg.seats, + stg.flights, + stg.segments, + stg.boarding_passes; + +-- 2. Слой ODS (Текущее состояние, SCD1) +TRUNCATE ods.bookings, + ods.tickets, + ods.airports, + ods.airplanes, + ods.routes, + ods.seats, + ods.flights, + ods.segments, + ods.boarding_passes; + +-- 3. Слой DDS (Схема "Звезда", SCD1/SCD2) +TRUNCATE dds.dim_calendar, + dds.dim_airports, + dds.dim_airplanes, + dds.dim_tariffs, + dds.dim_passengers, + dds.dim_routes, + dds.fact_flight_sales; + +-- 4. Слой DM (Витрины данных) +TRUNCATE dm.sales_report;