Исправлены критические замечания

This commit is contained in:
2026-01-18 21:33:46 +03:00
parent 762d539d52
commit ba74b521dc
11 changed files with 143 additions and 90 deletions
+19 -1
View File
@@ -85,7 +85,25 @@ make bookings-init
- проверяет количество строк в том же окне инкремента, а также ссылочную целостность и обязательные поля;
- при проблемах делает `RAISE EXCEPTION`, чтобы DAG падал “красным”.
6) `finish_summary`
6) Справочники (full load)
Каждый справочник загружается “снэпшотом” (все строки) и затем проверяется DQ-скриптом:
- `load_airports_to_stg``check_airports_dq` (`sql/stg/airports_load.sql`, `sql/stg/airports_dq.sql`)
- `load_airplanes_to_stg``check_airplanes_dq` (`sql/stg/airplanes_load.sql`, `sql/stg/airplanes_dq.sql`)
- `load_routes_to_stg``check_routes_dq` (`sql/stg/routes_load.sql`, `sql/stg/routes_dq.sql`)
- `load_seats_to_stg``check_seats_dq` (`sql/stg/seats_load.sql`, `sql/stg/seats_dq.sql`)
7) Транзакции
- `load_flights_to_stg``check_flights_dq` (инкремент по `scheduled_departure`)
- `load_segments_to_stg``check_segments_dq` (инкремент по `book_date` через tickets/bookings)
- `load_boarding_passes_to_stg``check_boarding_passes_dq` (full snapshot)
Важно: для некоторых таблиц “пустое окно инкремента” считается ошибкой (DQ делает `RAISE EXCEPTION`),
а для `boarding_passes` DQ может быть пропущена, если в источнике 0 строк.
8) `finish_summary`
- логирует краткую сводку в конце запуска.
+25 -32
View File
@@ -26,52 +26,46 @@
---
## 2) Критичные замечания (исправить перед тем, как показывать как эталон)
## 2) Критичные замечания (статус на текущий момент)
### 2.1. Некорректные утверждения про MPP и co-location (вводят студентов в заблуждение)
### 2.1. Некорректные утверждения про MPP и co-location (статус: исправлено)
В Greenplum производительность JOIN сильно зависит от распределения данных по сегментам.
Если ключ распределения двух таблиц совпадает с ключом JOIN — часто удаётся обойтись без перераспределения данных (motion).
Проблема: в некоторых DDL-комментариях сейчас обещается co-location там, где его не будет.
Раньше в некоторых DDL-комментариях обещалась co-location там, где её не будет.
Это педагогически опасно: студент запоминает неверную модель, а потом “не понимает”, почему запросы медленные.
Примеры мест, которые стоит скорректировать:
- `sql/stg/flights_ddl.sql`: `stg.flights` распределена по `flight_id`, а `stg.boarding_passes` — по `ticket_no`,
поэтому “co-location flights и boarding_passes при JOIN по flight_id” не выполняется.
- `sql/stg/routes_ddl.sql`: распределение `stg.routes` по `route_no` не даёт co-location с `stg.airports` (которая по `airport_code`)
и `stg.airplanes` (которая по `airplane_code`) при типичных JOIN’ах.
Что сделано:
- DDL-комментарии приведены к честной формулировке “ключ выбран так-то, но JOIN по другим ключам может требовать motion”.
- Исправлены места, где co-location заявлялась ошибочно (в т.ч. `routes`, `flights`, `segments`, `boarding_passes`).
Рекомендация: либо исправить распределение (если это действительно важно для учебного кейса),
либо **честно переписать комментарии**: “ключ выбран так-то, но JOIN по другим ключам может требовать motion”.
Файлы: `sql/stg/routes_ddl.sql`, `sql/stg/flights_ddl.sql`, `sql/stg/segments_ddl.sql`, `sql/stg/boarding_passes_ddl.sql`.
### 2.2. DQ-проверки ссылочной целостности иногда “смотрят в историю”, а не в текущий батч
### 2.2. DQ-проверки ссылочной целостности: “текущий батч” vs “вся история” (статус: зафиксировано и частично усилено)
Часть DQ-скриптов проверяет наличие “родительских” записей в таблице **без фильтра `batch_id`**.
При append-only истории это может скрыть проблемы текущей загрузки:
родитель был загружен в прошлом батче → проверка пройдёт, даже если текущий батч родителя не загрузил.
Пример:
- `sql/stg/routes_dq.sql` проверяет airports/airplanes без ограничения на `batch_id`.
Что сделано для справочников (snapshot), которые загружаются каждый запуск:
- `routes_dq.sql`: проверка airports/airplanes стала батч-строгой (`batch_id = текущий батч`).
- `seats_dq.sql`: проверка airplanes стала батч-строгой (`batch_id = текущий батч`).
- `flights_dq.sql`: проверка routes стала батч-строгой (`batch_id = текущий батч`).
Рекомендация: для snapshot-таблиц (справочники и boarding_passes) использовать батч-строгую проверку:
“в текущем `batch_id` все ссылки указывают на строки текущего `batch_id`”.
Это лучше учит идее “консистентность батча” и упрощает отладку.
Почему не всё делаем батч-строго:
- Для инкрементальных таблиц (например, `segments`) ссылки могут указывать на данные,
загруженные в предыдущих батчах → там корректнее проверять “существует в STG вообще”, а не “существует в текущем батче”.
### 2.3. Smoke-тесты DAG’ов есть, но почти не проверяют граф
### 2.3. Smoke-тесты DAG’ов (статус: исправлено)
В `tests/test_dags_smoke.py` новые тесты в основном проверяют “таски существуют” через `dag.has_task(...)`.
Как учебный пример теста это слабовато: студент видит тест, но не понимает, что именно он защищает.
Что сделано:
- Тесты усилены: теперь проверяются ключевые зависимости графа через `get_direct_relatives("downstream")`.
Рекомендация: тестировать зависимости так же, как это уже сделано для `csv_to_greenplum`
(через `dag.get_task(...).get_direct_relatives("downstream")`).
### 2.4. Документация по DAG (статус: синхронизировано базово)
### 2.4. Документация по DAG отстаёт от реальной логики
`docs/bookings_to_gp_stage.md` описывает только загрузку `bookings` и `tickets`,
но DAG теперь загружает ещё 7 таблиц (справочники и транзакции).
Рекомендация: обновить документ, чтобы студент мог запустить пайплайн “по инструкции” без сюрпризов.
Что сделано:
- `docs/bookings_to_gp_stage.md` обновлён так, чтобы отражать текущий набор таблиц и шагов пайплайна.
---
@@ -139,9 +133,8 @@ assert airports_load in tickets_dq.get_direct_relatives("downstream")
## 5) Чек-лист “готово как эталон”
- [ ] В DDL-комментариях нет неверных обещаний про co-location/уникальность ключей.
- [x] В DDL-комментариях нет неверных обещаний про co-location/уникальность ключей.
- [ ] Для DQ определена и описана политика “0 строк”: где fail, где skip.
- [ ] DQ ссылочной целостности не маскирует проблемы текущего батча (batch-строгие проверки там, где это уместно).
- [ ] `docs/bookings_to_gp_stage.md` соответствует фактическому DAG.
- [ ] Smoke-тесты проверяют хотя бы критические зависимости графа.
- [x] DQ ссылочной целостности не маскирует проблемы текущего батча (batch-строгие проверки там, где это уместно).
- [x] `docs/bookings_to_gp_stage.md` соответствует фактическому DAG.
- [x] Smoke-тесты проверяют хотя бы критические зависимости графа.
+14 -9
View File
@@ -105,16 +105,17 @@
#### stg.tickets (транзакции, инкремент)
- **Источник:** `bookings.tickets` (через PXF)
- **Ключ распределения:** `ticket_no`
- **Ключ распределения:** `book_ref`
- **Примечание:** JOIN `tickets``segments`/`boarding_passes` по `ticket_no` может требовать motion (ключи распределения разные).
- **Бизнес-колонки:**
- `ticket_no TEXT` - номер билета
- `book_ref TEXT` - номер бронирования
- `passenger_id TEXT` - идентификатор пассажира
- `passenger_name TEXT` - имя пассажира
- `contact_data TEXT` - контактные данные (JSONB)
- `outbound TEXT` - направление (в источнике boolean)
- **Технические колонки:** `src_created_at_ts` (из book_date через bookings), `load_dttm`, `batch_id`
- **Стратегия загрузки:** Инкремент по `book_date` (через bookings)
- **DQ проверки:** count (окно инкремента), дубликаты ticket_no, NULL обязательных полей, ссылочная целостность
- **DQ проверки:** count (окно инкремента), дубликаты ticket_no, NULL обязательных полей, пустой passenger_name, ссылочная целостность (bookings)
#### stg.airports (справочник, full load)
- **Источник:** `bookings.airports_data` (через PXF)
@@ -145,6 +146,7 @@
#### stg.routes (справочник, full load)
- **Источник:** `bookings.routes` (через PXF)
- **Ключ распределения:** `route_no`
- **Примечание:** JOIN по `departure_airport`/`arrival_airport`/`airplane_code` может требовать motion (ключи распределения разные).
- **Бизнес-колонки:**
- `route_no TEXT` - номер маршрута
- `validity TEXT` - период действия (из tstzrange)
@@ -156,7 +158,7 @@
- `duration TEXT` - длительность
- **Технические колонки:** `src_created_at_ts`, `load_dttm`, `batch_id`
- **Стратегия загрузки:** Full load
- **DQ проверки:** count, дубликаты (route_no, validity), NULL обязательных полей, ссылочная целостность
- **DQ проверки:** count, дубликаты (route_no, validity), NULL обязательных полей, ссылочная целостность (batch_id = текущий батч)
#### stg.seats (справочник, full load)
- **Источник:** `bookings.seats` (через PXF)
@@ -167,11 +169,12 @@
- `fare_conditions TEXT` - класс обслуживания
- **Технические колонки:** `src_created_at_ts`, `load_dttm`, `batch_id`
- **Стратегия загрузки:** Full load
- **DQ проверки:** count, дубликаты (airplane_code, seat_no), NULL обязательных полей, ссылочная целостность
- **DQ проверки:** count, дубликаты (airplane_code, seat_no), NULL обязательных полей, ссылочная целостность (batch_id = текущий батч)
#### stg.flights (транзакции, инкремент)
- **Источник:** `bookings.flights` (через PXF)
- **Ключ распределения:** `flight_id`
- **Примечание:** JOIN с таблицами, распределёнными по другим ключам, может требовать motion.
- **Бизнес-колонки:**
- `flight_id TEXT` - идентификатор рейса
- `route_no TEXT` - номер маршрута
@@ -182,11 +185,12 @@
- `actual_arrival TEXT` - фактическое время прилёта
- **Технические колонки:** `src_created_at_ts` (=scheduled_departure), `load_dttm`, `batch_id`
- **Стратегия загрузки:** Инкремент по `scheduled_departure`
- **DQ проверки:** count (окно инкремента), дубликаты flight_id, NULL обязательных полей, ссылочная целостность
- **DQ проверки:** count (окно инкремента), дубликаты flight_id, NULL обязательных полей, ссылочная целостность (routes, batch_id = текущий батч)
#### stg.segments (транзакции, инкремент)
- **Источник:** `bookings.segments` (через PXF)
- **Ключ распределения:** `ticket_no` (co-location с tickets)
- **Ключ распределения:** `ticket_no` (co-location с boarding_passes)
- **Примечание:** JOIN `segments``tickets` по `ticket_no` может требовать motion (stg.tickets распределена по `book_ref`).
- **Бизнес-колонки:**
- `ticket_no TEXT` - номер билета
- `flight_id TEXT` - идентификатор рейса
@@ -194,11 +198,12 @@
- `price TEXT` - цена
- **Технические колонки:** `src_created_at_ts` (из book_date через tickets), `load_dttm`, `batch_id`
- **Стратегия загрузки:** Инкремент по `book_date` (через tickets)
- **DQ проверки:** count (окно инкремента), дубликаты (ticket_no, flight_id), NULL обязательных полей, ссылочная целостность
- **DQ проверки:** count (окно инкремента), дубликаты (ticket_no, flight_id), NULL обязательных полей, ссылочная целостность (tickets, flights)
#### stg.boarding_passes (транзакции, full snapshot)
- **Источник:** `bookings.boarding_passes` (через PXF)
- **Ключ распределения:** `ticket_no` (co-location с tickets/segments)
- **Ключ распределения:** `ticket_no` (co-location с segments)
- **Примечание:** JOIN `boarding_passes``tickets` по `ticket_no` может требовать motion (stg.tickets распределена по `book_ref`).
- **Бизнес-колонки:**
- `ticket_no TEXT` - номер билета
- `flight_id TEXT` - идентификатор рейса
+3 -3
View File
@@ -32,7 +32,7 @@ WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1)
-- Ключ распределения: ticket_no
-- Обоснование: ticket_no — это основной бизнес-ключ для билетов.
-- Использование ticket_no обеспечивает:
-- 1. Co-location данных boarding_passes и tickets при JOIN по ticket_no
-- 2. Co-location данных boarding_passes и segments при JOIN по ticket_no
-- 3. Равномерное распределение данных по сегментам (ticket_no имеет высокую кардинальность)
-- 1. Co-location данных boarding_passes и segments при JOIN по ticket_no
-- 2. Равномерное распределение данных по сегментам (ticket_no имеет высокую кардинальность)
-- Примечание: stg.tickets распределена по book_ref, поэтому JOIN boarding_passes ↔ tickets по ticket_no может требовать motion.
DISTRIBUTED BY (ticket_no);
+1 -1
View File
@@ -38,5 +38,5 @@ WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1)
-- Использование flight_id обеспечивает:
-- 1. Равномерное распределение данных по сегментам (flight_id имеет высокую кардинальность)
-- 2. Оптимизацию запросов, которые фильтруют или группируют по flight_id
-- 3. Co-location данных flights и boarding_passes при JOIN по flight_id
-- Примечание: JOIN с таблицами, распределёнными по другим ключам, может требовать motion.
DISTRIBUTED BY (flight_id);
+3 -1
View File
@@ -76,7 +76,9 @@ BEGIN
SELECT COUNT(*)
INTO v_orphan_route_count
FROM stg.flights AS f
LEFT JOIN stg.routes AS r ON f.route_no = r.route_no
LEFT JOIN stg.routes AS r
ON f.route_no = r.route_no
AND r.batch_id = v_batch_id
WHERE f.batch_id = v_batch_id
AND r.route_no IS NULL;
+3 -3
View File
@@ -37,9 +37,9 @@ CREATE TABLE IF NOT EXISTS stg.routes (
)
WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1)
-- Ключ распределения: route_no
-- Обоснование: route_no — это уникальный идентификатор маршрута.
-- Обоснование: route_no — бизнес-идентификатор маршрута и часто используется в фильтрах/джойнах.
-- Использование route_no обеспечивает:
-- 1. Равномерное распределение данных по сегментам (route_no имеет высокую кардинальность)
-- 2. Co-location данных routes с airports и airplanes при JOIN
-- 3. Оптимизацию запросов, которые фильтруют или группируют по route_no
-- 2. Оптимизацию запросов, которые фильтруют или группируют по route_no
-- Примечание: JOIN по airport_code/airplane_code может требовать перераспределения данных (motion).
DISTRIBUTED BY (route_no);
+9 -3
View File
@@ -67,7 +67,9 @@ BEGIN
SELECT COUNT(*)
INTO v_orphan_airports_count
FROM stg.routes AS r
LEFT JOIN stg.airports AS da ON r.departure_airport = da.airport_code
LEFT JOIN stg.airports AS da
ON r.departure_airport = da.airport_code
AND da.batch_id = v_batch_id
WHERE r.batch_id = v_batch_id
AND da.airport_code IS NULL;
@@ -82,7 +84,9 @@ BEGIN
SELECT COUNT(*)
INTO v_orphan_airports_count
FROM stg.routes AS r
LEFT JOIN stg.airports AS aa ON r.arrival_airport = aa.airport_code
LEFT JOIN stg.airports AS aa
ON r.arrival_airport = aa.airport_code
AND aa.batch_id = v_batch_id
WHERE r.batch_id = v_batch_id
AND aa.airport_code IS NULL;
@@ -97,7 +101,9 @@ BEGIN
SELECT COUNT(*)
INTO v_orphan_airplanes_count
FROM stg.routes AS r
LEFT JOIN stg.airplanes AS a ON r.airplane_code = a.airplane_code
LEFT JOIN stg.airplanes AS a
ON r.airplane_code = a.airplane_code
AND a.batch_id = v_batch_id
WHERE r.batch_id = v_batch_id
AND a.airplane_code IS NULL;
+3 -1
View File
@@ -65,7 +65,9 @@ BEGIN
SELECT COUNT(*)
INTO v_orphan_airplanes_count
FROM stg.seats AS s
LEFT JOIN stg.airplanes AS a ON s.airplane_code = a.airplane_code
LEFT JOIN stg.airplanes AS a
ON s.airplane_code = a.airplane_code
AND a.batch_id = v_batch_id
WHERE s.batch_id = v_batch_id
AND a.airplane_code IS NULL;
+3 -3
View File
@@ -30,7 +30,7 @@ WITH (appendonly=true, orientation=row, compresstype=zlib, compresslevel=1)
-- Ключ распределения: ticket_no
-- Обоснование: ticket_no — это основной бизнес-ключ для билетов.
-- Использование ticket_no обеспечивает:
-- 1. Co-location данных segments и tickets при JOIN по ticket_no
-- 2. Co-location данных segments и boarding_passes при JOIN по ticket_no
-- 3. Равномерное распределение данных по сегментам (ticket_no имеет высокую кардинальность)
-- 1. Co-location данных segments и boarding_passes при JOIN по ticket_no
-- 2. Равномерное распределение данных по сегментам (ticket_no имеет высокую кардинальность)
-- Примечание: stg.tickets распределена по book_ref, поэтому JOIN segments ↔ tickets по ticket_no может требовать motion.
DISTRIBUTED BY (ticket_no);
+60 -33
View File
@@ -89,17 +89,25 @@ def test_bookings_stg_ddl_dag_structure():
}
assert expected_tasks.issubset(dag.task_dict.keys())
# Проверка линейных зависимостей
# Справочники создаются после bookings/tickets
assert dag.has_task("apply_stg_bookings_ddl")
assert dag.has_task("apply_stg_tickets_ddl")
assert dag.has_task("apply_stg_airports_ddl")
assert dag.has_task("apply_stg_airplanes_ddl")
assert dag.has_task("apply_stg_routes_ddl")
assert dag.has_task("apply_stg_seats_ddl")
assert dag.has_task("apply_stg_flights_ddl")
assert dag.has_task("apply_stg_segments_ddl")
assert dag.has_task("apply_stg_boarding_passes_ddl")
# Линейные зависимости: bookings/tickets → справочники → транзакции
t_bookings = dag.get_task("apply_stg_bookings_ddl")
t_tickets = dag.get_task("apply_stg_tickets_ddl")
t_airports = dag.get_task("apply_stg_airports_ddl")
t_airplanes = dag.get_task("apply_stg_airplanes_ddl")
t_routes = dag.get_task("apply_stg_routes_ddl")
t_seats = dag.get_task("apply_stg_seats_ddl")
t_flights = dag.get_task("apply_stg_flights_ddl")
t_segments = dag.get_task("apply_stg_segments_ddl")
t_boarding = dag.get_task("apply_stg_boarding_passes_ddl")
assert t_tickets in t_bookings.get_direct_relatives("downstream")
assert t_airports in t_tickets.get_direct_relatives("downstream")
assert t_airplanes in t_airports.get_direct_relatives("downstream")
assert t_routes in t_airplanes.get_direct_relatives("downstream")
assert t_seats in t_routes.get_direct_relatives("downstream")
assert t_flights in t_seats.get_direct_relatives("downstream")
assert t_segments in t_flights.get_direct_relatives("downstream")
assert t_boarding in t_segments.get_direct_relatives("downstream")
def test_bookings_to_gp_stage_dag_structure():
@@ -130,25 +138,44 @@ def test_bookings_to_gp_stage_dag_structure():
}
assert expected_tasks.issubset(dag.task_dict.keys())
# Проверка линейных зависимостей
# bookings/tickets → справочники → транзакции → финальный лог
assert dag.has_task("generate_bookings_day")
assert dag.has_task("load_bookings_to_stg")
assert dag.has_task("check_row_counts")
assert dag.has_task("load_tickets_to_stg")
assert dag.has_task("check_tickets_dq")
assert dag.has_task("load_airports_to_stg")
assert dag.has_task("check_airports_dq")
assert dag.has_task("load_airplanes_to_stg")
assert dag.has_task("check_airplanes_dq")
assert dag.has_task("load_routes_to_stg")
assert dag.has_task("check_routes_dq")
assert dag.has_task("load_seats_to_stg")
assert dag.has_task("check_seats_dq")
assert dag.has_task("load_flights_to_stg")
assert dag.has_task("check_flights_dq")
assert dag.has_task("load_segments_to_stg")
assert dag.has_task("check_segments_dq")
assert dag.has_task("load_boarding_passes_to_stg")
assert dag.has_task("check_boarding_passes_dq")
assert dag.has_task("finish_summary")
# Линейные зависимости: bookings/tickets → справочники → транзакции → финальный лог
t_generate = dag.get_task("generate_bookings_day")
t_bookings = dag.get_task("load_bookings_to_stg")
t_bookings_dq = dag.get_task("check_row_counts")
t_tickets = dag.get_task("load_tickets_to_stg")
t_tickets_dq = dag.get_task("check_tickets_dq")
t_airports = dag.get_task("load_airports_to_stg")
t_airports_dq = dag.get_task("check_airports_dq")
t_airplanes = dag.get_task("load_airplanes_to_stg")
t_airplanes_dq = dag.get_task("check_airplanes_dq")
t_routes = dag.get_task("load_routes_to_stg")
t_routes_dq = dag.get_task("check_routes_dq")
t_seats = dag.get_task("load_seats_to_stg")
t_seats_dq = dag.get_task("check_seats_dq")
t_flights = dag.get_task("load_flights_to_stg")
t_flights_dq = dag.get_task("check_flights_dq")
t_segments = dag.get_task("load_segments_to_stg")
t_segments_dq = dag.get_task("check_segments_dq")
t_boarding = dag.get_task("load_boarding_passes_to_stg")
t_boarding_dq = dag.get_task("check_boarding_passes_dq")
t_finish = dag.get_task("finish_summary")
assert t_bookings in t_generate.get_direct_relatives("downstream")
assert t_bookings_dq in t_bookings.get_direct_relatives("downstream")
assert t_tickets in t_bookings_dq.get_direct_relatives("downstream")
assert t_tickets_dq in t_tickets.get_direct_relatives("downstream")
assert t_airports in t_tickets_dq.get_direct_relatives("downstream")
assert t_airports_dq in t_airports.get_direct_relatives("downstream")
assert t_airplanes in t_airports_dq.get_direct_relatives("downstream")
assert t_airplanes_dq in t_airplanes.get_direct_relatives("downstream")
assert t_routes in t_airplanes_dq.get_direct_relatives("downstream")
assert t_routes_dq in t_routes.get_direct_relatives("downstream")
assert t_seats in t_routes_dq.get_direct_relatives("downstream")
assert t_seats_dq in t_seats.get_direct_relatives("downstream")
assert t_flights in t_seats_dq.get_direct_relatives("downstream")
assert t_flights_dq in t_flights.get_direct_relatives("downstream")
assert t_segments in t_flights_dq.get_direct_relatives("downstream")
assert t_segments_dq in t_segments.get_direct_relatives("downstream")
assert t_boarding in t_segments_dq.get_direct_relatives("downstream")
assert t_boarding_dq in t_boarding.get_direct_relatives("downstream")
assert t_finish in t_boarding_dq.get_direct_relatives("downstream")