From 15081384f4d84d17bcf9f699c0dd2095cb0d27a5 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sun, 18 Jan 2026 10:37:59 +0300 Subject: [PATCH] =?UTF-8?q?=D0=97=D0=B0=D0=BC=D0=B5=D1=87=D0=B0=D0=BD?= =?UTF-8?q?=D0=B8=D1=8F=20=D0=BE=D1=82=20codex?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- plans/stg_layer_implementation_plan.md | 52 ++++++++++++++++---------- 1 file changed, 33 insertions(+), 19 deletions(-) diff --git a/plans/stg_layer_implementation_plan.md b/plans/stg_layer_implementation_plan.md index c149f6c..80a69a6 100644 --- a/plans/stg_layer_implementation_plan.md +++ b/plans/stg_layer_implementation_plan.md @@ -6,20 +6,23 @@ ## Обзор задачи -Согласно [`docs/internal/db_schema.md`](../docs/internal/db_schema.md:31), STG слой реализован частично (2 из 9 таблиц: bookings, tickets). Необходимо реализовать оставшиеся 7 таблиц. +Согласно [`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) | **Инкремент по дате** | Больший объём, необходимость отслеживания изменений | +| Транзакции (flights, segments) | **Инкремент** | Больший объём; выбираем максимально естественное опорное поле | +| Транзакции (boarding_passes) | **Full snapshot** | В источнике строки создаются и обновляются со временем, простого инкремента без усложнений нет | + +> Важно: под **Full load** в STG подразумеваем «сняли слепок и дописали в append-only таблицу с `batch_id`», а не `TRUNCATE + INSERT`. Это даёт простую идемпотентность (по `batch_id`) и сохраняет историю загрузок. ### Оценка размера справочников | Справочник | Примерный размер | Оценка | |-------------|------------------|---------| -| **airports** | ~700-800 аэропортов | **Маленький** | +| **airports** | ~5K-6K аэропортов | **Маленький** | | **airplanes** | ~10 моделей самолётов | **Крошечный** | | **seats** | ~1700-2000 записей | **Маленький** | | **routes** | Ожидается несколько тысяч | **Маленький/Средний** | @@ -36,7 +39,7 @@ | `bookings.seats` | `stg.seats` | Справочник | Full | - | | `bookings.flights` | `stg.flights` | Транзакции | Инкремент | `scheduled_departure` | | `bookings.segments` | `stg.segments` | Транзакции | Инкремент | `book_date` (через tickets) | -| `bookings.boarding_passes` | `stg.boarding_passes` | Транзакции | Инкремент | `book_date` (через tickets) | +| `bookings.boarding_passes` | `stg.boarding_passes` | Транзакции | Full (snapshot) | - | ## Паттерн реализации (на основе bookings/tickets) @@ -69,12 +72,16 @@ CREATE TABLE IF NOT EXISTS stg.{table} ( -- бизнес-колонки как TEXT src_created_at_ts TIMESTAMP, load_dttm TIMESTAMP NOT NULL DEFAULT now(), - batch_id TEXT NOT NULL + batch_id TEXT ) WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1) DISTRIBUTED BY ({distribution_key}); ``` +Примечание по PXF/JDBC типам: для «сложных» типов Postgres (например, `jsonb`, `point`, массивы, `tstzrange`) +чаще всего проще и надёжнее объявлять колонки во внешней таблице как `TEXT`, чтобы избежать несовместимостей +драйвера/маппинга типов. Внутренний STG всё равно хранит бизнес-поля как `TEXT`. + ### Общая структура LOAD файла (Full load для справочников) ```sql @@ -146,6 +153,9 @@ ANALYZE stg.{table}; ### Общая структура DQ файла +Для **инкрементальных** таблиц сравниваем окно инкремента (по `src_created_at_ts`) между источником и STG. +Для **full snapshot** таблиц (справочники и `boarding_passes`) обычно достаточно сравнить общее количество строк в источнике с количеством строк, загруженных в текущий `batch_id`, плюс проверить дубликаты/NULL/ссылочную целостность. + ```sql -- Проверки качества данных для {table} @@ -293,13 +303,13 @@ END $$; - `duration TEXT` - Тех.колонки: `src_created_at_ts`, `load_dttm`, `batch_id` - Распределение: `DISTRIBUTED BY (route_no)` -- Обоснование: `route_no` — это уникальный идентификатор маршрута +- Обоснование: `route_no` — логический идентификатор маршрута; он нужен для JOIN с flights по `route_no` **Загрузка**: Full **DQ проверки**: - Count между источником и STG -- Дубликаты `route_no` +- Дубликаты `(route_no, validity)` - NULL обязательных полей (route_no, departure_airport, arrival_airport, airplane_code) - Ссылочная целостность на airports (departure_airport, arrival_airport) - Ссылочная целостность на airplanes (airplane_code) @@ -355,6 +365,8 @@ END $$; - NULL обязательных полей (flight_id, route_no, status, scheduled_departure) - Ссылочная целостность на routes (route_no) +> Примечание: `flights.status/actual_*` в источнике могут меняться со временем. Для учебного STG можно принять допущение "insert-only" (снимаем слепок на момент загрузки), либо усложнить и перезагружать скользящее окно по датам вылета. + ### 6. segments (транзакции, инкремент) **Внешняя таблица**: `stg.segments_ext` @@ -381,7 +393,7 @@ END $$; - Ссылочная целостность на tickets (ticket_no) - Ссылочная целостность на flights (flight_id) -### 7. boarding_passes (транзакции, инкремент) +### 7. boarding_passes (транзакции, full snapshot) **Внешняя таблица**: `stg.boarding_passes_ext` - Поля: `ticket_no`, `flight_id`, `seat_no`, `boarding_no`, `boarding_time` @@ -395,11 +407,11 @@ END $$; - `boarding_no TEXT` - `boarding_time TEXT` - Тех.колонки: `src_created_at_ts`, `load_dttm`, `batch_id` -- `src_created_at_ts` = берётся из `bookings.book_date` через JOIN с tickets +- `src_created_at_ts` = `now()` (в этой таблице нет удобного поля для инкремента, потому что строки могут создаваться и обновляться со временем) - Распределение: `DISTRIBUTED BY (ticket_no)` - Обоснование: co-location с tickets/segments для оптимизации JOIN -**Загрузка**: Инкремент по `book_date` (как в tickets) +**Загрузка**: Full snapshot (все строки при каждом запуске) **DQ проверки**: - Count между источником и STG @@ -408,6 +420,8 @@ END $$; - Ссылочная целостность на tickets (ticket_no) - Ссылочная целостность на segments (ticket_no, flight_id) +> Примечание: в источнике `boarding_passes` строки сначала создаются при CHECK-IN (без `boarding_time`), а потом обновляются при BOARDING. Поэтому инкремент "по времени" без усложнений будет пропускать часть событий и/или изменения. Для учебного стенда самый стабильный вариант — снимать полный слепок. + ## Обновление существующих DAG ### `airflow/dags/bookings_stg_ddl.py` @@ -454,7 +468,7 @@ apply_stg_segments_ddl = PostgresOperator( apply_stg_boarding_passes_ddl = PostgresOperator( task_id="apply_stg_boarding_passes_ddl", postgres_conn_id=GREENPLUM_CONN_ID, - sql="stg/boardings_ddl.sql", + sql="stg/boarding_passes_ddl.sql", ) ``` @@ -544,13 +558,13 @@ check_segments_dq = PostgresOperator( load_boarding_passes_to_stg = PostgresOperator( task_id="load_boarding_passes_to_stg", postgres_conn_id=GREENPLUM_CONN_ID, - sql="stg/boardings_load.sql", + 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/boardings_dq.sql", + sql="stg/boarding_passes_dq.sql", ) ``` @@ -577,7 +591,7 @@ check_boarding_passes_dq = PostgresOperator( \i stg/seats_ddl.sql \i stg/flights_ddl.sql \i stg/segments_ddl.sql -\i stg/boardings_ddl.sql +\i stg/boarding_passes_ddl.sql ``` ### Добавление тестов @@ -639,7 +653,7 @@ def test_bookings_to_gp_stage_dag_structure(): ### Обновление документации -Обновить статус в [`docs/internal/db_schema.md`](../docs/internal/db_schema.md:31) с "2 из 9" на "9 из 9". +Обновить статус в [`docs/internal/db_schema.md`](../docs/internal/db_schema.md) с "2 из 9" на "9 из 9". Добавить описание новых таблиц в документацию. @@ -685,7 +699,7 @@ graph TB - [ ] `sql/stg/seats_ddl.sql` - [ ] `sql/stg/flights_ddl.sql` - [ ] `sql/stg/segments_ddl.sql` - - [ ] `sql/stg/boardings_ddl.sql` + - [ ] `sql/stg/boarding_passes_ddl.sql` - [ ] Создать файлы LOAD для новых таблиц (7 файлов) - [ ] `sql/stg/airports_load.sql` - [ ] `sql/stg/airplanes_load.sql` @@ -693,7 +707,7 @@ graph TB - [ ] `sql/stg/seats_load.sql` - [ ] `sql/stg/flights_load.sql` - [ ] `sql/stg/segments_load.sql` - - [ ] `sql/stg/boardings_load.sql` + - [ ] `sql/stg/boarding_passes_load.sql` - [ ] Создать файлы DQ для новых таблиц (7 файлов) - [ ] `sql/stg/airports_dq.sql` - [ ] `sql/stg/airplanes_dq.sql` @@ -701,7 +715,7 @@ graph TB - [ ] `sql/stg/seats_dq.sql` - [ ] `sql/stg/flights_dq.sql` - [ ] `sql/stg/segments_dq.sql` - - [ ] `sql/stg/boardings_dq.sql` + - [ ] `sql/stg/boarding_passes_dq.sql` - [ ] Обновить DAG `bookings_stg_ddl.py` - [ ] Обновить DAG `bookings_to_gp_stage.py` - [ ] Обновить `sql/ddl_gp.sql` @@ -721,7 +735,7 @@ graph TB ## Связанные документы -- [`docs/internal/db_schema.md`](../docs/internal/db_schema.md) - Общая схема Б DWH +- [`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